Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
332 changes: 330 additions & 2 deletions continuityos/gate/monotonic_anchor.py
Original file line number Diff line number Diff line change
Expand Up @@ -17,6 +17,7 @@
DOMAIN = "continuityos.governance.execution-anchor.v1"
COMMITMENT_SCHEMA = "continuityos.governance.execution-anchor.commitment.v1"
RECEIPT_SCHEMA = "continuityos.governance.execution-anchor.receipt.v1"
RECONCILIATION_PLAN_SCHEMA = "continuityos.governance.execution-anchor.reconciliation-plan.v1"
EVENT_KIND = "monotonic_anchor"
PHASE_GENESIS = "GENESIS"
PHASE_STARTED = "EXECUTION_STARTED"
Expand Down Expand Up @@ -247,7 +248,8 @@ def _event_payload(row) -> dict[str, Any]:
return value

def validate_global(
self, ledger, *, _allow_pending_boundary: tuple[str, str] | None = None
self, ledger, *, _allow_pending_boundary: tuple[str, str] | None = None,
_allow_hardware_frontier_mismatch: bool = False,
) -> dict[str, Any]:
verification = ledger.verify()
if not verification.get("ok"):
Expand Down Expand Up @@ -352,7 +354,10 @@ def validate_global(
raise MonotonicAnchorError(
"monotonic anchor hardware identity substitution detected"
)
if provider["observed_digest"] != latest["observed_anchor_digest"]:
if (
provider["observed_digest"] != latest["observed_anchor_digest"]
and not _allow_hardware_frontier_mismatch
):
raise MonotonicAnchorError(
"local governance state differs from monotonic hardware frontier"
)
Expand Down Expand Up @@ -668,6 +673,329 @@ def _record_execution_terminal_locked(
)
return self._extend_and_persist(ledger, receipt)

@staticmethod
def _reconciliation_plan_hash(plan: dict[str, Any]) -> str:
body = dict(plan)
body.pop("plan_hash", None)
return hashlib.sha256(_canonical(body)).hexdigest()

@classmethod
def _finalize_reconciliation_plan(cls, plan: dict[str, Any]) -> dict[str, Any]:
result = dict(plan)
result["plan_hash"] = cls._reconciliation_plan_hash(result)
return result

def plan_offline_reconciliation(
self, ledger, *, preflight_hash: str
) -> dict[str, Any]:
"""Read-only deterministic recovery plan under the external witness lock."""
with self._require_r14_witness(ledger).locked():
return self._plan_offline_reconciliation_locked(
ledger, preflight_hash=preflight_hash
)

def _plan_offline_reconciliation_locked(
self, ledger, *, preflight_hash: str
) -> dict[str, Any]:
if not _is_nonzero_sha256(preflight_hash):
raise MonotonicAnchorError("reconciliation preflight hash is invalid")
row = ledger._attempt_row(preflight_hash)
if row is None:
raise MonotonicAnchorError("reconciliation execution attempt is missing")
try:
attempt = ledger._validate_attempt_row(row)
except Exception as exc:
raise MonotonicAnchorError(
"reconciliation execution attempt is invalid"
) from exc

phase = attempt["phase"]
if phase not in {"ATTEMPT_STARTED", "TERMINAL"}:
raise MonotonicAnchorError(
"execution attempt has no monotonic boundary to reconcile"
)

if phase == "ATTEMPT_STARTED":
boundary_phase = PHASE_STARTED
boundary_hash = attempt.get("execution_started_hash")
boundary_kind = "execution_started"
terminal_kind = None
else:
if not attempt.get("execution_started_hash"):
state = self.validate_global(
ledger, _allow_hardware_frontier_mismatch=True
)
hardware = self._read_provider()
action = (
"ALREADY_RECONCILED"
if hardware["observed_digest"] == state["observed_anchor_digest"]
else "HOLD_DIVERGED"
)
return self._finalize_reconciliation_plan({
"schema": RECONCILIATION_PLAN_SCHEMA,
"domain": DOMAIN,
"state_id": state["state_id"],
"preflight_hash": preflight_hash,
"binding_sha256": attempt["binding_sha256"],
"attempt_phase": phase,
"boundary_phase": None,
"boundary_kind": None,
"boundary_hash": None,
"terminal_kind": attempt.get("terminal_kind"),
"action": action,
"receipt": None,
"local_anchor_generation": state["anchor_generation"],
"local_anchor_digest": state["observed_anchor_digest"],
"hardware_observed_digest": hardware["observed_digest"],
"note": "terminal attempt crossed no subprocess boundary",
})
boundary_phase = PHASE_TERMINAL
boundary_hash = attempt.get("terminal_hash")
boundary_kind = attempt.get("terminal_kind")
terminal_kind = attempt.get("terminal_kind")

if not _is_nonzero_sha256(boundary_hash):
raise MonotonicAnchorError("reconciliation boundary hash is invalid")

# First validate all durable local evidence while allowing the one exact
# boundary to be pending and allowing hardware to be one deterministic
# digest ahead. This remains read-only.
state = self.validate_global(
ledger,
_allow_pending_boundary=(boundary_phase, boundary_hash),
_allow_hardware_frontier_mismatch=True,
)
hardware = self._read_provider()
anchored = state["request_receipts"].get(preflight_hash) or {}
existing = (
anchored.get("started")
if boundary_phase == PHASE_STARTED
else anchored.get("terminal")
)

if existing is not None:
if hardware["observed_digest"] != state["observed_anchor_digest"]:
action = "HOLD_DIVERGED"
note = "hardware differs from the latest durable monotonic receipt"
elif (
boundary_phase == PHASE_STARTED
and attempt["phase"] == "ATTEMPT_STARTED"
and anchored.get("terminal") is None
):
action = "OUTCOME_EVIDENCE_REQUIRED"
note = (
"started receipt is durable but no verified terminal outcome exists; "
"do not infer or fabricate subprocess outcome"
)
else:
action = "ALREADY_RECONCILED"
note = "exact monotonic boundary receipt is already durable"
return self._finalize_reconciliation_plan({
"schema": RECONCILIATION_PLAN_SCHEMA,
"domain": DOMAIN,
"state_id": state["state_id"],
"preflight_hash": preflight_hash,
"binding_sha256": attempt["binding_sha256"],
"attempt_phase": attempt["phase"],
"boundary_phase": boundary_phase,
"boundary_kind": boundary_kind,
"boundary_hash": boundary_hash,
"terminal_kind": terminal_kind,
"action": action,
"receipt": existing,
"local_anchor_generation": state["anchor_generation"],
"local_anchor_digest": state["observed_anchor_digest"],
"hardware_observed_digest": hardware["observed_digest"],
"note": note,
})

frontier = self._frontier(ledger)
if frontier["event_hash"] != boundary_hash:
raise MonotonicAnchorError(
"unreconciled monotonic boundary is not the current ledger frontier"
)
event = ledger.event(boundary_hash)
if event is None or event.get("kind") != boundary_kind:
raise MonotonicAnchorError(
"reconciliation boundary does not identify the expected ledger event"
)
event_payload = event.get("payload")
if not isinstance(event_payload, dict):
raise MonotonicAnchorError("reconciliation boundary payload is invalid")
if event_payload.get("preflight_hash") != preflight_hash:
raise MonotonicAnchorError("reconciliation boundary preflight mismatch")
if boundary_phase == PHASE_TERMINAL:
started = anchored.get("started")
if started is None:
raise MonotonicAnchorError(
"terminal reconciliation requires a durable started receipt"
)
if started["binding_sha256"] != attempt["binding_sha256"]:
raise MonotonicAnchorError(
"terminal reconciliation started binding mismatch"
)

receipt = self._build_receipt(
state_id=state["state_id"],
generation=state["anchor_generation"] + 1,
phase=boundary_phase,
preflight_hash=preflight_hash,
binding_sha256=attempt["binding_sha256"],
terminal_kind=terminal_kind,
frontier=frontier,
previous_digest=state["observed_anchor_digest"],
hardware=state,
)
previous_digest = receipt["previous_anchor_digest"]
expected_digest = receipt["observed_anchor_digest"]
observed = hardware["observed_digest"]
if observed == previous_digest:
action = "NEEDS_EXACTLY_ONE_EXTEND"
note = "hardware is at the previous durable digest"
elif observed == expected_digest:
action = "APPEND_EXACT_RECEIPT_ONLY"
note = (
"hardware already reached the deterministic next digest; "
"a second extend is forbidden"
)
else:
action = "HOLD_DIVERGED"
note = (
"hardware is neither at the previous durable digest nor the "
"deterministic expected next digest"
)

return self._finalize_reconciliation_plan({
"schema": RECONCILIATION_PLAN_SCHEMA,
"domain": DOMAIN,
"state_id": state["state_id"],
"preflight_hash": preflight_hash,
"binding_sha256": attempt["binding_sha256"],
"attempt_phase": attempt["phase"],
"boundary_phase": boundary_phase,
"boundary_kind": boundary_kind,
"boundary_hash": boundary_hash,
"terminal_kind": terminal_kind,
"action": action,
"receipt": receipt,
"local_anchor_generation": state["anchor_generation"],
"local_anchor_digest": state["observed_anchor_digest"],
"hardware_observed_digest": observed,
"note": note,
})

def apply_offline_reconciliation(
self, ledger, *, preflight_hash: str, plan_hash: str
) -> dict[str, Any]:
"""Apply one exact pre-planned recovery action; never retries automatically."""
if not _is_nonzero_sha256(plan_hash):
raise MonotonicAnchorError("reconciliation plan hash is invalid")
with self._require_r14_witness(ledger).locked():
plan = self._plan_offline_reconciliation_locked(
ledger, preflight_hash=preflight_hash
)
if plan["plan_hash"] != plan_hash:
raise MonotonicAnchorError(
"reconciliation plan hash differs from current authority state"
)
action = plan["action"]
if action == "ALREADY_RECONCILED":
return {
"status": action,
"plan_hash": plan_hash,
"receipt_hash": None,
"anchor_generation": plan["local_anchor_generation"],
"observed_anchor_digest": plan["local_anchor_digest"],
}
if action in {"HOLD_DIVERGED", "OUTCOME_EVIDENCE_REQUIRED"}:
raise MonotonicAnchorError(
f"reconciliation cannot apply while plan action is {action}"
)
if action not in {
"NEEDS_EXACTLY_ONE_EXTEND", "APPEND_EXACT_RECEIPT_ONLY"
}:
raise MonotonicAnchorError("unsupported reconciliation action")

receipt = plan.get("receipt")
receipt = self._require_receipt(receipt)
self._validate_receipt_crypto(receipt)
frontier = self._frontier(ledger)
if frontier["event_hash"] != receipt["ledger_event_hash"]:
raise MonotonicAnchorError(
"reconciliation boundary is no longer the current ledger frontier"
)
before = self._read_provider()
if (
before["nv_public_sha256"] != receipt["nv_public_sha256"]
or before["nv_name_sha256"] != receipt["nv_name_sha256"]
):
raise MonotonicAnchorError(
"reconciliation hardware identity changed before apply"
)

if action == "NEEDS_EXACTLY_ONE_EXTEND":
if before["observed_digest"] != receipt["previous_anchor_digest"]:
raise MonotonicAnchorError(
"reconciliation hardware moved before the planned extend"
)
try:
returned = _require_snapshot(self.provider.extend(
expected_previous_digest=receipt["previous_anchor_digest"],
commitment_sha256=receipt["commitment_sha256"],
))
except MonotonicAnchorError:
raise
except Exception as exc:
raise MonotonicAnchorError(
f"reconciliation extend failed: {type(exc).__name__}: {exc}"
) from exc
if (
returned["nv_public_sha256"] != receipt["nv_public_sha256"]
or returned["nv_name_sha256"] != receipt["nv_name_sha256"]
or returned["observed_digest"] != receipt["observed_anchor_digest"]
):
raise MonotonicAnchorError(
"reconciliation extend readback mismatch"
)
fresh = self._read_provider()
if fresh != returned:
raise MonotonicAnchorError(
"reconciliation hardware changed after extend"
)
else:
if before["observed_digest"] != receipt["observed_anchor_digest"]:
raise MonotonicAnchorError(
"append-only reconciliation hardware digest changed"
)

# No retry loop: if this durable append fails after hardware advance,
# the operator must generate a fresh plan. That fresh plan will
# deterministically classify APPEND_EXACT_RECEIPT_ONLY.
try:
receipt_hash = ledger.append(EVENT_KIND, receipt)
except Exception as exc:
raise MonotonicAnchorError(
"reconciliation receipt was not durably recorded"
) from exc

after = self._read_provider()
if (
after["nv_public_sha256"] != receipt["nv_public_sha256"]
or after["nv_name_sha256"] != receipt["nv_name_sha256"]
or after["observed_digest"] != receipt["observed_anchor_digest"]
):
raise MonotonicAnchorError(
"reconciliation hardware frontier changed after receipt"
)
validated = self.validate_global(ledger)
return {
"status": "RECONCILED",
"plan_hash": plan_hash,
"receipt_hash": receipt_hash,
"anchor_generation": validated["anchor_generation"],
"observed_anchor_digest": validated["observed_anchor_digest"],
}

def require_terminal_receipt(
self, ledger, *, preflight_hash: str, binding_sha256: str, attempt: dict[str, Any]
) -> dict[str, Any]:
Expand Down
Loading
Loading