diff --git a/README.md b/README.md index 24f2255..241d94a 100644 --- a/README.md +++ b/README.md @@ -34,8 +34,9 @@ object-storage contract under opaque Campaign-owned keys. Templates returns a bounded artifact or a Files-managed artifact for printable output. Database job rows retain the expected hashes and provenance. Workers resolve and verify the frozen evidence before delivery. Build failure compensates objects written -before database commit, and retention keeps database references when object -deletion fails so cleanup can be retried. +before database commit. Retention uses a fenced forward-recovery operation and +independently verifies both artifact absence and the committed locator update; +partial or unobservable cleanup remains visible in Ops. ## Dependencies diff --git a/docs/CAMPAIGN_DELIVERY_RUNBOOK.md b/docs/CAMPAIGN_DELIVERY_RUNBOOK.md index b69ee80..7eed86c 100644 --- a/docs/CAMPAIGN_DELIVERY_RUNBOOK.md +++ b/docs/CAMPAIGN_DELIVERY_RUNBOOK.md @@ -156,8 +156,11 @@ before attempting delivery. business fields. - A build failure deletes objects written before the database transaction can commit. -- Retention clears metadata only after object deletion succeeds; an unavailable - backend leaves the reference in place for a later retry. +- Retention starts a job-fenced forward-recovery operation before deletion, + commits metadata changes in the Campaign-owned boundary, and independently + probes the original locator before reporting success. An unavailable backend + leaves an outcome-unknown operation; a deletion/metadata mismatch becomes + recovery-required in Ops. - A hard process loss between object creation and metadata commit can leave an orphan object. Reconcile only within the Campaign build prefix and verify that no job references the object before deleting it. diff --git a/src/govoplan_campaign/backend/manifest.py b/src/govoplan_campaign/backend/manifest.py index 3fc0279..f8e334b 100644 --- a/src/govoplan_campaign/backend/manifest.py +++ b/src/govoplan_campaign/backend/manifest.py @@ -1052,7 +1052,7 @@ manifest = ModuleManifest( id="campaigns.reference.shared-build-artifacts", title="Operate Campaign build artifacts across workers", summary="Generated messages use shared object storage and are verified before delivery.", - body="Campaign stores generated EML under opaque shared object keys and records expected size, SHA-256 digest, and Message-ID in each job. A worker may run on another node and verifies that evidence before delivery. Before object or Files-owned output effects, a fenced Core recovery operation records the canonical build request, validated-version evidence, and reserved object prefix. Before a real Mail, Postbox, or print effect, a separate job-fenced operation records immutable message and recipient digests and later verifies the authoritative Campaign/channel attempt state. Definitive rejection is distinct from accepted, outcome-unknown, and recovery-required work in Ops. Object-only failures prove compensation; managed Files output and uncertain cleanup remain explicit forward-recovery work. Retention retains metadata when deletion fails. Never copy or edit runtime object keys as business data.", + body="Campaign stores generated EML under opaque shared object keys and records expected size, SHA-256 digest, and Message-ID in each job. A worker may run on another node and verifies that evidence before delivery. Before object or Files-owned output effects, a fenced Core recovery operation records the canonical build request, validated-version evidence, and reserved object prefix. Before a real Mail, Postbox, or print effect, a separate job-fenced operation records immutable message and recipient digests and later verifies the authoritative Campaign/channel attempt state. Definitive rejection is distinct from accepted, outcome-unknown, and recovery-required work in Ops. Object-only failures prove compensation; managed Files output and uncertain cleanup remain explicit forward-recovery work. Retention commits locator changes and independently verifies artifact absence; partial cleanup remains recovery-required. Never copy or edit runtime object keys as business data.", layer="evidence", documentation_types=("admin",), audience=("campaign_operator", "platform_operator", "release_reviewer"), diff --git a/src/govoplan_campaign/backend/retention.py b/src/govoplan_campaign/backend/retention.py index c3c13dd..1f3600d 100644 --- a/src/govoplan_campaign/backend/retention.py +++ b/src/govoplan_campaign/backend/retention.py @@ -1,6 +1,7 @@ from __future__ import annotations import copy +from dataclasses import dataclass import hashlib import json from datetime import datetime, timedelta, timezone @@ -9,12 +10,26 @@ 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, CampaignVersion, JobImapStatus, JobQueueStatus from govoplan_campaign.backend.runtime import get_settings @@ -36,6 +51,15 @@ FINAL_EML_SEND_STATUSES = { } +@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 @@ -113,6 +137,203 @@ def _apply_raw_json_retention( 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, *, @@ -120,6 +341,7 @@ def _apply_eml_retention( 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, @@ -127,6 +349,7 @@ def _apply_eml_retention( "files_deleted": 0, "files_missing": 0, "delete_failed": 0, + "recovery_blocked": 0, "skipped_not_final": 0, } jobs = ( @@ -153,10 +376,22 @@ def _apply_eml_retention( result["eligible"] += 1 if dry_run: continue - if job.eml_storage_key: - active_storage = storage or configured_storage_backend( + 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) @@ -219,8 +454,42 @@ def apply_campaign_retention( now: datetime, policy_for_campaign_id: Callable[[str | None], object], ) -> dict[str, dict[str, int]]: - return { - "raw_campaign_json": _apply_raw_json_retention(session, dry_run=dry_run, now=now, policy_for_campaign_id=policy_for_campaign_id), - "generated_eml": _apply_eml_retention(session, dry_run=dry_run, now=now, policy_for_campaign_id=policy_for_campaign_id), - "stored_report_detail": _apply_report_detail_retention(session, dry_run=dry_run, now=now, policy_for_campaign_id=policy_for_campaign_id), - } + 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 diff --git a/tests/test_retention_recovery.py b/tests/test_retention_recovery.py new file mode 100644 index 0000000..b1ff6dc --- /dev/null +++ b/tests/test_retention_recovery.py @@ -0,0 +1,150 @@ +from __future__ import annotations + +from contextlib import AbstractContextManager +from datetime import datetime, timezone +from types import SimpleNamespace +from unittest.mock import Mock + +import pytest + +from govoplan_core.core.recovery import RecoveryStatus +from govoplan_campaign.backend import retention + + +class _EvidenceSession(AbstractContextManager): + def __init__(self, job: object | None) -> None: + self.job = job + + def __enter__(self): + return self + + def __exit__(self, *_args) -> None: + return None + + def get(self, _model, _object_id): + return self.job + + +def _recovery(*, storage, storage_key: str | None, local_path: str | None): + return retention._GeneratedArtifactRecovery( + operation=Mock(), + job_id="job-1", + storage_key=storage_key, + local_path=local_path, + storage=storage, + ) + + +@pytest.mark.parametrize( + ("outcome", "expected_method", "expected_status"), + [ + ("succeeded", "succeed", None), + ("failed", "reject", None), + ("outcome_unknown", "unresolved", RecoveryStatus.OUTCOME_UNKNOWN), + ("recovery_required", "unresolved", RecoveryStatus.RECOVERY_REQUIRED), + ], +) +def test_retention_recovery_maps_verified_artifact_state( + monkeypatch, + outcome: str, + expected_method: str, + expected_status: RecoveryStatus | None, +) -> None: + recovery = _recovery( + storage=Mock(), + storage_key="campaigns/message.eml", + local_path=None, + ) + monkeypatch.setattr( + retention, + "_generated_artifact_recovery_evidence", + lambda _recovery: ( + outcome, + { + "verified": outcome != "outcome_unknown", + "checks": {"artifact_locations_probed": True}, + }, + ), + ) + + retention._finish_generated_artifact_recovery(recovery) + + method = getattr(recovery.operation, expected_method) + method.assert_called_once() + if expected_status is not None: + assert method.call_args.kwargs["status"] == expected_status + for other in {"succeed", "reject", "unresolved"} - {expected_method}: + getattr(recovery.operation, other).assert_not_called() + + +@pytest.mark.parametrize( + ("object_exists", "metadata_key", "expected"), + [ + (False, None, "succeeded"), + (True, "campaigns/message.eml", "failed"), + (False, "campaigns/message.eml", "recovery_required"), + (True, None, "recovery_required"), + ], +) +def test_retention_recovery_compares_storage_and_database_independently( + monkeypatch, + object_exists: bool, + metadata_key: str | None, + expected: str, +) -> None: + storage = Mock() + storage.exists.return_value = object_exists + job = SimpleNamespace( + eml_storage_key=metadata_key, + eml_local_path=None, + ) + monkeypatch.setattr( + retention, + "get_database", + lambda: SimpleNamespace( + SessionLocal=lambda: _EvidenceSession(job), + ), + ) + + outcome, evidence = retention._generated_artifact_recovery_evidence( + _recovery( + storage=storage, + storage_key="campaigns/message.eml", + local_path=None, + ) + ) + + assert outcome == expected + assert evidence["verified"] is True + + +def test_campaign_retention_commits_before_recovery_verification(monkeypatch) -> None: + events: list[str] = [] + session = Mock() + session.commit.side_effect = lambda: events.append("commit") + recovery = Mock() + + monkeypatch.setattr(retention, "_apply_raw_json_retention", lambda *_args, **_kwargs: {}) + monkeypatch.setattr(retention, "_apply_report_detail_retention", lambda *_args, **_kwargs: {}) + + def apply_eml(*_args, **kwargs): + kwargs["recovery_operations"].append(recovery) + return {"metadata_cleared": 1} + + monkeypatch.setattr(retention, "_apply_eml_retention", apply_eml) + monkeypatch.setattr( + retention, + "_finish_generated_artifact_recoveries", + lambda _recoveries: events.append("verify"), + ) + + result = retention.apply_campaign_retention( + session, + dry_run=False, + now=datetime.now(timezone.utc), + policy_for_campaign_id=lambda _campaign_id: object(), + ) + + assert result["generated_eml"] == {"metadata_cleared": 1} + assert events == ["commit", "verify"] + session.rollback.assert_not_called()