diff --git a/docs/CAMPAIGN_BUILD_RECOVERY.md b/docs/CAMPAIGN_BUILD_RECOVERY.md new file mode 100644 index 0000000..6291f36 --- /dev/null +++ b/docs/CAMPAIGN_BUILD_RECOVERY.md @@ -0,0 +1,26 @@ +# Campaign Build Recovery + +Campaign build is fenced per tenant and Campaign version. Before rendering or +writing generated artifacts, Campaign commits a Core recovery operation with +the canonical source and validation hashes, the runtime fence, and a reserved +opaque object prefix. A repeated idempotency key can replay only a verified +successful build; it cannot start a second active build. + +Generated EML and bounded print output are written to shared storage and checked +for exact size and SHA-256 content. Campaign then commits its jobs and execution +snapshot, compares the stored object and database manifests, and records +verified success. A managed Files output makes the operation forward-recoverable +because Campaign cannot undo a Files-owned artifact; object-only builds use +explicit compensation. + +If the Campaign database transaction fails, Campaign deletes every object it +recorded and verifies absence before recording recovered state. Failed deletion, +an unavailable storage check, process loss, or superseded-object cleanup failure +leaves a recovery-required operation visible in Ops. Do not retry such an +operation as a normal build. Verify its checkpoint chain and reserved prefix, +then reconcile it through the owning-module procedure. The orphan inventory +reconciler tracked in Campaign issue 91 will automate that bounded inspection. + +No recovery checkpoint contains message bodies, recipients, credentials, or +resolved provider secrets. Object keys remain restricted diagnostics rather +than Campaign business data. diff --git a/src/govoplan_campaign/backend/manifest.py b/src/govoplan_campaign/backend/manifest.py index ce9efed..202a58d 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. Failed builds compensate newly written objects; 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. Object-only failures prove compensation; managed Files output and uncertain cleanup remain explicit forward-recovery or recovery-required work in Ops. Retention retains metadata when deletion fails. Never copy or edit runtime object keys as business data.", layer="evidence", documentation_types=("admin",), audience=("campaign_operator", "platform_operator", "release_reviewer"), @@ -1074,15 +1074,20 @@ manifest = ModuleManifest( href="govoplan-campaign/docs/CAMPAIGN_DELIVERY_RUNBOOK.md", kind="repository", ), + DocumentationLink( + label="Campaign build recovery", + href="govoplan-campaign/docs/CAMPAIGN_BUILD_RECOVERY.md", + kind="repository", + ), ), related_modules=("files", "ops", "mail"), metadata={ "kind": "reference", "route": "/campaigns/queue", "screen": "Campaign operator queue", - "verification": "Build on one replica, deliver from another, verify object digest/size evidence, and exercise storage failure during build and retention.", + "verification": "Build on one replica, verify the recovery checkpoint chain and database/object manifests, deliver from another, then exercise process loss, stale fencing, storage failure, and retention cleanup.", "limitations": [ - "A hard process loss between object creation and database commit can leave an orphan object until an inventory reconciler removes it.", + "A hard process loss between object creation and database commit remains recovery-required under its reserved prefix until the inventory reconciler verifies or removes it.", "Database, object storage, and encryption keys require a coordinated deployment backup and restore procedure.", ], }, diff --git a/src/govoplan_campaign/backend/persistence/campaigns.py b/src/govoplan_campaign/backend/persistence/campaigns.py index bf6b1f3..1412201 100644 --- a/src/govoplan_campaign/backend/persistence/campaigns.py +++ b/src/govoplan_campaign/backend/persistence/campaigns.py @@ -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 diff --git a/src/govoplan_campaign/backend/routes/versions.py b/src/govoplan_campaign/backend/routes/versions.py index ac9304c..03659e9 100644 --- a/src/govoplan_campaign/backend/routes/versions.py +++ b/src/govoplan_campaign/backend/routes/versions.py @@ -1,9 +1,10 @@ from __future__ import annotations import hashlib +import json from urllib.parse import quote -from fastapi import APIRouter, Depends, Header, HTTPException, Query, Response, status +from fastapi import APIRouter, Depends, Header, HTTPException, Query, Request, Response, status from sqlalchemy.orm import Session from govoplan_campaign.backend.schemas import ( @@ -22,10 +23,21 @@ from govoplan_campaign.backend.schemas import ( from govoplan_core.auth import ApiPrincipal, has_scope, require_scope from govoplan_core.audit.logging import audit_from_principal from govoplan_core.core.object_storage import StorageBackendError +from govoplan_core.core.recovery import ( + RecoveryGuaranteeError, + RecoveryMode, + RecoveryPlan, +) +from govoplan_core.core.recovery_runtime import ( + RecoveryOperationBusy, + RecoveryOperationStateConflict, + begin_durable_recovery_operation, +) from govoplan_campaign.backend.db.models import ( CampaignVersion, ) -from govoplan_core.db.session import get_session +from govoplan_core.db.session import get_database, get_session +from govoplan_core.server.runtime_agent import application_runtime_identity from govoplan_campaign.backend.response_security import ( public_campaign_payload, ) @@ -74,6 +86,47 @@ from govoplan_campaign.backend.routes.attachments import ( router = APIRouter(prefix="/campaigns", tags=["campaigns"]) +def _canonical_sha256(value: object) -> str: + return hashlib.sha256( + json.dumps( + value, + ensure_ascii=False, + sort_keys=True, + separators=(",", ":"), + default=str, + ).encode("utf-8") + ).hexdigest() + + +def _campaign_build_recovery_plan(raw_json: dict[str, object]) -> RecoveryPlan: + delivery = raw_json.get("delivery") + print_config = delivery.get("print") if isinstance(delivery, dict) else None + persists_managed_output = bool( + isinstance(print_config, dict) + and print_config.get("persist_to_files", True) + ) + if persists_managed_output: + return RecoveryPlan( + mode=RecoveryMode.FORWARD_RECOVERY, + preconditions=("validated Campaign version is locked",), + forward_recovery_steps=( + "reuse the Templates render idempotency key", + "reconcile the managed output and Campaign build manifests", + ), + verification_steps=( + "compare committed Campaign jobs with generated object evidence", + ), + ) + return RecoveryPlan( + mode=RecoveryMode.COMPENSATION, + preconditions=("validated Campaign version is locked",), + compensation_steps=("delete every object under the reserved build prefix",), + verification_steps=( + "compare committed Campaign jobs with generated object evidence", + ), + ) + + @router.get("/{campaign_id}/versions/{version_id}/print-output/download") def download_print_output( campaign_id: str, @@ -728,6 +781,7 @@ def validate_version( @router.post("/versions/{version_id}/build") def build_version( version_id: str, + request: Request, payload: BuildCampaignRequest | None = None, session: Session = Depends(get_session), principal: ApiPrincipal = Depends(require_scope("campaigns:campaign:build")), @@ -737,14 +791,80 @@ def build_version( _require_mail_profile_use_if_needed( principal, version.raw_json if isinstance(version.raw_json, dict) else {} ) + raw_json = version.raw_json if isinstance(version.raw_json, dict) else {} + write_eml = payload.write_eml if payload else True + source_sha256 = _canonical_sha256(raw_json) + validation_sha256 = _canonical_sha256( + version.validation_summary + if isinstance(version.validation_summary, dict) + else {} + ) + idempotency_key = ( + payload.idempotency_key + if payload and payload.idempotency_key + else ( + f"campaign-build:{version.id}:{source_sha256}:" + f"{validation_sha256}:{int(write_eml)}" + ) + ) try: + identity = application_runtime_identity(request.app) + recovery_start = begin_durable_recovery_operation( + get_database().SessionLocal, + identity=identity, + module_id="campaigns", + operation_type="build-artifacts", + idempotency_key=idempotency_key, + request={ + "tenant_id": principal.tenant_id, + "campaign_id": version.campaign_id, + "version_id": version.id, + "source_sha256": source_sha256, + "validation_sha256": validation_sha256, + "write_eml": write_eml, + }, + recovery_plan=_campaign_build_recovery_plan(raw_json), + precondition_evidence={ + "campaign_version_id": version.id, + "source_sha256": source_sha256, + "validation_sha256": validation_sha256, + "locked_at": version.locked_at.isoformat() + if version.locked_at + else None, + "workflow_state": version.workflow_state, + }, + lease_resource_key=( + f"campaign:build:{principal.tenant_id}:{version.id}" + ), + lease_ttl_seconds=30 * 60, + resource_type="campaign_version", + resource_id=version.id, + metadata={"actor_account_id": principal.account_id}, + ) + if recovery_start.replayed: + session.refresh(version) + if not isinstance(version.build_summary, dict): + raise RecoveryGuaranteeError( + "A completed build operation has no committed Campaign summary" + ) + return public_campaign_payload( + version.build_summary, + include_diagnostics=has_scope( + principal, "campaigns:diagnostic:read" + ), + ) + recovery_operation = recovery_start.operation + if recovery_operation is None: # pragma: no cover - guarded by replay branch + raise RecoveryGuaranteeError("Campaign build authority was not created") result = build_campaign_version( session, tenant_id=principal.tenant_id, version_id=version_id, - write_eml=payload.write_eml if payload else True, + write_eml=write_eml, user_id=principal.user.id, principal=principal, + recovery_operation=recovery_operation, + build_id=recovery_start.operation_id, ) audit_from_principal( session, @@ -753,8 +873,9 @@ def build_version( object_type="campaign_version", object_id=version_id, details={ - "write_eml": payload.write_eml if payload else True, + "write_eml": write_eml, "built_count": result.get("built_count"), + "recovery_operation_id": recovery_start.operation_id, }, commit=True, ) @@ -766,6 +887,21 @@ def build_version( raise HTTPException( status_code=status.HTTP_404_NOT_FOUND, detail=str(exc) ) from exc + except (RecoveryOperationBusy, RecoveryOperationStateConflict) as exc: + raise HTTPException( + status_code=status.HTTP_409_CONFLICT, + detail=str(exc), + ) from exc + except RecoveryGuaranteeError as exc: + raise HTTPException( + status_code=status.HTTP_503_SERVICE_UNAVAILABLE, + detail=str(exc), + ) from exc + except RuntimeError as exc: + raise HTTPException( + status_code=status.HTTP_503_SERVICE_UNAVAILABLE, + detail="Campaign build coordination is unavailable", + ) from exc except Exception as exc: raise HTTPException( status_code=status.HTTP_422_UNPROCESSABLE_CONTENT, detail=str(exc) diff --git a/src/govoplan_campaign/backend/schemas.py b/src/govoplan_campaign/backend/schemas.py index 5e56bb8..b1a6a81 100644 --- a/src/govoplan_campaign/backend/schemas.py +++ b/src/govoplan_campaign/backend/schemas.py @@ -677,6 +677,7 @@ class BuildCampaignRequest(BaseModel): model_config = ConfigDict(extra="forbid") write_eml: bool = True + idempotency_key: str | None = Field(default=None, min_length=1, max_length=200) class ApiKeyCreateRequest(BaseModel): diff --git a/tests/test_campaign_build_recovery.py b/tests/test_campaign_build_recovery.py new file mode 100644 index 0000000..6a8ae60 --- /dev/null +++ b/tests/test_campaign_build_recovery.py @@ -0,0 +1,94 @@ +from __future__ import annotations + +import hashlib +from pathlib import Path + +import pytest + +from govoplan_core.core.object_storage import LocalFilesystemStorageBackend +from govoplan_core.core.recovery import RecoveryMode +from govoplan_campaign.backend.persistence.campaigns import ( + CampaignPersistenceError, + _StoredEmlArtifact, + _build_storage_expectations, + _delete_storage_keys, + _verify_build_storage_manifest, + _verify_storage_keys_absent, +) +from govoplan_campaign.backend.routes.versions import _campaign_build_recovery_plan + + +def test_object_only_and_managed_output_builds_use_distinct_recovery_modes() -> None: + assert _campaign_build_recovery_plan({}).mode == RecoveryMode.COMPENSATION + assert ( + _campaign_build_recovery_plan( + {"delivery": {"print": {"persist_to_files": True}}} + ).mode + == RecoveryMode.FORWARD_RECOVERY + ) + assert ( + _campaign_build_recovery_plan( + {"delivery": {"print": {"persist_to_files": False}}} + ).mode + == RecoveryMode.COMPENSATION + ) + + +def test_generated_object_manifest_verifies_exact_bytes(tmp_path: Path) -> None: + storage = LocalFilesystemStorageBackend(tmp_path) + payload = b"Message-ID: \r\n\r\nbody" + key = "campaign-artifacts/tenant/campaign/version/build/00000001.eml" + storage.put_bytes(key, payload) + artifact = _StoredEmlArtifact( + storage_key=key, + size_bytes=len(payload), + sha256=hashlib.sha256(payload).hexdigest(), + message_id_header="", + ) + + evidence = _verify_build_storage_manifest( + storage, + _build_storage_expectations( + stored_eml_by_index={1: artifact}, + print_outputs_by_index={}, + ), + ) + + assert evidence["object_count"] == 1 + assert evidence["total_bytes"] == len(payload) + storage.put_bytes(key, b"tampered") + with pytest.raises(CampaignPersistenceError, match="does not match"): + _verify_build_storage_manifest( + storage, + _build_storage_expectations( + stored_eml_by_index={1: artifact}, + print_outputs_by_index={}, + ), + ) + + +def test_compensation_requires_verified_object_absence(tmp_path: Path) -> None: + storage = LocalFilesystemStorageBackend(tmp_path) + key = "campaign-artifacts/tenant/campaign/version/build/object" + storage.put_bytes(key, b"payload") + assert _delete_storage_keys(storage, [key]) == [] + assert _verify_storage_keys_absent(storage, [key]) == (True, 0) + + +class _UnremovableStorage: + name = "unremovable" + + def delete(self, _key: str) -> None: + from govoplan_core.core.object_storage import StorageBackendError + + raise StorageBackendError("unavailable") + + def exists(self, _key: str) -> bool: + return True + + +def test_failed_compensation_remains_observable() -> None: + storage = _UnremovableStorage() + key = "campaign-artifacts/tenant/campaign/version/build/object" + assert _delete_storage_keys(storage, [key]) == [key] # type: ignore[arg-type] + assert _verify_storage_keys_absent(storage, [key]) == (False, 1) # type: ignore[arg-type] diff --git a/tests/test_imap_append_idempotency.py b/tests/test_imap_append_idempotency.py index 24c7ba4..45e9687 100644 --- a/tests/test_imap_append_idempotency.py +++ b/tests/test_imap_append_idempotency.py @@ -289,6 +289,7 @@ def test_imap_reconciliation_preserves_attempt_and_only_not_appended_is_retryabl imap_claimed_at=datetime.now(timezone.utc), imap_claim_token="claim-1", last_error="unknown", + delivery_provenance={}, ) attempt = SimpleNamespace( id="attempt-1", diff --git a/tests/test_route_registration.py b/tests/test_route_registration.py index a0e06d8..34b8357 100644 --- a/tests/test_route_registration.py +++ b/tests/test_route_registration.py @@ -38,7 +38,7 @@ def test_campaign_router_composes_every_workflow_operation_once() -> None: actual = _operation_keys(router) assert actual == expected - assert len(actual) == 70 + assert len(actual) == 71 assert not [operation for operation, count in Counter(actual).items() if count > 1] diff --git a/webui/src/api/campaigns.ts b/webui/src/api/campaigns.ts index bac075e..9186647 100644 --- a/webui/src/api/campaigns.ts +++ b/webui/src/api/campaigns.ts @@ -1209,7 +1209,7 @@ writeEml = true) : Promise> { return apiFetch>(settings, `/api/v1/campaigns/versions/${versionId}/build`, { method: "POST", - body: JSON.stringify({ write_eml: writeEml }) + body: JSON.stringify({ write_eml: writeEml, idempotency_key: crypto.randomUUID() }) }); }