diff --git a/api/changelog.d/scan-summary-deadlocks.fixed.md b/api/changelog.d/scan-summary-deadlocks.fixed.md new file mode 100644 index 0000000000..0fcb4479bd --- /dev/null +++ b/api/changelog.d/scan-summary-deadlocks.fixed.md @@ -0,0 +1 @@ +`scan-summary` aggregation now upserts summaries in deterministic conflict-key order, preventing PostgreSQL deadlocks during concurrent reaggregation diff --git a/api/src/backend/tasks/jobs/scan.py b/api/src/backend/tasks/jobs/scan.py index 8290c45a7a..39db32abe4 100644 --- a/api/src/backend/tasks/jobs/scan.py +++ b/api/src/backend/tasks/jobs/scan.py @@ -1464,7 +1464,7 @@ def aggregate_findings(tenant_id: str, scan_id: str): ) with rls_transaction(tenant_id): - scan_aggregations = { + scan_aggregations = [ ScanSummary( tenant_id=tenant_id, scan_id=scan_id, @@ -1489,9 +1489,18 @@ def aggregate_findings(tenant_id: str, scan_id: str): for agg in aggregation if agg["resources__service"] is not None and agg["resources__region"] is not None - } - # Upsert so re-runs (post-mute reaggregation) don't trip - # `unique_scan_summary`; race-safe under concurrent writers. + ] + # Needed sort so concurrent upserts acquire locks consistently + scan_aggregations.sort( + key=lambda summary: ( + summary.tenant_id, + summary.scan_id, + summary.check_id, + summary.service, + summary.severity, + summary.region, + ) + ) ScanSummary.objects.bulk_create( scan_aggregations, batch_size=3000, diff --git a/api/src/backend/tasks/tests/test_scan.py b/api/src/backend/tasks/tests/test_scan.py index f49ca6655b..2a66985276 100644 --- a/api/src/backend/tasks/tests/test_scan.py +++ b/api/src/backend/tasks/tests/test_scan.py @@ -3652,6 +3652,95 @@ class TestAggregateFindings: regions = {s.region for s in summaries} assert regions == {"us-east-1", "us-west-2"} + @patch("tasks.jobs.scan.Finding.objects.filter") + @patch("tasks.jobs.scan.ScanSummary.objects.bulk_create") + @patch("tasks.jobs.scan.rls_transaction") + def test_aggregate_findings_orders_upserts_by_conflict_key( + self, mock_rls_transaction, mock_bulk_create, mock_findings_filter + ): + """Scan summaries must use a stable lock order for concurrent upserts.""" + tenant_id = str(uuid.uuid4()) + scan_id = str(uuid.uuid4()) + counts = { + "fail": 1, + "_pass": 0, + "muted_count": 0, + "total": 1, + "new": 1, + "changed": 0, + "unchanged": 0, + "fail_new": 1, + "fail_changed": 0, + "pass_new": 0, + "pass_changed": 0, + "muted_new": 0, + "muted_changed": 0, + } + + mock_queryset = MagicMock() + mock_queryset.values.return_value = mock_queryset + mock_queryset.annotate.return_value = [ + { + "check_id": "check-b", + "resources__service": "s3", + "severity": "high", + "resources__region": "us-east-1", + **counts, + }, + { + "check_id": "check-a", + "resources__service": "sqs", + "severity": "high", + "resources__region": "us-east-1", + **counts, + }, + { + "check_id": "check-a", + "resources__service": "s3", + "severity": "medium", + "resources__region": "us-east-1", + **counts, + }, + { + "check_id": "check-a", + "resources__service": "s3", + "severity": "high", + "resources__region": "us-west-2", + **counts, + }, + { + "check_id": "check-a", + "resources__service": "s3", + "severity": "high", + "resources__region": "us-east-1", + **counts, + }, + ] + + ctx = MagicMock() + ctx.__enter__.return_value = None + ctx.__exit__.return_value = False + mock_rls_transaction.return_value = ctx + mock_findings_filter.return_value = mock_queryset + + aggregate_findings(tenant_id, scan_id) + + summaries = mock_bulk_create.call_args.args[0] + assert isinstance(summaries, list) + conflict_keys = [ + ( + str(summary.tenant_id), + str(summary.scan_id), + summary.check_id, + summary.service, + summary.severity, + summary.region, + ) + for summary in summaries + ] + assert len(conflict_keys) == 5 + assert conflict_keys == sorted(conflict_keys) + @patch("tasks.jobs.scan.Finding.objects.filter") @patch("tasks.jobs.scan.ScanSummary.objects.bulk_create") @patch("tasks.jobs.scan.rls_transaction")