From b79df3ddad68275caded3de73e481caf67375f37 Mon Sep 17 00:00:00 2001 From: bitmaster162 Date: Fri, 18 Sep 2026 18:41:03 +0700 Subject: [PATCH] R24 add witnessed replay rollback tamper protection --- .../trusted_human_approval_production.py | 45 + continuityos/witnessed_approval_replay.py | 821 ++++++++++++++++++ docs/R24_WITNESSED_REPLAY_ROLLBACK_TAMPER.md | 122 +++ tests/test_witnessed_approval_replay.py | 634 ++++++++++++++ 4 files changed, 1622 insertions(+) create mode 100644 continuityos/witnessed_approval_replay.py create mode 100644 docs/R24_WITNESSED_REPLAY_ROLLBACK_TAMPER.md create mode 100644 tests/test_witnessed_approval_replay.py diff --git a/continuityos/trusted_human_approval_production.py b/continuityos/trusted_human_approval_production.py index 8cffb61..9c36d1f 100644 --- a/continuityos/trusted_human_approval_production.py +++ b/continuityos/trusted_human_approval_production.py @@ -11,6 +11,9 @@ from .durable_approval_replay import SQLiteApprovalReplayGuard from .multi_host_approval_replay import MULTI_HOST, PostgresApprovalReplayAuthority +from .witnessed_approval_replay import ( + ROLLBACK_PROTECTION, PostgresWitnessedApprovalReplayAuthority, +) from .persistent_governance_store import resolve_persistent_governance_store_path from .trusted_human_approval import HumanApprovalResult, verify_and_consume_human_approval @@ -53,6 +56,20 @@ def build_multi_host_production_replay_authority( return authority +def build_witnessed_multi_host_production_replay_authority( + *, replay_dsn: str, replay_witness: Any, replay_namespace: str = "human-approval", +) -> PostgresWitnessedApprovalReplayAuthority: + """Build the R24 rollback/tamper-aware MULTI_HOST replay authority.""" + authority = PostgresWitnessedApprovalReplayAuthority( + replay_dsn, witness=replay_witness, namespace=replay_namespace + ) + if authority.replay_scope != MULTI_HOST: + raise RuntimeError("production human approval: multi-host replay scope drift") + if authority.rollback_protection != ROLLBACK_PROTECTION: + raise RuntimeError("production human approval: rollback protection drift") + return authority + + def verify_and_consume_human_approval_production( *, replay_db_path: str | Path, @@ -127,9 +144,37 @@ def verify_and_consume_human_approval_production_multi_host( ) +def verify_and_consume_human_approval_production_multi_host_witnessed( + *, replay_dsn: str, replay_witness: Any, replay_namespace: str, + execution_host_count: int, request_receipt: Any, approval_envelope: Any, + trusted_key_registry: Any, pinned_registry_sha256: str, repository: str, + current_base_sha: str, current_head_sha: str, current_tree_sha: str, + now_unix: int, +) -> HumanApprovalResult: + """Verify one approval with R24 external-witness rollback protection.""" + if _execution_host_count(execution_host_count) < 2: + raise ValueError( + "production human approval: witnessed multi-host binding requires topology > 1" + ) + authority = build_witnessed_multi_host_production_replay_authority( + replay_dsn=replay_dsn, replay_witness=replay_witness, + replay_namespace=replay_namespace, + ) + return verify_and_consume_human_approval( + request_receipt=request_receipt, approval_envelope=approval_envelope, + trusted_key_registry=trusted_key_registry, + pinned_registry_sha256=pinned_registry_sha256, repository=repository, + current_base_sha=current_base_sha, current_head_sha=current_head_sha, + current_tree_sha=current_tree_sha, now_unix=now_unix, + replay_guard=authority, + ) + + __all__ = [ "DEFAULT_BUSY_TIMEOUT_MS", "SINGLE_HOST", "MULTI_HOST", "build_production_replay_guard", "build_multi_host_production_replay_authority", + "build_witnessed_multi_host_production_replay_authority", "verify_and_consume_human_approval_production", "verify_and_consume_human_approval_production_multi_host", + "verify_and_consume_human_approval_production_multi_host_witnessed", ] diff --git a/continuityos/witnessed_approval_replay.py b/continuityos/witnessed_approval_replay.py new file mode 100644 index 0000000..c0e498b --- /dev/null +++ b/continuityos/witnessed_approval_replay.py @@ -0,0 +1,821 @@ +"""Witnessed multi-host replay with rollback/tamper detection and recovery. + +R24 composes PostgreSQL replay state with an external append-only witness. The +witness is authoritative for monotonic claim history; PostgreSQL can be rebuilt +from witnessed records after a database snapshot rollback. No witness backend +is silently substituted or auto-provisioned. +""" +from __future__ import annotations + +import importlib +import json +from typing import Any, Callable, Mapping + +from .multi_host_approval_replay import ( + ALREADY_CONSUMED, + BACKEND, + CLAIMED, + CONFLICT, + MULTI_HOST, + SCHEMA as R23_SCHEMA, + MultiHostApprovalReplayError, + _approval_id, + _canonical_json, + _close_quietly, + _dsn, + _hex64, + _namespace, + _receipt, + _sha256_text, +) + +SCHEMA = "continuityos.witnessed_approval_replay/v1" +STATE_SCHEMA = "continuityos.replay_witness_state/v1" +RECORD_SCHEMA = "continuityos.replay_witness_record/v1" +ROLLBACK_PROTECTION = "EXTERNAL_APPEND_ONLY_WITNESS" +WITNESS_SCOPE = "EXTERNAL_APPEND_ONLY" +GENESIS_DOMAIN = "continuityos.replay_witness_genesis/v1" + + +def _generation(value: object) -> int: + if type(value) is not int or value < 0: + raise MultiHostApprovalReplayError("witnessed replay: invalid generation") + return value + + +def _genesis_head(namespace: str) -> str: + return _sha256_text(_canonical_json({ + "schema": GENESIS_DOMAIN, + "namespace": namespace, + })) + + +def _state(namespace: str, generation: int, head_sha256: str) -> dict[str, Any]: + return { + "schema": STATE_SCHEMA, + "namespace": namespace, + "generation": _generation(generation), + "head_sha256": _hex64("head_sha256", head_sha256), + } + + +def _require_state(value: Any, *, namespace: str) -> dict[str, Any]: + if type(value) is not dict or set(value) != { + "schema", "namespace", "generation", "head_sha256" + }: + raise MultiHostApprovalReplayError("witnessed replay: witness state invalid") + if value["schema"] != STATE_SCHEMA or value["namespace"] != namespace: + raise MultiHostApprovalReplayError("witnessed replay: witness state invalid") + return _state(namespace, value["generation"], value["head_sha256"]) + + +def _record_core( + *, + namespace: str, + generation: int, + previous_head_sha256: str, + approval_id: str, + subject: dict[str, Any], + subject_sha256: str, + nonce: str, + digest_sha256: str, +) -> dict[str, Any]: + return { + "schema": RECORD_SCHEMA, + "namespace": namespace, + "generation": _generation(generation), + "previous_head_sha256": _hex64( + "previous_head_sha256", previous_head_sha256 + ), + "approval_id": _approval_id(approval_id), + "subject": dict(subject), + "subject_sha256": _hex64("subject_sha256", subject_sha256), + "nonce": _hex64("nonce", nonce), + "digest_sha256": _hex64("digest_sha256", digest_sha256), + } + + +def _build_record( + *, + namespace: str, + generation: int, + previous_head_sha256: str, + approval_id: str, + subject: dict[str, Any], + nonce: str, + digest_sha256: str, +) -> dict[str, Any]: + subject_json = _canonical_json(subject) + core = _record_core( + namespace=namespace, + generation=generation, + previous_head_sha256=previous_head_sha256, + approval_id=approval_id, + subject=dict(subject), + subject_sha256=_sha256_text(subject_json), + nonce=nonce, + digest_sha256=digest_sha256, + ) + head = _sha256_text(_canonical_json(core)) + return {**core, "head_sha256": head, "record_id": "rwr_" + head} + + +def _require_record( + value: Any, + *, + namespace: str, + expected_generation: int | None = None, + expected_previous_head: str | None = None, +) -> dict[str, Any]: + if type(value) is not dict or set(value) != { + "schema", "namespace", "generation", "previous_head_sha256", + "approval_id", "subject", "subject_sha256", "nonce", + "digest_sha256", "head_sha256", "record_id", + }: + raise MultiHostApprovalReplayError("witnessed replay: witness record invalid") + if value["schema"] != RECORD_SCHEMA or value["namespace"] != namespace: + raise MultiHostApprovalReplayError("witnessed replay: witness record invalid") + if type(value["subject"]) is not dict: + raise MultiHostApprovalReplayError("witnessed replay: witness record invalid") + rebuilt = _build_record( + namespace=namespace, + generation=value["generation"], + previous_head_sha256=value["previous_head_sha256"], + approval_id=value["approval_id"], + subject=dict(value["subject"]), + nonce=value["nonce"], + digest_sha256=value["digest_sha256"], + ) + if rebuilt != value: + raise MultiHostApprovalReplayError("witnessed replay: witness record tampered") + if expected_generation is not None and value["generation"] != expected_generation: + raise MultiHostApprovalReplayError("witnessed replay: witness generation gap") + if ( + expected_previous_head is not None + and value["previous_head_sha256"] != expected_previous_head + ): + raise MultiHostApprovalReplayError("witnessed replay: witness chain mismatch") + return dict(value) + + +class PostgresWitnessedApprovalReplayAuthority: + """R24 PostgreSQL replay authority backed by an external monotonic witness.""" + + replay_scope = MULTI_HOST + backend = BACKEND + rollback_protection = ROLLBACK_PROTECTION + + def __init__( + self, + dsn: str, + *, + witness: Any, + namespace: str = "human-approval", + connect: Callable[[str], Any] | None = None, + ) -> None: + self.dsn = _dsn(dsn) + self.namespace = _namespace(namespace) + if getattr(witness, "witness_scope", None) != WITNESS_SCOPE: + raise MultiHostApprovalReplayError( + "witnessed replay: external append-only witness required" + ) + for name in ("current_state", "records_after", "append_record"): + if not callable(getattr(witness, name, None)): + raise MultiHostApprovalReplayError( + "witnessed replay: witness contract invalid" + ) + self.witness = witness + if connect is None: + try: + module = importlib.import_module("psycopg") + connect = module.connect + except (ImportError, AttributeError) as exc: + raise MultiHostApprovalReplayError( + "witnessed replay: psycopg backend unavailable" + ) from exc + self._connect = connect + self._initialize() + self.synchronize() + + def _open(self) -> Any: + try: + return self._connect(self.dsn) + except Exception as exc: + raise MultiHostApprovalReplayError( + "witnessed replay: PostgreSQL connection failed" + ) from exc + + def _initialize(self) -> None: + con = self._open() + cur = None + try: + cur = con.cursor() + cur.execute( + "CREATE TABLE IF NOT EXISTS continuityos_replay_meta (" + "key TEXT PRIMARY KEY, value TEXT NOT NULL)" + ) + cur.execute( + "CREATE TABLE IF NOT EXISTS continuityos_approval_claims (" + "namespace TEXT NOT NULL, approval_id TEXT NOT NULL, " + "subject_json TEXT NOT NULL, subject_sha256 TEXT NOT NULL, " + "nonce TEXT NOT NULL, digest_sha256 TEXT NOT NULL, " + "PRIMARY KEY(namespace, approval_id))" + ) + cur.execute( + "CREATE TABLE IF NOT EXISTS continuityos_replay_witness_state (" + "namespace TEXT PRIMARY KEY, generation BIGINT NOT NULL, " + "head_sha256 TEXT NOT NULL)" + ) + cur.execute( + "CREATE TABLE IF NOT EXISTS continuityos_replay_witness_journal (" + "namespace TEXT NOT NULL, generation BIGINT NOT NULL, " + "previous_head_sha256 TEXT NOT NULL, head_sha256 TEXT NOT NULL, " + "approval_id TEXT NOT NULL, subject_json TEXT NOT NULL, " + "subject_sha256 TEXT NOT NULL, nonce TEXT NOT NULL, " + "digest_sha256 TEXT NOT NULL, " + "PRIMARY KEY(namespace, generation), " + "UNIQUE(namespace, approval_id))" + ) + cur.execute( + "INSERT INTO continuityos_replay_meta(key, value) " + "VALUES('schema', %s) ON CONFLICT (key) DO NOTHING", + (R23_SCHEMA,), + ) + cur.execute( + "INSERT INTO continuityos_replay_meta(key, value) " + "VALUES('witness_schema', %s) ON CONFLICT (key) DO NOTHING", + (SCHEMA,), + ) + self._require_schema_meta(cur) + cur.execute( + "SELECT generation, head_sha256 " + "FROM continuityos_replay_witness_state WHERE namespace=%s", + (self.namespace,), + ) + state_row = cur.fetchone() + if state_row is None: + cur.execute( + "SELECT COUNT(*) FROM continuityos_approval_claims " + "WHERE namespace=%s", + (self.namespace,), + ) + count_row = cur.fetchone() + if count_row is None or type(count_row[0]) is not int: + raise MultiHostApprovalReplayError( + "witnessed replay: claim count unavailable" + ) + if count_row[0] != 0: + raise MultiHostApprovalReplayError( + "witnessed replay: explicit R23 migration required" + ) + cur.execute( + "INSERT INTO continuityos_replay_witness_state(" + "namespace, generation, head_sha256) VALUES(%s, %s, %s)", + (self.namespace, 0, _genesis_head(self.namespace)), + ) + else: + _state(self.namespace, state_row[0], state_row[1]) + con.commit() + except MultiHostApprovalReplayError: + try: + con.rollback() + except Exception: + pass + raise + except Exception as exc: + try: + con.rollback() + except Exception: + pass + raise MultiHostApprovalReplayError( + "witnessed replay: schema initialization failed" + ) from exc + finally: + _close_quietly(cur) + _close_quietly(con) + + def _require_schema_meta(self, cur: Any) -> None: + cur.execute( + "SELECT value FROM continuityos_replay_meta WHERE key='schema'" + ) + if cur.fetchone() != (R23_SCHEMA,): + raise MultiHostApprovalReplayError( + "witnessed replay: R23 schema identity mismatch" + ) + cur.execute( + "SELECT value FROM continuityos_replay_meta WHERE key='witness_schema'" + ) + if cur.fetchone() != (SCHEMA,): + raise MultiHostApprovalReplayError( + "witnessed replay: R24 schema identity mismatch" + ) + + def _witness_current(self) -> dict[str, Any]: + try: + value = self.witness.current_state(self.namespace) + except Exception as exc: + raise MultiHostApprovalReplayError( + "witnessed replay: witness unavailable" + ) from exc + return _require_state(value, namespace=self.namespace) + + def _witness_records_after(self, generation: int) -> list[dict[str, Any]]: + try: + value = self.witness.records_after(self.namespace, generation) + except Exception as exc: + raise MultiHostApprovalReplayError( + "witnessed replay: witness history unavailable" + ) from exc + if type(value) is not list: + raise MultiHostApprovalReplayError( + "witnessed replay: witness history invalid" + ) + return [dict(item) if type(item) is dict else item for item in value] + + def _append_witness( + self, + *, + expected_generation: int, + expected_head_sha256: str, + record: dict[str, Any], + ) -> dict[str, Any]: + try: + value = self.witness.append_record( + self.namespace, + expected_generation, + expected_head_sha256, + dict(record), + ) + except Exception as exc: + raise MultiHostApprovalReplayError( + "witnessed replay: witness append failed" + ) from exc + observed = _require_record( + value, + namespace=self.namespace, + expected_generation=record["generation"], + expected_previous_head=record["previous_head_sha256"], + ) + if observed != record: + raise MultiHostApprovalReplayError( + "witnessed replay: witness append readback mismatch" + ) + return observed + + def _read_state_for_update(self, cur: Any) -> dict[str, Any]: + cur.execute( + "SELECT generation, head_sha256 " + "FROM continuityos_replay_witness_state " + "WHERE namespace=%s FOR UPDATE", + (self.namespace,), + ) + row = cur.fetchone() + if row is None: + raise MultiHostApprovalReplayError( + "witnessed replay: database state missing" + ) + return _state(self.namespace, row[0], row[1]) + + def _audit_database( + self, cur: Any, state_value: Mapping[str, Any] + ) -> list[dict[str, Any]]: + generation = _generation(state_value["generation"]) + expected_head = _genesis_head(self.namespace) + cur.execute( + "SELECT generation, previous_head_sha256, head_sha256, approval_id, " + "subject_json, subject_sha256, nonce, digest_sha256 " + "FROM continuityos_replay_witness_journal " + "WHERE namespace=%s ORDER BY generation", + (self.namespace,), + ) + rows = list(cur.fetchall()) + if len(rows) != generation: + raise MultiHostApprovalReplayError( + "witnessed replay: database journal length mismatch" + ) + records: list[dict[str, Any]] = [] + journal_claims: dict[str, tuple[str, str, str, str]] = {} + for index, row in enumerate(rows, start=1): + ( + row_generation, previous_head, row_head, approval_id, + subject_json, subject_sha, nonce, digest, + ) = row + if row_generation != index or previous_head != expected_head: + raise MultiHostApprovalReplayError( + "witnessed replay: database journal chain mismatch" + ) + try: + subject = json.loads(subject_json) + except Exception as exc: + raise MultiHostApprovalReplayError( + "witnessed replay: database journal subject invalid" + ) from exc + if type(subject) is not dict or _canonical_json(subject) != subject_json: + raise MultiHostApprovalReplayError( + "witnessed replay: database journal subject invalid" + ) + record = _build_record( + namespace=self.namespace, + generation=row_generation, + previous_head_sha256=previous_head, + approval_id=approval_id, + subject=dict(subject), + nonce=nonce, + digest_sha256=digest, + ) + if record["head_sha256"] != row_head or record["subject_sha256"] != subject_sha: + raise MultiHostApprovalReplayError( + "witnessed replay: database journal tampered" + ) + records.append(record) + journal_claims[approval_id] = ( + subject_json, subject_sha, nonce, digest + ) + expected_head = row_head + if expected_head != state_value["head_sha256"]: + raise MultiHostApprovalReplayError( + "witnessed replay: database state head mismatch" + ) + cur.execute( + "SELECT approval_id, subject_json, subject_sha256, nonce, digest_sha256 " + "FROM continuityos_approval_claims WHERE namespace=%s ORDER BY approval_id", + (self.namespace,), + ) + claim_rows = list(cur.fetchall()) + claims = { + row[0]: (row[1], row[2], row[3], row[4]) for row in claim_rows + } + if len(claim_rows) != generation or claims != journal_claims: + raise MultiHostApprovalReplayError( + "witnessed replay: database claim set tampered" + ) + return records + + def _apply_record(self, cur: Any, record: Mapping[str, Any]) -> None: + subject_json = _canonical_json(record["subject"]) + expected_claim = ( + subject_json, + record["subject_sha256"], + record["nonce"], + record["digest_sha256"], + ) + cur.execute( + "INSERT INTO continuityos_approval_claims(" + "namespace, approval_id, subject_json, subject_sha256, nonce, digest_sha256" + ") VALUES(%s, %s, %s, %s, %s, %s) " + "ON CONFLICT (namespace, approval_id) DO NOTHING " + "RETURNING subject_json, subject_sha256, nonce, digest_sha256", + ( + self.namespace, + record["approval_id"], + *expected_claim, + ), + ) + inserted = cur.fetchone() + if inserted is None: + cur.execute( + "SELECT subject_json, subject_sha256, nonce, digest_sha256 " + "FROM continuityos_approval_claims " + "WHERE namespace=%s AND approval_id=%s", + (self.namespace, record["approval_id"]), + ) + inserted = cur.fetchone() + if inserted is None or tuple(inserted) != expected_claim: + raise MultiHostApprovalReplayError( + "witnessed replay: recovery claim conflict" + ) + expected_journal = ( + record["previous_head_sha256"], + record["head_sha256"], + record["approval_id"], + subject_json, + record["subject_sha256"], + record["nonce"], + record["digest_sha256"], + ) + cur.execute( + "INSERT INTO continuityos_replay_witness_journal(" + "namespace, generation, previous_head_sha256, head_sha256, approval_id, " + "subject_json, subject_sha256, nonce, digest_sha256" + ") VALUES(%s, %s, %s, %s, %s, %s, %s, %s, %s) " + "ON CONFLICT (namespace, generation) DO NOTHING " + "RETURNING previous_head_sha256, head_sha256, approval_id, " + "subject_json, subject_sha256, nonce, digest_sha256", + ( + self.namespace, + record["generation"], + *expected_journal, + ), + ) + journal_row = cur.fetchone() + if journal_row is None: + cur.execute( + "SELECT previous_head_sha256, head_sha256, approval_id, " + "subject_json, subject_sha256, nonce, digest_sha256 " + "FROM continuityos_replay_witness_journal " + "WHERE namespace=%s AND generation=%s", + (self.namespace, record["generation"]), + ) + journal_row = cur.fetchone() + if journal_row is None or tuple(journal_row) != expected_journal: + raise MultiHostApprovalReplayError( + "witnessed replay: recovery journal conflict" + ) + cur.execute( + "UPDATE continuityos_replay_witness_state " + "SET generation=%s, head_sha256=%s " + "WHERE namespace=%s AND generation=%s AND head_sha256=%s " + "RETURNING generation, head_sha256", + ( + record["generation"], + record["head_sha256"], + self.namespace, + record["generation"] - 1, + record["previous_head_sha256"], + ), + ) + if cur.fetchone() != (record["generation"], record["head_sha256"]): + raise MultiHostApprovalReplayError( + "witnessed replay: database state CAS failed" + ) + + def _snapshot_database(self) -> dict[str, Any]: + con = self._open() + cur = None + try: + cur = con.cursor() + cur.execute("SET TRANSACTION ISOLATION LEVEL SERIALIZABLE") + self._require_schema_meta(cur) + value = self._read_state_for_update(cur) + self._audit_database(cur, value) + con.commit() + return dict(value) + except MultiHostApprovalReplayError: + try: + con.rollback() + except Exception: + pass + raise + except Exception as exc: + try: + con.rollback() + except Exception: + pass + raise MultiHostApprovalReplayError( + "witnessed replay: database snapshot failed" + ) from exc + finally: + _close_quietly(cur) + _close_quietly(con) + + def _read_claim_row(self, approval_id: str) -> tuple[str, str, str, str] | None: + con = self._open() + cur = None + try: + cur = con.cursor() + cur.execute("SET TRANSACTION ISOLATION LEVEL SERIALIZABLE") + self._require_schema_meta(cur) + cur.execute( + "SELECT subject_json, subject_sha256, nonce, digest_sha256 " + "FROM continuityos_approval_claims " + "WHERE namespace=%s AND approval_id=%s", + (self.namespace, approval_id), + ) + row = cur.fetchone() + con.commit() + return None if row is None else tuple(row) + except MultiHostApprovalReplayError: + try: + con.rollback() + except Exception: + pass + raise + except Exception as exc: + try: + con.rollback() + except Exception: + pass + raise MultiHostApprovalReplayError( + "witnessed replay: claim read failed" + ) from exc + finally: + _close_quietly(cur) + _close_quietly(con) + + def _receipt_for_existing( + self, + *, + approval_id: str, + subject: dict[str, Any], + subject_sha256: str, + nonce: str, + digest_sha256: str, + existing: tuple[str, str, str, str], + ) -> dict[str, Any]: + same = existing == ( + _canonical_json(subject), subject_sha256, nonce, digest_sha256 + ) + return _receipt( + namespace=self.namespace, + approval_id=approval_id, + subject=dict(subject), + subject_sha256=subject_sha256, + nonce=nonce, + digest_sha256=digest_sha256, + status=ALREADY_CONSUMED if same else CONFLICT, + existing=(existing[1], existing[2], existing[3]), + ) + + def synchronize(self) -> dict[str, Any]: + for _attempt in range(8): + db_snapshot = self._snapshot_database() + + # External witness I/O is deliberately outside every PostgreSQL + # transaction/row lock. The witness CAS is the global monotonic + # serializer; PostgreSQL is revalidated before applying its prefix. + witness_state = self._witness_current() + if witness_state["generation"] < db_snapshot["generation"]: + raise MultiHostApprovalReplayError( + "witnessed replay: database ahead of witness" + ) + + needed = witness_state["generation"] - db_snapshot["generation"] + history: list[dict[str, Any]] = [] + if needed: + observed_history = self._witness_records_after( + db_snapshot["generation"] + ) + if len(observed_history) < needed: + raise MultiHostApprovalReplayError( + "witnessed replay: witness history incomplete" + ) + history = observed_history[:needed] + + con = self._open() + cur = None + try: + cur = con.cursor() + cur.execute("SET TRANSACTION ISOLATION LEVEL SERIALIZABLE") + self._require_schema_meta(cur) + locked_state = self._read_state_for_update(cur) + self._audit_database(cur, locked_state) + if locked_state != db_snapshot: + con.rollback() + continue + + if witness_state["generation"] == locked_state["generation"]: + if witness_state["head_sha256"] != locked_state["head_sha256"]: + raise MultiHostApprovalReplayError( + "witnessed replay: database/witness head mismatch" + ) + con.commit() + return dict(locked_state) + + expected_generation = locked_state["generation"] + expected_head = locked_state["head_sha256"] + for raw in history: + record = _require_record( + raw, + namespace=self.namespace, + expected_generation=expected_generation + 1, + expected_previous_head=expected_head, + ) + self._apply_record(cur, record) + expected_generation = record["generation"] + expected_head = record["head_sha256"] + + if ( + expected_generation != witness_state["generation"] + or expected_head != witness_state["head_sha256"] + ): + raise MultiHostApprovalReplayError( + "witnessed replay: witness history incomplete" + ) + recovered = _state( + self.namespace, expected_generation, expected_head + ) + self._audit_database(cur, recovered) + con.commit() + return recovered + except MultiHostApprovalReplayError: + try: + con.rollback() + except Exception: + pass + raise + except Exception as exc: + try: + con.rollback() + except Exception: + pass + raise MultiHostApprovalReplayError( + "witnessed replay: synchronization failed" + ) from exc + finally: + _close_quietly(cur) + _close_quietly(con) + + raise MultiHostApprovalReplayError( + "witnessed replay: database synchronization contention" + ) + + def claim_once( + self, + *, + approval_id: str, + subject: dict[str, Any], + nonce: str, + digest_sha256: str, + ) -> dict[str, Any]: + identifier = _approval_id(approval_id) + subject_value = dict(subject) + subject_json = _canonical_json(subject_value) + subject_sha256 = _sha256_text(subject_json) + nonce_value = _hex64("nonce", nonce) + digest_value = _hex64("digest_sha256", digest_sha256) + + last_append_error: MultiHostApprovalReplayError | None = None + for _attempt in range(8): + state_value = self.synchronize() + existing = self._read_claim_row(identifier) + if existing is not None: + return self._receipt_for_existing( + approval_id=identifier, + subject=subject_value, + subject_sha256=subject_sha256, + nonce=nonce_value, + digest_sha256=digest_value, + existing=existing, + ) + + record = _build_record( + namespace=self.namespace, + generation=state_value["generation"] + 1, + previous_head_sha256=state_value["head_sha256"], + approval_id=identifier, + subject=subject_value, + nonce=nonce_value, + digest_sha256=digest_value, + ) + + try: + self._append_witness( + expected_generation=state_value["generation"], + expected_head_sha256=state_value["head_sha256"], + record=record, + ) + except MultiHostApprovalReplayError as exc: + last_append_error = exc + # CAS loss to another writer is recoverable. An ambiguous + # append outcome is never promoted to CLAIMED: after recovery + # an exact row is ALREADY_CONSUMED, preserving one-success. + try: + self.synchronize() + existing = self._read_claim_row(identifier) + except MultiHostApprovalReplayError: + raise exc + if existing is not None: + return self._receipt_for_existing( + approval_id=identifier, + subject=subject_value, + subject_sha256=subject_sha256, + nonce=nonce_value, + digest_sha256=digest_value, + existing=existing, + ) + continue + + # The append succeeded outside the DB critical section. Reconcile + # the authoritative witness prefix, then prove our exact row exists. + self.synchronize() + existing = self._read_claim_row(identifier) + expected = (subject_json, subject_sha256, nonce_value, digest_value) + if existing != expected: + raise MultiHostApprovalReplayError( + "witnessed replay: appended claim recovery mismatch" + ) + return _receipt( + namespace=self.namespace, + approval_id=identifier, + subject=subject_value, + subject_sha256=subject_sha256, + nonce=nonce_value, + digest_sha256=digest_value, + status=CLAIMED, + ) + + if last_append_error is not None: + raise last_append_error + raise MultiHostApprovalReplayError( + "witnessed replay: claim contention exceeded" + ) + + +__all__ = [ + "SCHEMA", + "STATE_SCHEMA", + "RECORD_SCHEMA", + "ROLLBACK_PROTECTION", + "WITNESS_SCOPE", + "PostgresWitnessedApprovalReplayAuthority", +] diff --git a/docs/R24_WITNESSED_REPLAY_ROLLBACK_TAMPER.md b/docs/R24_WITNESSED_REPLAY_ROLLBACK_TAMPER.md new file mode 100644 index 0000000..53f1b8d --- /dev/null +++ b/docs/R24_WITNESSED_REPLAY_ROLLBACK_TAMPER.md @@ -0,0 +1,122 @@ +# R24 — Witnessed Replay Rollback / Tamper Detection + +R24 hardens the R23 multi-host replay authority against database rollback, +snapshot restore, silent row tampering, and unwitnessed direct writes. + +## Threat model + +R23 proves atomic multi-host claim-once behavior inside one shared PostgreSQL +state. It does not detect a later restoration of that database to an older +valid snapshot. A restored snapshot can make a previously consumed approval +appear unused. + +R24 adds an external append-only witness in a separate rollback domain. +If PostgreSQL and the witness can be rolled back together, R24 cannot prove +monotonic history and must not be described as rollback-resistant. + +## Protocol + +For each namespace R24 maintains: + +- PostgreSQL witnessed state: generation plus head_sha256; +- PostgreSQL witnessed journal: one chained record per consumed approval; +- the existing R23 approval claim set; +- an external append-only witness carrying the authoritative record chain. + +Each witness record binds the exact approval ID, canonical subject, subject +SHA-256, nonce, approval-envelope digest, previous head, generation, and new +head. + +Before every claim R24 performs a full integrity sweep: + +1. replay journal generations from genesis and recompute every chained head; +2. require the journal head to equal PostgreSQL witnessed state; +3. require the complete approval claim set to equal the journal bindings; +4. compare PostgreSQL state with the external witness state. + +Database-ahead-of-witness, same-generation head mismatch, journal tamper, +claim-row tamper, or direct unwitnessed claims fail closed. + +## Crash-safe ordering and recovery + +The external witness append is committed before the PostgreSQL claim/journal +transaction. This ordering intentionally makes the witness authoritative. + +If the process crashes after witness append but before PostgreSQL commit, the +witness is ahead. On the next initialization or claim, R24 validates the +witness chain and replays missing witnessed records into PostgreSQL. The +approval remains consumed; a crash cannot reopen it. + +If PostgreSQL is restored to an older snapshot, the same recovery path rebuilds +the missing suffix from the witness. If PostgreSQL contains state not present +in the witness, R24 fails closed instead of guessing which side is correct. + +## Migration boundary + +R24 does not silently bless pre-existing R23 claims. If a namespace contains +R23 claims but has no R24 witnessed state, initialization fails with +explicit R23 migration required. + +A future bounded migration procedure must create independently reviewed +genesis or migration evidence. Auto-hashing the current database into a new +witness would turn a compromised snapshot into trusted history and is +therefore forbidden. + +## Production binding + +R23 production entrypoints remain available for historical compatibility but +are rollback-unqualified. R24 adds a separate witnessed multi-host production +entrypoint requiring an explicit external witness object with +witness_scope=EXTERNAL_APPEND_ONLY. + +There is no fallback from the witnessed path to ordinary R23 PostgreSQL, +SQLite, or memory. + +## Authority boundary + +R24 only governs replay ownership and recovery of replay evidence. It does not +merge pull requests, deploy software, execute runtime actions, trade, access +wallets, or grant capital authority. + +## Qualification boundary + +Current tests use a deterministic transactional PostgreSQL fake plus an +independent append-only witness fake. They cover concurrent claim-once, +database rollback recovery, crash-after-witness-append recovery, database +tamper, direct unwitnessed writes, witness rollback, and witness-record tamper. + +This is not live production qualification. Production claims require: + +- a real external append-only witness in a rollback domain independent from + PostgreSQL; +- one shared PostgreSQL service; +- at least two independent execution hosts or process domains; +- injected crash and network-partition tests around witness append and DB commit; +- restore-from-snapshot testing against real PostgreSQL; +- operator recovery evidence. + +The current full integrity sweep is O(n) in witnessed claims per claim. That is +an intentional safety-first design for low-volume Human approval replay state, +not a claim of unbounded production scalability. + +## External I/O and availability boundary + +No witness network call is executed while a PostgreSQL transaction or +FOR UPDATE row lock is held. R24 snapshots and audits PostgreSQL, releases the +critical section, performs witness I/O, then reacquires PostgreSQL state and +requires the snapshot to remain current before applying any witnessed suffix. + +The witness CAS is the monotonic serializer across hosts. PostgreSQL is a +recoverable materialization of that witnessed history. + +A single external witness is intentionally fail-closed and is therefore an +availability dependency. R24 core does not implement witness transport, +timeouts, retries, replication, metrics, paging, or failover. A production +witness adapter must provide bounded network operations, explicit health and +latency telemetry, durable append/read semantics, and an operator recovery +procedure. High availability or a documented RTO/RPO is required before +calling the witnessed path production-qualified. + +Witness rotation is also not automatic. Replacing or re-keying the witness +requires a separately reviewed continuity/migration procedure that preserves +the existing chain; silently starting a fresh witness is forbidden. diff --git a/tests/test_witnessed_approval_replay.py b/tests/test_witnessed_approval_replay.py new file mode 100644 index 0000000..4787693 --- /dev/null +++ b/tests/test_witnessed_approval_replay.py @@ -0,0 +1,634 @@ +from __future__ import annotations + +import copy +import threading +from concurrent.futures import ThreadPoolExecutor + +import pytest + +import continuityos.witnessed_approval_replay as replay +from continuityos.multi_host_approval_replay import ( + ALREADY_CONSUMED, CLAIMED, MultiHostApprovalReplayError, +) + +NS = "human-approval" +APPROVAL_ID = "hap_" + "a" * 64 +NONCE = "1" * 64 +DIGEST = "2" * 64 +SUBJECT = { + "repository": "bitmaster162/continuityos", + "baseline_sha": "b" * 40, + "candidate_sha": "c" * 40, + "candidate_tree_sha": "d" * 40, +} + + +class DbState: + def __init__(self) -> None: + self.lock = threading.RLock() + self.meta: dict[str, str] = {} + self.claims: dict[tuple[str, str], tuple[str, str, str, str]] = {} + self.states: dict[str, tuple[int, str]] = {} + self.journal: dict[tuple[str, int], tuple[str, str, str, str, str, str, str]] = {} + self.fail_next_commit = False + self.active_transactions = 0 + + +class FakeConnection: + def __init__(self, state: DbState) -> None: + self.state = state + self.locked = False + self.meta = {} + self.claims = {} + self.states = {} + self.journal = {} + + def cursor(self): + if not self.locked: + self.state.lock.acquire() + self.locked = True + self.state.active_transactions += 1 + self.meta = copy.deepcopy(self.state.meta) + self.claims = copy.deepcopy(self.state.claims) + self.states = copy.deepcopy(self.state.states) + self.journal = copy.deepcopy(self.state.journal) + return FakeCursor(self) + + def commit(self) -> None: + if self.state.fail_next_commit: + self.state.fail_next_commit = False + self._release() + raise RuntimeError("simulated crash before database commit") + self.state.meta = copy.deepcopy(self.meta) + self.state.claims = copy.deepcopy(self.claims) + self.state.states = copy.deepcopy(self.states) + self.state.journal = copy.deepcopy(self.journal) + self._release() + + def rollback(self) -> None: + self._release() + + def close(self) -> None: + self._release() + + def _release(self) -> None: + if self.locked: + self.locked = False + self.state.active_transactions -= 1 + self.state.lock.release() + + +class FakeCursor: + def __init__(self, con: FakeConnection) -> None: + self.con = con + self.rows: list[tuple] = [] + + def execute(self, sql: str, params=None) -> None: + q = " ".join(sql.lower().split()) + p = tuple(params or ()) + self.rows = [] + if q.startswith("set transaction isolation level"): + return + if q.startswith("create table"): + return + if q.startswith("insert into continuityos_replay_meta"): + key = "witness_schema" if "witness_schema" in q else "schema" + self.con.meta.setdefault(key, p[0]) + return + if q.startswith("select value from continuityos_replay_meta"): + key = "witness_schema" if "witness_schema" in q else "schema" + if key in self.con.meta: + self.rows = [(self.con.meta[key],)] + return + if q.startswith("select generation, head_sha256 from continuityos_replay_witness_state"): + row = self.con.states.get(p[0]) + self.rows = [] if row is None else [row] + return + if q.startswith("select count(*) from continuityos_approval_claims"): + self.rows = [(sum(1 for ns, _ in self.con.claims if ns == p[0]),)] + return + if q.startswith("insert into continuityos_replay_witness_state"): + self.con.states[p[0]] = (p[1], p[2]) + return + if q.startswith("select generation, previous_head_sha256"): + ns = p[0] + items = sorted( + ((gen, row) for (row_ns, gen), row in self.con.journal.items() if row_ns == ns), + key=lambda item: item[0], + ) + self.rows = [(gen, *row) for gen, row in items] + return + if q.startswith("select approval_id, subject_json"): + ns = p[0] + items = sorted( + ((aid, row) for (row_ns, aid), row in self.con.claims.items() if row_ns == ns), + key=lambda item: item[0], + ) + self.rows = [(aid, *row) for aid, row in items] + return + if q.startswith("insert into continuityos_approval_claims"): + key = (p[0], p[1]) + row = (p[2], p[3], p[4], p[5]) + if key not in self.con.claims: + self.con.claims[key] = row + self.rows = [row] + return + if q.startswith("select subject_json"): + row = self.con.claims.get((p[0], p[1])) + self.rows = [] if row is None else [row] + return + if q.startswith("insert into continuityos_replay_witness_journal"): + key = (p[0], p[1]) + row = (p[2], p[3], p[4], p[5], p[6], p[7], p[8]) + if key not in self.con.journal: + self.con.journal[key] = row + self.rows = [row] + return + if q.startswith("select previous_head_sha256"): + row = self.con.journal.get((p[0], p[1])) + self.rows = [] if row is None else [row] + return + if q.startswith("update continuityos_replay_witness_state"): + new_generation, new_head, ns, expected_generation, expected_head = p + if self.con.states.get(ns) == (expected_generation, expected_head): + self.con.states[ns] = (new_generation, new_head) + self.rows = [(new_generation, new_head)] + return + raise AssertionError(f"unexpected SQL: {q}") + + def fetchone(self): + return self.rows.pop(0) if self.rows else None + + def fetchall(self): + rows = list(self.rows) + self.rows = [] + return rows + + def close(self) -> None: + pass + + +def connect_factory(state: DbState): + return lambda dsn: FakeConnection(state) + + +class FakeWitness: + witness_scope = replay.WITNESS_SCOPE + + def __init__(self, db: DbState | None = None) -> None: + self.lock = threading.RLock() + self.records: dict[str, list[dict]] = {} + self.db = db + self.fail_db_commit_after_append = False + self.require_no_db_transaction = False + + def current_state(self, namespace: str) -> dict: + if self.require_no_db_transaction and self.db is not None: + assert self.db.active_transactions == 0 + with self.lock: + records = self.records.get(namespace, []) + if not records: + return replay._state(namespace, 0, replay._genesis_head(namespace)) + last = records[-1] + return replay._state(namespace, last["generation"], last["head_sha256"]) + + def records_after(self, namespace: str, generation: int) -> list[dict]: + if self.require_no_db_transaction and self.db is not None: + assert self.db.active_transactions == 0 + with self.lock: + return copy.deepcopy(self.records.get(namespace, [])[generation:]) + + def append_record( + self, namespace: str, expected_generation: int, + expected_head_sha256: str, record: dict, + ) -> dict: + if self.require_no_db_transaction and self.db is not None: + assert self.db.active_transactions == 0 + with self.lock: + current = self.current_state(namespace) + if ( + current["generation"] != expected_generation + or current["head_sha256"] != expected_head_sha256 + ): + raise RuntimeError("witness CAS conflict") + expected = replay._require_record( + record, + namespace=namespace, + expected_generation=expected_generation + 1, + expected_previous_head=expected_head_sha256, + ) + self.records.setdefault(namespace, []).append(copy.deepcopy(expected)) + if self.fail_db_commit_after_append: + self.fail_db_commit_after_append = False + assert self.db is not None + self.db.fail_next_commit = True + return copy.deepcopy(expected) + + +def authority(db: DbState, witness: FakeWitness): + return replay.PostgresWitnessedApprovalReplayAuthority( + "postgresql://r24-test", + witness=witness, + namespace=NS, + connect=connect_factory(db), + ) + + +def claim( + value: replay.PostgresWitnessedApprovalReplayAuthority, + *, + approval_id: str = APPROVAL_ID, + subject: dict | None = None, +): + return value.claim_once( + approval_id=approval_id, + subject=dict(subject or SUBJECT), + nonce=NONCE, + digest_sha256=DIGEST, + ) + + +def test_claim_once_and_replay_are_r23_receipt_compatible(): + db = DbState() + witness = FakeWitness(db) + first = authority(db, witness) + one = claim(first) + two = claim(authority(db, witness)) + assert one["status"] == CLAIMED + assert two["status"] == ALREADY_CONSUMED + assert one["schema"] == "continuityos.multi_host_replay_claim/v1" + assert db.states[NS][0] == 1 + assert len(witness.records[NS]) == 1 + + +def test_concurrent_authorities_append_exactly_one_witness_record(): + db = DbState() + witness = FakeWitness(db) + first = authority(db, witness) + second = authority(db, witness) + + def run(index: int): + return claim(first if index % 2 == 0 else second) + + with ThreadPoolExecutor(max_workers=12) as pool: + receipts = list(pool.map(run, range(24))) + statuses = [item["status"] for item in receipts] + assert statuses.count(CLAIMED) == 1 + assert statuses.count(ALREADY_CONSUMED) == 23 + assert len(witness.records[NS]) == 1 + + +def test_database_snapshot_rollback_recovers_from_witness(): + db = DbState() + witness = FakeWitness(db) + assert claim(authority(db, witness))["status"] == CLAIMED + with db.lock: + db.claims.clear() + db.journal.clear() + db.states[NS] = (0, replay._genesis_head(NS)) + recovered = authority(db, witness) + assert db.states[NS][0] == 1 + assert claim(recovered)["status"] == ALREADY_CONSUMED + + +def test_crash_after_witness_append_before_db_commit_recovers(): + db = DbState() + witness = FakeWitness(db) + guard = authority(db, witness) + witness.fail_db_commit_after_append = True + with pytest.raises(MultiHostApprovalReplayError, match="database snapshot failed|synchronization failed|claim failed"): + claim(guard) + assert witness.current_state(NS)["generation"] == 1 + assert db.states[NS][0] == 0 + recovered = authority(db, witness) + assert db.states[NS][0] == 1 + assert claim(recovered)["status"] == ALREADY_CONSUMED + + +def test_claim_row_tamper_fails_closed(): + db = DbState() + witness = FakeWitness(db) + assert claim(authority(db, witness))["status"] == CLAIMED + with db.lock: + old = db.claims[(NS, APPROVAL_ID)] + db.claims[(NS, APPROVAL_ID)] = (old[0], old[1], old[2], "f" * 64) + with pytest.raises(MultiHostApprovalReplayError, match="claim set tampered"): + authority(db, witness) + + +def test_direct_unwitnessed_claim_fails_closed(): + db = DbState() + witness = FakeWitness(db) + authority(db, witness) + subject_json = replay._canonical_json(SUBJECT) + with db.lock: + db.claims[(NS, APPROVAL_ID)] = ( + subject_json, replay._sha256_text(subject_json), NONCE, DIGEST + ) + with pytest.raises(MultiHostApprovalReplayError, match="claim set tampered"): + authority(db, witness) + + +def test_database_ahead_of_witness_fails_closed(): + db = DbState() + witness = FakeWitness(db) + assert claim(authority(db, witness))["status"] == CLAIMED + with witness.lock: + witness.records[NS] = [] + with pytest.raises(MultiHostApprovalReplayError, match="database ahead of witness"): + authority(db, witness) + + +def test_witness_record_tamper_fails_closed_during_recovery(): + db = DbState() + witness = FakeWitness(db) + assert claim(authority(db, witness))["status"] == CLAIMED + with db.lock: + db.claims.clear() + db.journal.clear() + db.states[NS] = (0, replay._genesis_head(NS)) + with witness.lock: + witness.records[NS][0]["digest_sha256"] = "e" * 64 + with pytest.raises(MultiHostApprovalReplayError, match="witness record tampered"): + authority(db, witness) + + +def test_existing_r23_claims_require_explicit_migration(): + db = DbState() + witness = FakeWitness(db) + subject_json = replay._canonical_json(SUBJECT) + db.claims[(NS, APPROVAL_ID)] = ( + subject_json, replay._sha256_text(subject_json), NONCE, DIGEST + ) + with pytest.raises(MultiHostApprovalReplayError, match="explicit R23 migration required"): + authority(db, witness) + + +def test_witness_contract_is_mandatory(): + class BadWitness: + pass + + db = DbState() + with pytest.raises( + MultiHostApprovalReplayError, match="external append-only witness required" + ): + replay.PostgresWitnessedApprovalReplayAuthority( + "postgresql://r24-test", witness=BadWitness(), + namespace=NS, connect=connect_factory(db), + ) + + +def test_module_has_no_executor_or_authority_widening_surface(): + from pathlib import Path + + source = Path(replay.__file__).read_text(encoding="utf-8") + assert "subprocess" not in source + assert "merge_pull_request" not in source + assert "can_trade" not in source + assert "capital_permission" not in source + + +def test_production_builder_requires_witnessed_markers(monkeypatch): + import continuityos.trusted_human_approval_production as production + + class Stub: + replay_scope = replay.MULTI_HOST + rollback_protection = replay.ROLLBACK_PROTECTION + + marker = Stub() + monkeypatch.setattr( + production, "PostgresWitnessedApprovalReplayAuthority", + lambda dsn, witness, namespace: marker, + ) + assert production.build_witnessed_multi_host_production_replay_authority( + replay_dsn="postgresql://shared", + replay_witness=object(), + replay_namespace=NS, + ) is marker + + +def test_witnessed_production_binding_requires_multi_host_topology(monkeypatch): + import continuityos.trusted_human_approval_production as production + + monkeypatch.setattr( + production, + "build_witnessed_multi_host_production_replay_authority", + lambda **kwargs: (_ for _ in ()).throw(AssertionError("must not build")), + ) + with pytest.raises(ValueError, match="requires topology > 1"): + production.verify_and_consume_human_approval_production_multi_host_witnessed( + replay_dsn="postgresql://shared", replay_witness=object(), + replay_namespace=NS, execution_host_count=1, + request_receipt={}, approval_envelope={}, trusted_key_registry={}, + pinned_registry_sha256="0" * 64, repository="bitmaster162/continuityos", + current_base_sha="b" * 40, current_head_sha="c" * 40, + current_tree_sha="d" * 40, now_unix=1, + ) + + +def test_witnessed_production_binding_passes_exact_authority(monkeypatch): + import continuityos.trusted_human_approval_production as production + + marker = object() + monkeypatch.setattr( + production, + "build_witnessed_multi_host_production_replay_authority", + lambda **kwargs: marker, + ) + monkeypatch.setattr( + production, "verify_and_consume_human_approval", + lambda **kwargs: kwargs["replay_guard"], + ) + result = production.verify_and_consume_human_approval_production_multi_host_witnessed( + replay_dsn="postgresql://shared", replay_witness=object(), + replay_namespace=NS, execution_host_count=2, + request_receipt={}, approval_envelope={}, trusted_key_registry={}, + pinned_registry_sha256="0" * 64, repository="bitmaster162/continuityos", + current_base_sha="b" * 40, current_head_sha="c" * 40, + current_tree_sha="d" * 40, now_unix=1, + ) + assert result is marker + + +def test_live_r23_schema_meta_tamper_fails_closed(): + db = DbState() + witness = FakeWitness(db) + guard = authority(db, witness) + with db.lock: + db.meta["schema"] = "wrong-schema" + with pytest.raises(MultiHostApprovalReplayError, match="R23 schema identity mismatch"): + guard.synchronize() + + +def test_live_r24_schema_meta_tamper_fails_closed(): + db = DbState() + witness = FakeWitness(db) + guard = authority(db, witness) + with db.lock: + db.meta["witness_schema"] = "wrong-schema" + with pytest.raises(MultiHostApprovalReplayError, match="R24 schema identity mismatch"): + guard.synchronize() + + +def test_r17_replay_stays_denied_after_database_rollback_recovery(): + import importlib.util + from pathlib import Path + from continuityos.trusted_human_approval import verify_and_consume_human_approval + + fixture_path = Path(__file__).with_name("test_trusted_human_approval.py") + spec = importlib.util.spec_from_file_location("_r17_fixture_r24", fixture_path) + fixture = importlib.util.module_from_spec(spec) + assert spec.loader is not None + spec.loader.exec_module(fixture) + + request = fixture.request() + private, registry, pin = fixture.key_material() + envelope = fixture.signed_envelope(request, private) + db = DbState() + witness = FakeWitness(db) + result = verify_and_consume_human_approval( + request_receipt=request, + approval_envelope=envelope, + trusted_key_registry=registry, + pinned_registry_sha256=pin, + repository=fixture.REPOSITORY, + current_base_sha=fixture.BASE, + current_head_sha=fixture.HEAD, + current_tree_sha=fixture.TREE, + now_unix=fixture.NOW, + replay_guard=authority(db, witness), + ) + assert result.replay_claim_receipt["status"] == CLAIMED + + with db.lock: + db.claims.clear() + db.journal.clear() + db.states[NS] = (0, replay._genesis_head(NS)) + + recovered = authority(db, witness) + with pytest.raises(ValueError, match="approval replay detected"): + verify_and_consume_human_approval( + request_receipt=request, + approval_envelope=envelope, + trusted_key_registry=registry, + pinned_registry_sha256=pin, + repository=fixture.REPOSITORY, + current_base_sha=fixture.BASE, + current_head_sha=fixture.HEAD, + current_tree_sha=fixture.TREE, + now_unix=fixture.NOW, + replay_guard=recovered, + ) + + + +def test_external_witness_io_occurs_outside_database_transaction(): + db = DbState() + witness = FakeWitness(db) + witness.require_no_db_transaction = True + guard = authority(db, witness) + assert claim(guard)["status"] == CLAIMED + assert claim(guard)["status"] == ALREADY_CONSUMED + + +def _claim_id(ch: str) -> str: + return "hap_" + ch * 64 + + +def test_truncated_witness_history_fails_closed(): + db = DbState() + witness = FakeWitness(db) + guard = authority(db, witness) + assert claim(guard, approval_id=_claim_id("a"))["status"] == CLAIMED + assert claim(guard, approval_id=_claim_id("b"))["status"] == CLAIMED + with db.lock: + db.claims.clear() + db.journal.clear() + db.states[NS] = (0, replay._genesis_head(NS)) + original = witness.records_after + witness.records_after = lambda namespace, generation: original( + namespace, generation + )[:1] + with pytest.raises( + MultiHostApprovalReplayError, match="witness history incomplete" + ): + authority(db, witness) + + +def test_reordered_witness_history_fails_closed(): + db = DbState() + witness = FakeWitness(db) + guard = authority(db, witness) + assert claim(guard, approval_id=_claim_id("a"))["status"] == CLAIMED + assert claim(guard, approval_id=_claim_id("b"))["status"] == CLAIMED + with db.lock: + db.claims.clear() + db.journal.clear() + db.states[NS] = (0, replay._genesis_head(NS)) + original = witness.records_after + witness.records_after = lambda namespace, generation: list( + reversed(original(namespace, generation)) + ) + with pytest.raises( + MultiHostApprovalReplayError, match="witness generation gap" + ): + authority(db, witness) + + +def test_ambiguous_append_response_never_promotes_to_claimed(): + class AmbiguousWitness(FakeWitness): + def __init__(self, db): + super().__init__(db) + self.once = True + + def append_record(self, *args, **kwargs): + value = super().append_record(*args, **kwargs) + if self.once: + self.once = False + raise RuntimeError("response lost after durable append") + return value + + db = DbState() + witness = AmbiguousWitness(db) + receipt = claim(authority(db, witness)) + assert receipt["status"] == ALREADY_CONSUMED + assert len(witness.records[NS]) == 1 + assert db.states[NS][0] == 1 + + +def test_mixed_synchronize_and_claim_concurrency_is_consistent(): + db = DbState() + witness = FakeWitness(db) + first = authority(db, witness) + second = authority(db, witness) + + def run(index: int): + if index % 3 == 0: + return ("sync", (first if index % 2 == 0 else second).synchronize()) + return ("claim", claim(first if index % 2 == 0 else second)) + + with ThreadPoolExecutor(max_workers=12) as pool: + results = list(pool.map(run, range(30))) + receipts = [value for kind, value in results if kind == "claim"] + statuses = [item["status"] for item in receipts] + assert statuses.count(CLAIMED) == 1 + assert statuses.count(ALREADY_CONSUMED) == len(receipts) - 1 + assert witness.current_state(NS)["generation"] == 1 + assert db.states[NS][0] == 1 + + +def test_namespaces_are_isolated(): + db = DbState() + witness = FakeWitness(db) + left = replay.PostgresWitnessedApprovalReplayAuthority( + "postgresql://r24-test", witness=witness, namespace="human-approval-a", + connect=connect_factory(db), + ) + right = replay.PostgresWitnessedApprovalReplayAuthority( + "postgresql://r24-test", witness=witness, namespace="human-approval-b", + connect=connect_factory(db), + ) + assert claim(left)["status"] == CLAIMED + assert claim(right)["status"] == CLAIMED + assert witness.current_state("human-approval-a")["generation"] == 1 + assert witness.current_state("human-approval-b")["generation"] == 1