241 lines
9.0 KiB
Python
241 lines
9.0 KiB
Python
from __future__ import annotations
|
|
|
|
from datetime import UTC, datetime
|
|
from unittest.mock import 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,
|
|
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__,
|
|
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
|