from __future__ import annotations from datetime import datetime, timedelta, timezone from unittest.mock import patch from sqlalchemy import Column, MetaData, String, Table, create_engine, select from sqlalchemy.orm import sessionmaker import pytest from govoplan_core.core.recovery import ( RecoveryCheckpoint, RecoveryGuaranteeError, RecoveryMode, RecoveryOperation, RecoveryPlan, RecoveryStatus, verify_recovery_evidence_chain, ) from govoplan_core.core.recovery_runtime import ( RecoveryOperationBusy, RecoveryOperationStateConflict, begin_durable_recovery_operation, claim_durable_recovery_operation, ) from govoplan_core.core.runtime_coordination import DistributedLease, RuntimeIdentity from govoplan_core.db.base import Base def _fixture(): engine = create_engine("sqlite+pysqlite:///:memory:") Base.metadata.create_all( engine, tables=[ DistributedLease.__table__, RecoveryOperation.__table__, RecoveryCheckpoint.__table__, ], ) return engine, sessionmaker(bind=engine, expire_on_commit=False) def _identity(node: str, incarnation: str) -> RuntimeIdentity: return RuntimeIdentity( installation_id="installation-1", node_id=node, incarnation=incarnation, role="worker", software_version="test", composition_hash="a" * 64, ) def _start(factory, identity, *, key: str = "build-1"): return begin_durable_recovery_operation( factory, identity=identity, module_id="campaigns", operation_type="build-artifacts", idempotency_key=key, request={"version_id": "version-1", "write_eml": True}, recovery_plan=RecoveryPlan( mode=RecoveryMode.COMPENSATION, preconditions=("validated version is locked",), compensation_steps=("delete build object prefix",), verification_steps=("compare database and object manifests",), ), precondition_evidence={"validation_sha256": "b" * 64}, lease_resource_key="campaign:build:version-1", resource_type="campaign_version", resource_id="version-1", ) def _start_atomic(factory, identity, *, key: str = "sync-1"): return begin_durable_recovery_operation( factory, identity=identity, module_id="connectors", operation_type="read-snapshot", idempotency_key=key, request={"provider_id": "provider-1", "cursor": "revision-1"}, recovery_plan=RecoveryPlan( mode=RecoveryMode.ATOMIC, preconditions=("the provider read is non-mutating",), verification_steps=("compare the committed projection",), ), precondition_evidence={"provider_mutation": False}, lease_resource_key="connectors:provider-1", ) def test_durable_operation_commits_before_caller_effect_and_replays_success() -> None: engine, factory = _fixture() try: started = _start(factory, _identity("worker-1", "incarnation-1")) assert started.operation is not None with factory() as session: persisted = session.get(RecoveryOperation, started.operation_id) assert persisted is not None assert persisted.status == RecoveryStatus.RUNNING.value assert persisted.checkpoint_count == 3 started.operation.checkpoint( kind="object-prefix-reserved", summary="Build object prefix reserved", evidence={"prefix": "campaign-artifacts/build-1/"}, ) started.operation.succeed( evidence={ "verified": True, "checks": {"database_manifest": "matched", "object_manifest": "matched"}, } ) replay = _start(factory, _identity("worker-2", "incarnation-2")) assert replay.replayed is True assert replay.operation is None with factory() as session: assert verify_recovery_evidence_chain(session, started.operation_id) finally: engine.dispose() def test_atomic_terminal_commits_domain_rows_and_recovery_evidence_together() -> None: engine, factory = _fixture() metadata = MetaData() projection = Table( "test_recovery_projection", metadata, Column("id", String(36), primary_key=True), ) metadata.create_all(engine) try: started = _start_atomic(factory, _identity("worker-1", "incarnation-1")) assert started.operation is not None with factory() as session: session.execute(projection.insert().values(id="projection-1")) started.operation.commit_atomic_success( session, evidence={ "verified": True, "checks": {"projection_id": "projection-1"}, }, ) with factory() as session: assert session.scalar(select(projection.c.id)) == "projection-1" operation = session.get(RecoveryOperation, started.operation_id) assert operation is not None assert operation.status == RecoveryStatus.SUCCEEDED.value assert verify_recovery_evidence_chain(session, operation.id) finally: engine.dispose() def test_failed_atomic_commit_rolls_back_domain_and_terminal_checkpoint() -> None: engine, factory = _fixture() metadata = MetaData() projection = Table( "test_recovery_projection_rollback", metadata, Column("id", String(36), primary_key=True), ) metadata.create_all(engine) try: started = _start_atomic(factory, _identity("worker-1", "incarnation-1")) assert started.operation is not None with factory() as session: session.execute(projection.insert().values(id="rolled-back")) with ( patch.object(session, "commit", side_effect=RuntimeError("commit failed")), pytest.raises(RuntimeError, match="commit failed"), ): started.operation.commit_atomic_success( session, evidence={ "verified": True, "checks": {"projection_id": "rolled-back"}, }, ) with factory() as session: assert session.execute(select(projection.c.id)).all() == [] operation = session.get(RecoveryOperation, started.operation_id) assert operation is not None assert operation.status == RecoveryStatus.RUNNING.value assert verify_recovery_evidence_chain(session, operation.id) finally: engine.dispose() def test_same_fence_cannot_start_duplicate_running_operation() -> None: engine, factory = _fixture() try: identity = _identity("worker-1", "incarnation-1") started = _start(factory, identity) with pytest.raises(RecoveryOperationStateConflict, match="already running"): _start(factory, identity) assert started.operation is not None started.operation.release_unresolved() finally: engine.dispose() def test_verified_provider_rejection_is_terminal_without_recovery() -> None: engine, factory = _fixture() try: started = _start(factory, _identity("worker-1", "incarnation-1")) assert started.operation is not None started.operation.reject( summary="Provider definitively rejected the request", evidence={ "verified": True, "provider_outcome": "rejected", "checks": {"provider_response": "definitive-rejection"}, }, ) with factory() as session: operation = session.get(RecoveryOperation, started.operation_id) assert operation is not None assert operation.status == RecoveryStatus.REJECTED.value assert operation.completed_at is not None assert verify_recovery_evidence_chain(session, operation.id) finally: engine.dispose() def test_other_runtime_cannot_use_an_active_fence() -> None: engine, factory = _fixture() try: started = _start(factory, _identity("worker-1", "incarnation-1")) with pytest.raises(RecoveryOperationBusy): _start(factory, _identity("worker-2", "incarnation-2"), key="build-2") assert started.operation is not None started.operation.release_unresolved() finally: engine.dispose() def test_expired_crash_fence_is_taken_over_as_recovery_required() -> None: engine, factory = _fixture() try: started = _start(factory, _identity("worker-1", "incarnation-1")) assert started.operation is not None with factory() as session: lease = session.execute(select(DistributedLease)).scalar_one() lease.expires_at = datetime.now(timezone.utc) - timedelta(seconds=1) session.add(lease) session.commit() recovery = claim_durable_recovery_operation( factory, identity=_identity("worker-2", "incarnation-2"), operation_id=started.operation_id, ) with factory() as session: operation = session.get(RecoveryOperation, started.operation_id) assert operation is not None assert operation.status == RecoveryStatus.RECOVERY_REQUIRED.value assert operation.fencing_token == 2 assert verify_recovery_evidence_chain(session, operation.id) recovery.compensate( failure_summary="worker stopped during object publication", failure_evidence={"object_prefix": "campaign-artifacts/build-1/"}, recovery_evidence={ "verified": True, "checks": {"object_prefix_empty": True}, }, ) finally: engine.dispose() def test_tampered_checkpoint_blocks_verified_success() -> None: engine, factory = _fixture() try: started = _start(factory, _identity("worker-1", "incarnation-1")) assert started.operation is not None with factory() as session: checkpoint = session.execute( select(RecoveryCheckpoint).order_by(RecoveryCheckpoint.sequence) ).scalars().first() assert checkpoint is not None checkpoint.summary = "tampered" session.add(checkpoint) session.commit() with pytest.raises(RecoveryGuaranteeError, match="chain verification failed"): started.operation.succeed( evidence={"verified": True, "checks": {"objects": "matched"}} ) finally: engine.dispose() @pytest.mark.parametrize( ("effect_occurred", "expected_status"), [ (True, RecoveryStatus.SUCCEEDED.value), (False, RecoveryStatus.RECOVERED.value), ], ) def test_unknown_provider_outcome_can_be_resolved_from_external_evidence( effect_occurred: bool, expected_status: str, ) -> None: engine, factory = _fixture() try: started = _start(factory, _identity("worker-1", "incarnation-1")) assert started.operation is not None started.operation.unresolved( status=RecoveryStatus.OUTCOME_UNKNOWN, summary="Provider outcome is unknown", evidence={"effect_started": True}, failure_summary="Inspect the provider before retrying", ) recovery = claim_durable_recovery_operation( factory, identity=_identity("worker-2", "incarnation-2"), operation_id=started.operation_id, ) recovery.resolve_unknown( effect_occurred=effect_occurred, summary="Operator verified the provider outcome", evidence={ "verified": True, "checks": {"provider_evidence": "case-1"}, "effect_occurred": effect_occurred, "reference": "case-1", }, ) with factory() as session: operation = session.get(RecoveryOperation, started.operation_id) assert operation is not None assert operation.status == expected_status assert verify_recovery_evidence_chain(session, operation.id) finally: engine.dispose()