Files
govoplan-campaign/tests/test_campaign_scheduling.py
T

531 lines
21 KiB
Python

from __future__ import annotations
from contextlib import nullcontext
from datetime import UTC, datetime
from types import SimpleNamespace
from unittest.mock import Mock, patch
from sqlalchemy import Column, String, Table, create_engine
from sqlalchemy.orm import Session, sessionmaker
from sqlalchemy.orm.attributes import flag_modified
from govoplan_campaign.backend.campaign.scheduling import (
campaign_schedule_source_snapshot,
canonical_configuration_hash,
dispatch_due_campaign_schedules,
next_schedule_fire,
)
from govoplan_campaign.backend.db.models import (
Campaign,
CampaignJob,
CampaignSchedule,
CampaignScheduleOccurrence,
CampaignShare,
CampaignVersion,
)
from govoplan_core.db.base import Base
from govoplan_core.core.change_sequence import ChangeSequenceEntry
def _access_table(name: str) -> Table:
existing = Base.metadata.tables.get(name)
if existing is not None:
return existing
return Table(
name,
Base.metadata,
Column("id", String(36), primary_key=True),
)
def _create_generated_campaign(session: Session, **kwargs):
raw = kwargs["raw_json"]
metadata = raw["campaign"]
campaign = Campaign(
tenant_id=kwargs["tenant_id"],
created_by_user_id=kwargs["user_id"],
owner_user_id=kwargs["user_id"],
external_id=metadata["id"],
name=metadata["name"],
status="draft",
)
session.add(campaign)
session.flush()
version = CampaignVersion(
campaign_id=campaign.id,
version_number=1,
raw_json=raw,
)
session.add(version)
session.flush()
campaign.current_version_id = version.id
return campaign, version
class TestCampaignScheduling:
def setup_method(self):
self.engine = create_engine("sqlite+pysqlite:///:memory:")
users = _access_table("access_users")
groups = _access_table("access_groups")
Base.metadata.create_all(
self.engine,
tables=[
users,
groups,
Campaign.__table__,
CampaignVersion.__table__,
CampaignJob.__table__,
CampaignShare.__table__,
CampaignSchedule.__table__,
CampaignScheduleOccurrence.__table__,
ChangeSequenceEntry.__table__,
],
)
self.SessionLocal = sessionmaker(
bind=self.engine,
class_=Session,
expire_on_commit=False,
)
configuration = {
"version": "1.0",
"campaign": {"id": "source", "name": "Monthly notice"},
}
snapshot = campaign_schedule_source_snapshot(
configuration=configuration,
campaign_settings={"retention": "sealed"},
mail_profile_policy={"profile_id": "profile-1"},
shares=[],
)
with self.SessionLocal() as session:
user_values = {"id": "user-1"}
if "tenant_id" in users.c:
user_values.update(
tenant_id="tenant-1",
account_id="account-1",
email="user-1@example.test",
)
session.execute(users.insert().values(**user_values))
campaign = Campaign(
id="campaign-1",
tenant_id="tenant-1",
created_by_user_id="user-1",
owner_user_id="user-1",
external_id="source",
name="Monthly notice",
status="sent",
current_version_id="version-1",
)
version = CampaignVersion(
id="version-1",
campaign_id=campaign.id,
version_number=1,
workflow_state="completed",
raw_json=configuration,
)
schedule = CampaignSchedule(
id="schedule-1",
tenant_id="tenant-1",
campaign_id=campaign.id,
source_version_id=version.id,
created_by_user_id="user-1",
name="Monthly notice",
recurrence_kind="daily",
interval_count=1,
timezone="Europe/Berlin",
starts_at=datetime(2026, 8, 7, 8, tzinfo=UTC),
next_fire_at=datetime(2026, 8, 7, 8, tzinfo=UTC),
max_occurrences=2,
copy_options={
"include_recipients": True,
"include_files": True,
"include_shares": False,
"include_policies": True,
"include_mail_profile": True,
},
source_snapshot=snapshot,
source_snapshot_hash=canonical_configuration_hash(snapshot),
)
session.add_all((campaign, version, schedule))
session.commit()
def teardown_method(self):
self.engine.dispose()
def test_due_occurrences_prepare_distinct_drafts_and_complete_bound(self):
with self.SessionLocal() as session, patch(
"govoplan_campaign.backend.campaign.scheduling.create_campaign_version_from_json",
side_effect=lambda *args, **kwargs: _create_generated_campaign(session, **kwargs),
), patch("govoplan_campaign.backend.campaign.scheduling.audit_event"):
first = dispatch_due_campaign_schedules(
session,
tenant_id="tenant-1",
now=datetime(2026, 8, 7, 8, tzinfo=UTC),
)
session.commit()
assert first["prepared"] == 1
schedule = session.get(CampaignSchedule, "schedule-1")
assert schedule is not None
assert schedule.active is True
assert schedule.occurrence_count == 1
assert schedule.next_fire_at is not None
second = dispatch_due_campaign_schedules(
session,
tenant_id="tenant-1",
now=datetime(2026, 8, 8, 8, tzinfo=UTC),
)
session.commit()
assert second["prepared"] == 1
assert schedule.active is False
assert schedule.next_fire_at is None
occurrences = session.query(CampaignScheduleOccurrence).all()
assert len(occurrences) == 2
assert len({item.generated_campaign_id for item in occurrences}) == 2
assert session.get(Campaign, "campaign-1").status == "sent"
def test_snapshot_integrity_failure_pauses_schedule_for_operator(self):
with self.SessionLocal() as session:
schedule = session.get(CampaignSchedule, "schedule-1")
assert schedule is not None
schedule.source_snapshot["configuration"]["campaign"]["name"] = "Tampered"
flag_modified(schedule, "source_snapshot")
session.commit()
with patch("govoplan_campaign.backend.campaign.scheduling.audit_event"):
result = dispatch_due_campaign_schedules(
session,
tenant_id="tenant-1",
now=datetime(2026, 8, 7, 8, tzinfo=UTC),
)
session.commit()
assert result["failed"] == 1
assert schedule.active is False
assert "integrity" in (schedule.last_error or "")
occurrence = session.query(CampaignScheduleOccurrence).one()
assert occurrence.status == "failed"
def test_monthly_recurrence_clamps_end_of_month(self):
result = next_schedule_fire(
datetime(2026, 1, 31, 9, tzinfo=UTC),
recurrence_kind="monthly",
interval_count=1,
timezone_name="UTC",
)
assert result == datetime(2026, 2, 28, 9, tzinfo=UTC)
def test_occurrence_uses_sealed_policy_state_and_advances_revision(self):
with self.SessionLocal() as session:
source = session.get(Campaign, "campaign-1")
assert source is not None
source.settings = {"retention": "changed-after-scheduling"}
source.mail_profile_policy = {"profile_id": "profile-2"}
session.commit()
with patch(
"govoplan_campaign.backend.campaign.scheduling.create_campaign_version_from_json",
side_effect=lambda *args, **kwargs: _create_generated_campaign(
session,
**kwargs,
),
), patch("govoplan_campaign.backend.campaign.scheduling.audit_event"):
result = dispatch_due_campaign_schedules(
session,
tenant_id="tenant-1",
now=datetime(2026, 8, 7, 8, tzinfo=UTC),
)
session.commit()
assert result["prepared"] == 1
schedule = session.get(CampaignSchedule, "schedule-1")
assert schedule is not None
generated = session.get(Campaign, schedule.last_campaign_id)
assert generated is not None
assert generated.settings == {"retention": "sealed"}
assert generated.mail_profile_policy == {"profile_id": "profile-1"}
assert schedule.resource_revision == 2
def test_autonomous_occurrences_allocate_commands_once_and_complete_bound(self):
context = SimpleNamespace(
snapshot=SimpleNamespace(
mail_profile_id="profile-1",
smtp_transport_revision="transport-1",
smtp_server_id="smtp-1",
smtp_credential_id="credential-1",
),
message_bytes=b"From: Sender <sender@example.test>\r\nTo: one@example.test\r\n\r\nHello",
envelope_from="sender@example.test",
envelope_recipients=["one@example.test"],
)
job = SimpleNamespace(
id="job-1",
resolved_recipients={"from": {"email": "sender@example.test"}},
)
mail = Mock()
mail.durable_delivery_available = True
mail.delivery_command_summary.return_value = {
"id": "command-1",
"status": "accepted",
"accepted_count": 1,
"refused_count": 0,
"failure_code": None,
}
mail.submit_delivery_command.side_effect = [
{"id": "command-1", "status": "pending", "duplicate": False},
{"id": "command-2", "status": "pending", "duplicate": False},
]
with self.SessionLocal() as session:
schedule = session.get(CampaignSchedule, "schedule-1")
assert schedule is not None
schedule.delivery_mode = "autonomous"
schedule.approved_execution_snapshot_hash = "a" * 64
session.commit()
validation = {
"execution_snapshot_hash": "a" * 64,
"approval_request_id": "approval-1",
"approval_subject_digest": "b" * 64,
"job_count": 1,
"job_manifest_sha256": "c" * 64,
}
patches = (
patch(
"govoplan_campaign.backend.campaign.scheduling.validate_autonomous_schedule_source",
return_value=validation,
),
patch(
"govoplan_campaign.backend.campaign.scheduling._autonomous_source_jobs",
return_value=[job],
),
patch(
"govoplan_campaign.backend.campaign.scheduling._send_job_delivery_context",
return_value=context,
),
patch(
"govoplan_campaign.backend.campaign.scheduling._synchronous_smtp_batch_manager",
return_value=nullcontext(None),
),
patch(
"govoplan_campaign.backend.campaign.scheduling.mail_integration",
return_value=mail,
),
patch("govoplan_campaign.backend.campaign.scheduling.audit_event"),
)
with patches[0], patches[1], patches[2], patches[3], patches[4], patches[5]:
first = dispatch_due_campaign_schedules(
session,
tenant_id="tenant-1",
now=datetime(2026, 8, 7, 8, tzinfo=UTC),
)
session.commit()
second = dispatch_due_campaign_schedules(
session,
tenant_id="tenant-1",
now=datetime(2026, 8, 8, 8, tzinfo=UTC),
)
session.commit()
assert first["autonomous_prepared"] == 1
assert second["autonomous_prepared"] == 1
assert schedule.active is False
occurrences = (
session.query(CampaignScheduleOccurrence)
.order_by(CampaignScheduleOccurrence.scheduled_for)
.all()
)
assert [item.status for item in occurrences] == ["accepted", "prepared"]
assert [item.delivery_command_ids for item in occurrences] == [
["command-1"],
["command-2"],
]
assert len({item.idempotency_key for item in occurrences}) == 2
assert [item.recovery_state for item in occurrences] == [
"complete",
"pending",
]
assert occurrences[0].evidence["source_campaign_id"] == "campaign-1"
assert occurrences[0].evidence["source_version_id"] == "version-1"
assert (
occurrences[0].evidence["source_snapshot_hash"]
== schedule.source_snapshot_hash
)
assert mail.submit_delivery_command.call_count == 2
def test_autonomous_unknown_outcome_pauses_without_resubmission(self):
mail = Mock()
mail.durable_delivery_available = True
mail.delivery_command_summary.return_value = {
"id": "command-1",
"status": "outcome_unknown",
"accepted_count": 0,
"refused_count": 0,
"failure_code": "smtp_outcome_unknown",
}
with self.SessionLocal() as session:
schedule = session.get(CampaignSchedule, "schedule-1")
assert schedule is not None
schedule.delivery_mode = "autonomous"
occurrence = CampaignScheduleOccurrence(
tenant_id="tenant-1",
schedule_id=schedule.id,
scheduled_for=datetime(2026, 8, 6, 8, tzinfo=UTC),
status="prepared",
idempotency_key="occurrence-1",
delivery_command_ids=["command-1"],
recovery_state="pending",
)
session.add(occurrence)
session.commit()
with patch(
"govoplan_campaign.backend.campaign.scheduling.mail_integration",
return_value=mail,
), patch(
"govoplan_campaign.backend.campaign.scheduling._notify_schedule_operator"
) as notify:
result = dispatch_due_campaign_schedules(
session,
tenant_id="tenant-1",
now=datetime(2026, 8, 6, 9, tzinfo=UTC),
)
session.commit()
assert result["refreshed"]["uncertain"] == 1
assert result["selected"] == 0
assert occurrence.status == "uncertain"
assert occurrence.recovery_state == "operator_required"
assert schedule.active is False
assert schedule.last_outcome == "uncertain"
notify.assert_called_once()
mail.submit_delivery_command.assert_not_called()
def test_autonomous_pending_occurrence_defers_the_next_delivery(self):
mail = Mock()
mail.durable_delivery_available = True
mail.delivery_command_summary.return_value = {
"id": "command-1",
"status": "pending",
"accepted_count": 0,
"refused_count": 0,
"failure_code": None,
}
with self.SessionLocal() as session:
schedule = session.get(CampaignSchedule, "schedule-1")
assert schedule is not None
schedule.delivery_mode = "autonomous"
session.add(
CampaignScheduleOccurrence(
tenant_id="tenant-1",
schedule_id=schedule.id,
scheduled_for=datetime(2026, 8, 6, 8, tzinfo=UTC),
status="prepared",
idempotency_key="occurrence-1",
delivery_command_ids=["command-1"],
recovery_state="pending",
)
)
session.commit()
with patch(
"govoplan_campaign.backend.campaign.scheduling.mail_integration",
return_value=mail,
):
result = dispatch_due_campaign_schedules(
session,
tenant_id="tenant-1",
now=datetime(2026, 8, 7, 8, tzinfo=UTC),
)
assert result["refreshed"]["checked"] == 1
assert result["deferred"] == 1
assert result["autonomous_prepared"] == 0
assert schedule.active is True
assert schedule.occurrence_count == 0
mail.submit_delivery_command.assert_not_called()
def test_missing_mail_recovery_capability_pauses_an_open_occurrence(self):
mail = Mock()
mail.durable_delivery_available = False
with self.SessionLocal() as session:
schedule = session.get(CampaignSchedule, "schedule-1")
assert schedule is not None
schedule.delivery_mode = "autonomous"
occurrence = CampaignScheduleOccurrence(
tenant_id="tenant-1",
schedule_id=schedule.id,
scheduled_for=datetime(2026, 8, 6, 8, tzinfo=UTC),
status="prepared",
idempotency_key="occurrence-1",
delivery_command_ids=["command-1"],
recovery_state="pending",
)
session.add(occurrence)
session.commit()
with patch(
"govoplan_campaign.backend.campaign.scheduling.mail_integration",
return_value=mail,
), patch(
"govoplan_campaign.backend.campaign.scheduling._notify_schedule_operator"
) as notify:
result = dispatch_due_campaign_schedules(
session,
tenant_id="tenant-1",
now=datetime(2026, 8, 7, 8, tzinfo=UTC),
)
assert result["refreshed"]["uncertain"] == 1
assert occurrence.status == "uncertain"
assert occurrence.evidence["recovery_reason"] == (
"mail_delivery_outbox_unavailable"
)
assert schedule.active is False
notify.assert_called_once()
def test_autonomous_source_requires_an_explicit_approval(self):
with self.SessionLocal() as session, patch(
"govoplan_campaign.backend.campaign.scheduling.campaign_approval_gate",
return_value=None,
):
campaign = session.get(Campaign, "campaign-1")
version = session.get(CampaignVersion, "version-1")
assert campaign is not None and version is not None
from govoplan_campaign.backend.campaign.scheduling import (
validate_autonomous_schedule_source,
)
try:
validate_autonomous_schedule_source(
session,
campaign=campaign,
version=version,
)
except RuntimeError as exc:
assert "explicit Approval request" in str(exc)
else: # pragma: no cover - defensive assertion
raise AssertionError("Autonomous source validation unexpectedly passed")
def test_duplicate_occurrence_recovers_schedule_without_another_effect(self):
with self.SessionLocal() as session:
schedule = session.get(CampaignSchedule, "schedule-1")
assert schedule is not None
recorded = CampaignScheduleOccurrence(
tenant_id="tenant-1",
schedule_id=schedule.id,
scheduled_for=schedule.next_fire_at,
status="prepared",
idempotency_key="existing-key",
recovery_state="pending",
)
session.add(recorded)
session.commit()
with patch("govoplan_campaign.backend.campaign.scheduling.audit_event"):
result = dispatch_due_campaign_schedules(
session,
tenant_id="tenant-1",
now=datetime(2026, 8, 7, 8, tzinfo=UTC),
)
session.commit()
assert result["duplicates"] == 1
assert result["failed"] == 0
assert schedule.occurrence_count == 1
assert schedule.next_fire_at == datetime(2026, 8, 8, 8, tzinfo=UTC)
assert session.query(CampaignScheduleOccurrence).count() == 1