diff --git a/src/sentry/integrations/cursor_origin/repository_events.py b/src/sentry/integrations/cursor_origin/repository_events.py index 7b93074be41d..4d5ec12e36bd 100644 --- a/src/sentry/integrations/cursor_origin/repository_events.py +++ b/src/sentry/integrations/cursor_origin/repository_events.py @@ -4,10 +4,12 @@ from collections.abc import Mapping, Sequence from typing import Any +from sentry.constants import ObjectStatus from sentry.integrations.cursor_origin.constants import CURSOR_ORIGIN_WEB_BASE_URL from sentry.integrations.cursor_origin.handlers import WebhookEventHandler from sentry.integrations.cursor_origin.repository import active_repositories from sentry.integrations.cursor_origin.webhook_types import ( + RepositoryDeletedEvent, RepositoryMetadataEvent, RepositorySnapshot, ) @@ -15,11 +17,21 @@ RpcIntegration, RpcOrganizationIntegration, ) +from sentry.integrations.services.repository import repository_service +from sentry.integrations.source_code_management.repo_audit import log_repo_change +from sentry.integrations.source_code_management.sync_repos import DISABLE_ACTIVITY_CUTOFF_DAYS +from sentry.integrations.types import IntegrationProviderSlug from sentry.integrations.utils.metrics import IntegrationWebhookEventType +from sentry.models.organization import Organization from sentry.models.repository import Repository +from sentry.organizations.services.organization.serial import serialize_rpc_organization +from sentry.plugins.providers.integration_repository import get_integration_repository_provider +from sentry.utils import metrics logger = logging.getLogger("sentry.integrations.cursor_origin") +PROVIDER = f"integrations:{IntegrationProviderSlug.CURSOR_ORIGIN.value}" + class RepositoryMetadataUpdatedHandler(WebhookEventHandler): """Apply a default-branch change. The name is repaired for every event.""" @@ -84,3 +96,125 @@ def store_default_branch(repo: Repository, snapshot: RepositorySnapshot, deliver }, ) repo.update(config={**repo.config, "default_branch": snapshot.default_branch}) + + +def _active_org_integrations( + integration: RpcIntegration, + org_integrations: Sequence[RpcOrganizationIntegration], + delivery_id: str, +) -> list[RpcOrganizationIntegration]: + if integration.status != ObjectStatus.ACTIVE: + logger.info( + "cursor_origin.repository.inactive_integration", + extra={"delivery_id": delivery_id, "integration_id": integration.id}, + ) + return [] + return [oi for oi in org_integrations if oi.status == ObjectStatus.ACTIVE] + + +class RepositoryCreatedHandler(WebhookEventHandler): + EVENT_TYPE = IntegrationWebhookEventType.INBOUND_SYNC + + def __call__( + self, + payload: Mapping[str, Any], + delivery_id: str, + integration: RpcIntegration, + org_integrations: Sequence[RpcOrganizationIntegration], + ) -> None: + snapshot = RepositoryMetadataEvent.from_payload(payload).repository + provider = get_integration_repository_provider(integration) + config = { + "name": snapshot.full_name, + "external_id": snapshot.id, + "default_branch": snapshot.default_branch, + "integration_id": integration.id, + } + + org_integrations = _active_org_integrations(integration, org_integrations, delivery_id) + for organization in Organization.objects.filter( + id__in=[oi.organization_id for oi in org_integrations] + ): + created, reactivated, _ = provider.create_repositories( + configs=[config], organization=serialize_rpc_organization(organization) + ) + if created: + repository_service.auto_link_repos_by_name( + organization_id=organization.id, repo_ids=[repo.id for repo in created] + ) + for repo in created: + log_repo_change( + event_name="REPO_ADDED", + organization_id=organization.id, + repo=repo, + source="Cursor Origin webhook", + provider=integration.provider, + ) + for repo in reactivated: + log_repo_change( + event_name="REPO_ENABLED", + organization_id=organization.id, + repo=repo, + source="Cursor Origin webhook", + provider=integration.provider, + ) + + +class RepositoryDeletedHandler(WebhookEventHandler): + EVENT_TYPE = IntegrationWebhookEventType.INBOUND_SYNC + + def __call__( + self, + payload: Mapping[str, Any], + delivery_id: str, + integration: RpcIntegration, + org_integrations: Sequence[RpcOrganizationIntegration], + ) -> None: + external_id = RepositoryDeletedEvent.from_payload(payload).repository.id + for org_integration in _active_org_integrations(integration, org_integrations, delivery_id): + organization_id = org_integration.organization_id + if repository_service.find_recently_active_repo_external_ids( + organization_id=organization_id, + integration_id=integration.id, + provider=PROVIDER, + external_ids=[external_id], + cutoff_days=DISABLE_ACTIVITY_CUTOFF_DAYS, + ): + logger.info( + "cursor_origin.repository.disable_skipped_due_to_activity", + extra={ + "delivery_id": delivery_id, + "organization_id": organization_id, + "external_id": external_id, + "cutoff_days": DISABLE_ACTIVITY_CUTOFF_DAYS, + }, + ) + metrics.incr( + "cursor_origin.repository.disable_skipped_due_to_activity", sample_rate=1.0 + ) + continue + + active = [ + repo + for repo in repository_service.get_repositories( + organization_id=organization_id, + integration_id=integration.id, + providers=[PROVIDER], + external_id=external_id, + ) + if repo.status == ObjectStatus.ACTIVE + ] + repository_service.disable_repositories_by_external_ids( + organization_id=organization_id, + integration_id=integration.id, + provider=PROVIDER, + external_ids=[external_id], + ) + for repo in active: + log_repo_change( + event_name="REPO_DISABLED", + organization_id=organization_id, + repo=repo, + source="Cursor Origin webhook", + provider=integration.provider, + ) diff --git a/src/sentry/integrations/cursor_origin/webhook.py b/src/sentry/integrations/cursor_origin/webhook.py index db7525eb70ad..7ff96a3680f0 100644 --- a/src/sentry/integrations/cursor_origin/webhook.py +++ b/src/sentry/integrations/cursor_origin/webhook.py @@ -35,6 +35,8 @@ from sentry.integrations.cursor_origin.pull_request import PullRequestLifecycleHandler from sentry.integrations.cursor_origin.push import RepositoryPushedHandler from sentry.integrations.cursor_origin.repository_events import ( + RepositoryCreatedHandler, + RepositoryDeletedHandler, RepositoryMetadataUpdatedHandler, refresh_repository_name, ) @@ -155,6 +157,8 @@ def verify_delivery(request: HttpRequest, body: bytes) -> Verification: "pull_request.metadata.updated": PullRequestLifecycleHandler, "pull_request.published": PullRequestLifecycleHandler, "pull_request.reopened": PullRequestLifecycleHandler, + "repository.created": RepositoryCreatedHandler, + "repository.deleted": RepositoryDeletedHandler, "repository.metadata.updated": RepositoryMetadataUpdatedHandler, "repository.pushed": RepositoryPushedHandler, } diff --git a/src/sentry/integrations/cursor_origin/webhook_types.py b/src/sentry/integrations/cursor_origin/webhook_types.py index 5df8623e72b8..61808e630e8b 100644 --- a/src/sentry/integrations/cursor_origin/webhook_types.py +++ b/src/sentry/integrations/cursor_origin/webhook_types.py @@ -208,3 +208,14 @@ def from_payload(cls, payload: Mapping[str, Any]) -> InstallationEvent: return cls.parse_obj(payload) except ValidationError as e: raise OriginPayloadError(str(e)) from e + + +class RepositoryDeletedEvent(OriginModel): + repository: Repository + + @classmethod + def from_payload(cls, payload: Mapping[str, Any]) -> RepositoryDeletedEvent: + try: + return cls.parse_obj(payload) + except ValidationError as e: + raise OriginPayloadError(str(e)) from e diff --git a/src/sentry/plugins/providers/integration_repository.py b/src/sentry/plugins/providers/integration_repository.py index bf93cfe8107d..9a2d4f98a201 100644 --- a/src/sentry/plugins/providers/integration_repository.py +++ b/src/sentry/plugins/providers/integration_repository.py @@ -293,7 +293,8 @@ def create_repositories( Returns (created, reactivated, missing) — newly created repos, repos that were reactivated or updated from a hidden/unlinked state, and repo configs that could not be created because a repository with that configuration - already exists. + already exists. A repo that was already active has its config refreshed but + is not reported as reactivated. """ external_id_to_repo_config: dict[str, RepositoryConfig] = {} for config in configs: @@ -301,6 +302,7 @@ def create_repositories( external_id_to_repo_config[result["external_id"]] = result repos_to_update: list[RpcRepository] = [] + refreshed_repos: list[RpcRepository] = [] created_repos: list[RpcRepository] = [] transferred_repos: list[RpcRepository] = [] @@ -368,7 +370,10 @@ def create_repositories( missing_repos.append(repo_config) # We anticipate to only update one repository, but we update any duplicates as well. for repo in repositories: - repos_to_update.append(self._apply_repo_config(repo, repo_config)) + if repo.status == ObjectStatus.ACTIVE: + refreshed_repos.append(self._apply_repo_config(repo, repo_config)) + else: + repos_to_update.append(self._apply_repo_config(repo, repo_config)) continue # if we don't find the repo on this integration, the unique constraint was hit by a @@ -394,12 +399,12 @@ def create_repositories( ) missing_repos.append(repo_config) - if repos_to_update: + if repos_to_update or refreshed_repos: repository_service.update_repositories( organization_id=organization.id, - updates=repos_to_update, + updates=repos_to_update + refreshed_repos, ) - for repo in repos_to_update: + for repo in repos_to_update + refreshed_repos: self.on_create_repository(repo, organization) return created_repos, repos_to_update + transferred_repos, missing_repos diff --git a/tests/sentry/integrations/cursor_origin/test_repository_events.py b/tests/sentry/integrations/cursor_origin/test_repository_events.py index 91b19a0bf2d3..200c0ca8a192 100644 --- a/tests/sentry/integrations/cursor_origin/test_repository_events.py +++ b/tests/sentry/integrations/cursor_origin/test_repository_events.py @@ -7,6 +7,8 @@ from sentry.constants import ObjectStatus from sentry.integrations.cursor_origin.repository_events import ( + RepositoryCreatedHandler, + RepositoryDeletedHandler, RepositoryMetadataUpdatedHandler, refresh_repository_name, ) @@ -15,6 +17,7 @@ RepositoryMetadataEvent, ) from sentry.integrations.services.integration import integration_service +from sentry.models.commit import Commit from sentry.models.repository import Repository from sentry.testutils.cases import TestCase from sentry.testutils.silo import cell_silo_test @@ -199,3 +202,184 @@ def test_another_repository_is_left_alone(self) -> None: self._refresh(payload) assert self._reloaded().name == REPO + + +@cell_silo_test +class RepositoryDeletedHandlerTest(TestCase): + def setUp(self) -> None: + super().setUp() + self.integration = self.create_integration( + organization=self.organization, + provider="cursor_origin", + name="acme", + external_id=INSTALLATION_ID, + status=ObjectStatus.ACTIVE, + ) + self.repo = Repository.objects.create( + organization_id=self.organization.id, + name=REPO, + provider="integrations:cursor_origin", + integration_id=self.integration.id, + external_id=REPO_EXTERNAL_ID, + config={"name": REPO}, + ) + context = integration_service.organization_contexts( + provider="cursor_origin", external_id=INSTALLATION_ID + ) + assert context.integration is not None + self.rpc_integration = context.integration + self.org_integrations = context.organization_integrations + + def _handle(self, payload: dict[str, Any]) -> None: + RepositoryDeletedHandler()( + payload, DELIVERY_ID, self.rpc_integration, self.org_integrations + ) + + def test_a_deleted_repository_is_disabled(self) -> None: + self._handle({"repository": {"id": REPO_EXTERNAL_ID, "name": "rocket"}}) + + assert Repository.objects.get(id=self.repo.id).status == ObjectStatus.DISABLED + + def test_a_recently_active_repository_is_left_active(self) -> None: + """GitHub skips disabling a repository with activity in the last 30 days.""" + Commit.objects.create( + organization_id=self.organization.id, repository_id=self.repo.id, key="abc" + ) + + self._handle({"repository": {"id": REPO_EXTERNAL_ID, "name": "rocket"}}) + + assert Repository.objects.get(id=self.repo.id).status == ObjectStatus.ACTIVE + + def test_another_repository_is_left_active(self) -> None: + self._handle({"repository": {"id": "r_01other", "name": "other"}}) + + assert Repository.objects.get(id=self.repo.id).status == ObjectStatus.ACTIVE + + def test_a_payload_with_no_repository_id_is_refused(self) -> None: + with pytest.raises(OriginPayloadError, match="repository -> id"): + self._handle({"repository": {"name": "rocket"}}) + + def test_an_inactive_installation_disables_nothing(self) -> None: + """After a suspension, a late delivery must not act on the repositories.""" + self.rpc_integration = self.rpc_integration.copy(update={"status": ObjectStatus.DISABLED}) + + self._handle({"repository": {"id": REPO_EXTERNAL_ID, "name": "rocket"}}) + + assert Repository.objects.get(id=self.repo.id).status == ObjectStatus.ACTIVE + + def test_an_organization_being_uninstalled_is_skipped(self) -> None: + self.org_integrations = [ + oi.copy(update={"status": ObjectStatus.PENDING_DELETION}) + for oi in self.org_integrations + ] + + self._handle({"repository": {"id": REPO_EXTERNAL_ID, "name": "rocket"}}) + + assert Repository.objects.get(id=self.repo.id).status == ObjectStatus.ACTIVE + + +@cell_silo_test +class RepositoryCreatedHandlerTest(TestCase): + def setUp(self) -> None: + super().setUp() + self.integration = self.create_integration( + organization=self.organization, + provider="cursor_origin", + name="acme", + external_id=INSTALLATION_ID, + status=ObjectStatus.ACTIVE, + ) + context = integration_service.organization_contexts( + provider="cursor_origin", external_id=INSTALLATION_ID + ) + assert context.integration is not None + self.rpc_integration = context.integration + self.org_integrations = context.organization_integrations + + def _handle(self) -> None: + RepositoryCreatedHandler()( + _payload(), DELIVERY_ID, self.rpc_integration, self.org_integrations + ) + + def _repositories(self) -> list[Repository]: + return list( + Repository.objects.filter( + organization_id=self.organization.id, external_id=REPO_EXTERNAL_ID + ) + ) + + def test_a_new_repository_is_added(self) -> None: + self._handle() + + [repo] = self._repositories() + assert repo.name == REPO + assert repo.url == f"{WEB}/{REPO}" + assert repo.integration_id == self.integration.id + assert repo.config == {"name": REPO, "default_branch": "main"} + + def test_a_repository_left_by_an_earlier_installation_is_reactivated(self) -> None: + """GitHub audits a reactivated repository as `REPO_ENABLED` rather than added.""" + unlinked = Repository.objects.create( + organization_id=self.organization.id, + name=REPO, + provider="integrations:cursor_origin", + external_id=REPO_EXTERNAL_ID, + integration_id=None, + status=ObjectStatus.DISABLED, + config={"name": REPO}, + ) + + with mock.patch( + "sentry.integrations.cursor_origin.repository_events.log_repo_change" + ) as mock_log: + self._handle() + + repo = Repository.objects.get(id=unlinked.id) + assert repo.status == ObjectStatus.ACTIVE + assert repo.integration_id == self.integration.id + assert [call.kwargs["event_name"] for call in mock_log.call_args_list] == ["REPO_ENABLED"] + + def test_a_redelivery_adds_and_audits_nothing(self) -> None: + with mock.patch( + "sentry.integrations.cursor_origin.repository_events.log_repo_change" + ) as mock_log: + self._handle() + self._handle() + + assert len(self._repositories()) == 1 + assert [call.kwargs["event_name"] for call in mock_log.call_args_list] == ["REPO_ADDED"] + + def test_a_repository_the_sync_already_added_is_not_audited(self) -> None: + Repository.objects.create( + organization_id=self.organization.id, + name=REPO, + provider="integrations:cursor_origin", + external_id=REPO_EXTERNAL_ID, + integration_id=self.integration.id, + config={"name": REPO, "default_branch": "main"}, + ) + + with mock.patch( + "sentry.integrations.cursor_origin.repository_events.log_repo_change" + ) as mock_log: + self._handle() + + assert len(self._repositories()) == 1 + assert mock_log.call_count == 0 + + def test_an_inactive_installation_adds_nothing(self) -> None: + self.rpc_integration = self.rpc_integration.copy(update={"status": ObjectStatus.DISABLED}) + + self._handle() + + assert self._repositories() == [] + + def test_an_organization_being_uninstalled_is_skipped(self) -> None: + self.org_integrations = [ + oi.copy(update={"status": ObjectStatus.PENDING_DELETION}) + for oi in self.org_integrations + ] + + self._handle() + + assert self._repositories() == [] diff --git a/tests/sentry/middleware/integrations/parsers/test_cursor_origin.py b/tests/sentry/middleware/integrations/parsers/test_cursor_origin.py index 1a243f70593c..d473aec82230 100644 --- a/tests/sentry/middleware/integrations/parsers/test_cursor_origin.py +++ b/tests/sentry/middleware/integrations/parsers/test_cursor_origin.py @@ -129,7 +129,7 @@ def test_a_repository_event_is_forwarded_to_the_cells(self) -> None: def test_an_event_no_handler_reads_is_dropped(self) -> None: self._integration() - response = self._parser(**_headers("repository.created")).get_response() + response = self._parser(**_headers("repository.check_run.created")).get_response() assert response.status_code == status.HTTP_202_ACCEPTED assert_no_webhook_payloads() diff --git a/tests/sentry/plugins/test_integration_repository.py b/tests/sentry/plugins/test_integration_repository.py index c293fcb0a3d4..29a0658d5b00 100644 --- a/tests/sentry/plugins/test_integration_repository.py +++ b/tests/sentry/plugins/test_integration_repository.py @@ -236,3 +236,53 @@ def test_create_repository__only_activates_hidden_repo(self, get_jwt: MagicMock) assert exc_info.value.existing_repo is None repo.refresh_from_db() assert repo.status == ObjectStatus.PENDING_DELETION + + +class CreateRepositoriesTest(TestCase): + def setUp(self) -> None: + super().setUp() + self.integration = self.create_integration( + organization=self.organization, provider="github", external_id="654321" + ) + self.config = { + "identifier": "getsentry/sentry", + "external_id": "654321", + "integration_id": self.integration.id, + } + + @cached_property + def provider(self) -> GitHubRepositoryProvider: + return GitHubRepositoryProvider("integrations:github") + + def _existing(self, status: int) -> Repository: + return Repository.objects.create( + name="getsentry/old-name", + provider="integrations:github", + organization_id=self.organization.id, + integration_id=self.integration.id, + external_id="654321", + status=status, + ) + + def test_an_active_repository_is_refreshed_but_not_reactivated(self) -> None: + repo = self._existing(ObjectStatus.ACTIVE) + + created, reactivated, _ = self.provider.create_repositories( + configs=[self.config], organization=self.organization + ) + + assert (created, reactivated) == ([], []) + repo.refresh_from_db() + assert repo.name == "getsentry/sentry" + + def test_a_disabled_repository_is_reactivated(self) -> None: + repo = self._existing(ObjectStatus.DISABLED) + + created, reactivated, _ = self.provider.create_repositories( + configs=[self.config], organization=self.organization + ) + + assert created == [] + assert [r.id for r in reactivated] == [repo.id] + repo.refresh_from_db() + assert repo.status == ObjectStatus.ACTIVE