965 lines
35 KiB
Python
965 lines
35 KiB
Python
from __future__ import annotations
|
|
|
|
import calendar
|
|
import copy
|
|
import hashlib
|
|
import json
|
|
from collections.abc import Mapping
|
|
from datetime import UTC, datetime, timedelta
|
|
from email import policy
|
|
from email.parser import BytesParser
|
|
from zoneinfo import ZoneInfo, ZoneInfoNotFoundError
|
|
|
|
from sqlalchemy.orm import Session
|
|
|
|
from govoplan_campaign.backend.campaign.copying import campaign_copy_configuration
|
|
from govoplan_campaign.backend.db.models import (
|
|
Campaign,
|
|
CampaignJob,
|
|
CampaignSchedule,
|
|
CampaignScheduleOccurrence,
|
|
CampaignShare,
|
|
CampaignVersion,
|
|
JobBuildStatus,
|
|
)
|
|
from govoplan_campaign.backend.approval_gate import (
|
|
assert_campaign_approval,
|
|
campaign_approval_gate,
|
|
)
|
|
from govoplan_campaign.backend.campaign.models import DeliveryChannelPolicy
|
|
from govoplan_campaign.backend.integrations import mail_integration
|
|
from govoplan_campaign.backend.persistence.campaigns import (
|
|
create_campaign_version_from_json,
|
|
)
|
|
from govoplan_campaign.backend.sending.execution import ensure_execution_snapshot
|
|
from govoplan_campaign.backend.sending.jobs import (
|
|
_from_header_from_job,
|
|
_send_job_delivery_context,
|
|
_single_job_validation_allowed,
|
|
_synchronous_smtp_batch_manager,
|
|
)
|
|
from govoplan_core.audit.logging import audit_event
|
|
|
|
|
|
RECURRENCE_KINDS = frozenset({"once", "daily", "weekly", "monthly"})
|
|
SCHEDULE_SOURCE_SCHEMA = "govoplan.campaign.schedule-source.v1"
|
|
SCHEDULE_DELIVERY_MODES = frozenset({"manual", "autonomous"})
|
|
|
|
|
|
def canonical_configuration_hash(value: Mapping[str, object]) -> str:
|
|
return hashlib.sha256(
|
|
json.dumps(
|
|
value,
|
|
ensure_ascii=True,
|
|
sort_keys=True,
|
|
separators=(",", ":"),
|
|
).encode("utf-8")
|
|
).hexdigest()
|
|
|
|
|
|
def campaign_schedule_source_snapshot(
|
|
*,
|
|
configuration: Mapping[str, object],
|
|
campaign_settings: Mapping[str, object],
|
|
mail_profile_policy: Mapping[str, object],
|
|
shares: list[Mapping[str, object]],
|
|
) -> dict[str, object]:
|
|
"""Seal every selected source domain so worker execution cannot drift."""
|
|
|
|
return {
|
|
"schema": SCHEDULE_SOURCE_SCHEMA,
|
|
"configuration": copy.deepcopy(dict(configuration)),
|
|
"campaign_settings": copy.deepcopy(dict(campaign_settings)),
|
|
"mail_profile_policy": copy.deepcopy(dict(mail_profile_policy)),
|
|
"shares": [copy.deepcopy(dict(item)) for item in shares],
|
|
}
|
|
|
|
|
|
def next_schedule_fire(
|
|
scheduled_for: datetime,
|
|
*,
|
|
recurrence_kind: str,
|
|
interval_count: int,
|
|
timezone_name: str,
|
|
) -> datetime | None:
|
|
if recurrence_kind == "once":
|
|
return None
|
|
if recurrence_kind not in RECURRENCE_KINDS:
|
|
raise ValueError(f"Unsupported campaign recurrence: {recurrence_kind}")
|
|
if interval_count < 1:
|
|
raise ValueError("Campaign recurrence interval must be positive")
|
|
try:
|
|
zone = ZoneInfo(timezone_name)
|
|
except ZoneInfoNotFoundError as exc:
|
|
raise ValueError(f"Unknown campaign schedule timezone: {timezone_name}") from exc
|
|
local = _as_utc(scheduled_for).astimezone(zone)
|
|
if recurrence_kind == "daily":
|
|
upcoming = local + timedelta(days=interval_count)
|
|
elif recurrence_kind == "weekly":
|
|
upcoming = local + timedelta(weeks=interval_count)
|
|
else:
|
|
month_index = local.year * 12 + local.month - 1 + interval_count
|
|
year, month_offset = divmod(month_index, 12)
|
|
month = month_offset + 1
|
|
day = min(local.day, calendar.monthrange(year, month)[1])
|
|
upcoming = local.replace(year=year, month=month, day=day)
|
|
return upcoming.astimezone(UTC)
|
|
|
|
|
|
def dispatch_due_campaign_schedules(
|
|
session: Session,
|
|
*,
|
|
tenant_id: str | None = None,
|
|
now: datetime | None = None,
|
|
limit: int = 50,
|
|
) -> dict[str, object]:
|
|
observed_at = _as_utc(now or datetime.now(UTC))
|
|
refreshed = refresh_autonomous_schedule_outcomes(
|
|
session,
|
|
tenant_id=tenant_id,
|
|
now=observed_at,
|
|
)
|
|
query = session.query(CampaignSchedule).filter(
|
|
CampaignSchedule.active.is_(True),
|
|
CampaignSchedule.next_fire_at.is_not(None),
|
|
CampaignSchedule.next_fire_at <= observed_at,
|
|
)
|
|
if tenant_id is not None:
|
|
query = query.filter(CampaignSchedule.tenant_id == tenant_id)
|
|
schedules = (
|
|
query.order_by(CampaignSchedule.next_fire_at.asc(), CampaignSchedule.id.asc())
|
|
.with_for_update(skip_locked=True)
|
|
.limit(max(1, min(limit, 250)))
|
|
.all()
|
|
)
|
|
result: dict[str, object] = {
|
|
"selected": len(schedules),
|
|
"prepared": 0,
|
|
"autonomous_prepared": 0,
|
|
"failed": 0,
|
|
"completed": 0,
|
|
"coalesced": 0,
|
|
"duplicates": 0,
|
|
"deferred": 0,
|
|
"campaign_ids": [],
|
|
"operator_actions": [],
|
|
"refreshed": refreshed,
|
|
}
|
|
for schedule in schedules:
|
|
scheduled_for = _as_utc(schedule.next_fire_at or observed_at)
|
|
if schedule.delivery_mode == "autonomous" and _has_open_occurrence(
|
|
session, schedule_id=schedule.id
|
|
):
|
|
result["deferred"] = int(result["deferred"]) + 1
|
|
continue
|
|
try:
|
|
with session.begin_nested():
|
|
if schedule.delivery_mode == "autonomous":
|
|
_occurrence, skipped = _prepare_autonomous_occurrence(
|
|
session,
|
|
schedule=schedule,
|
|
scheduled_for=scheduled_for,
|
|
observed_at=observed_at,
|
|
)
|
|
campaign_id = schedule.campaign_id
|
|
result["autonomous_prepared"] = (
|
|
int(result["autonomous_prepared"]) + 1
|
|
)
|
|
else:
|
|
campaign, _version, skipped = _prepare_occurrence(
|
|
session,
|
|
schedule=schedule,
|
|
scheduled_for=scheduled_for,
|
|
observed_at=observed_at,
|
|
)
|
|
campaign_id = campaign.id
|
|
result["prepared"] = int(result["prepared"]) + 1
|
|
result["coalesced"] = int(result["coalesced"]) + skipped
|
|
result["campaign_ids"].append(campaign_id) # type: ignore[union-attr]
|
|
if not schedule.active:
|
|
result["completed"] = int(result["completed"]) + 1
|
|
except Exception as exc: # noqa: BLE001 - persist bounded operator evidence
|
|
session.expire_all()
|
|
recorded = (
|
|
session.query(CampaignScheduleOccurrence)
|
|
.filter(
|
|
CampaignScheduleOccurrence.schedule_id == schedule.id,
|
|
CampaignScheduleOccurrence.scheduled_for == scheduled_for,
|
|
)
|
|
.one_or_none()
|
|
)
|
|
if recorded is not None:
|
|
result["duplicates"] = int(result["duplicates"]) + 1
|
|
if (
|
|
schedule.active
|
|
and schedule.next_fire_at is not None
|
|
and _as_utc(schedule.next_fire_at) == scheduled_for
|
|
and recorded.status not in {"failed", "uncertain"}
|
|
):
|
|
_advance_schedule(
|
|
session,
|
|
schedule=schedule,
|
|
occurrence=recorded,
|
|
scheduled_for=scheduled_for,
|
|
observed_at=observed_at,
|
|
sequence=schedule.occurrence_count + 1,
|
|
)
|
|
continue
|
|
session.add(
|
|
CampaignScheduleOccurrence(
|
|
tenant_id=schedule.tenant_id,
|
|
schedule_id=schedule.id,
|
|
scheduled_for=scheduled_for,
|
|
status="failed",
|
|
idempotency_key=_occurrence_idempotency_key(
|
|
schedule.id, scheduled_for
|
|
),
|
|
error=str(exc)[:4000],
|
|
recovery_state="failed",
|
|
evidence={"delivery_mode": schedule.delivery_mode},
|
|
last_checked_at=observed_at,
|
|
)
|
|
)
|
|
schedule.active = False
|
|
schedule.last_error = str(exc)[:4000]
|
|
schedule.last_outcome = "failed"
|
|
schedule.last_recovery_state = "operator_required"
|
|
schedule.resource_revision += 1
|
|
session.add(schedule)
|
|
result["failed"] = int(result["failed"]) + 1
|
|
result["operator_actions"].append( # type: ignore[union-attr]
|
|
{
|
|
"schedule_id": schedule.id,
|
|
"campaign_id": schedule.campaign_id,
|
|
"reason": "draft_preparation_failed",
|
|
"delivery_mode": schedule.delivery_mode,
|
|
}
|
|
)
|
|
_notify_schedule_operator(
|
|
session,
|
|
schedule=schedule,
|
|
reason="policy_or_systemic_preflight_failed",
|
|
)
|
|
session.flush()
|
|
return result
|
|
|
|
|
|
def _prepare_occurrence(
|
|
session: Session,
|
|
*,
|
|
schedule: CampaignSchedule,
|
|
scheduled_for: datetime,
|
|
observed_at: datetime,
|
|
) -> tuple[Campaign, CampaignVersion, int]:
|
|
existing = (
|
|
session.query(CampaignScheduleOccurrence)
|
|
.filter(
|
|
CampaignScheduleOccurrence.schedule_id == schedule.id,
|
|
CampaignScheduleOccurrence.scheduled_for == scheduled_for,
|
|
)
|
|
.one_or_none()
|
|
)
|
|
if existing is not None:
|
|
raise RuntimeError("Campaign schedule occurrence was already recorded")
|
|
|
|
source_campaign = session.get(Campaign, schedule.campaign_id)
|
|
if source_campaign is None or source_campaign.tenant_id != schedule.tenant_id:
|
|
raise RuntimeError("Campaign schedule source is no longer available")
|
|
source_version = session.get(CampaignVersion, schedule.source_version_id)
|
|
if source_version is None or source_version.campaign_id != source_campaign.id:
|
|
raise RuntimeError("Campaign schedule source version is no longer available")
|
|
if canonical_configuration_hash(schedule.source_snapshot) != schedule.source_snapshot_hash:
|
|
raise RuntimeError("Campaign schedule source snapshot integrity check failed")
|
|
snapshot = _schedule_snapshot(schedule.source_snapshot)
|
|
|
|
sequence = schedule.occurrence_count + 1
|
|
external_id = _scheduled_external_id(
|
|
source_campaign.external_id,
|
|
schedule.id,
|
|
sequence,
|
|
)
|
|
local_date = scheduled_for.astimezone(ZoneInfo(schedule.timezone)).date().isoformat()
|
|
generated_name = f"{schedule.name} - {local_date}"
|
|
raw_json = campaign_copy_configuration(
|
|
snapshot["configuration"],
|
|
schedule.copy_options,
|
|
)
|
|
metadata = raw_json.get("campaign")
|
|
if not isinstance(metadata, dict):
|
|
raise RuntimeError("Campaign schedule snapshot has no campaign metadata")
|
|
metadata["id"] = external_id
|
|
metadata["name"] = generated_name
|
|
metadata["mode"] = "draft"
|
|
|
|
generated_campaign, generated_version = create_campaign_version_from_json(
|
|
session,
|
|
tenant_id=schedule.tenant_id,
|
|
user_id=schedule.created_by_user_id,
|
|
raw_json=raw_json,
|
|
source_filename=None,
|
|
source_base_path=schedule.source_base_path,
|
|
commit=False,
|
|
)
|
|
if bool(schedule.copy_options.get("include_policies", True)):
|
|
generated_campaign.settings = copy.deepcopy(snapshot["campaign_settings"])
|
|
if bool(schedule.copy_options.get("include_mail_profile", True)):
|
|
generated_campaign.mail_profile_policy = copy.deepcopy(
|
|
snapshot["mail_profile_policy"]
|
|
)
|
|
if bool(schedule.copy_options.get("include_shares", False)):
|
|
_copy_snapshot_shares(
|
|
session,
|
|
schedule=schedule,
|
|
generated_campaign=generated_campaign,
|
|
shares=snapshot["shares"],
|
|
)
|
|
|
|
occurrence = CampaignScheduleOccurrence(
|
|
tenant_id=schedule.tenant_id,
|
|
schedule_id=schedule.id,
|
|
scheduled_for=scheduled_for,
|
|
status="prepared",
|
|
idempotency_key=_occurrence_idempotency_key(schedule.id, scheduled_for),
|
|
generated_campaign_id=generated_campaign.id,
|
|
generated_version_id=generated_version.id,
|
|
recovery_state="none",
|
|
evidence={"delivery_mode": "manual"},
|
|
last_checked_at=observed_at,
|
|
)
|
|
session.add(occurrence)
|
|
session.flush()
|
|
schedule.last_campaign_id = generated_campaign.id
|
|
schedule.last_outcome = "prepared"
|
|
schedule.last_recovery_state = "none"
|
|
coalesced = _advance_schedule(
|
|
session,
|
|
schedule=schedule,
|
|
occurrence=occurrence,
|
|
scheduled_for=scheduled_for,
|
|
observed_at=observed_at,
|
|
sequence=sequence,
|
|
)
|
|
audit_event(
|
|
session,
|
|
tenant_id=schedule.tenant_id,
|
|
user_id=schedule.created_by_user_id,
|
|
action="campaign.schedule.draft_prepared",
|
|
object_type="campaign_schedule",
|
|
object_id=schedule.id,
|
|
details={
|
|
"source_campaign_id": source_campaign.id,
|
|
"source_version_id": source_version.id,
|
|
"scheduled_for": scheduled_for.isoformat(),
|
|
"generated_campaign_id": generated_campaign.id,
|
|
"generated_version_id": generated_version.id,
|
|
"occurrence": sequence,
|
|
"coalesced_missed_intervals": coalesced,
|
|
"delivery_started": False,
|
|
},
|
|
commit=False,
|
|
)
|
|
return generated_campaign, generated_version, coalesced
|
|
|
|
|
|
def validate_autonomous_schedule_source(
|
|
session: Session,
|
|
*,
|
|
campaign: Campaign,
|
|
version: CampaignVersion,
|
|
) -> dict[str, object]:
|
|
"""Validate the exact immutable execution that an autonomous schedule reuses."""
|
|
|
|
gate = campaign_approval_gate(version)
|
|
if gate is None:
|
|
raise RuntimeError(
|
|
"Autonomous delivery requires an explicit Approval request for the built source version."
|
|
)
|
|
assert_campaign_approval(session, tenant_id=campaign.tenant_id, version=version)
|
|
snapshot = ensure_execution_snapshot(session, version)
|
|
snapshot_hash = str(version.execution_snapshot_hash or "")
|
|
if len(snapshot_hash) != 64:
|
|
raise RuntimeError("The approved Campaign execution snapshot is incomplete.")
|
|
jobs = _autonomous_source_jobs(
|
|
session,
|
|
tenant_id=campaign.tenant_id,
|
|
campaign_id=campaign.id,
|
|
version=version,
|
|
)
|
|
mail = mail_integration()
|
|
if not mail.durable_delivery_available:
|
|
raise RuntimeError(
|
|
"Autonomous delivery requires Mail's durable delivery-command outbox."
|
|
)
|
|
if not snapshot.mail_profile_id or not snapshot.smtp_transport_revision:
|
|
raise RuntimeError(
|
|
"The approved Campaign execution has no immutable Mail transport evidence."
|
|
)
|
|
summary = mail.campaign_profile_delivery_summary(
|
|
session,
|
|
tenant_id=campaign.tenant_id,
|
|
campaign_id=campaign.id,
|
|
profile_id=snapshot.mail_profile_id,
|
|
smtp_server_id=snapshot.smtp_server_id,
|
|
smtp_credential_id=snapshot.smtp_credential_id,
|
|
)
|
|
if not summary.get("smtp_available"):
|
|
raise RuntimeError("The approved Campaign Mail transport is unavailable.")
|
|
if summary.get("smtp_transport_revision") != snapshot.smtp_transport_revision:
|
|
raise RuntimeError(
|
|
"The Campaign Mail transport changed after approval; rebuild and approve a new source version."
|
|
)
|
|
return {
|
|
"execution_snapshot_hash": snapshot_hash,
|
|
"approval_request_id": str(gate.get("request_id") or ""),
|
|
"approval_subject_digest": str(gate.get("subject_digest") or ""),
|
|
"job_count": len(jobs),
|
|
"job_manifest_sha256": canonical_configuration_hash(
|
|
{"jobs": [{"id": job.id, "eml_sha256": job.eml_sha256} for job in jobs]}
|
|
),
|
|
}
|
|
|
|
|
|
def _autonomous_source_jobs(
|
|
session: Session,
|
|
*,
|
|
tenant_id: str,
|
|
campaign_id: str,
|
|
version: CampaignVersion,
|
|
) -> list[CampaignJob]:
|
|
jobs = (
|
|
session.query(CampaignJob)
|
|
.filter(
|
|
CampaignJob.tenant_id == tenant_id,
|
|
CampaignJob.campaign_id == campaign_id,
|
|
CampaignJob.campaign_version_id == version.id,
|
|
)
|
|
.order_by(CampaignJob.entry_index.asc(), CampaignJob.id.asc())
|
|
.all()
|
|
)
|
|
if not jobs:
|
|
raise RuntimeError(
|
|
"Autonomous delivery requires a built source version with recipient jobs."
|
|
)
|
|
for job in jobs:
|
|
if job.build_status != JobBuildStatus.BUILT.value:
|
|
raise RuntimeError(
|
|
"Autonomous delivery requires every source message to be built."
|
|
)
|
|
if not _single_job_validation_allowed(version, job, include_warnings=True):
|
|
raise RuntimeError(
|
|
"Autonomous delivery requires every source message to pass its reviewed recipient and attachment gates."
|
|
)
|
|
if DeliveryChannelPolicy(job.delivery_channel_policy) != DeliveryChannelPolicy.MAIL:
|
|
raise RuntimeError(
|
|
"Autonomous schedules currently support Mail-only delivery; use manual mode for hybrid, Postbox, or print delivery."
|
|
)
|
|
return jobs
|
|
|
|
|
|
def _prepare_autonomous_occurrence(
|
|
session: Session,
|
|
*,
|
|
schedule: CampaignSchedule,
|
|
scheduled_for: datetime,
|
|
observed_at: datetime,
|
|
) -> tuple[CampaignScheduleOccurrence, int]:
|
|
existing = (
|
|
session.query(CampaignScheduleOccurrence)
|
|
.filter(
|
|
CampaignScheduleOccurrence.schedule_id == schedule.id,
|
|
CampaignScheduleOccurrence.scheduled_for == scheduled_for,
|
|
)
|
|
.one_or_none()
|
|
)
|
|
if existing is not None:
|
|
raise RuntimeError("Campaign schedule occurrence was already recorded")
|
|
if canonical_configuration_hash(schedule.source_snapshot) != schedule.source_snapshot_hash:
|
|
raise RuntimeError("Campaign schedule source snapshot integrity check failed")
|
|
campaign = session.get(Campaign, schedule.campaign_id)
|
|
version = session.get(CampaignVersion, schedule.source_version_id)
|
|
if campaign is None or campaign.tenant_id != schedule.tenant_id:
|
|
raise RuntimeError("Campaign schedule source is no longer available")
|
|
if version is None or version.campaign_id != campaign.id:
|
|
raise RuntimeError("Campaign schedule source version is no longer available")
|
|
validation = validate_autonomous_schedule_source(
|
|
session,
|
|
campaign=campaign,
|
|
version=version,
|
|
)
|
|
if (
|
|
not schedule.approved_execution_snapshot_hash
|
|
or validation["execution_snapshot_hash"]
|
|
!= schedule.approved_execution_snapshot_hash
|
|
):
|
|
raise RuntimeError(
|
|
"The approved Campaign execution changed after the autonomous schedule was created."
|
|
)
|
|
jobs = _autonomous_source_jobs(
|
|
session,
|
|
tenant_id=schedule.tenant_id,
|
|
campaign_id=campaign.id,
|
|
version=version,
|
|
)
|
|
occurrence_key = _occurrence_idempotency_key(schedule.id, scheduled_for)
|
|
occurrence = CampaignScheduleOccurrence(
|
|
tenant_id=schedule.tenant_id,
|
|
schedule_id=schedule.id,
|
|
scheduled_for=scheduled_for,
|
|
status="preparing",
|
|
idempotency_key=occurrence_key,
|
|
recovery_state="prepared",
|
|
evidence={
|
|
"delivery_mode": "autonomous",
|
|
"source_campaign_id": campaign.id,
|
|
"source_version_id": version.id,
|
|
"source_snapshot_hash": schedule.source_snapshot_hash,
|
|
**validation,
|
|
},
|
|
last_checked_at=observed_at,
|
|
)
|
|
session.add(occurrence)
|
|
session.flush()
|
|
|
|
contexts = {job.id: _send_job_delivery_context(session, job) for job in jobs}
|
|
with _synchronous_smtp_batch_manager(session, jobs=jobs, contexts=contexts):
|
|
pass
|
|
|
|
mail = mail_integration()
|
|
commands: list[dict[str, object]] = []
|
|
for job in jobs:
|
|
context = contexts[job.id]
|
|
if context.envelope_from is None or not context.envelope_recipients:
|
|
raise RuntimeError("A frozen Campaign message has no delivery envelope.")
|
|
message = BytesParser(policy=policy.default).parsebytes(context.message_bytes)
|
|
commands.append(
|
|
mail.submit_delivery_command(
|
|
session,
|
|
tenant_id=schedule.tenant_id,
|
|
command_type="campaign_schedule_occurrence",
|
|
source_module="campaigns",
|
|
source_resource_type="campaign",
|
|
source_resource_id=campaign.id,
|
|
source_version_id=version.id,
|
|
idempotency_key=f"{occurrence_key}:{job.id}",
|
|
profile_id=context.snapshot.mail_profile_id,
|
|
message_bytes=context.message_bytes,
|
|
envelope_from=context.envelope_from,
|
|
envelope_recipients=context.envelope_recipients,
|
|
from_header=_from_header_from_job(job) or str(message.get("From") or ""),
|
|
expected_smtp_transport_revision=(
|
|
context.snapshot.smtp_transport_revision or ""
|
|
),
|
|
smtp_server_id=context.snapshot.smtp_server_id,
|
|
smtp_credential_id=context.snapshot.smtp_credential_id,
|
|
created_by_user_id=schedule.created_by_user_id,
|
|
)
|
|
)
|
|
occurrence.delivery_command_ids = [str(item["id"]) for item in commands]
|
|
occurrence.status = "prepared"
|
|
occurrence.recovery_state = "pending"
|
|
occurrence.evidence = {
|
|
**occurrence.evidence,
|
|
"command_count": len(commands),
|
|
"duplicate_command_count": sum(bool(item.get("duplicate")) for item in commands),
|
|
"command_status_counts": _status_counts(commands),
|
|
}
|
|
occurrence.last_checked_at = observed_at
|
|
sequence = schedule.occurrence_count + 1
|
|
schedule.last_campaign_id = campaign.id
|
|
schedule.last_outcome = "prepared"
|
|
schedule.last_recovery_state = "pending"
|
|
coalesced = _advance_schedule(
|
|
session,
|
|
schedule=schedule,
|
|
occurrence=occurrence,
|
|
scheduled_for=scheduled_for,
|
|
observed_at=observed_at,
|
|
sequence=sequence,
|
|
)
|
|
audit_event(
|
|
session,
|
|
tenant_id=schedule.tenant_id,
|
|
user_id=schedule.created_by_user_id,
|
|
action="campaign.schedule.delivery_prepared",
|
|
object_type="campaign_schedule_occurrence",
|
|
object_id=occurrence.id,
|
|
details={
|
|
"schedule_id": schedule.id,
|
|
"campaign_id": campaign.id,
|
|
"source_version_id": version.id,
|
|
"scheduled_for": scheduled_for.isoformat(),
|
|
"occurrence_idempotency_key": occurrence_key,
|
|
"delivery_command_count": len(commands),
|
|
"execution_snapshot_hash": validation["execution_snapshot_hash"],
|
|
"approval_request_id": validation["approval_request_id"],
|
|
"coalesced_missed_intervals": coalesced,
|
|
},
|
|
commit=False,
|
|
)
|
|
return occurrence, coalesced
|
|
|
|
|
|
def _advance_schedule(
|
|
session: Session,
|
|
*,
|
|
schedule: CampaignSchedule,
|
|
occurrence: CampaignScheduleOccurrence,
|
|
scheduled_for: datetime,
|
|
observed_at: datetime,
|
|
sequence: int,
|
|
) -> int:
|
|
schedule.occurrence_count = sequence
|
|
schedule.last_fired_at = scheduled_for
|
|
schedule.last_error = None
|
|
next_fire = next_schedule_fire(
|
|
scheduled_for,
|
|
recurrence_kind=schedule.recurrence_kind,
|
|
interval_count=schedule.interval_count,
|
|
timezone_name=schedule.timezone,
|
|
)
|
|
coalesced = 0
|
|
while next_fire is not None and next_fire <= observed_at:
|
|
session.add(
|
|
CampaignScheduleOccurrence(
|
|
tenant_id=schedule.tenant_id,
|
|
schedule_id=schedule.id,
|
|
scheduled_for=next_fire,
|
|
status="superseded",
|
|
idempotency_key=_occurrence_idempotency_key(schedule.id, next_fire),
|
|
recovery_state="superseded",
|
|
evidence={
|
|
"delivery_mode": schedule.delivery_mode,
|
|
"reason": "coalesced_missed_interval",
|
|
"superseded_by_occurrence_id": occurrence.id,
|
|
},
|
|
last_checked_at=observed_at,
|
|
)
|
|
)
|
|
next_fire = next_schedule_fire(
|
|
next_fire,
|
|
recurrence_kind=schedule.recurrence_kind,
|
|
interval_count=schedule.interval_count,
|
|
timezone_name=schedule.timezone,
|
|
)
|
|
coalesced += 1
|
|
if (
|
|
next_fire is None
|
|
or sequence >= schedule.max_occurrences
|
|
or (schedule.ends_at is not None and next_fire > _as_utc(schedule.ends_at))
|
|
):
|
|
schedule.active = False
|
|
schedule.next_fire_at = None
|
|
else:
|
|
schedule.next_fire_at = next_fire
|
|
schedule.resource_revision += 1
|
|
session.add(schedule)
|
|
return coalesced
|
|
|
|
|
|
def refresh_autonomous_schedule_outcomes(
|
|
session: Session,
|
|
*,
|
|
tenant_id: str | None = None,
|
|
now: datetime | None = None,
|
|
) -> dict[str, int]:
|
|
observed_at = _as_utc(now or datetime.now(UTC))
|
|
query = session.query(CampaignScheduleOccurrence).filter(
|
|
CampaignScheduleOccurrence.status.in_(("prepared", "uncertain")),
|
|
)
|
|
if tenant_id is not None:
|
|
query = query.filter(CampaignScheduleOccurrence.tenant_id == tenant_id)
|
|
counts = {
|
|
"checked": 0,
|
|
"accepted": 0,
|
|
"uncertain": 0,
|
|
"failed": 0,
|
|
"skipped": 0,
|
|
}
|
|
mail = mail_integration()
|
|
if not mail.durable_delivery_available:
|
|
for occurrence in query.order_by(
|
|
CampaignScheduleOccurrence.created_at
|
|
).limit(250):
|
|
if not occurrence.delivery_command_ids:
|
|
continue
|
|
counts["checked"] += 1
|
|
counts["uncertain"] += 1
|
|
_mark_occurrence_uncertain(
|
|
session,
|
|
occurrence=occurrence,
|
|
observed_at=observed_at,
|
|
reason="mail_delivery_outbox_unavailable",
|
|
)
|
|
return counts
|
|
for occurrence in query.order_by(CampaignScheduleOccurrence.created_at).limit(250):
|
|
if not occurrence.delivery_command_ids:
|
|
continue
|
|
summaries: list[dict[str, object]] = []
|
|
try:
|
|
summaries = [
|
|
mail.delivery_command_summary(
|
|
session,
|
|
tenant_id=occurrence.tenant_id,
|
|
command_id=command_id,
|
|
)
|
|
for command_id in occurrence.delivery_command_ids
|
|
]
|
|
except Exception:
|
|
counts["checked"] += 1
|
|
counts["uncertain"] += 1
|
|
_mark_occurrence_uncertain(
|
|
session,
|
|
occurrence=occurrence,
|
|
observed_at=observed_at,
|
|
reason="mail_delivery_status_unavailable",
|
|
)
|
|
continue
|
|
counts["checked"] += 1
|
|
outcome, recovery_state = _aggregate_command_outcome(summaries)
|
|
previous_outcome = occurrence.status
|
|
previous_recovery_state = occurrence.recovery_state
|
|
occurrence.status = outcome
|
|
occurrence.recovery_state = recovery_state
|
|
occurrence.last_checked_at = observed_at
|
|
occurrence.evidence = {
|
|
**(occurrence.evidence or {}),
|
|
"command_status_counts": _status_counts(summaries),
|
|
"accepted_recipient_count": sum(
|
|
int(item.get("accepted_count") or 0) for item in summaries
|
|
),
|
|
"refused_recipient_count": sum(
|
|
int(item.get("refused_count") or 0) for item in summaries
|
|
),
|
|
"failure_codes": sorted(
|
|
{
|
|
str(item["failure_code"])
|
|
for item in summaries
|
|
if item.get("failure_code")
|
|
}
|
|
),
|
|
}
|
|
schedule = session.get(CampaignSchedule, occurrence.schedule_id)
|
|
if schedule is not None:
|
|
schedule.last_outcome = outcome
|
|
schedule.last_recovery_state = recovery_state
|
|
transitioned_to_operator_required = (
|
|
outcome in {"uncertain", "failed"}
|
|
and (
|
|
previous_outcome != outcome
|
|
or previous_recovery_state != recovery_state
|
|
or schedule.active
|
|
)
|
|
)
|
|
if transitioned_to_operator_required:
|
|
schedule.active = False
|
|
schedule.last_error = (
|
|
"Autonomous delivery needs operator review; automatic recurrence is paused."
|
|
)
|
|
schedule.resource_revision += 1
|
|
_notify_schedule_operator(
|
|
session,
|
|
schedule=schedule,
|
|
reason=f"delivery_{outcome}",
|
|
)
|
|
session.add(schedule)
|
|
session.add(occurrence)
|
|
if outcome in counts:
|
|
counts[outcome] += 1
|
|
return counts
|
|
|
|
|
|
def _mark_occurrence_uncertain(
|
|
session: Session,
|
|
*,
|
|
occurrence: CampaignScheduleOccurrence,
|
|
observed_at: datetime,
|
|
reason: str,
|
|
) -> None:
|
|
previous_outcome = occurrence.status
|
|
previous_recovery_state = occurrence.recovery_state
|
|
occurrence.status = "uncertain"
|
|
occurrence.recovery_state = "operator_required"
|
|
occurrence.last_checked_at = observed_at
|
|
occurrence.evidence = {
|
|
**(occurrence.evidence or {}),
|
|
"recovery_reason": reason,
|
|
}
|
|
schedule = session.get(CampaignSchedule, occurrence.schedule_id)
|
|
if schedule is not None:
|
|
transitioned = (
|
|
previous_outcome != "uncertain"
|
|
or previous_recovery_state != "operator_required"
|
|
or schedule.active
|
|
)
|
|
schedule.active = False
|
|
schedule.last_outcome = "uncertain"
|
|
schedule.last_recovery_state = "operator_required"
|
|
schedule.last_error = (
|
|
"Autonomous delivery status is unavailable; automatic recurrence is paused."
|
|
)
|
|
if transitioned:
|
|
schedule.resource_revision += 1
|
|
_notify_schedule_operator(
|
|
session,
|
|
schedule=schedule,
|
|
reason=reason,
|
|
)
|
|
session.add(schedule)
|
|
session.add(occurrence)
|
|
|
|
|
|
def _has_open_occurrence(session: Session, *, schedule_id: str) -> bool:
|
|
rows = (
|
|
session.query(CampaignScheduleOccurrence.delivery_command_ids)
|
|
.filter(
|
|
CampaignScheduleOccurrence.schedule_id == schedule_id,
|
|
CampaignScheduleOccurrence.status == "prepared",
|
|
)
|
|
.limit(1000)
|
|
.all()
|
|
)
|
|
return any(bool(command_ids) for (command_ids,) in rows)
|
|
|
|
|
|
def _aggregate_command_outcome(
|
|
summaries: list[dict[str, object]],
|
|
) -> tuple[str, str]:
|
|
statuses = {str(item.get("status") or "") for item in summaries}
|
|
if statuses and statuses <= {"accepted", "reconciled_accepted"}:
|
|
return "accepted", "complete"
|
|
if statuses and statuses <= {"reconciled_not_accepted"}:
|
|
return "skipped", "reconciled"
|
|
if statuses & {"outcome_unknown", "in_progress"}:
|
|
return "uncertain", "operator_required"
|
|
if statuses & {"permanent_failure", "partially_refused", "reconciled_not_accepted"}:
|
|
return "failed", "operator_required"
|
|
return "prepared", "pending"
|
|
|
|
|
|
def _status_counts(items: list[dict[str, object]]) -> dict[str, int]:
|
|
result: dict[str, int] = {}
|
|
for item in items:
|
|
status = str(item.get("status") or "unknown")
|
|
result[status] = result.get(status, 0) + 1
|
|
return result
|
|
|
|
|
|
def _occurrence_idempotency_key(schedule_id: str, scheduled_for: datetime) -> str:
|
|
return f"campaign-schedule:{schedule_id}:{_as_utc(scheduled_for).isoformat()}"
|
|
|
|
|
|
def _notify_schedule_operator(
|
|
session: Session,
|
|
*,
|
|
schedule: CampaignSchedule,
|
|
reason: str,
|
|
) -> None:
|
|
from govoplan_core.core.notifications import (
|
|
NotificationDispatchRequest,
|
|
notification_dispatch_provider,
|
|
)
|
|
from govoplan_campaign.backend.runtime import get_registry
|
|
|
|
provider = notification_dispatch_provider(get_registry())
|
|
if provider is None:
|
|
return
|
|
try:
|
|
provider.enqueue_notification(
|
|
session,
|
|
NotificationDispatchRequest(
|
|
tenant_id=schedule.tenant_id,
|
|
source_module="campaigns",
|
|
source_resource_type="campaign_schedule",
|
|
source_resource_id=schedule.id,
|
|
event_kind="campaign.schedule.operator_required",
|
|
channel="inbox",
|
|
recipient_type="user" if schedule.created_by_user_id else None,
|
|
recipient_id=schedule.created_by_user_id,
|
|
subject=f"Campaign schedule paused: {schedule.name}",
|
|
body_text=(
|
|
"Autonomous Campaign delivery was paused before another occurrence. "
|
|
"Review its recovery evidence before resuming."
|
|
),
|
|
action_url=f"/campaigns/{schedule.campaign_id}",
|
|
priority=2,
|
|
payload={"schedule_id": schedule.id, "reason": reason},
|
|
),
|
|
enqueue_delivery=False,
|
|
)
|
|
except Exception:
|
|
return
|
|
|
|
|
|
def _copy_snapshot_shares(
|
|
session: Session,
|
|
*,
|
|
schedule: CampaignSchedule,
|
|
generated_campaign: Campaign,
|
|
shares: object,
|
|
) -> None:
|
|
if not isinstance(shares, list):
|
|
raise RuntimeError("Campaign schedule share snapshot is invalid")
|
|
for source in shares:
|
|
if not isinstance(source, Mapping):
|
|
raise RuntimeError("Campaign schedule share snapshot is invalid")
|
|
target_type = str(source.get("target_type") or "")
|
|
target_id = str(source.get("target_id") or "")
|
|
permission = str(source.get("permission") or "read")
|
|
if not target_type or not target_id:
|
|
raise RuntimeError("Campaign schedule share snapshot is incomplete")
|
|
session.add(
|
|
CampaignShare(
|
|
tenant_id=schedule.tenant_id,
|
|
campaign_id=generated_campaign.id,
|
|
target_type=target_type,
|
|
target_id=target_id,
|
|
permission=permission,
|
|
created_by_user_id=schedule.created_by_user_id,
|
|
)
|
|
)
|
|
|
|
|
|
def _schedule_snapshot(value: object) -> dict[str, object]:
|
|
if not isinstance(value, Mapping) or value.get("schema") != SCHEDULE_SOURCE_SCHEMA:
|
|
raise RuntimeError("Campaign schedule source snapshot schema is invalid")
|
|
configuration = value.get("configuration")
|
|
settings = value.get("campaign_settings")
|
|
mail_policy = value.get("mail_profile_policy")
|
|
shares = value.get("shares")
|
|
if (
|
|
not isinstance(configuration, Mapping)
|
|
or not isinstance(settings, Mapping)
|
|
or not isinstance(mail_policy, Mapping)
|
|
or not isinstance(shares, list)
|
|
):
|
|
raise RuntimeError("Campaign schedule source snapshot is incomplete")
|
|
return {
|
|
"configuration": dict(configuration),
|
|
"campaign_settings": dict(settings),
|
|
"mail_profile_policy": dict(mail_policy),
|
|
"shares": shares,
|
|
}
|
|
|
|
|
|
def _scheduled_external_id(source: str, schedule_id: str, sequence: int) -> str:
|
|
suffix = f"-scheduled-{schedule_id[:8]}-{sequence}"
|
|
return f"{source[:255 - len(suffix)]}{suffix}"
|
|
|
|
|
|
def _as_utc(value: datetime) -> datetime:
|
|
if value.tzinfo is None:
|
|
return value.replace(tzinfo=UTC)
|
|
return value.astimezone(UTC)
|
|
|
|
|
|
__all__ = [
|
|
"RECURRENCE_KINDS",
|
|
"SCHEDULE_SOURCE_SCHEMA",
|
|
"campaign_schedule_source_snapshot",
|
|
"canonical_configuration_hash",
|
|
"dispatch_due_campaign_schedules",
|
|
"next_schedule_fire",
|
|
"refresh_autonomous_schedule_outcomes",
|
|
"validate_autonomous_schedule_source",
|
|
]
|