test(campaign): cover bounded delivery modes

This commit is contained in:
2026-07-22 09:00:35 +02:00
parent aae4ef5952
commit b0282ebff2

View File

@@ -4,7 +4,9 @@ from types import SimpleNamespace
from unittest.mock import Mock, patch from unittest.mock import Mock, patch
import pytest import pytest
from fastapi import HTTPException
from govoplan_campaign.backend import router
from govoplan_campaign.backend.delivery_policy import ( from govoplan_campaign.backend.delivery_policy import (
CampaignDeliveryPolicyError, CampaignDeliveryPolicyError,
DEFAULT_SYNCHRONOUS_SEND_MAX_RECIPIENT_JOBS, DEFAULT_SYNCHRONOUS_SEND_MAX_RECIPIENT_JOBS,
@@ -17,10 +19,14 @@ from govoplan_campaign.backend.db.models import (
JobValidationStatus, JobValidationStatus,
) )
from govoplan_campaign.backend.sending.jobs import ( from govoplan_campaign.backend.sending.jobs import (
QueueCampaignResult,
SynchronousSendRejected, SynchronousSendRejected,
_ensure_synchronous_send_count_allowed, _ensure_synchronous_send_count_allowed,
_preflight_synchronous_send_batch, _preflight_synchronous_send_batch,
queue_campaign_jobs,
send_campaign_now,
synchronous_send_candidate_jobs, synchronous_send_candidate_jobs,
synchronous_send_options,
) )
@@ -35,6 +41,9 @@ class _PolicySession:
def _version() -> SimpleNamespace: def _version() -> SimpleNamespace:
return SimpleNamespace( return SimpleNamespace(
id="version-1", id="version-1",
locked_at=object(),
published_at=None,
validation_summary={"ok": True},
build_summary={"build_token": "build-1"}, build_summary={"build_token": "build-1"},
editor_state={ editor_state={
"review_send": { "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: def _job(job_id: str, **overrides: object) -> SimpleNamespace:
values: dict[str, object] = { values: dict[str, object] = {
"id": job_id, "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 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: def test_batch_preflight_checks_every_message_before_provider_effects() -> None:
jobs = [_job("one"), _job("two")] jobs = [_job("one"), _job("two")]
contexts = { contexts = {