Files
govoplan-campaign/src/govoplan_campaign/backend/retention.py
T

517 lines
18 KiB
Python

from __future__ import annotations
import copy
from dataclasses import dataclass
import hashlib
import json
from datetime import datetime, timedelta, timezone
from pathlib import Path
from typing import Any, Callable
from sqlalchemy.orm import Session
from govoplan_core.core.recovery import (
RecoveryGuaranteeError,
RecoveryMode,
RecoveryPlan,
RecoveryStatus,
)
from govoplan_core.core.recovery_runtime import (
DurableRecoveryOperation,
RecoveryOperationBusy,
RecoveryOperationStateConflict,
begin_durable_recovery_operation,
)
from govoplan_core.core.object_storage import (
StorageBackend,
StorageBackendError,
StorageObjectMissing,
configured_storage_backend,
)
from govoplan_core.core.runtime_coordination import process_runtime_identity
from govoplan_core.db.session import get_database
from govoplan_core.settings import settings as core_settings
from govoplan_campaign.backend.db.models import (
CampaignJob,
CampaignSchedule,
CampaignVersion,
JobImapStatus,
JobQueueStatus,
)
from govoplan_campaign.backend.runtime import get_settings
FINAL_VERSION_STATES = {
"completed",
"partially_completed",
"outcome_unknown",
"failed",
"cancelled",
"archived",
}
FINAL_EML_SEND_STATUSES = {
"smtp_accepted",
"sent",
"failed_permanent",
"cancelled",
"outcome_unknown",
}
@dataclass(frozen=True, slots=True)
class _GeneratedArtifactRecovery:
operation: DurableRecoveryOperation
job_id: str
storage_key: str | None
local_path: str | None
storage: StorageBackend | None
def _cutoff(days: int | None, *, now: datetime) -> datetime | None:
if days is None:
return None
return now - timedelta(days=days)
def _is_before_cutoff(value: datetime | None, cutoff: datetime | None) -> bool:
if value is None or cutoff is None:
return False
candidate = value.replace(tzinfo=timezone.utc) if value.tzinfo is None else value
return candidate < cutoff
def _json_sha256(value: Any) -> str:
payload = json.dumps(value, sort_keys=True, ensure_ascii=False, default=str).encode("utf-8")
return hashlib.sha256(payload).hexdigest()
def _summary_was_redacted(summary: Any) -> bool:
return isinstance(summary, dict) and bool(summary.get("_retention", {}).get("report_detail_redacted"))
def _redacted_report_summary(summary: Any, *, now: datetime) -> tuple[dict[str, Any] | None, bool]:
if not isinstance(summary, dict):
return summary, False
if _summary_was_redacted(summary):
return summary, False
detail_keys = {"messages", "issues", "recent_failures", "jobs"}
if not any(key in summary for key in detail_keys):
return summary, False
next_summary = copy.deepcopy(summary)
for key in detail_keys:
next_summary.pop(key, None)
retention = dict(next_summary.get("_retention") or {})
retention.update({"report_detail_redacted": True, "redacted_at": now.isoformat()})
next_summary["_retention"] = retention
return next_summary, True
def _apply_raw_json_retention(
session: Session,
*,
dry_run: bool,
now: datetime,
policy_for_campaign_id: Callable[[str | None], object],
) -> dict[str, int]:
result = {"eligible": 0, "redacted": 0, "skipped_not_final": 0, "already_redacted": 0}
versions = session.query(CampaignVersion).order_by(CampaignVersion.updated_at.asc()).all()
for version in versions:
policy = policy_for_campaign_id(version.campaign_id)
cutoff = _cutoff(0 if not policy.store_raw_campaign_json else policy.raw_campaign_json_retention_days, now=now)
if not _is_before_cutoff(version.updated_at, cutoff):
continue
raw_json = version.raw_json if isinstance(version.raw_json, dict) else {}
if raw_json.get("_retention", {}).get("raw_json_redacted"):
result["already_redacted"] += 1
continue
if version.workflow_state not in FINAL_VERSION_STATES:
result["skipped_not_final"] += 1
continue
result["eligible"] += 1
if dry_run:
continue
snapshot = version.execution_snapshot if isinstance(version.execution_snapshot, dict) else {}
version.raw_json = {
"version": version.schema_version or "1.0",
"_retention": {
"raw_json_redacted": True,
"redacted_at": now.isoformat(),
"original_sha256": snapshot.get("campaign_json_sha256") or _json_sha256(raw_json),
},
}
session.add(version)
result["redacted"] += 1
return result
def _artifact_locator_sha256(
*,
storage_key: str | None,
local_path: str | None,
) -> str:
return _json_sha256(
{
"storage_key": storage_key,
"local_path": local_path,
}
)
def _begin_generated_artifact_recovery(
*,
job: CampaignJob,
storage: StorageBackend | None,
) -> _GeneratedArtifactRecovery | None:
storage_key = str(job.eml_storage_key) if job.eml_storage_key else None
local_path = str(job.eml_local_path) if job.eml_local_path else None
locator_sha256 = _artifact_locator_sha256(
storage_key=storage_key,
local_path=local_path,
)
try:
started = begin_durable_recovery_operation(
get_database().SessionLocal,
identity=process_runtime_identity(),
module_id="campaigns",
operation_type="generated-artifact-retention",
idempotency_key=(
f"campaign-retention:{job.id}:{locator_sha256[:32]}"
),
request={
"tenant_id": job.tenant_id,
"campaign_id": job.campaign_id,
"version_id": job.campaign_version_id,
"job_id": job.id,
"message_sha256": job.eml_sha256,
"artifact_locator_sha256": locator_sha256,
},
recovery_plan=RecoveryPlan(
mode=RecoveryMode.FORWARD_RECOVERY,
preconditions=(
"the Campaign job is terminal and outside its retention window",
"no IMAP append or delivery outcome remains unresolved",
),
forward_recovery_steps=(
"verify whether each recorded artifact still exists",
"clear the database locator only after absence is established",
),
verification_steps=(
"reload the Campaign job through an independent session",
"probe every original object or local-development path",
),
),
precondition_evidence={
"job_id": job.id,
"queue_status": job.queue_status,
"send_status": job.send_status,
"imap_status": job.imap_status,
"message_sha256": job.eml_sha256,
"artifact_locator_sha256": locator_sha256,
},
lease_resource_key=f"campaign:retention:{job.tenant_id}:{job.id}",
lease_ttl_seconds=15 * 60,
resource_type="campaign_job",
resource_id=job.id,
metadata={
"resources": [
"postgresql",
"object-storage" if storage_key else "local-development-storage",
],
},
)
except (RecoveryOperationBusy, RecoveryOperationStateConflict):
return None
if started.replayed or started.operation is None:
return None
return _GeneratedArtifactRecovery(
operation=started.operation,
job_id=job.id,
storage_key=storage_key,
local_path=local_path,
storage=storage,
)
def _generated_artifact_recovery_evidence(
recovery: _GeneratedArtifactRecovery,
) -> tuple[str, dict[str, Any]]:
probes: dict[str, bool | None] = {}
if recovery.storage_key:
try:
if recovery.storage is None:
raise StorageBackendError("Artifact storage is unavailable")
probes["object_missing"] = not recovery.storage.exists(
recovery.storage_key
)
except (StorageBackendError, OSError):
probes["object_missing"] = None
if recovery.local_path:
try:
probes["local_path_missing"] = not Path(recovery.local_path).exists()
except OSError:
probes["local_path_missing"] = None
with get_database().SessionLocal() as evidence_session:
job = evidence_session.get(CampaignJob, recovery.job_id)
job_present = job is not None
metadata_cleared = bool(
job is None
or (
(
not recovery.storage_key
or job.eml_storage_key != recovery.storage_key
)
and (
not recovery.local_path
or job.eml_local_path != recovery.local_path
)
)
)
metadata_intact = bool(
job is not None
and job.eml_storage_key == recovery.storage_key
and job.eml_local_path == recovery.local_path
)
probe_values = tuple(probes.values())
probe_verified = bool(probe_values) and all(
value is not None for value in probe_values
)
artifacts_absent = probe_verified and all(value is True for value in probe_values)
artifacts_intact = probe_verified and all(value is False for value in probe_values)
evidence = {
"verified": probe_verified,
"checks": {
"job_state_reloaded": True,
"artifact_locations_probed": probe_verified,
},
"job_present": job_present,
"metadata_cleared": metadata_cleared,
"metadata_intact": metadata_intact,
"artifact_probes": probes,
}
if not probe_verified:
return "outcome_unknown", evidence
if artifacts_absent and metadata_cleared:
return "succeeded", evidence
if artifacts_intact and metadata_intact:
return "failed", evidence
return "recovery_required", evidence
def _finish_generated_artifact_recovery(
recovery: _GeneratedArtifactRecovery,
) -> None:
outcome, evidence = _generated_artifact_recovery_evidence(recovery)
if outcome == "succeeded":
recovery.operation.succeed(evidence=evidence)
elif outcome == "failed":
recovery.operation.reject(
summary="Generated Campaign artifacts were not deleted",
evidence=evidence,
)
elif outcome == "outcome_unknown":
recovery.operation.unresolved(
status=RecoveryStatus.OUTCOME_UNKNOWN,
summary="Generated artifact deletion could not be verified",
evidence=evidence,
failure_summary="Artifact storage availability prevented verification",
)
else:
recovery.operation.unresolved(
status=RecoveryStatus.RECOVERY_REQUIRED,
summary="Generated artifact retention is only partially complete",
evidence=evidence,
failure_summary="Artifact and Campaign metadata state require reconciliation",
)
def _finish_generated_artifact_recoveries(
recoveries: list[_GeneratedArtifactRecovery],
) -> None:
failures: list[Exception] = []
for recovery in recoveries:
try:
_finish_generated_artifact_recovery(recovery)
except Exception as exc: # preserve every operation's chance to close
failures.append(exc)
if failures:
raise RecoveryGuaranteeError(
f"{len(failures)} Campaign retention recovery operation(s) could not be finalized"
) from failures[0]
def _apply_eml_retention(
session: Session,
*,
dry_run: bool,
now: datetime,
policy_for_campaign_id: Callable[[str | None], object],
storage: StorageBackend | None = None,
recovery_operations: list[_GeneratedArtifactRecovery] | None = None,
) -> dict[str, int]:
result = {
"eligible": 0,
"metadata_cleared": 0,
"files_deleted": 0,
"files_missing": 0,
"delete_failed": 0,
"recovery_blocked": 0,
"skipped_not_final": 0,
"skipped_schedule_source": 0,
}
protected_source_versions = {
str(version_id)
for (version_id,) in (
session.query(CampaignSchedule.source_version_id)
.filter(
CampaignSchedule.delivery_mode == "autonomous",
CampaignSchedule.next_fire_at.is_not(None),
)
.all()
)
}
jobs = (
session.query(CampaignJob)
.filter((CampaignJob.eml_local_path.is_not(None)) | (CampaignJob.eml_storage_key.is_not(None)))
.order_by(CampaignJob.updated_at.asc())
.all()
)
for job in jobs:
if getattr(job, "campaign_version_id", None) in protected_source_versions:
result["skipped_schedule_source"] += 1
continue
policy = policy_for_campaign_id(job.campaign_id)
cutoff = _cutoff(policy.generated_eml_retention_days, now=now)
if not _is_before_cutoff(job.updated_at, cutoff):
continue
if job.queue_status in {JobQueueStatus.QUEUED.value, JobQueueStatus.SENDING.value} or job.send_status not in FINAL_EML_SEND_STATUSES:
result["skipped_not_final"] += 1
continue
if job.imap_status in {
JobImapStatus.PENDING.value,
JobImapStatus.APPENDING.value,
JobImapStatus.OUTCOME_UNKNOWN.value,
}:
result["skipped_not_final"] += 1
continue
result["eligible"] += 1
if dry_run:
continue
active_storage = storage
if job.eml_storage_key and active_storage is None:
active_storage = configured_storage_backend(
get_settings() or core_settings
)
if recovery_operations is not None:
recovery = _begin_generated_artifact_recovery(
job=job,
storage=active_storage,
)
if recovery is None:
result["recovery_blocked"] += 1
continue
recovery_operations.append(recovery)
if job.eml_storage_key:
assert active_storage is not None
try:
if active_storage.exists(job.eml_storage_key):
active_storage.delete(job.eml_storage_key)
result["files_deleted"] += 1
else:
result["files_missing"] += 1
except StorageObjectMissing:
result["files_missing"] += 1
except StorageBackendError:
result["delete_failed"] += 1
continue
if job.eml_local_path:
path = Path(job.eml_local_path)
if path.exists():
path.unlink()
result["files_deleted"] += 1
else:
result["files_missing"] += 1
job.eml_local_path = None
job.eml_storage_key = None
session.add(job)
result["metadata_cleared"] += 1
return result
def _apply_report_detail_retention(
session: Session,
*,
dry_run: bool,
now: datetime,
policy_for_campaign_id: Callable[[str | None], object],
) -> dict[str, int]:
result = {"eligible_versions": 0, "summaries_redacted": 0, "already_redacted": 0}
versions = session.query(CampaignVersion).order_by(CampaignVersion.updated_at.asc()).all()
for version in versions:
policy = policy_for_campaign_id(version.campaign_id)
cutoff = _cutoff(policy.stored_report_detail_retention_days, now=now)
if not _is_before_cutoff(version.updated_at, cutoff):
continue
next_validation, validation_changed = _redacted_report_summary(version.validation_summary, now=now)
next_build, build_changed = _redacted_report_summary(version.build_summary, now=now)
if not validation_changed and not build_changed:
if _summary_was_redacted(version.validation_summary) or _summary_was_redacted(version.build_summary):
result["already_redacted"] += 1
continue
result["eligible_versions"] += 1
if dry_run:
continue
version.validation_summary = next_validation
version.build_summary = next_build
session.add(version)
result["summaries_redacted"] += int(validation_changed) + int(build_changed)
return result
def apply_campaign_retention(
session: Session,
*,
dry_run: bool,
now: datetime,
policy_for_campaign_id: Callable[[str | None], object],
) -> dict[str, dict[str, int]]:
recoveries: list[_GeneratedArtifactRecovery] = []
try:
# Start external-effect fences before queries for database-only
# redaction can autoflush unrelated changes in the caller session.
generated_eml = _apply_eml_retention(
session,
dry_run=dry_run,
now=now,
policy_for_campaign_id=policy_for_campaign_id,
recovery_operations=None if dry_run else recoveries,
)
result = {
"raw_campaign_json": _apply_raw_json_retention(
session,
dry_run=dry_run,
now=now,
policy_for_campaign_id=policy_for_campaign_id,
),
"generated_eml": generated_eml,
"stored_report_detail": _apply_report_detail_retention(
session,
dry_run=dry_run,
now=now,
policy_for_campaign_id=policy_for_campaign_id,
),
}
if not dry_run:
# External artifact deletion and its locator update form one
# module-owned recovery boundary. The outer Policy audit commits
# separately after Campaign has verified this boundary.
session.commit()
except Exception:
session.rollback()
if recoveries:
_finish_generated_artifact_recoveries(recoveries)
raise
if recoveries:
_finish_generated_artifact_recoveries(recoveries)
return result