From 4774a8025cea3fd0962a23741e6d0241d4062663 Mon Sep 17 00:00:00 2001 From: Albrecht Degering Date: Wed, 22 Jul 2026 20:30:00 +0200 Subject: [PATCH] fix(campaign): roll back rejected immediate queues --- src/govoplan_campaign/backend/router.py | 4 ++ src/govoplan_campaign/backend/sending/jobs.py | 12 ++++- tests/test_synchronous_delivery_policy.py | 49 ++++++++++++++++++- 3 files changed, 63 insertions(+), 2 deletions(-) diff --git a/src/govoplan_campaign/backend/router.py b/src/govoplan_campaign/backend/router.py index 5a6c7cd..73b31ee 100644 --- a/src/govoplan_campaign/backend/router.py +++ b/src/govoplan_campaign/backend/router.py @@ -3608,6 +3608,10 @@ def send_campaign_now_endpoint( ) return SendCampaignNowResponse(result=response_result) except SynchronousSendRejected as exc: + # A synchronous request stages queue state before the all-message + # preflight can run. Rejecting that preflight must not leave work + # eligible for a background worker when no provider effect occurred. + session.rollback() audit_from_principal( session, principal, diff --git a/src/govoplan_campaign/backend/sending/jobs.py b/src/govoplan_campaign/backend/sending/jobs.py index 7741b78..20e4f8b 100644 --- a/src/govoplan_campaign/backend/sending/jobs.py +++ b/src/govoplan_campaign/backend/sending/jobs.py @@ -622,6 +622,7 @@ def _persist_campaign_queue( version: CampaignVersion, queued: list[CampaignJob], delivery_mode: str, + commit: bool = True, ) -> None: if queued: previous_status = campaign.status @@ -640,7 +641,8 @@ def _persist_campaign_queue( version_id=version.id, ) session.add(campaign) - session.commit() + if commit: + session.commit() def _enqueue_campaign_jobs(queued: list[CampaignJob], *, enabled: bool) -> int: @@ -661,6 +663,7 @@ def queue_campaign_jobs( include_warnings: bool = True, dry_run: bool = False, delivery_mode: str | None = None, + commit_queue: bool = True, ) -> QueueCampaignResult: """Move queueable DB jobs to QUEUED and optionally enqueue Celery tasks.""" @@ -694,6 +697,7 @@ def queue_campaign_jobs( version=version, queued=queued, delivery_mode=selected_delivery_mode, + commit=commit_queue, ) enqueued_count = _enqueue_campaign_jobs( queued, @@ -766,6 +770,7 @@ def send_campaign_now( enqueue_celery=False, dry_run=dry_run, delivery_mode=DELIVERY_MODE_SYNCHRONOUS, + commit_queue=False, ) if dry_run: return SendCampaignNowResult( @@ -802,6 +807,11 @@ def send_campaign_now( jobs=jobs, policy=synchronous_policy, ) + # Queue state and its inbox notification become durable only after every + # message and the selected transport revision have passed preflight. This + # preserves late-ack recovery without leaving rejected work eligible for a + # background worker. + session.commit() results: list[dict[str, Any]] = [] sent_count = 0 diff --git a/tests/test_synchronous_delivery_policy.py b/tests/test_synchronous_delivery_policy.py index 7dbe096..301c4f1 100644 --- a/tests/test_synchronous_delivery_policy.py +++ b/tests/test_synchronous_delivery_policy.py @@ -195,7 +195,10 @@ def test_post_queue_growth_is_rejected_before_batch_or_provider_preflight() -> N patch("govoplan_campaign.backend.sending.jobs._ensure_campaign_execution_snapshot"), patch("govoplan_campaign.backend.sending.jobs.effective_synchronous_send_policy", return_value=policy), patch("govoplan_campaign.backend.sending.jobs._campaign_jobs_for_queue", return_value=initial_jobs), - patch("govoplan_campaign.backend.sending.jobs.queue_campaign_jobs", return_value=queued), + patch( + "govoplan_campaign.backend.sending.jobs.queue_campaign_jobs", + return_value=queued, + ) as queue, patch("govoplan_campaign.backend.sending.jobs._campaign_jobs_for_version", return_value=post_queue_jobs), patch("govoplan_campaign.backend.sending.jobs._preflight_synchronous_send_batch") as batch_preflight, ): @@ -207,6 +210,7 @@ def test_post_queue_growth_is_rejected_before_batch_or_provider_preflight() -> N ) assert rejected.value.eligible_count == 3 + assert queue.call_args.kwargs["commit_queue"] is False batch_preflight.assert_not_called() @@ -251,6 +255,7 @@ def test_asynchronous_mode_matches_actual_worker_availability( assert result.worker_queue_available is workers_available assert result.enqueued_count == expected_enqueued assert persist.call_args.kwargs["delivery_mode"] == expected_mode + assert persist.call_args.kwargs["commit"] is True assert enqueue.call_args.kwargs["enabled"] is workers_available @@ -356,3 +361,45 @@ def test_batch_preflight_checks_every_message_before_provider_effects() -> None: assert state_preflight.call_count == 2 assert input_preflight.call_count == 2 provider.send_campaign_email_bytes.assert_not_called() + + +def test_rejected_synchronous_preflight_rolls_back_staged_queue_before_audit() -> None: + session = Mock() + campaign = SimpleNamespace(id="campaign-1", current_version_id="version-1") + version = SimpleNamespace( + id="version-1", + raw_json={}, + locked_at=object(), + validation_summary={"ok": True}, + build_summary={"built_count": 1}, + ) + rejection = SynchronousSendRejected( + "Preflight rejected the staged send.", + reason="batch_preflight_failed", + eligible_count=1, + ) + + with ( + patch.object(router, "_get_campaign_for_principal"), + patch.object(router, "_require_permission"), + patch.object(router, "_get_campaign_for_tenant", return_value=campaign), + patch.object(router, "_get_version_for_tenant", return_value=version), + patch.object(router, "_require_mail_profile_use_if_needed"), + patch.object(router, "is_user_locked_version", return_value=False), + patch.object(router, "send_campaign_now", side_effect=rejection), + patch.object(router, "audit_from_principal") as audit, + pytest.raises(HTTPException) as rejected, + ): + router.send_campaign_now_endpoint( + "campaign-1", + session=session, + principal=_Principal( + "campaigns:campaign:send", "campaigns:recipient:read" + ), # type: ignore[arg-type] + ) + + assert rejected.value.status_code == 422 + session.rollback.assert_called_once_with() + audit.assert_called_once() + assert audit.call_args.kwargs["action"] == "campaign.send_now_rejected" + assert audit.call_args.kwargs["commit"] is True