fix(campaign): roll back rejected immediate queues
This commit is contained in:
@@ -3608,6 +3608,10 @@ def send_campaign_now_endpoint(
|
|||||||
)
|
)
|
||||||
return SendCampaignNowResponse(result=response_result)
|
return SendCampaignNowResponse(result=response_result)
|
||||||
except SynchronousSendRejected as exc:
|
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(
|
audit_from_principal(
|
||||||
session,
|
session,
|
||||||
principal,
|
principal,
|
||||||
|
|||||||
@@ -622,6 +622,7 @@ def _persist_campaign_queue(
|
|||||||
version: CampaignVersion,
|
version: CampaignVersion,
|
||||||
queued: list[CampaignJob],
|
queued: list[CampaignJob],
|
||||||
delivery_mode: str,
|
delivery_mode: str,
|
||||||
|
commit: bool = True,
|
||||||
) -> None:
|
) -> None:
|
||||||
if queued:
|
if queued:
|
||||||
previous_status = campaign.status
|
previous_status = campaign.status
|
||||||
@@ -640,7 +641,8 @@ def _persist_campaign_queue(
|
|||||||
version_id=version.id,
|
version_id=version.id,
|
||||||
)
|
)
|
||||||
session.add(campaign)
|
session.add(campaign)
|
||||||
session.commit()
|
if commit:
|
||||||
|
session.commit()
|
||||||
|
|
||||||
|
|
||||||
def _enqueue_campaign_jobs(queued: list[CampaignJob], *, enabled: bool) -> int:
|
def _enqueue_campaign_jobs(queued: list[CampaignJob], *, enabled: bool) -> int:
|
||||||
@@ -661,6 +663,7 @@ def queue_campaign_jobs(
|
|||||||
include_warnings: bool = True,
|
include_warnings: bool = True,
|
||||||
dry_run: bool = False,
|
dry_run: bool = False,
|
||||||
delivery_mode: str | None = None,
|
delivery_mode: str | None = None,
|
||||||
|
commit_queue: bool = True,
|
||||||
) -> QueueCampaignResult:
|
) -> QueueCampaignResult:
|
||||||
"""Move queueable DB jobs to QUEUED and optionally enqueue Celery tasks."""
|
"""Move queueable DB jobs to QUEUED and optionally enqueue Celery tasks."""
|
||||||
|
|
||||||
@@ -694,6 +697,7 @@ def queue_campaign_jobs(
|
|||||||
version=version,
|
version=version,
|
||||||
queued=queued,
|
queued=queued,
|
||||||
delivery_mode=selected_delivery_mode,
|
delivery_mode=selected_delivery_mode,
|
||||||
|
commit=commit_queue,
|
||||||
)
|
)
|
||||||
enqueued_count = _enqueue_campaign_jobs(
|
enqueued_count = _enqueue_campaign_jobs(
|
||||||
queued,
|
queued,
|
||||||
@@ -766,6 +770,7 @@ def send_campaign_now(
|
|||||||
enqueue_celery=False,
|
enqueue_celery=False,
|
||||||
dry_run=dry_run,
|
dry_run=dry_run,
|
||||||
delivery_mode=DELIVERY_MODE_SYNCHRONOUS,
|
delivery_mode=DELIVERY_MODE_SYNCHRONOUS,
|
||||||
|
commit_queue=False,
|
||||||
)
|
)
|
||||||
if dry_run:
|
if dry_run:
|
||||||
return SendCampaignNowResult(
|
return SendCampaignNowResult(
|
||||||
@@ -802,6 +807,11 @@ def send_campaign_now(
|
|||||||
jobs=jobs,
|
jobs=jobs,
|
||||||
policy=synchronous_policy,
|
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]] = []
|
results: list[dict[str, Any]] = []
|
||||||
sent_count = 0
|
sent_count = 0
|
||||||
|
|||||||
@@ -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._ensure_campaign_execution_snapshot"),
|
||||||
patch("govoplan_campaign.backend.sending.jobs.effective_synchronous_send_policy", return_value=policy),
|
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._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._campaign_jobs_for_version", return_value=post_queue_jobs),
|
||||||
patch("govoplan_campaign.backend.sending.jobs._preflight_synchronous_send_batch") as batch_preflight,
|
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 rejected.value.eligible_count == 3
|
||||||
|
assert queue.call_args.kwargs["commit_queue"] is False
|
||||||
batch_preflight.assert_not_called()
|
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.worker_queue_available is workers_available
|
||||||
assert result.enqueued_count == expected_enqueued
|
assert result.enqueued_count == expected_enqueued
|
||||||
assert persist.call_args.kwargs["delivery_mode"] == expected_mode
|
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
|
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 state_preflight.call_count == 2
|
||||||
assert input_preflight.call_count == 2
|
assert input_preflight.call_count == 2
|
||||||
provider.send_campaign_email_bytes.assert_not_called()
|
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
|
||||||
|
|||||||
Reference in New Issue
Block a user