Adopt recovery ledger for Campaign builds

This commit is contained in:
2026-08-03 03:03:00 +02:00
parent c6bbdae2e1
commit dd09b06c47
9 changed files with 592 additions and 103 deletions
@@ -22,6 +22,8 @@ from govoplan_core.core.object_storage import (
StorageBackendError,
configured_storage_backend,
)
from govoplan_core.core.recovery import RecoveryStatus
from govoplan_core.core.recovery_runtime import DurableRecoveryOperation
from govoplan_core.core.templates import TemplateRenderRequest
from govoplan_core.settings import settings as core_settings
from govoplan_campaign.backend.db.models import (
@@ -160,6 +162,7 @@ def _persist_built_eml_artifacts(
version_id: str,
build_id: str,
built_messages: list[Any],
written_storage_keys: list[str] | None = None,
) -> dict[int, _StoredEmlArtifact]:
artifacts: dict[int, _StoredEmlArtifact] = {}
try:
@@ -196,6 +199,8 @@ def _persist_built_eml_artifacts(
sha256=digest,
message_id_header=str(message_id) if message_id else None,
)
if written_storage_keys is not None:
written_storage_keys.append(storage_key)
message.eml_path = None
message.eml_size_bytes = len(payload)
except Exception:
@@ -207,14 +212,16 @@ def _persist_built_eml_artifacts(
return artifacts
def _delete_storage_keys(storage: StorageBackend, keys: list[str]) -> None:
def _delete_storage_keys(storage: StorageBackend, keys: list[str]) -> list[str]:
failed: list[str] = []
for key in keys:
try:
storage.delete(key)
except StorageBackendError:
# A committed build remains authoritative. Reconciliation can
# remove an orphaned superseded object later.
continue
# Keep the failure observable so the caller can require recovery
# instead of claiming that best-effort deletion compensated it.
failed.append(key)
return failed
def _next_version_number(session: Session, campaign_id: str) -> int:
@@ -1135,6 +1142,7 @@ def _resolve_built_print_outputs(
config: CampaignConfig,
built_messages: list[Any],
entries_by_index: dict[int, Any],
written_storage_keys: list[str] | None = None,
) -> dict[int, dict[str, Any]]:
printable: list[tuple[Any, Any, dict[str, Any]]] = []
for built in built_messages:
@@ -1229,6 +1237,8 @@ def _resolve_built_print_outputs(
raise CampaignPersistenceError(
f"Printable Campaign output could not be persisted: {exc}"
) from exc
if written_storage_keys is not None:
written_storage_keys.append(storage_key)
artifact["storage_key"] = storage_key
artifact["download_path"] = (
f"/api/v1/campaigns/{version.campaign_id}/versions/{version.id}/"
@@ -1471,6 +1481,123 @@ def _apply_campaign_build_state(
campaign.status = CampaignStatus.VALIDATED.value
def _build_storage_expectations(
*,
stored_eml_by_index: dict[int, _StoredEmlArtifact],
print_outputs_by_index: dict[int, dict[str, Any]],
) -> dict[str, dict[str, Any]]:
expectations = {
artifact.storage_key: {
"size_bytes": artifact.size_bytes,
"sha256": artifact.sha256,
"kind": "message/rfc822",
}
for artifact in stored_eml_by_index.values()
}
for output in print_outputs_by_index.values():
artifact = output.get("artifact") if isinstance(output, dict) else None
if not isinstance(artifact, dict) or not artifact.get("storage_key"):
continue
expectations[str(artifact["storage_key"])] = {
"size_bytes": int(output.get("output_size_bytes") or 0),
"sha256": str(output.get("output_sha256") or ""),
"kind": "print-output",
}
return expectations
def _verify_build_storage_manifest(
storage: StorageBackend,
expectations: dict[str, dict[str, Any]],
) -> dict[str, Any]:
manifest: list[dict[str, Any]] = []
total_bytes = 0
for key in sorted(expectations):
expected = expectations[key]
expected_size = int(expected["size_bytes"])
expected_sha256 = str(expected["sha256"])
try:
info = storage.stat(key)
digest = hashlib.sha256()
observed_size = 0
for chunk in storage.iter_bytes(key):
digest.update(chunk)
observed_size += len(chunk)
except StorageBackendError as exc:
raise CampaignPersistenceError(
"A generated Campaign artifact could not be verified in shared storage"
) from exc
observed_sha256 = digest.hexdigest()
if (
info.size_bytes != expected_size
or observed_size != expected_size
or observed_sha256 != expected_sha256
):
raise CampaignPersistenceError(
"A generated Campaign artifact does not match its build evidence"
)
manifest.append(
{
"key_sha256": hashlib.sha256(key.encode("utf-8")).hexdigest(),
"size_bytes": expected_size,
"sha256": expected_sha256,
"kind": expected["kind"],
}
)
total_bytes += expected_size
return {
"object_count": len(manifest),
"total_bytes": total_bytes,
"manifest_sha256": _canonical_sha256(manifest),
}
def _verify_storage_keys_absent(
storage: StorageBackend,
keys: list[str],
) -> tuple[bool, int]:
remaining = 0
try:
for key in sorted(set(keys)):
if storage.exists(key):
remaining += 1
except StorageBackendError:
return False, max(1, remaining)
return remaining == 0, remaining
def _persisted_build_manifest(
session: Session,
*,
version_id: str,
expected_storage_keys: list[str],
) -> dict[str, Any]:
jobs = (
session.query(CampaignJob)
.filter(CampaignJob.campaign_version_id == version_id)
.order_by(CampaignJob.entry_index)
.all()
)
persisted_keys: set[str] = set()
for job in jobs:
if job.eml_storage_key:
persisted_keys.add(str(job.eml_storage_key))
output = job.resolved_print_output
artifact = output.get("artifact") if isinstance(output, dict) else None
if isinstance(artifact, dict) and artifact.get("storage_key"):
persisted_keys.add(str(artifact["storage_key"]))
expected_keys = set(expected_storage_keys)
if persisted_keys != expected_keys:
raise CampaignPersistenceError(
"Persisted Campaign jobs do not match the generated artifact manifest"
)
return {
"job_count": len(jobs),
"referenced_object_count": len(persisted_keys),
"reference_manifest_sha256": _canonical_sha256(sorted(persisted_keys)),
}
def build_campaign_version(
session: Session,
*,
@@ -1479,6 +1606,8 @@ def build_campaign_version(
write_eml: bool = True,
user_id: str | None = None,
principal: ApiPrincipal | None = None,
recovery_operation: DurableRecoveryOperation | None = None,
build_id: str | None = None,
) -> dict[str, Any]:
version, snapshot_path, config = load_version_config(session, version_id)
campaign = session.get(Campaign, version.campaign_id)
@@ -1502,99 +1631,111 @@ def build_campaign_version(
files = files_integration()
storage = _object_storage()
build_id = uuid4().hex
with TemporaryDirectory(prefix="govoplan-campaign-build-") as output_directory:
with files.prepared_campaign_snapshot(
session,
tenant_id=tenant_id,
campaign_id=campaign.id,
raw_json=version.raw_json if isinstance(version.raw_json, dict) else {},
include_bytes=True,
prefix="govoplan-managed-build-",
) as prepared:
managed_raw = load_campaign_json(prepared.path)
managed_config = load_campaign_config_from_json(
session,
tenant_id=tenant_id,
raw_json=managed_raw,
campaign_id=campaign.id,
)
result = build_campaign_messages(
managed_config,
campaign_file=prepared.path,
output_dir=Path(output_directory),
write_eml=write_eml,
)
files.annotate_built_messages_with_managed_files(
result.built_messages,
prepared.managed_files_by_local_path,
)
entries_by_index = {
index: entry
for index, entry in enumerate(
load_campaign_entries(
managed_config,
campaign_file=prepared.path,
),
start=1,
)
}
delivery_provenance_by_index = {
index: dict(entry.distribution_source or {})
for index, entry in entries_by_index.items()
}
resolved_postbox_targets_by_index = _resolve_built_postbox_targets(
session,
tenant_id=tenant_id,
config=managed_config,
built_messages=result.built_messages,
entries_by_index=entries_by_index,
)
resolved_print_outputs_by_index = _resolve_built_print_outputs(
session,
storage=storage,
tenant_id=tenant_id,
build_id=build_id,
version=version,
principal=principal,
config=managed_config,
built_messages=result.built_messages,
entries_by_index=entries_by_index,
)
_prepare_built_calendar_invitations(
version=version,
config=managed_config,
built_messages=result.built_messages,
entries_by_index=entries_by_index,
delivery_provenance_by_index=delivery_provenance_by_index,
user_id=user_id,
)
new_print_storage_keys = sorted(
{
str(artifact["storage_key"])
for output in resolved_print_outputs_by_index.values()
if isinstance(output, dict)
for artifact in [output.get("artifact")]
if isinstance(artifact, dict) and artifact.get("storage_key")
}
effective_build_id = build_id or uuid4().hex
storage_prefix = (
f"campaign-artifacts/{tenant_id}/{campaign.id}/{version.id}/"
f"{effective_build_id}/"
)
written_storage_keys: list[str] = []
new_storage_keys: list[str] = []
old_storage_keys: list[str] = []
domain_committed = False
if recovery_operation is not None:
recovery_operation.checkpoint(
kind="object-prefix-reserved",
summary="The Campaign build reserved its reconciliation prefix",
evidence={
"storage_prefix": storage_prefix,
"campaign_version_id": version.id,
},
)
try:
try:
with TemporaryDirectory(prefix="govoplan-campaign-build-") as output_directory:
with files.prepared_campaign_snapshot(
session,
tenant_id=tenant_id,
campaign_id=campaign.id,
raw_json=version.raw_json
if isinstance(version.raw_json, dict)
else {},
include_bytes=True,
prefix="govoplan-managed-build-",
) as prepared:
managed_raw = load_campaign_json(prepared.path)
managed_config = load_campaign_config_from_json(
session,
tenant_id=tenant_id,
raw_json=managed_raw,
campaign_id=campaign.id,
)
result = build_campaign_messages(
managed_config,
campaign_file=prepared.path,
output_dir=Path(output_directory),
write_eml=write_eml,
)
files.annotate_built_messages_with_managed_files(
result.built_messages,
prepared.managed_files_by_local_path,
)
entries_by_index = {
index: entry
for index, entry in enumerate(
load_campaign_entries(
managed_config,
campaign_file=prepared.path,
),
start=1,
)
}
delivery_provenance_by_index = {
index: dict(entry.distribution_source or {})
for index, entry in entries_by_index.items()
}
resolved_postbox_targets_by_index = _resolve_built_postbox_targets(
session,
tenant_id=tenant_id,
config=managed_config,
built_messages=result.built_messages,
entries_by_index=entries_by_index,
)
resolved_print_outputs_by_index = _resolve_built_print_outputs(
session,
storage=storage,
tenant_id=tenant_id,
build_id=effective_build_id,
version=version,
principal=principal,
config=managed_config,
built_messages=result.built_messages,
entries_by_index=entries_by_index,
written_storage_keys=written_storage_keys,
)
_prepare_built_calendar_invitations(
version=version,
config=managed_config,
built_messages=result.built_messages,
entries_by_index=entries_by_index,
delivery_provenance_by_index=delivery_provenance_by_index,
user_id=user_id,
)
stored_eml_by_index = _persist_built_eml_artifacts(
storage=storage,
tenant_id=tenant_id,
campaign_id=campaign.id,
version_id=version.id,
build_id=build_id,
build_id=effective_build_id,
built_messages=result.built_messages,
written_storage_keys=written_storage_keys,
)
except Exception:
_delete_storage_keys(storage, new_print_storage_keys)
raise
new_storage_keys = [
*new_print_storage_keys,
*(item.storage_key for item in stored_eml_by_index.values()),
]
try:
new_storage_keys = sorted(set(written_storage_keys))
storage_manifest = _verify_build_storage_manifest(
storage,
_build_storage_expectations(
stored_eml_by_index=stored_eml_by_index,
print_outputs_by_index=resolved_print_outputs_by_index,
),
)
report_json = _campaign_build_report(result, files)
report_json["built_by_user_id"] = user_id
if resolved_print_outputs_by_index:
@@ -1667,9 +1808,94 @@ def build_campaign_version(
session.add(version)
session.add(campaign)
session.commit()
except Exception:
session.rollback()
_delete_storage_keys(storage, new_storage_keys)
domain_committed = True
database_manifest = _persisted_build_manifest(
session,
version_id=version.id,
expected_storage_keys=new_storage_keys,
)
superseded_delete_failures = _delete_storage_keys(storage, old_storage_keys)
if recovery_operation is not None:
recovery_operation.checkpoint(
kind="build-manifest-verified",
summary="Generated objects and committed Campaign rows were compared",
evidence={
"storage": storage_manifest,
"database": database_manifest,
"storage_prefix": storage_prefix,
},
)
if superseded_delete_failures:
recovery_operation.unresolved(
status=RecoveryStatus.RECOVERY_REQUIRED,
summary="The new build committed but superseded objects require cleanup",
evidence={
"failed_object_count": len(superseded_delete_failures),
"failed_manifest_sha256": _canonical_sha256(
sorted(superseded_delete_failures)
),
},
failure_summary=(
"Campaign build succeeded, but superseded object cleanup "
"requires reconciliation"
),
)
else:
recovery_operation.succeed(
evidence={
"verified": True,
"checks": {
"storage_manifest": storage_manifest,
"database_manifest": database_manifest,
"superseded_cleanup": "complete",
},
}
)
return report_json
except Exception as exc:
if not domain_committed:
session.rollback()
deletion_failures = _delete_storage_keys(
storage,
sorted(set(written_storage_keys)),
)
absent, remaining_count = _verify_storage_keys_absent(
storage,
written_storage_keys,
)
if recovery_operation is not None and not recovery_operation.closed:
if deletion_failures or not absent:
recovery_operation.unresolved(
status=RecoveryStatus.RECOVERY_REQUIRED,
summary="Campaign build compensation could not be verified",
evidence={
"storage_prefix": storage_prefix,
"written_object_count": len(set(written_storage_keys)),
"delete_failure_count": len(deletion_failures),
"remaining_object_count": remaining_count,
"exception_type": type(exc).__name__,
},
failure_summary=(
"Generated Campaign objects may remain after a failed build"
),
)
else:
recovery_operation.compensate(
failure_summary="Campaign build failed before database commit",
failure_evidence={
"storage_prefix": storage_prefix,
"written_object_count": len(set(written_storage_keys)),
"exception_type": type(exc).__name__,
},
recovery_evidence={
"verified": True,
"checks": {
"database_transaction": "rolled-back",
"written_objects_absent": True,
},
},
)
elif recovery_operation is not None and not recovery_operation.closed:
recovery_operation.release_unresolved()
raise
_delete_storage_keys(storage, old_storage_keys)
return report_json