Files
govoplan-core/tests/test_recovery_runtime.py
T

421 lines
15 KiB
Python

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_verified_external_success_commits_projection_and_evidence_together() -> None:
engine, factory = _fixture()
metadata = MetaData()
projection = Table(
"test_verified_external_projection",
metadata,
Column("id", String(36), primary_key=True),
)
metadata.create_all(engine)
try:
started = _start(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_verified_success(
session,
evidence={
"verified": True,
"checks": {
"provider_result": "accepted",
"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.mode == RecoveryMode.COMPENSATION.value
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()
def test_unknown_resolution_commits_domain_projection_and_evidence_together() -> None:
engine, factory = _fixture()
metadata = MetaData()
projection = Table(
"test_unknown_resolution_projection",
metadata,
Column("id", String(36), primary_key=True),
)
metadata.create_all(engine)
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,
)
with factory() as session:
session.execute(projection.insert().values(id="confirmed-effect"))
recovery.commit_unknown_resolution(
session,
effect_occurred=True,
summary="Operator verified the provider outcome",
evidence={
"verified": True,
"checks": {"provider_evidence": "case-1"},
"effect_occurred": True,
},
)
with factory() as session:
assert session.scalar(select(projection.c.id)) == "confirmed-effect"
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()