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 \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