From b0282ebff2d2a1c0f755b3f27849aa33f59ea7bf Mon Sep 17 00:00:00 2001 From: Albrecht Degering Date: Wed, 22 Jul 2026 09:00:35 +0200 Subject: [PATCH] test(campaign): cover bounded delivery modes --- tests/test_synchronous_delivery_policy.py | 187 ++++++++++++++++++++++ 1 file changed, 187 insertions(+) diff --git a/tests/test_synchronous_delivery_policy.py b/tests/test_synchronous_delivery_policy.py index d31d52e..7dbe096 100644 --- a/tests/test_synchronous_delivery_policy.py +++ b/tests/test_synchronous_delivery_policy.py @@ -4,7 +4,9 @@ from types import SimpleNamespace from unittest.mock import Mock, patch import pytest +from fastapi import HTTPException +from govoplan_campaign.backend import router from govoplan_campaign.backend.delivery_policy import ( CampaignDeliveryPolicyError, DEFAULT_SYNCHRONOUS_SEND_MAX_RECIPIENT_JOBS, @@ -17,10 +19,14 @@ from govoplan_campaign.backend.db.models import ( JobValidationStatus, ) from govoplan_campaign.backend.sending.jobs import ( + QueueCampaignResult, SynchronousSendRejected, _ensure_synchronous_send_count_allowed, _preflight_synchronous_send_batch, + queue_campaign_jobs, + send_campaign_now, synchronous_send_candidate_jobs, + synchronous_send_options, ) @@ -35,6 +41,9 @@ class _PolicySession: def _version() -> SimpleNamespace: return SimpleNamespace( id="version-1", + locked_at=object(), + published_at=None, + validation_summary={"ok": True}, build_summary={"build_token": "build-1"}, editor_state={ "review_send": { @@ -46,6 +55,16 @@ def _version() -> SimpleNamespace: ) +class _Principal: + def __init__(self, *scopes: str) -> None: + self.scopes = set(scopes) + self.tenant_id = "tenant-1" + self.user = SimpleNamespace(id="user-1") + + def has(self, scope: str) -> bool: + return scope in self.scopes + + def _job(job_id: str, **overrides: object) -> SimpleNamespace: values: dict[str, object] = { "id": job_id, @@ -133,6 +152,174 @@ def test_limit_and_zero_count_reject_before_delivery() -> None: assert oversized.value.audit_details()["synchronous_send_policy"]["max_recipient_jobs"] == 2 +def test_exact_synchronous_policy_boundary_is_allowed() -> None: + policy = effective_synchronous_send_policy( + _PolicySession(), # type: ignore[arg-type] + tenant_id="tenant-1", + environ={"GOVOPLAN_CAMPAIGN_SYNCHRONOUS_SEND_MAX_RECIPIENTS": "2"}, + ) + + _ensure_synchronous_send_count_allowed(2, policy=policy) + + +def test_post_queue_growth_is_rejected_before_batch_or_provider_preflight() -> None: + policy = effective_synchronous_send_policy( + _PolicySession(), # type: ignore[arg-type] + tenant_id="tenant-1", + environ={"GOVOPLAN_CAMPAIGN_SYNCHRONOUS_SEND_MAX_RECIPIENTS": "2"}, + ) + campaign = SimpleNamespace(id="campaign-1", current_version_id="version-1") + initial_jobs = [_job("one"), _job("two")] + post_queue_jobs = [ + _job( + job_id, + queue_status=JobQueueStatus.QUEUED.value, + send_status=JobSendStatus.QUEUED.value, + ) + for job_id in ("one", "two", "concurrent") + ] + queued = QueueCampaignResult( + campaign_id="campaign-1", + version_id="version-1", + queued_count=2, + skipped_count=0, + blocked_count=0, + enqueued_count=0, + delivery_mode="synchronous", + ) + + with ( + patch("govoplan_campaign.backend.sending.jobs._get_campaign_for_tenant", return_value=campaign), + patch("govoplan_campaign.backend.sending.jobs._get_current_version", return_value=_version()), + patch("govoplan_campaign.backend.sending.jobs._ensure_version_validated_and_locked"), + 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._campaign_jobs_for_version", return_value=post_queue_jobs), + patch("govoplan_campaign.backend.sending.jobs._preflight_synchronous_send_batch") as batch_preflight, + ): + with pytest.raises(SynchronousSendRejected, match="above the effective") as rejected: + send_campaign_now( + object(), # type: ignore[arg-type] + tenant_id="tenant-1", + campaign_id="campaign-1", + ) + + assert rejected.value.eligible_count == 3 + batch_preflight.assert_not_called() + + +@pytest.mark.parametrize( + ("workers_available", "expected_mode", "expected_enqueued"), + ((False, "database_queue", 0), (True, "worker_queue", 1)), +) +def test_asynchronous_mode_matches_actual_worker_availability( + workers_available: bool, + expected_mode: str, + expected_enqueued: int, +) -> None: + campaign = SimpleNamespace(id="campaign-1") + version = _version() + job = _job("one") + + with ( + patch("govoplan_campaign.backend.sending.jobs._celery_enabled", return_value=workers_available), + patch("govoplan_campaign.backend.sending.jobs._get_campaign_for_tenant", return_value=campaign), + patch("govoplan_campaign.backend.sending.jobs._get_current_version", return_value=version), + patch("govoplan_campaign.backend.sending.jobs._ensure_version_validated_and_locked"), + patch("govoplan_campaign.backend.sending.jobs._ensure_campaign_execution_snapshot"), + patch("govoplan_campaign.backend.sending.jobs._campaign_jobs_for_queue", return_value=[job]), + patch( + "govoplan_campaign.backend.sending.jobs._select_campaign_jobs_for_queue", + return_value=([job], 0, 0), + ), + patch("govoplan_campaign.backend.sending.jobs._persist_campaign_queue") as persist, + patch( + "govoplan_campaign.backend.sending.jobs._enqueue_campaign_jobs", + return_value=expected_enqueued, + ) as enqueue, + ): + result = queue_campaign_jobs( + object(), # type: ignore[arg-type] + tenant_id="tenant-1", + campaign_id="campaign-1", + enqueue_celery=True, + ) + + assert result.delivery_mode == expected_mode + assert result.worker_queue_available is workers_available + assert result.enqueued_count == expected_enqueued + assert persist.call_args.kwargs["delivery_mode"] == expected_mode + assert enqueue.call_args.kwargs["enabled"] is workers_available + + +def test_invalid_policy_disables_synchronous_mode_without_hiding_queue_availability() -> None: + campaign = SimpleNamespace(id="campaign-1") + version = _version() + with ( + patch("govoplan_campaign.backend.sending.jobs._get_campaign_for_tenant", return_value=campaign), + patch("govoplan_campaign.backend.sending.jobs._get_version_for_campaign", return_value=version), + patch("govoplan_campaign.backend.sending.jobs._campaign_jobs_for_version", return_value=[_job("one")]), + patch("govoplan_campaign.backend.sending.jobs._celery_enabled", return_value=True), + patch( + "govoplan_campaign.backend.sending.jobs.effective_synchronous_send_policy", + side_effect=CampaignDeliveryPolicyError("invalid deployment value"), + ), + ): + options = synchronous_send_options( + object(), # type: ignore[arg-type] + tenant_id="tenant-1", + campaign_id="campaign-1", + ) + + assert options["worker_queue_available"] is True + assert options["synchronous_send"]["allowed"] is False + assert options["synchronous_send"]["reason"] == "policy_configuration_invalid" + assert options["synchronous_send"]["policy"] == {} + + +@pytest.mark.parametrize( + ("path", "required_scope"), + ( + ("/campaigns/{campaign_id}/send-now", "campaigns:campaign:send"), + ("/campaigns/{campaign_id}/queue", "campaigns:campaign:queue"), + ), +) +def test_delivery_endpoints_require_their_mode_permission_and_recipient_authority( + path: str, + required_scope: str, +) -> None: + route = next(item for item in router.router.routes if item.path == path) + dependency = next(item for item in route.dependant.dependencies if item.name == "principal") + + with pytest.raises(HTTPException) as missing_mode_permission: + dependency.call(_Principal("campaigns:recipient:read")) + assert missing_mode_permission.value.status_code == 403 + allowed = _Principal(required_scope, "campaigns:recipient:read") + assert dependency.call(allowed) is allowed + + mode_only = _Principal(required_scope) + with ( + patch.object(router, "_get_campaign_for_principal"), + pytest.raises(HTTPException) as missing_recipient_authority, + ): + if required_scope == "campaigns:campaign:send": + router.send_campaign_now_endpoint( + "campaign-1", + session=Mock(), + principal=mode_only, # type: ignore[arg-type] + ) + else: + router.queue_campaign( + "campaign-1", + session=Mock(), + principal=mode_only, # type: ignore[arg-type] + ) + assert missing_recipient_authority.value.status_code == 403 + assert "campaigns:recipient:read" in missing_recipient_authority.value.detail + + def test_batch_preflight_checks_every_message_before_provider_effects() -> None: jobs = [_job("one"), _job("two")] contexts = {