Fence generated artifact retention

This commit is contained in:
2026-08-03 03:54:58 +02:00
parent 9da03090a7
commit 50ce8b0acb
5 changed files with 435 additions and 12 deletions
+3 -2
View File
@@ -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
+5 -2
View File
@@ -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.
+1 -1
View File
@@ -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"),
+276 -7
View File
@@ -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
+150
View File
@@ -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()