Commit verified recovery projections atomically
This commit is contained in:
@@ -738,6 +738,13 @@ effects, transitions partial/unknown outcomes honestly, and records verified
|
|||||||
completion or recovery. Plaintext secrets must never enter recovery metadata or
|
completion or recovery. Plaintext secrets must never enter recovery metadata or
|
||||||
evidence.
|
evidence.
|
||||||
|
|
||||||
|
For a conclusive external result, modules may commit their local success
|
||||||
|
projection and the verified terminal checkpoint in one database transaction via
|
||||||
|
`DurableRecoveryOperation.commit_verified_success`. This does not make the
|
||||||
|
external provider effect atomic. It prevents a local `succeeded` state from
|
||||||
|
becoming authoritative when the recovery evidence chain is damaged or the
|
||||||
|
terminal checkpoint cannot commit.
|
||||||
|
|
||||||
## Install, Uninstall, And Catalogs
|
## Install, Uninstall, And Catalogs
|
||||||
|
|
||||||
Core owns the install plan, signed catalog validation, license entitlement
|
Core owns the install plan, signed catalog validation, license entitlement
|
||||||
|
|||||||
@@ -113,12 +113,35 @@ class DurableRecoveryOperation:
|
|||||||
) -> None:
|
) -> None:
|
||||||
"""Commit domain writes and verified success in one DB transaction."""
|
"""Commit domain writes and verified success in one DB transaction."""
|
||||||
|
|
||||||
self._commit_atomic_terminal(
|
self._commit_terminal(
|
||||||
session,
|
session,
|
||||||
status=RecoveryStatus.SUCCEEDED,
|
status=RecoveryStatus.SUCCEEDED,
|
||||||
summary="Operation effects and authoritative state were verified",
|
summary="Operation effects and authoritative state were verified",
|
||||||
kind="verified-success",
|
kind="verified-success",
|
||||||
evidence=evidence,
|
evidence=evidence,
|
||||||
|
require_atomic_mode=True,
|
||||||
|
)
|
||||||
|
|
||||||
|
def commit_verified_success(
|
||||||
|
self,
|
||||||
|
session: Session,
|
||||||
|
*,
|
||||||
|
evidence: dict[str, Any],
|
||||||
|
) -> None:
|
||||||
|
"""Commit a verified success projection and checkpoint together.
|
||||||
|
|
||||||
|
Non-atomic operations use this only after their external effect has a
|
||||||
|
conclusive provider result. It does not make that effect atomic; it
|
||||||
|
prevents local success from outrunning its durable verification.
|
||||||
|
"""
|
||||||
|
|
||||||
|
self._commit_terminal(
|
||||||
|
session,
|
||||||
|
status=RecoveryStatus.SUCCEEDED,
|
||||||
|
summary="Operation effects and authoritative state were verified",
|
||||||
|
kind="verified-success",
|
||||||
|
evidence=evidence,
|
||||||
|
require_atomic_mode=False,
|
||||||
)
|
)
|
||||||
|
|
||||||
def commit_atomic_failure(
|
def commit_atomic_failure(
|
||||||
@@ -130,12 +153,13 @@ class DurableRecoveryOperation:
|
|||||||
) -> None:
|
) -> None:
|
||||||
"""Commit domain failure evidence and the terminal state atomically."""
|
"""Commit domain failure evidence and the terminal state atomically."""
|
||||||
|
|
||||||
self._commit_atomic_terminal(
|
self._commit_terminal(
|
||||||
session,
|
session,
|
||||||
status=RecoveryStatus.FAILED,
|
status=RecoveryStatus.FAILED,
|
||||||
summary=summary,
|
summary=summary,
|
||||||
kind="verified-failure",
|
kind="verified-failure",
|
||||||
evidence=evidence,
|
evidence=evidence,
|
||||||
|
require_atomic_mode=True,
|
||||||
)
|
)
|
||||||
|
|
||||||
def commit_atomic_rejection(
|
def commit_atomic_rejection(
|
||||||
@@ -147,12 +171,13 @@ class DurableRecoveryOperation:
|
|||||||
) -> None:
|
) -> None:
|
||||||
"""Commit a definitive rejection and its domain evidence atomically."""
|
"""Commit a definitive rejection and its domain evidence atomically."""
|
||||||
|
|
||||||
self._commit_atomic_terminal(
|
self._commit_terminal(
|
||||||
session,
|
session,
|
||||||
status=RecoveryStatus.REJECTED,
|
status=RecoveryStatus.REJECTED,
|
||||||
summary=summary,
|
summary=summary,
|
||||||
kind="verified-rejection",
|
kind="verified-rejection",
|
||||||
evidence=evidence,
|
evidence=evidence,
|
||||||
|
require_atomic_mode=True,
|
||||||
)
|
)
|
||||||
|
|
||||||
def fail(self, *, summary: str, evidence: dict[str, Any]) -> None:
|
def fail(self, *, summary: str, evidence: dict[str, Any]) -> None:
|
||||||
@@ -375,7 +400,7 @@ class DurableRecoveryOperation:
|
|||||||
self.lease_claim = claim
|
self.lease_claim = claim
|
||||||
return operation, claim
|
return operation, claim
|
||||||
|
|
||||||
def _commit_atomic_terminal(
|
def _commit_terminal(
|
||||||
self,
|
self,
|
||||||
session: Session,
|
session: Session,
|
||||||
*,
|
*,
|
||||||
@@ -383,16 +408,20 @@ class DurableRecoveryOperation:
|
|||||||
summary: str,
|
summary: str,
|
||||||
kind: str,
|
kind: str,
|
||||||
evidence: dict[str, Any],
|
evidence: dict[str, Any],
|
||||||
|
require_atomic_mode: bool,
|
||||||
) -> None:
|
) -> None:
|
||||||
if status not in {
|
if status not in {
|
||||||
RecoveryStatus.SUCCEEDED,
|
RecoveryStatus.SUCCEEDED,
|
||||||
RecoveryStatus.FAILED,
|
RecoveryStatus.FAILED,
|
||||||
RecoveryStatus.REJECTED,
|
RecoveryStatus.REJECTED,
|
||||||
}:
|
}:
|
||||||
raise ValueError("Unsupported atomic terminal recovery status")
|
raise ValueError("Unsupported terminal recovery status")
|
||||||
try:
|
try:
|
||||||
operation, claim = self._locked_and_renewed(session)
|
operation, claim = self._locked_and_renewed(session)
|
||||||
if operation.mode != RecoveryMode.ATOMIC.value:
|
if (
|
||||||
|
require_atomic_mode
|
||||||
|
and operation.mode != RecoveryMode.ATOMIC.value
|
||||||
|
):
|
||||||
raise RecoveryGuaranteeError(
|
raise RecoveryGuaranteeError(
|
||||||
"Atomic terminal commits require an atomic recovery plan"
|
"Atomic terminal commits require an atomic recovery plan"
|
||||||
)
|
)
|
||||||
|
|||||||
@@ -153,6 +153,42 @@ def test_atomic_terminal_commits_domain_rows_and_recovery_evidence_together() ->
|
|||||||
engine.dispose()
|
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:
|
def test_failed_atomic_commit_rolls_back_domain_and_terminal_checkpoint() -> None:
|
||||||
engine, factory = _fixture()
|
engine, factory = _fixture()
|
||||||
metadata = MetaData()
|
metadata = MetaData()
|
||||||
|
|||||||
Reference in New Issue
Block a user