From a9dac499374247503be1f49d148ebced7e162c6b Mon Sep 17 00:00:00 2001 From: bitmaster162 Date: Sat, 19 Sep 2026 17:50:02 +0700 Subject: [PATCH 1/6] R25 add replay production qualification harness --- continuityos/witnessed_approval_replay.py | 155 ++-- docs/R25_REPLAY_PRODUCTION_QUALIFICATION.md | 155 ++++ tests/test_r25_replay_qualification.py | 145 +++ tests/test_witnessed_approval_replay.py | 38 + tools/r25_replay_qualification.py | 979 ++++++++++++++++++++ tools/r25_witness_service.py | 281 ++++++ 6 files changed, 1693 insertions(+), 60 deletions(-) create mode 100644 docs/R25_REPLAY_PRODUCTION_QUALIFICATION.md create mode 100644 tests/test_r25_replay_qualification.py create mode 100644 tools/r25_replay_qualification.py create mode 100644 tools/r25_witness_service.py diff --git a/continuityos/witnessed_approval_replay.py b/continuityos/witnessed_approval_replay.py index c0e498b..7a0f80b 100644 --- a/continuityos/witnessed_approval_replay.py +++ b/continuityos/witnessed_approval_replay.py @@ -2,8 +2,10 @@ 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 witnessed records after a database snapshot rollback. R25 adds bounded +retry hardening for PostgreSQL serialization/deadlock aborts discovered by the +real qualification harness. No witness backend is silently substituted or +auto-provisioned. """ from __future__ import annotations @@ -35,6 +37,23 @@ ROLLBACK_PROTECTION = "EXTERNAL_APPEND_ONLY_WITNESS" WITNESS_SCOPE = "EXTERNAL_APPEND_ONLY" GENESIS_DOMAIN = "continuityos.replay_witness_genesis/v1" +_RETRYABLE_TRANSACTION_SQLSTATES = frozenset({"40001", "40P01"}) + + +def _is_retryable_transaction_error(exc: BaseException) -> bool: + current: BaseException | None = exc + seen: set[int] = set() + for _depth in range(8): + if current is None or id(current) in seen: + break + seen.add(id(current)) + sqlstate = getattr(current, "sqlstate", None) + if sqlstate is None: + sqlstate = getattr(current, "pgcode", None) + if sqlstate in _RETRYABLE_TRANSACTION_SQLSTATES: + return True + current = current.__cause__ or current.__context__ + return False def _generation(value: object) -> int: @@ -540,67 +559,81 @@ def _apply_record(self, cur: Any, record: Mapping[str, Any]) -> None: ) 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: + for _attempt in range(8): + con = self._open() + cur = None try: - con.rollback() - except Exception: - pass - raise MultiHostApprovalReplayError( - "witnessed replay: database snapshot failed" - ) from exc - finally: - _close_quietly(cur) - _close_quietly(con) + 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 + if _is_retryable_transaction_error(exc): + continue + raise MultiHostApprovalReplayError( + "witnessed replay: database snapshot failed" + ) from exc + finally: + _close_quietly(cur) + _close_quietly(con) + raise MultiHostApprovalReplayError( + "witnessed replay: database snapshot contention" + ) - 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: + def _read_claim_row( + self, approval_id: str + ) -> tuple[str, str, str, str] | None: + for _attempt in range(8): + con = self._open() + cur = None try: - con.rollback() - except Exception: - pass - raise MultiHostApprovalReplayError( - "witnessed replay: claim read failed" - ) from exc - finally: - _close_quietly(cur) - _close_quietly(con) + 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 + if _is_retryable_transaction_error(exc): + continue + raise MultiHostApprovalReplayError( + "witnessed replay: claim read failed" + ) from exc + finally: + _close_quietly(cur) + _close_quietly(con) + raise MultiHostApprovalReplayError( + "witnessed replay: claim read contention" + ) def _receipt_for_existing( self, @@ -708,6 +741,8 @@ def synchronize(self) -> dict[str, Any]: con.rollback() except Exception: pass + if _is_retryable_transaction_error(exc): + continue raise MultiHostApprovalReplayError( "witnessed replay: synchronization failed" ) from exc diff --git a/docs/R25_REPLAY_PRODUCTION_QUALIFICATION.md b/docs/R25_REPLAY_PRODUCTION_QUALIFICATION.md new file mode 100644 index 0000000..eea8e23 --- /dev/null +++ b/docs/R25_REPLAY_PRODUCTION_QUALIFICATION.md @@ -0,0 +1,155 @@ +# R25 — Replay Production Qualification Harness + +R25 qualifies the R24 witnessed multi-host replay protocol against a real +PostgreSQL server and an external durable witness process instead of relying +only on deterministic fake backends. + +R25 is a qualification and hardening increment. It does not add merge, +deployment, runtime-execution, trading, wallet, or capital authority. + +## Baseline + +R25 starts from master +57c90dd4c7f8ec2a47a9d8b8a48e7d9fb1a11eaa. + +That baseline already includes: + +- R23 durable multi-host PostgreSQL CAS; +- R24 external witnessed rollback/tamper protection; +- Remote Commander PRs #197, #198, and #200. + +The Remote Commander files are disjoint from the replay qualification files. + +## Real qualification environment + +The local qualification harness uses: + +- Docker Engine 29.7.2; +- PostgreSQL 17-alpine in a dedicated container and named volume; +- psycopg 3.3.6 through the project multi-host-replay optional dependency; +- a separate Python 3.13 Alpine witness container and separate named volume; +- host ports published only on 127.0.0.1; +- at least two independent worker subprocesses. + +The witness fixture persists canonical JSONL records and fsyncs each append +before acknowledging it. It implements compare-and-set against exact +generation and head SHA-256. A deterministic post-fsync barrier allows the +controller to kill a worker in the crash window after witness durability but +before PostgreSQL reconciliation. + +The witness fixture is intentionally not a production witness service. It is +loopback HTTP without authentication or TLS and reports production_ready=false. + +## Fault matrix + +A green R25 local-container receipt requires all of these cases: + +1. Concurrent double claim + - eight independent workers consume the same approval concurrently; + - exactly one returns CLAIMED; + - exactly seven return ALREADY_CONSUMED; + - PostgreSQL contains one claim and one witnessed generation. + +2. PostgreSQL snapshot rollback + - capture a real pg_dump before the claim; + - consume the approval; + - restore PostgreSQL to the pre-claim database snapshot; + - the next claim rebuilds the missing suffix from the external witness; + - the result is ALREADY_CONSUMED, never a reopened capability. + +3. Crash after durable witness append + - the witness fsyncs the claim; + - the witness signals the deterministic post-append barrier; + - the controller kills the worker before PostgreSQL reconciliation; + - PostgreSQL still shows generation zero before recovery; + - the next worker reconstructs the claim and returns ALREADY_CONSUMED. + +4. Witness outage + - stop the witness container; + - replay authority fails closed; + - there is no fallback to ordinary PostgreSQL, SQLite, or memory; + - after witness recovery a fresh claim can proceed. + +5. PostgreSQL outage + - stop the PostgreSQL container; + - replay authority fails closed with a bounded PostgreSQL connection error; + - after PostgreSQL recovery a fresh claim can proceed. + +6. Witness rollback + - capture witness bytes before a claim; + - consume the approval so PostgreSQL advances; + - stop the witness and restore its old log bytes; + - restart the witness; + - synchronize fails closed because the database is ahead of the witness. + +7. PostgreSQL row tamper + - consume an approval; + - directly alter its stored digest in PostgreSQL; + - synchronize fails closed on the claim-set integrity audit. + +## R25 finding: real SERIALIZABLE contention + +The first real-PostgreSQL concurrency run exposed a production-hardening gap +that deterministic fake tests did not reproduce. + +With two concurrent workers, PostgreSQL SERIALIZABLE could abort one +reconciliation transaction. No double spend occurred: the other worker +recovered the witnessed claim and returned ALREADY_CONSUMED. However, the +desired response invariant of one CLAIMED plus one ALREADY_CONSUMED was not +met because the aborted transaction was surfaced as a generic fail-closed +synchronization error. + +R25 hardens witnessed replay with bounded retries for PostgreSQL SQLSTATE: + +- 40001 — serialization_failure; +- 40P01 — deadlock_detected. + +Only those retryable transaction failures are retried. Witness errors, +integrity errors, tamper evidence, connection failures, and all unknown +database failures remain fail-closed. + +The existing eight-attempt contention bound remains the outer limit. Exhaustion +still returns a replay contention error rather than widening authority. + +## Qualification status + +A successful harness run returns: + +LOCAL_CONTAINER_QUALIFICATION_GREEN + +This does **not** mean production_qualified_multi_host=true. + +The receipt must explicitly keep production_qualified_multi_host=false because: + +- PostgreSQL, witness, and workers still share one physical laptop; +- the qualification witness transport is unauthenticated loopback HTTP; +- service-stop tests are not packet-level network partitions; +- host-level rollback could correlate the two Docker named volumes. + +## Required next qualification + +A future promotion to real multi-host production qualification requires: + +- an independently administered witness in a different physical or failure + domain; +- at least two physical hosts or independent VMs for execution workers; +- authenticated and encrypted witness transport; +- packet-level partition and asymmetric network-fault injection; +- real operator witness rotation and recovery drill; +- explicit RTO/RPO and witness availability evidence. + +## Running locally + +The qualification must run from a clean R25 candidate for final evidence. + +On Windows, Docker Desktop launched through Remote Commander may need the +standard ProgramData environment value restored before starting Docker: + + ProgramData=C:\ProgramData + +The harness itself uses only ephemeral loopback ports, dynamically named +containers, and dynamically named Docker volumes. Unless --keep is supplied, +it removes its qualification containers and volumes in a finally block. + +The receipt is written outside the repository so producing evidence does not +mutate the candidate Git tree. diff --git a/tests/test_r25_replay_qualification.py b/tests/test_r25_replay_qualification.py new file mode 100644 index 0000000..9cf4794 --- /dev/null +++ b/tests/test_r25_replay_qualification.py @@ -0,0 +1,145 @@ +from __future__ import annotations + +import importlib.util +from pathlib import Path + +import pytest + +import continuityos.witnessed_approval_replay as replay + + +ROOT = Path(__file__).resolve().parents[1] + +if not ( + (ROOT / "tools" / "r25_witness_service.py").is_file() + and (ROOT / "tools" / "r25_replay_qualification.py").is_file() +): + pytest.skip( + "R25 qualification tools are repository-only and not packaged in the wheel", + allow_module_level=True, + ) + + +def _load(name: str, path: Path): + spec = importlib.util.spec_from_file_location(name, path) + assert spec is not None and spec.loader is not None + module = importlib.util.module_from_spec(spec) + spec.loader.exec_module(module) + return module + + +witness_service = _load( + "_r25_witness_service", ROOT / "tools" / "r25_witness_service.py" +) +qualification = _load( + "_r25_replay_qualification", ROOT / "tools" / "r25_replay_qualification.py" +) + + +def _record(namespace: str) -> dict: + return replay._build_record( + namespace=namespace, + generation=1, + previous_head_sha256=replay._genesis_head(namespace), + approval_id="hap_" + "a" * 64, + subject={"case": "r25-unit"}, + nonce="1" * 64, + digest_sha256="2" * 64, + ) + + +def _store(path: Path): + return witness_service.Store( + path, + pause_namespace=None, + pause_generation=None, + pause_marker=None, + release_marker=None, + pause_timeout=0.1, + ) + + +def test_qualification_witness_persists_and_reloads_fsynced_chain(tmp_path): + namespace = "r25-witness-persist" + path = tmp_path / "witness.jsonl" + record = _record(namespace) + store = _store(path) + + ok, observed = store.append( + namespace=namespace, + expected_generation=0, + expected_head_sha256=replay._genesis_head(namespace), + record=record, + ) + assert ok is True + assert observed == record + assert path.read_bytes().endswith(b"\n") + + reloaded = _store(path) + assert reloaded.current_state(namespace) == replay._state( + namespace, 1, record["head_sha256"] + ) + assert reloaded.records_after(namespace, 0) == [record] + + +def test_qualification_witness_cas_conflict_does_not_append(tmp_path): + namespace = "r25-witness-cas" + path = tmp_path / "witness.jsonl" + record = _record(namespace) + store = _store(path) + ok, _ = store.append( + namespace=namespace, + expected_generation=0, + expected_head_sha256=replay._genesis_head(namespace), + record=record, + ) + assert ok is True + before = path.read_bytes() + + ok, current = store.append( + namespace=namespace, + expected_generation=0, + expected_head_sha256=replay._genesis_head(namespace), + record=record, + ) + assert ok is False + assert current["generation"] == 1 + assert path.read_bytes() == before + + +def test_qualification_claim_identity_is_deterministic_and_scoped(): + one = qualification._claim_values("same") + two = qualification._claim_values("same") + other = qualification._claim_values("other") + + assert one == two + assert one != other + approval_id, subject, nonce, digest = one + assert approval_id.startswith("hap_") + assert len(approval_id) == 68 + assert len(nonce) == 64 + assert len(digest) == 64 + assert subject["qualification_case"] == "same" + + +def test_qualification_receipt_contract_cannot_claim_physical_multi_host(): + source = ( + ROOT / "tools" / "r25_replay_qualification.py" + ).read_text(encoding="utf-8") + assert '"production_qualified_multi_host": False' in source + assert '"same_physical_host": True' in source + assert '"worker_subprocess_domains": worker_count' in source + assert '"qualification_harness_sha256": harness_sha256' in source + assert '"witness_fixture_sha256": witness_fixture_sha256' in source + assert '"merge": False' in source + assert '"deploy": False' in source + assert '"can_trade": False' in source + assert '"capital_permission": "DENY"' in source + + +def test_qualification_witness_health_declares_not_production_ready(): + source = ( + ROOT / "tools" / "r25_witness_service.py" + ).read_text(encoding="utf-8") + assert '"production_ready": False' in source + assert 'durability": "jsonl_fsync"' in source diff --git a/tests/test_witnessed_approval_replay.py b/tests/test_witnessed_approval_replay.py index 4787693..440d7d2 100644 --- a/tests/test_witnessed_approval_replay.py +++ b/tests/test_witnessed_approval_replay.py @@ -23,6 +23,14 @@ } +class RetryableTransactionError(RuntimeError): + sqlstate = "40001" + + +class DeadlockTransactionError(RuntimeError): + sqlstate = "40P01" + + class DbState: def __init__(self) -> None: self.lock = threading.RLock() @@ -31,6 +39,7 @@ def __init__(self) -> None: 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.retryable_commit_failures = 0 self.active_transactions = 0 @@ -55,6 +64,10 @@ def cursor(self): return FakeCursor(self) def commit(self) -> None: + if self.state.retryable_commit_failures > 0: + self.state.retryable_commit_failures -= 1 + self._release() + raise RetryableTransactionError("simulated serialization failure") if self.state.fail_next_commit: self.state.fail_next_commit = False self._release() @@ -632,3 +645,28 @@ def test_namespaces_are_isolated(): assert claim(right)["status"] == CLAIMED assert witness.current_state("human-approval-a")["generation"] == 1 assert witness.current_state("human-approval-b")["generation"] == 1 + + +def test_retryable_transaction_error_detection_is_sqlstate_scoped(): + assert replay._is_retryable_transaction_error( + RetryableTransactionError("serialization") + ) + assert replay._is_retryable_transaction_error( + DeadlockTransactionError("deadlock") + ) + outer = RuntimeError("wrapper") + outer.__cause__ = RetryableTransactionError("nested serialization") + assert replay._is_retryable_transaction_error(outer) + assert not replay._is_retryable_transaction_error( + RuntimeError("ordinary failure") + ) + + +def test_snapshot_retries_retryable_serialization_failure(): + db = DbState() + witness = FakeWitness(db) + guard = authority(db, witness) + db.retryable_commit_failures = 1 + state = guard.synchronize() + assert state["generation"] == 0 + assert db.retryable_commit_failures == 0 diff --git a/tools/r25_replay_qualification.py b/tools/r25_replay_qualification.py new file mode 100644 index 0000000..5a66f99 --- /dev/null +++ b/tools/r25_replay_qualification.py @@ -0,0 +1,979 @@ +"""R25 real PostgreSQL + external witness qualification harness.""" +from __future__ import annotations + +import argparse +import hashlib +import json +import os +import socket +import subprocess +import sys +import tempfile +import time +import urllib.error +import urllib.parse +import urllib.request +from pathlib import Path +from typing import Any + +from continuityos.witnessed_approval_replay import ( + ALREADY_CONSUMED, + CLAIMED, + PostgresWitnessedApprovalReplayAuthority, + WITNESS_SCOPE, +) + +RECEIPT_SCHEMA = "continuityos.r25.replay_production_qualification/v1" +GREEN_STATUS = "LOCAL_CONTAINER_QUALIFICATION_GREEN" + + +def _canonical_json(value: object) -> str: + return json.dumps( + value, sort_keys=True, separators=(",", ":"), ensure_ascii=True + ) + + +def _hash(text: str) -> str: + return hashlib.sha256(text.encode("ascii")).hexdigest() + + +def _free_port() -> int: + with socket.socket(socket.AF_INET, socket.SOCK_STREAM) as sock: + sock.bind(("127.0.0.1", 0)) + return int(sock.getsockname()[1]) + + +def _run( + args: list[str], + *, + check: bool = True, + input_bytes: bytes | None = None, + text: bool = True, +) -> subprocess.CompletedProcess: + env = dict(os.environ) + if os.name == "nt": + env.setdefault("ProgramData", r"C:\ProgramData") + result = subprocess.run( + args, + input=input_bytes if not text else None, + stdout=subprocess.PIPE, + stderr=subprocess.PIPE, + text=text, + env=env, + check=False, + ) + if check and result.returncode != 0: + stdout = result.stdout if text else result.stdout.decode("utf-8", "replace") + stderr = result.stderr if text else result.stderr.decode("utf-8", "replace") + raise RuntimeError( + f"command failed ({result.returncode}): {args!r}\n" + f"stdout={stdout}\nstderr={stderr}" + ) + return result + + +def _docker(*args: str, check: bool = True) -> subprocess.CompletedProcess: + return _run(["docker", *args], check=check) + + +def _wait_http(url: str, timeout: float = 30.0) -> dict[str, Any]: + deadline = time.monotonic() + timeout + last: Exception | None = None + while time.monotonic() < deadline: + try: + with urllib.request.urlopen(url, timeout=1.0) as response: + value = json.loads(response.read().decode("ascii")) + if type(value) is dict: + return value + except Exception as exc: + last = exc + time.sleep(0.2) + raise RuntimeError(f"HTTP endpoint not ready: {url}: {last}") + + +def _wait_pg(container: str, timeout: float = 45.0) -> None: + deadline = time.monotonic() + timeout + while time.monotonic() < deadline: + result = _docker( + "exec", container, "pg_isready", "-U", "postgres", + "-d", "continuityos_r25", check=False, + ) + if result.returncode == 0: + return + time.sleep(0.5) + raise RuntimeError("PostgreSQL did not become ready") + + +class HttpWitness: + witness_scope = WITNESS_SCOPE + + def __init__(self, base_url: str, *, timeout: float = 2.0) -> None: + self.base_url = base_url.rstrip("/") + self.timeout = timeout + + def _get(self, path: str, query: dict[str, object]) -> Any: + encoded = urllib.parse.urlencode(query) + request = urllib.request.Request( + f"{self.base_url}{path}?{encoded}", method="GET" + ) + try: + with urllib.request.urlopen(request, timeout=self.timeout) as response: + return json.loads(response.read().decode("ascii")) + except Exception as exc: + raise RuntimeError(f"qualification witness unavailable: {exc}") from exc + + def current_state(self, namespace: str) -> dict: + value = self._get("/state", {"namespace": namespace}) + if type(value) is not dict: + raise RuntimeError("qualification witness state invalid") + return value + + def records_after(self, namespace: str, generation: int) -> list[dict]: + value = self._get("/records", {"namespace": namespace, "after": generation}) + if type(value) is not list: + raise RuntimeError("qualification witness history invalid") + return value + + def append_record( + self, namespace: str, expected_generation: int, + expected_head_sha256: str, record: dict, + ) -> dict: + payload = _canonical_json({ + "namespace": namespace, + "expected_generation": expected_generation, + "expected_head_sha256": expected_head_sha256, + "record": record, + }).encode("ascii") + request = urllib.request.Request( + f"{self.base_url}/append", data=payload, method="POST", + headers={"Content-Type": "application/json"}, + ) + try: + with urllib.request.urlopen(request, timeout=self.timeout) as response: + value = json.loads(response.read().decode("ascii")) + except urllib.error.HTTPError as exc: + body = exc.read().decode("ascii", "replace") + raise RuntimeError( + f"qualification witness append rejected: {exc.code}: {body}" + ) from exc + except Exception as exc: + raise RuntimeError( + f"qualification witness append unavailable: {exc}" + ) from exc + if type(value) is not dict: + raise RuntimeError("qualification witness append response invalid") + return value + + +def _claim_values(case: str) -> tuple[str, dict, str, str]: + return ( + "hap_" + _hash(f"approval:{case}"), + { + "repository": "bitmaster162/continuityos", + "baseline_sha": _hash(f"base:{case}")[:40], + "candidate_sha": _hash(f"head:{case}")[:40], + "candidate_tree_sha": _hash(f"tree:{case}")[:40], + "qualification_case": case, + }, + _hash(f"nonce:{case}"), + _hash(f"digest:{case}"), + ) + + +def _authority( + dsn: str, witness_url: str, namespace: str +) -> PostgresWitnessedApprovalReplayAuthority: + return PostgresWitnessedApprovalReplayAuthority( + dsn, witness=HttpWitness(witness_url), namespace=namespace + ) + + +def _worker_sync(args: argparse.Namespace) -> int: + try: + state = _authority(args.dsn, args.witness_url, args.namespace).synchronize() + print(_canonical_json({"ok": True, "state": state})) + return 0 + except Exception as exc: + print(_canonical_json({ + "ok": False, + "error": type(exc).__name__, + "detail": str(exc), + })) + return 2 + + +def _worker_claim(args: argparse.Namespace) -> int: + try: + approval_id, subject, nonce, digest = _claim_values(args.case) + receipt = _authority( + args.dsn, args.witness_url, args.namespace + ).claim_once( + approval_id=approval_id, + subject=subject, + nonce=nonce, + digest_sha256=digest, + ) + print(_canonical_json({"ok": True, "receipt": receipt})) + return 0 + except Exception as exc: + print(_canonical_json({ + "ok": False, + "error": type(exc).__name__, + "detail": str(exc), + })) + return 2 + + +def _worker_command( + python: str, + script: Path, + *, + mode: str, + dsn: str, + witness_url: str, + namespace: str, + case: str | None = None, +) -> list[str]: + command = [ + python, + str(script), + mode, + "--dsn", + dsn, + "--witness-url", + witness_url, + "--namespace", + namespace, + ] + if case is not None: + command.extend(["--case", case]) + return command + + +def _worker_run( + python: str, + script: Path, + *, + mode: str, + dsn: str, + witness_url: str, + namespace: str, + case: str | None = None, +) -> tuple[int, dict]: + command = _worker_command( + python, + script, + mode=mode, + dsn=dsn, + witness_url=witness_url, + namespace=namespace, + case=case, + ) + result = _run(command, check=False) + lines = [line for line in result.stdout.splitlines() if line.strip()] + if not lines: + return result.returncode, { + "ok": False, + "error": "NoWorkerReceipt", + "detail": result.stderr.strip(), + } + try: + value = json.loads(lines[-1]) + except json.JSONDecodeError: + value = { + "ok": False, + "error": "InvalidWorkerReceipt", + "detail": result.stdout[-2000:], + } + return result.returncode, value + + +def _spawn_worker( + python: str, + script: Path, + *, + mode: str, + dsn: str, + witness_url: str, + namespace: str, + case: str | None = None, +) -> subprocess.Popen: + command = _worker_command( + python, + script, + mode=mode, + dsn=dsn, + witness_url=witness_url, + namespace=namespace, + case=case, + ) + env = dict(os.environ) + if os.name == "nt": + env.setdefault("ProgramData", r"C:\ProgramData") + return subprocess.Popen( + command, + stdout=subprocess.PIPE, + stderr=subprocess.PIPE, + text=True, + env=env, + ) + + +def _collect_worker( + process: subprocess.Popen, timeout: float = 45.0 +) -> tuple[int, dict]: + stdout, stderr = process.communicate(timeout=timeout) + lines = [line for line in stdout.splitlines() if line.strip()] + if not lines: + return process.returncode or 0, { + "ok": False, + "error": "NoWorkerReceipt", + "detail": stderr.strip(), + } + try: + value = json.loads(lines[-1]) + except json.JSONDecodeError: + value = { + "ok": False, + "error": "InvalidWorkerReceipt", + "detail": stdout[-2000:], + } + return process.returncode or 0, value + + +def _status(receipt: dict) -> str | None: + if receipt.get("ok") is not True: + return None + inner = receipt.get("receipt") + return inner.get("status") if type(inner) is dict else None + + +def _db_namespace_state(dsn: str, namespace: str) -> dict[str, Any]: + import psycopg + + with psycopg.connect(dsn) as connection: + with connection.cursor() as cursor: + cursor.execute( + "SELECT generation, head_sha256 " + "FROM continuityos_replay_witness_state WHERE namespace=%s", + (namespace,), + ) + row = cursor.fetchone() + cursor.execute( + "SELECT COUNT(*) FROM continuityos_approval_claims " + "WHERE namespace=%s", + (namespace,), + ) + count = cursor.fetchone()[0] + return { + "generation": None if row is None else int(row[0]), + "head_sha256": None if row is None else str(row[1]), + "claim_count": int(count), + } + + +def _tamper_digest(dsn: str, namespace: str) -> None: + import psycopg + + with psycopg.connect(dsn) as connection: + with connection.cursor() as cursor: + cursor.execute( + "UPDATE continuityos_approval_claims " + "SET digest_sha256=%s WHERE namespace=%s", + ("f" * 64, namespace), + ) + if cursor.rowcount != 1: + raise RuntimeError("expected exactly one claim row to tamper") + connection.commit() + + +def _pg_version(dsn: str) -> str: + import psycopg + + with psycopg.connect(dsn) as connection: + with connection.cursor() as cursor: + cursor.execute("SELECT version()") + return str(cursor.fetchone()[0]) + + +def _snapshot_database(container: str) -> bytes: + result = _run( + [ + "docker", + "exec", + container, + "pg_dump", + "-U", + "postgres", + "-d", + "continuityos_r25", + "-Fc", + ], + text=False, + ) + return bytes(result.stdout) + + +def _restore_database(container: str, dump: bytes) -> None: + _docker( + "exec", + container, + "psql", + "-U", + "postgres", + "-d", + "postgres", + "-c", + "SELECT pg_terminate_backend(pid) FROM pg_stat_activity " + "WHERE datname='continuityos_r25' AND pid <> pg_backend_pid();", + ) + _docker( + "exec", + container, + "dropdb", + "-U", + "postgres", + "--if-exists", + "continuityos_r25", + ) + _docker("exec", container, "createdb", "-U", "postgres", "continuityos_r25") + result = _run( + [ + "docker", + "exec", + "-i", + container, + "pg_restore", + "-U", + "postgres", + "-d", + "continuityos_r25", + ], + input_bytes=dump, + text=False, + check=False, + ) + if result.returncode != 0: + raise RuntimeError( + "pg_restore failed: " + + result.stderr.decode("utf-8", "replace")[-4000:] + ) + + +def _volume_bytes(volume: str) -> bytes: + result = _run( + [ + "docker", + "run", + "--rm", + "-v", + f"{volume}:/data", + "python:3.13-alpine", + "sh", + "-c", + "cat /data/witness.jsonl 2>/dev/null || true", + ], + text=False, + ) + return bytes(result.stdout) + + +def _restore_volume_bytes(volume: str, content: bytes) -> None: + result = _run( + [ + "docker", + "run", + "--rm", + "-i", + "-v", + f"{volume}:/data", + "python:3.13-alpine", + "sh", + "-c", + "cat > /data/witness.jsonl", + ], + input_bytes=content, + text=False, + check=False, + ) + if result.returncode != 0: + raise RuntimeError( + "witness volume restore failed: " + + result.stderr.decode("utf-8", "replace")[-4000:] + ) + + +def _assert_worker_failure(code: int, value: dict, label: str) -> dict: + if code == 0 or value.get("ok") is not False: + raise AssertionError( + f"{label} did not fail closed: code={code} value={value}" + ) + return { + "status": "PASS", + "worker_exit": code, + "error": value.get("error"), + "detail": str(value.get("detail", ""))[:1000], + } + + +def _qualify(args: argparse.Namespace) -> int: + import psycopg + + repo = Path(args.repo_root).resolve() + script = Path(__file__).resolve() + witness_script = repo / "tools" / "r25_witness_service.py" + head = _run(["git", "-C", str(repo), "rev-parse", "HEAD"]).stdout.strip() + tree = _run( + ["git", "-C", str(repo), "rev-parse", "HEAD^{tree}"] + ).stdout.strip() + dirty = _run( + ["git", "-C", str(repo), "status", "--porcelain"] + ).stdout.strip() + if dirty and not args.allow_dirty: + raise RuntimeError("R25 qualification requires a clean Git worktree") + if args.expected_head and head != args.expected_head: + raise RuntimeError("Git HEAD does not match --expected-head") + + docker_version = _docker( + "version", "--format", "{{.Server.Version}}" + ).stdout.strip() + postgres_image = json.loads( + _docker("image", "inspect", "postgres:17-alpine").stdout + )[0] + witness_image = json.loads( + _docker("image", "inspect", "python:3.13-alpine").stdout + )[0] + harness_sha256 = hashlib.sha256(script.read_bytes()).hexdigest() + witness_fixture_sha256 = hashlib.sha256( + witness_script.read_bytes() + ).hexdigest() + run_id = f"{os.getpid()}-{int(time.time())}" + pg_name = f"continuityos-r25-pg-{run_id}" + witness_name = f"continuityos-r25-witness-{run_id}" + pg_volume = f"continuityos-r25-pg-{run_id}" + witness_volume = f"continuityos-r25-witness-{run_id}" + pg_port = _free_port() + witness_port = _free_port() + dsn = ( + f"postgresql://postgres@127.0.0.1:{pg_port}/continuityos_r25?connect_timeout=3" + ) + witness_url = f"http://127.0.0.1:{witness_port}" + work = ( + Path(args.work_dir).resolve() + if args.work_dir + else Path(tempfile.mkdtemp(prefix="continuityos-r25-")) + ) + work.mkdir(parents=True, exist_ok=True) + control = work / "control" + control.mkdir(parents=True, exist_ok=True) + pause_marker = control / "crash.marker" + release_marker = control / "crash.release" + crash_namespace = f"r25-crash-{run_id}" + cases: dict[str, Any] = {} + started_pg = False + started_witness = False + + def start_witness() -> None: + nonlocal started_witness + inspect = _docker("container", "inspect", witness_name, check=False) + if inspect.returncode == 0: + _docker("start", witness_name) + else: + _docker( + "run", + "-d", + "--name", + witness_name, + "-p", + f"127.0.0.1:{witness_port}:8080", + "-v", + f"{witness_volume}:/data", + "-v", + f"{witness_script}:/app/witness.py:ro", + "-v", + f"{control}:/control", + "python:3.13-alpine", + "python", + "/app/witness.py", + "--host", + "0.0.0.0", + "--port", + "8080", + "--data-file", + "/data/witness.jsonl", + "--pause-namespace", + crash_namespace, + "--pause-generation", + "1", + "--pause-marker", + "/control/crash.marker", + "--release-marker", + "/control/crash.release", + "--pause-timeout", + "60", + ) + started_witness = True + _wait_http(f"{witness_url}/health") + + def sync_namespace(namespace: str) -> dict: + code, value = _worker_run( + args.python, + script, + mode="sync", + dsn=dsn, + witness_url=witness_url, + namespace=namespace, + ) + if code != 0 or value.get("ok") is not True: + raise AssertionError(f"sync failed for {namespace}: {value}") + return value + + def claim_namespace(namespace: str, case: str) -> tuple[int, dict]: + return _worker_run( + args.python, + script, + mode="claim", + dsn=dsn, + witness_url=witness_url, + namespace=namespace, + case=case, + ) + + try: + _docker("volume", "create", pg_volume) + _docker("volume", "create", witness_volume) + _docker( + "run", + "-d", + "--name", + pg_name, + "-e", + "POSTGRES_HOST_AUTH_METHOD=trust", + "-e", + "POSTGRES_DB=continuityos_r25", + "-p", + f"127.0.0.1:{pg_port}:5432", + "-v", + f"{pg_volume}:/var/lib/postgresql/data", + "postgres:17-alpine", + ) + started_pg = True + _wait_pg(pg_name) + start_witness() + + postgres_version = _pg_version(dsn) + cases["preflight"] = { + "status": "PASS", + "docker_server": docker_version, + "postgres": postgres_version, + "psycopg": psycopg.__version__, + "witness_health": _wait_http(f"{witness_url}/health"), + } + + concurrency_ns = f"r25-concurrency-{run_id}" + sync_namespace(concurrency_ns) + worker_count = 8 + workers = [ + _spawn_worker( + args.python, + script, + mode="claim", + dsn=dsn, + witness_url=witness_url, + namespace=concurrency_ns, + case="concurrent", + ) + for _index in range(worker_count) + ] + results = [_collect_worker(worker) for worker in workers] + statuses = [_status(value) for _code, value in results] + if ( + statuses.count(CLAIMED) != 1 + or statuses.count(ALREADY_CONSUMED) != worker_count - 1 + ): + raise AssertionError( + f"concurrent claim invariant failed: {results}" + ) + cases["concurrent_double_claim"] = { + "status": "PASS", + "worker_processes": worker_count, + "claim_statuses": sorted(str(item) for item in statuses), + "db_state": _db_namespace_state(dsn, concurrency_ns), + } + + + snapshot_ns = f"r25-snapshot-{run_id}" + sync_namespace(snapshot_ns) + db_dump = _snapshot_database(pg_name) + code, first = claim_namespace(snapshot_ns, "snapshot") + if code != 0 or _status(first) != CLAIMED: + raise AssertionError(f"snapshot initial claim failed: {first}") + claimed_state = _db_namespace_state(dsn, snapshot_ns) + _restore_database(pg_name, db_dump) + restored_state = _db_namespace_state(dsn, snapshot_ns) + code, recovered = claim_namespace(snapshot_ns, "snapshot") + if code != 0 or _status(recovered) != ALREADY_CONSUMED: + raise AssertionError(f"snapshot recovery failed: {recovered}") + cases["postgres_snapshot_rollback"] = { + "status": "PASS", + "claimed_state": claimed_state, + "restored_pre_recovery_state": restored_state, + "recovered_status": _status(recovered), + "post_recovery_state": _db_namespace_state(dsn, snapshot_ns), + } + + sync_namespace(crash_namespace) + pause_marker.unlink(missing_ok=True) + release_marker.unlink(missing_ok=True) + crash_worker = _spawn_worker( + args.python, + script, + mode="claim", + dsn=dsn, + witness_url=witness_url, + namespace=crash_namespace, + case="crash-after-witness", + ) + deadline = time.monotonic() + 30.0 + while time.monotonic() < deadline and not pause_marker.exists(): + if crash_worker.poll() is not None: + raise AssertionError( + "crash worker exited before durable witness barrier" + ) + time.sleep(0.05) + if not pause_marker.exists(): + crash_worker.kill() + raise AssertionError("durable witness barrier was not reached") + crash_worker.kill() + crash_worker.wait(timeout=5) + pre_recovery_crash_state = _db_namespace_state( + dsn, crash_namespace + ) + release_marker.write_text("release", encoding="utf-8") + code, crash_recovered = claim_namespace( + crash_namespace, "crash-after-witness" + ) + if ( + code != 0 + or _status(crash_recovered) != ALREADY_CONSUMED + ): + raise AssertionError( + f"crash recovery failed: {crash_recovered}" + ) + cases["crash_after_witness_append"] = { + "status": "PASS", + "worker_killed_after_witness_fsync": True, + "db_before_recovery": pre_recovery_crash_state, + "recovered_status": _status(crash_recovered), + "db_after_recovery": _db_namespace_state( + dsn, crash_namespace + ), + } + + witness_outage_ns = f"r25-witness-outage-{run_id}" + sync_namespace(witness_outage_ns) + _docker("stop", witness_name) + started_witness = False + code, outage = claim_namespace( + witness_outage_ns, "witness-outage" + ) + cases["witness_outage_fail_closed"] = _assert_worker_failure( + code, outage, "witness outage" + ) + _docker("start", witness_name) + started_witness = True + _wait_http(f"{witness_url}/health") + code, after_outage = claim_namespace( + witness_outage_ns, "witness-outage" + ) + if code != 0 or _status(after_outage) != CLAIMED: + raise AssertionError( + f"post-witness-outage recovery failed: {after_outage}" + ) + cases["witness_outage_fail_closed"][ + "post_recovery_status" + ] = _status(after_outage) + + pg_outage_ns = f"r25-pg-outage-{run_id}" + sync_namespace(pg_outage_ns) + _docker("stop", pg_name) + started_pg = False + code, pg_outage = claim_namespace(pg_outage_ns, "pg-outage") + cases["postgres_outage_fail_closed"] = _assert_worker_failure( + code, pg_outage, "PostgreSQL outage" + ) + _docker("start", pg_name) + started_pg = True + _wait_pg(pg_name) + code, pg_after = claim_namespace(pg_outage_ns, "pg-outage") + if code != 0 or _status(pg_after) != CLAIMED: + raise AssertionError( + f"post-PostgreSQL-outage recovery failed: {pg_after}" + ) + cases["postgres_outage_fail_closed"][ + "post_recovery_status" + ] = _status(pg_after) + + + rollback_ns = f"r25-witness-rollback-{run_id}" + sync_namespace(rollback_ns) + witness_before = _volume_bytes(witness_volume) + code, rollback_claim = claim_namespace( + rollback_ns, "witness-rollback" + ) + if code != 0 or _status(rollback_claim) != CLAIMED: + raise AssertionError( + f"witness rollback setup failed: {rollback_claim}" + ) + _docker("stop", witness_name) + started_witness = False + _restore_volume_bytes(witness_volume, witness_before) + _docker("start", witness_name) + started_witness = True + _wait_http(f"{witness_url}/health") + code, rollback_detected = _worker_run( + args.python, + script, + mode="sync", + dsn=dsn, + witness_url=witness_url, + namespace=rollback_ns, + ) + cases["witness_rollback_fail_closed"] = _assert_worker_failure( + code, rollback_detected, "witness rollback" + ) + + tamper_ns = f"r25-db-tamper-{run_id}" + sync_namespace(tamper_ns) + code, tamper_claim = claim_namespace(tamper_ns, "db-tamper") + if code != 0 or _status(tamper_claim) != CLAIMED: + raise AssertionError( + f"tamper setup claim failed: {tamper_claim}" + ) + _tamper_digest(dsn, tamper_ns) + code, tamper_detected = _worker_run( + args.python, + script, + mode="sync", + dsn=dsn, + witness_url=witness_url, + namespace=tamper_ns, + ) + cases["postgres_row_tamper_fail_closed"] = _assert_worker_failure( + code, tamper_detected, "PostgreSQL row tamper" + ) + + receipt = { + "schema": RECEIPT_SCHEMA, + "status": GREEN_STATUS, + "git_head": head, + "git_tree": tree, + "environment": { + "docker_server": docker_version, + "postgres_image": "postgres:17-alpine", + "postgres_image_id": postgres_image["Id"], + "postgres_repo_digests": postgres_image.get("RepoDigests", []), + "python_witness_image": "python:3.13-alpine", + "python_witness_image_id": witness_image["Id"], + "python_witness_repo_digests": witness_image.get("RepoDigests", []), + "postgres_version": postgres_version, + "psycopg_version": psycopg.__version__, + "qualification_harness_sha256": harness_sha256, + "witness_fixture_sha256": witness_fixture_sha256, + "postgres_container_domain": True, + "witness_container_domain": True, + "worker_subprocess_domains": worker_count, + "separate_named_volumes": True, + "host_loopback_only_ports": True, + "same_physical_host": True, + }, + "cases": cases, + "authority": { + "merge": False, + "deploy": False, + "runtime_execution": False, + "can_trade": False, + "capital_permission": "DENY", + }, + "qualification_boundary": { + "production_qualified_multi_host": False, + "reasons": [ + "PostgreSQL, witness, and workers share one physical host", + "qualification witness is unauthenticated loopback HTTP", + "service-stop outage is not packet-level network partition", + "Docker-host rollback could correlate both named volumes", + ], + "required_next": [ + "independent physical or failure-domain witness", + "at least two physical or VM execution hosts", + "authenticated and encrypted witness transport", + "packet-level network partition injection", + "operator witness rotation and recovery drill", + ], + }, + } + receipt_path = Path(args.receipt).resolve() + receipt_path.parent.mkdir(parents=True, exist_ok=True) + receipt_path.write_text( + json.dumps(receipt, indent=2, sort_keys=True) + "\n", + encoding="utf-8", + ) + print(_canonical_json(receipt)) + return 0 + finally: + if not args.keep: + if ( + started_witness + or _docker( + "container", "inspect", witness_name, check=False + ).returncode == 0 + ): + _docker("rm", "-f", witness_name, check=False) + if ( + started_pg + or _docker("container", "inspect", pg_name, check=False).returncode == 0 + ): + _docker("rm", "-f", pg_name, check=False) + _docker("volume", "rm", "-f", witness_volume, check=False) + _docker("volume", "rm", "-f", pg_volume, check=False) + + +def _parser() -> argparse.ArgumentParser: + parser = argparse.ArgumentParser() + sub = parser.add_subparsers(dest="mode", required=True) + + sync = sub.add_parser("sync") + sync.add_argument("--dsn", required=True) + sync.add_argument("--witness-url", required=True) + sync.add_argument("--namespace", required=True) + + claim = sub.add_parser("claim") + claim.add_argument("--dsn", required=True) + claim.add_argument("--witness-url", required=True) + claim.add_argument("--namespace", required=True) + claim.add_argument("--case", required=True) + + qualify = sub.add_parser("qualify") + qualify.add_argument("--repo-root", default=".") + qualify.add_argument("--python", default=sys.executable) + qualify.add_argument("--expected-head") + qualify.add_argument("--receipt", required=True) + qualify.add_argument("--work-dir") + qualify.add_argument("--allow-dirty", action="store_true") + qualify.add_argument("--keep", action="store_true") + return parser + + +def main() -> int: + args = _parser().parse_args() + if args.mode == "sync": + return _worker_sync(args) + if args.mode == "claim": + return _worker_claim(args) + if args.mode == "qualify": + return _qualify(args) + raise AssertionError(args.mode) + + +if __name__ == "__main__": + raise SystemExit(main()) diff --git a/tools/r25_witness_service.py b/tools/r25_witness_service.py new file mode 100644 index 0000000..3901b66 --- /dev/null +++ b/tools/r25_witness_service.py @@ -0,0 +1,281 @@ +"""R25 qualification-only append-only witness service. + +This service is intentionally not a production witness implementation. It is a +loopback-published Docker qualification fixture with durable JSONL append+fsync +semantics and a deterministic post-append barrier for crash injection. +""" +from __future__ import annotations + +import argparse +import hashlib +import json +import os +import threading +import time +from http import HTTPStatus +from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer +from pathlib import Path +from urllib.parse import parse_qs, urlparse + +STATE_SCHEMA = "continuityos.replay_witness_state/v1" +GENESIS_DOMAIN = "continuityos.replay_witness_genesis/v1" + + +def _canonical_json(value: object) -> str: + return json.dumps( + value, sort_keys=True, separators=(",", ":"), ensure_ascii=True + ) + + +def _sha256_text(value: str) -> str: + return hashlib.sha256(value.encode("ascii")).hexdigest() + + +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: + return { + "schema": STATE_SCHEMA, + "namespace": namespace, + "generation": generation, + "head_sha256": head_sha256, + } + + +class Store: + def __init__( + self, + path: Path, + *, + pause_namespace: str | None, + pause_generation: int | None, + pause_marker: Path | None, + release_marker: Path | None, + pause_timeout: float, + ) -> None: + self.path = path + self.path.parent.mkdir(parents=True, exist_ok=True) + self.lock = threading.RLock() + self.records: dict[str, list[dict]] = {} + self.pause_namespace = pause_namespace + self.pause_generation = pause_generation + self.pause_marker = pause_marker + self.release_marker = release_marker + self.pause_timeout = pause_timeout + self._pause_used = False + if self.path.exists(): + with self.path.open("r", encoding="utf-8") as handle: + for lineno, raw in enumerate(handle, start=1): + raw = raw.strip() + if not raw: + continue + item = json.loads(raw) + if ( + type(item) is not dict + or set(item) != {"namespace", "record"} + or type(item["namespace"]) is not str + or type(item["record"]) is not dict + ): + raise RuntimeError( + f"invalid witness log entry at line {lineno}" + ) + self.records.setdefault(item["namespace"], []).append( + item["record"] + ) + + def current_state(self, namespace: str) -> dict: + with self.lock: + rows = self.records.get(namespace, []) + if not rows: + return _state(namespace, 0, _genesis_head(namespace)) + last = rows[-1] + return _state( + namespace, int(last["generation"]), str(last["head_sha256"]) + ) + + def records_after(self, namespace: str, generation: int) -> list[dict]: + with self.lock: + return [ + json.loads(_canonical_json(row)) + for row in self.records.get(namespace, [])[generation:] + ] + + def append( + self, + *, + namespace: str, + expected_generation: int, + expected_head_sha256: str, + record: dict, + ) -> tuple[bool, dict]: + with self.lock: + current = self.current_state(namespace) + if ( + current["generation"] != expected_generation + or current["head_sha256"] != expected_head_sha256 + ): + return False, current + if ( + record.get("namespace") != namespace + or record.get("generation") != expected_generation + 1 + or record.get("previous_head_sha256") != expected_head_sha256 + ): + raise ValueError("record does not extend expected witness state") + + item = {"namespace": namespace, "record": record} + encoded = (_canonical_json(item) + "\n").encode("ascii") + with self.path.open("ab", buffering=0) as handle: + handle.write(encoded) + os.fsync(handle.fileno()) + self.records.setdefault(namespace, []).append( + json.loads(_canonical_json(record)) + ) + + should_pause = ( + not self._pause_used + and self.pause_namespace == namespace + and self.pause_generation == record.get("generation") + ) + if should_pause: + self._pause_used = True + if self.pause_marker is not None: + self.pause_marker.parent.mkdir(parents=True, exist_ok=True) + self.pause_marker.write_text( + str(record.get("record_id", "")), encoding="utf-8" + ) + + if should_pause: + deadline = time.monotonic() + self.pause_timeout + while time.monotonic() < deadline: + if self.release_marker is not None and self.release_marker.exists(): + break + time.sleep(0.05) + + return True, json.loads(_canonical_json(record)) + + +class Handler(BaseHTTPRequestHandler): + server_version = "ContinuityOSR25Witness/1" + + @property + def store(self) -> Store: + return self.server.store # type: ignore[attr-defined] + + def log_message(self, fmt: str, *args: object) -> None: + return + + def _send(self, status: int, value: object) -> None: + payload = (_canonical_json(value) + "\n").encode("ascii") + self.send_response(status) + self.send_header("Content-Type", "application/json") + self.send_header("Content-Length", str(len(payload))) + self.end_headers() + try: + self.wfile.write(payload) + except (BrokenPipeError, ConnectionResetError): + pass + + def _query(self) -> tuple[str, dict[str, list[str]]]: + parsed = urlparse(self.path) + return parsed.path, parse_qs(parsed.query, strict_parsing=True) + + def do_GET(self) -> None: + try: + path, query = self._query() + if path == "/health": + self._send( + HTTPStatus.OK, + { + "ok": True, + "service": "continuityos-r25-qualification-witness", + "durability": "jsonl_fsync", + "production_ready": False, + }, + ) + return + if path == "/state": + namespace = query["namespace"][0] + self._send(HTTPStatus.OK, self.store.current_state(namespace)) + return + if path == "/records": + namespace = query["namespace"][0] + after = int(query["after"][0]) + if after < 0: + raise ValueError("after must be non-negative") + self._send( + HTTPStatus.OK, self.store.records_after(namespace, after) + ) + return + self._send(HTTPStatus.NOT_FOUND, {"error": "not found"}) + except Exception as exc: + self._send( + HTTPStatus.BAD_REQUEST, + {"error": type(exc).__name__, "detail": str(exc)}, + ) + + def do_POST(self) -> None: + try: + path, _query = self._query() + if path != "/append": + self._send(HTTPStatus.NOT_FOUND, {"error": "not found"}) + return + length = int(self.headers.get("Content-Length", "0")) + if length <= 0 or length > 1_000_000: + raise ValueError("invalid content length") + body = json.loads(self.rfile.read(length).decode("ascii")) + if type(body) is not dict: + raise ValueError("body must be an object") + ok, value = self.store.append( + namespace=str(body["namespace"]), + expected_generation=int(body["expected_generation"]), + expected_head_sha256=str(body["expected_head_sha256"]), + record=dict(body["record"]), + ) + if not ok: + self._send(HTTPStatus.CONFLICT, value) + return + self._send(HTTPStatus.OK, value) + except Exception as exc: + self._send( + HTTPStatus.BAD_REQUEST, + {"error": type(exc).__name__, "detail": str(exc)}, + ) + + +def main() -> int: + parser = argparse.ArgumentParser() + parser.add_argument("--host", default="0.0.0.0") + parser.add_argument("--port", type=int, default=8080) + parser.add_argument("--data-file", required=True) + parser.add_argument("--pause-namespace") + parser.add_argument("--pause-generation", type=int) + parser.add_argument("--pause-marker") + parser.add_argument("--release-marker") + parser.add_argument("--pause-timeout", type=float, default=60.0) + args = parser.parse_args() + + store = Store( + Path(args.data_file), + pause_namespace=args.pause_namespace, + pause_generation=args.pause_generation, + pause_marker=Path(args.pause_marker) if args.pause_marker else None, + release_marker=Path(args.release_marker) if args.release_marker else None, + pause_timeout=args.pause_timeout, + ) + server = ThreadingHTTPServer((args.host, args.port), Handler) + server.store = store # type: ignore[attr-defined] + try: + server.serve_forever(poll_interval=0.1) + except KeyboardInterrupt: + pass + finally: + server.server_close() + return 0 + + +if __name__ == "__main__": + raise SystemExit(main()) From 2be88a33ea0b97a33d17dde1ac92436c126e96cb Mon Sep 17 00:00:00 2001 From: bitmaster162 Date: Sat, 19 Sep 2026 18:07:33 +0700 Subject: [PATCH 2/6] R25 verify qualification witness chain on reload --- docs/R25_REPLAY_PRODUCTION_QUALIFICATION.md | 4 +- tests/test_r25_replay_qualification.py | 65 +++++++++++++++++ tools/r25_witness_service.py | 78 +++++++++++++++++---- 3 files changed, 131 insertions(+), 16 deletions(-) diff --git a/docs/R25_REPLAY_PRODUCTION_QUALIFICATION.md b/docs/R25_REPLAY_PRODUCTION_QUALIFICATION.md index eea8e23..036e2f7 100644 --- a/docs/R25_REPLAY_PRODUCTION_QUALIFICATION.md +++ b/docs/R25_REPLAY_PRODUCTION_QUALIFICATION.md @@ -33,7 +33,9 @@ The local qualification harness uses: The witness fixture persists canonical JSONL records and fsyncs each append before acknowledging it. It implements compare-and-set against exact -generation and head SHA-256. A deterministic post-fsync barrier allows the +generation and head SHA-256. On startup it revalidates each stored record's +schema, namespace, generation sequence, previous-head link, subject hash, +record hash, and record ID before admitting the log into memory. A deterministic post-fsync barrier allows the controller to kill a worker in the crash window after witness durability but before PostgreSQL reconciliation. diff --git a/tests/test_r25_replay_qualification.py b/tests/test_r25_replay_qualification.py index 9cf4794..42b350c 100644 --- a/tests/test_r25_replay_qualification.py +++ b/tests/test_r25_replay_qualification.py @@ -1,6 +1,7 @@ from __future__ import annotations import importlib.util +import json from pathlib import Path import pytest @@ -143,3 +144,67 @@ def test_qualification_witness_health_declares_not_production_ready(): ).read_text(encoding="utf-8") assert '"production_ready": False' in source assert 'durability": "jsonl_fsync"' in source + +def test_qualification_witness_reload_rejects_tampered_record_hash(tmp_path): + namespace = "r25-witness-reload-tamper" + path = tmp_path / "witness.jsonl" + record = _record(namespace) + store = _store(path) + ok, _ = store.append( + namespace=namespace, + expected_generation=0, + expected_head_sha256=replay._genesis_head(namespace), + record=record, + ) + assert ok is True + + item = json.loads(path.read_text(encoding="utf-8").strip()) + item["record"]["digest_sha256"] = "f" * 64 + path.write_text( + witness_service._canonical_json(item) + "\n", + encoding="utf-8", + ) + + with pytest.raises(RuntimeError, match="invalid witness record hash"): + _store(path) + + +def test_qualification_witness_reload_rejects_generation_gap(tmp_path): + namespace = "r25-witness-generation-gap" + path = tmp_path / "witness.jsonl" + record = _record(namespace) + store = _store(path) + ok, _ = store.append( + namespace=namespace, + expected_generation=0, + expected_head_sha256=replay._genesis_head(namespace), + record=record, + ) + assert ok is True + + item = json.loads(path.read_text(encoding="utf-8").strip()) + item["record"]["generation"] = 2 + path.write_text( + witness_service._canonical_json(item) + "\n", + encoding="utf-8", + ) + + with pytest.raises(RuntimeError, match="invalid witness generation sequence"): + _store(path) + + +def test_qualification_witness_append_rejects_invalid_record_before_write(tmp_path): + namespace = "r25-witness-append-tamper" + path = tmp_path / "witness.jsonl" + record = _record(namespace) + record["subject"]["case"] = "tampered-without-rehash" + store = _store(path) + + with pytest.raises(RuntimeError, match="invalid witness subject hash"): + store.append( + namespace=namespace, + expected_generation=0, + expected_head_sha256=replay._genesis_head(namespace), + record=record, + ) + assert not path.exists() or path.read_bytes() == b"" diff --git a/tools/r25_witness_service.py b/tools/r25_witness_service.py index 3901b66..77198ee 100644 --- a/tools/r25_witness_service.py +++ b/tools/r25_witness_service.py @@ -18,6 +18,7 @@ from urllib.parse import parse_qs, urlparse STATE_SCHEMA = "continuityos.replay_witness_state/v1" +RECORD_SCHEMA = "continuityos.replay_witness_record/v1" GENESIS_DOMAIN = "continuityos.replay_witness_genesis/v1" @@ -46,6 +47,44 @@ def _state(namespace: str, generation: int, head_sha256: str) -> dict: } +def _validate_record( + record: dict, + *, + namespace: str, + expected_generation: int, + expected_previous_head: str, +) -> dict: + expected_keys = { + "schema", "namespace", "generation", "previous_head_sha256", + "approval_id", "subject", "subject_sha256", "nonce", + "digest_sha256", "head_sha256", "record_id", + } + if type(record) is not dict or set(record) != expected_keys: + raise RuntimeError("invalid witness record shape") + if record["schema"] != RECORD_SCHEMA or record["namespace"] != namespace: + raise RuntimeError("invalid witness record identity") + if record["generation"] != expected_generation: + raise RuntimeError("invalid witness generation sequence") + if record["previous_head_sha256"] != expected_previous_head: + raise RuntimeError("invalid witness previous-head chain") + if type(record["subject"]) is not dict: + raise RuntimeError("invalid witness subject") + subject_sha = _sha256_text(_canonical_json(record["subject"])) + if record["subject_sha256"] != subject_sha: + raise RuntimeError("invalid witness subject hash") + core = { + key: record[key] + for key in expected_keys + if key not in {"head_sha256", "record_id"} + } + head = _sha256_text(_canonical_json(core)) + if record["head_sha256"] != head: + raise RuntimeError("invalid witness record hash") + if record["record_id"] != "rwr_" + head: + raise RuntimeError("invalid witness record id") + return json.loads(_canonical_json(record)) + + class Store: def __init__( self, @@ -83,9 +122,20 @@ def __init__( raise RuntimeError( f"invalid witness log entry at line {lineno}" ) - self.records.setdefault(item["namespace"], []).append( - item["record"] + namespace = item["namespace"] + rows = self.records.setdefault(namespace, []) + previous_head = ( + _genesis_head(namespace) + if not rows + else str(rows[-1]["head_sha256"]) + ) + validated = _validate_record( + item["record"], + namespace=namespace, + expected_generation=len(rows) + 1, + expected_previous_head=previous_head, ) + rows.append(validated) def current_state(self, namespace: str) -> dict: with self.lock: @@ -119,33 +169,31 @@ def append( or current["head_sha256"] != expected_head_sha256 ): return False, current - if ( - record.get("namespace") != namespace - or record.get("generation") != expected_generation + 1 - or record.get("previous_head_sha256") != expected_head_sha256 - ): - raise ValueError("record does not extend expected witness state") + validated_record = _validate_record( + record, + namespace=namespace, + expected_generation=expected_generation + 1, + expected_previous_head=expected_head_sha256, + ) - item = {"namespace": namespace, "record": record} + item = {"namespace": namespace, "record": validated_record} encoded = (_canonical_json(item) + "\n").encode("ascii") with self.path.open("ab", buffering=0) as handle: handle.write(encoded) os.fsync(handle.fileno()) - self.records.setdefault(namespace, []).append( - json.loads(_canonical_json(record)) - ) + self.records.setdefault(namespace, []).append(validated_record) should_pause = ( not self._pause_used and self.pause_namespace == namespace - and self.pause_generation == record.get("generation") + and self.pause_generation == validated_record.get("generation") ) if should_pause: self._pause_used = True if self.pause_marker is not None: self.pause_marker.parent.mkdir(parents=True, exist_ok=True) self.pause_marker.write_text( - str(record.get("record_id", "")), encoding="utf-8" + str(validated_record["record_id"]), encoding="utf-8" ) if should_pause: @@ -155,7 +203,7 @@ def append( break time.sleep(0.05) - return True, json.loads(_canonical_json(record)) + return True, json.loads(_canonical_json(validated_record)) class Handler(BaseHTTPRequestHandler): From 64034d22491d0f6a8c8fa46cb6ac0ac7a483fe28 Mon Sep 17 00:00:00 2001 From: bitmaster162 Date: Sat, 19 Sep 2026 19:12:02 +0700 Subject: [PATCH 3/6] R25: make replay retry budgets explicit --- continuityos/witnessed_approval_replay.py | 17 ++++++++++++----- 1 file changed, 12 insertions(+), 5 deletions(-) diff --git a/continuityos/witnessed_approval_replay.py b/continuityos/witnessed_approval_replay.py index 7a0f80b..db0ffcd 100644 --- a/continuityos/witnessed_approval_replay.py +++ b/continuityos/witnessed_approval_replay.py @@ -38,12 +38,15 @@ WITNESS_SCOPE = "EXTERNAL_APPEND_ONLY" GENESIS_DOMAIN = "continuityos.replay_witness_genesis/v1" _RETRYABLE_TRANSACTION_SQLSTATES = frozenset({"40001", "40P01"}) +_DB_TRANSACTION_RETRY_LIMIT = 8 +_LOGICAL_CONTENTION_LIMIT = 8 +_EXCEPTION_CHAIN_LIMIT = 8 def _is_retryable_transaction_error(exc: BaseException) -> bool: current: BaseException | None = exc seen: set[int] = set() - for _depth in range(8): + for _depth in range(_EXCEPTION_CHAIN_LIMIT): if current is None or id(current) in seen: break seen.add(id(current)) @@ -559,7 +562,7 @@ def _apply_record(self, cur: Any, record: Mapping[str, Any]) -> None: ) def _snapshot_database(self) -> dict[str, Any]: - for _attempt in range(8): + for _attempt in range(_DB_TRANSACTION_RETRY_LIMIT): con = self._open() cur = None try: @@ -589,6 +592,8 @@ def _snapshot_database(self) -> dict[str, Any]: finally: _close_quietly(cur) _close_quietly(con) + # Exhausting the DB transaction budget is terminal fail-closed for + # this call. Do not multiply it by re-entering an outer logical loop. raise MultiHostApprovalReplayError( "witnessed replay: database snapshot contention" ) @@ -596,7 +601,7 @@ def _snapshot_database(self) -> dict[str, Any]: def _read_claim_row( self, approval_id: str ) -> tuple[str, str, str, str] | None: - for _attempt in range(8): + for _attempt in range(_DB_TRANSACTION_RETRY_LIMIT): con = self._open() cur = None try: @@ -631,6 +636,8 @@ def _read_claim_row( finally: _close_quietly(cur) _close_quietly(con) + # Exhausting the DB transaction budget is terminal fail-closed for + # this call. The caller must not turn this into another retry budget. raise MultiHostApprovalReplayError( "witnessed replay: claim read contention" ) @@ -660,7 +667,7 @@ def _receipt_for_existing( ) def synchronize(self) -> dict[str, Any]: - for _attempt in range(8): + for _attempt in range(_LOGICAL_CONTENTION_LIMIT): db_snapshot = self._snapshot_database() # External witness I/O is deliberately outside every PostgreSQL @@ -770,7 +777,7 @@ def claim_once( digest_value = _hex64("digest_sha256", digest_sha256) last_append_error: MultiHostApprovalReplayError | None = None - for _attempt in range(8): + for _attempt in range(_LOGICAL_CONTENTION_LIMIT): state_value = self.synchronize() existing = self._read_claim_row(identifier) if existing is not None: From 2f0ab21250ecf80f6e4900c3ac52010198f12fc6 Mon Sep 17 00:00:00 2001 From: bitmaster162 Date: Sat, 19 Sep 2026 19:55:58 +0700 Subject: [PATCH 4/6] R25: cap nested retry budget and causal scope --- continuityos/witnessed_approval_replay.py | 29 +++++++++++++++++++---- tests/test_witnessed_approval_replay.py | 16 +++++++++++++ 2 files changed, 41 insertions(+), 4 deletions(-) diff --git a/continuityos/witnessed_approval_replay.py b/continuityos/witnessed_approval_replay.py index db0ffcd..5c57e9e 100644 --- a/continuityos/witnessed_approval_replay.py +++ b/continuityos/witnessed_approval_replay.py @@ -43,6 +43,17 @@ _EXCEPTION_CHAIN_LIMIT = 8 +class _RetryBudget: + def __init__(self, limit: int) -> None: + self.remaining = limit + + def take(self) -> bool: + if self.remaining <= 0: + return False + self.remaining -= 1 + return True + + def _is_retryable_transaction_error(exc: BaseException) -> bool: current: BaseException | None = exc seen: set[int] = set() @@ -55,7 +66,9 @@ def _is_retryable_transaction_error(exc: BaseException) -> bool: sqlstate = getattr(current, "pgcode", None) if sqlstate in _RETRYABLE_TRANSACTION_SQLSTATES: return True - current = current.__cause__ or current.__context__ + # Retry only through an explicit causal chain. Implicit __context__ + # can contain an unrelated prior serialization/deadlock exception. + current = current.__cause__ return False @@ -561,8 +574,11 @@ def _apply_record(self, cur: Any, record: Mapping[str, Any]) -> None: "witnessed replay: database state CAS failed" ) - def _snapshot_database(self) -> dict[str, Any]: - for _attempt in range(_DB_TRANSACTION_RETRY_LIMIT): + def _snapshot_database( + self, *, retry_budget: _RetryBudget | None = None + ) -> dict[str, Any]: + budget = retry_budget or _RetryBudget(_DB_TRANSACTION_RETRY_LIMIT) + while budget.take(): con = self._open() cur = None try: @@ -667,8 +683,13 @@ def _receipt_for_existing( ) def synchronize(self) -> dict[str, Any]: + # Share one snapshot retry budget across the whole logical operation so + # near-exhaustion cannot reset on every outer contention attempt. + snapshot_retry_budget = _RetryBudget(_DB_TRANSACTION_RETRY_LIMIT) for _attempt in range(_LOGICAL_CONTENTION_LIMIT): - db_snapshot = self._snapshot_database() + db_snapshot = self._snapshot_database( + retry_budget=snapshot_retry_budget + ) # External witness I/O is deliberately outside every PostgreSQL # transaction/row lock. The witness CAS is the global monotonic diff --git a/tests/test_witnessed_approval_replay.py b/tests/test_witnessed_approval_replay.py index 440d7d2..133aefd 100644 --- a/tests/test_witnessed_approval_replay.py +++ b/tests/test_witnessed_approval_replay.py @@ -657,11 +657,27 @@ def test_retryable_transaction_error_detection_is_sqlstate_scoped(): outer = RuntimeError("wrapper") outer.__cause__ = RetryableTransactionError("nested serialization") assert replay._is_retryable_transaction_error(outer) + contextual = RuntimeError("ordinary wrapper") + contextual.__context__ = RetryableTransactionError("implicit serialization context") + assert not replay._is_retryable_transaction_error(contextual) assert not replay._is_retryable_transaction_error( RuntimeError("ordinary failure") ) +def test_snapshot_retry_budget_is_shared_across_calls(): + db = DbState() + witness = FakeWitness(db) + guard = authority(db, witness) + budget = replay._RetryBudget(2) + db.retryable_commit_failures = 1 + state = guard._snapshot_database(retry_budget=budget) + assert state["generation"] == 0 + assert budget.remaining == 0 + with pytest.raises(MultiHostApprovalReplayError, match="database snapshot contention"): + guard._snapshot_database(retry_budget=budget) + + def test_snapshot_retries_retryable_serialization_failure(): db = DbState() witness = FakeWitness(db) From 614c3805b20278880e2ef6bbe46db28619eddd3e Mon Sep 17 00:00:00 2001 From: bitmaster162 Date: Sat, 19 Sep 2026 20:54:39 +0700 Subject: [PATCH 5/6] R25: charge snapshot budget only on retries --- continuityos/witnessed_approval_replay.py | 24 ++++++++++++----------- tests/test_witnessed_approval_replay.py | 13 ++++++++++-- 2 files changed, 24 insertions(+), 13 deletions(-) diff --git a/continuityos/witnessed_approval_replay.py b/continuityos/witnessed_approval_replay.py index 5c57e9e..f943560 100644 --- a/continuityos/witnessed_approval_replay.py +++ b/continuityos/witnessed_approval_replay.py @@ -577,8 +577,11 @@ def _apply_record(self, cur: Any, record: Mapping[str, Any]) -> None: def _snapshot_database( self, *, retry_budget: _RetryBudget | None = None ) -> dict[str, Any]: - budget = retry_budget or _RetryBudget(_DB_TRANSACTION_RETRY_LIMIT) - while budget.take(): + # The limit historically means at most 8 total transaction attempts: + # one initial attempt plus up to 7 retryable failures. Successful + # snapshots must not consume the shared retry budget. + budget = retry_budget or _RetryBudget(_DB_TRANSACTION_RETRY_LIMIT - 1) + while True: con = self._open() cur = None try: @@ -601,18 +604,17 @@ def _snapshot_database( except Exception: pass if _is_retryable_transaction_error(exc): - continue + if budget.take(): + continue + raise MultiHostApprovalReplayError( + "witnessed replay: database snapshot contention" + ) from exc raise MultiHostApprovalReplayError( "witnessed replay: database snapshot failed" ) from exc finally: _close_quietly(cur) _close_quietly(con) - # Exhausting the DB transaction budget is terminal fail-closed for - # this call. Do not multiply it by re-entering an outer logical loop. - raise MultiHostApprovalReplayError( - "witnessed replay: database snapshot contention" - ) def _read_claim_row( self, approval_id: str @@ -683,9 +685,9 @@ def _receipt_for_existing( ) def synchronize(self) -> dict[str, Any]: - # Share one snapshot retry budget across the whole logical operation so - # near-exhaustion cannot reset on every outer contention attempt. - snapshot_retry_budget = _RetryBudget(_DB_TRANSACTION_RETRY_LIMIT) + # Share one retryable-failure budget across the whole logical + # operation. Successful snapshot calls do not consume it. + snapshot_retry_budget = _RetryBudget(_DB_TRANSACTION_RETRY_LIMIT - 1) for _attempt in range(_LOGICAL_CONTENTION_LIMIT): db_snapshot = self._snapshot_database( retry_budget=snapshot_retry_budget diff --git a/tests/test_witnessed_approval_replay.py b/tests/test_witnessed_approval_replay.py index 133aefd..70a4c28 100644 --- a/tests/test_witnessed_approval_replay.py +++ b/tests/test_witnessed_approval_replay.py @@ -665,17 +665,26 @@ def test_retryable_transaction_error_detection_is_sqlstate_scoped(): ) -def test_snapshot_retry_budget_is_shared_across_calls(): +def test_snapshot_retry_budget_counts_failures_not_successful_calls(): db = DbState() witness = FakeWitness(db) guard = authority(db, witness) budget = replay._RetryBudget(2) + + for _ in range(replay._LOGICAL_CONTENTION_LIMIT + 1): + state = guard._snapshot_database(retry_budget=budget) + assert state["generation"] == 0 + assert budget.remaining == 2 + db.retryable_commit_failures = 1 state = guard._snapshot_database(retry_budget=budget) assert state["generation"] == 0 - assert budget.remaining == 0 + assert budget.remaining == 1 + + db.retryable_commit_failures = 2 with pytest.raises(MultiHostApprovalReplayError, match="database snapshot contention"): guard._snapshot_database(retry_budget=budget) + assert budget.remaining == 0 def test_snapshot_retries_retryable_serialization_failure(): From b05ac6da9b9abeedb0e087bad1d309e6804d87f6 Mon Sep 17 00:00:00 2001 From: bitmaster162 Date: Sat, 19 Sep 2026 21:11:34 +0700 Subject: [PATCH 6/6] R25: document frozen retry ceilings --- continuityos/witnessed_approval_replay.py | 11 +++++++++++ tests/test_witnessed_approval_replay.py | 16 ++++++++++++++++ 2 files changed, 27 insertions(+) diff --git a/continuityos/witnessed_approval_replay.py b/continuityos/witnessed_approval_replay.py index f943560..7a3855b 100644 --- a/continuityos/witnessed_approval_replay.py +++ b/continuityos/witnessed_approval_replay.py @@ -38,6 +38,8 @@ WITNESS_SCOPE = "EXTERNAL_APPEND_ONLY" GENESIS_DOMAIN = "continuityos.replay_witness_genesis/v1" _RETRYABLE_TRANSACTION_SQLSTATES = frozenset({"40001", "40P01"}) +# Review-frozen safety ceilings. These are intentionally not runtime +# configurable: the reviewed worst-case retry/inspection work stays bounded. _DB_TRANSACTION_RETRY_LIMIT = 8 _LOGICAL_CONTENTION_LIMIT = 8 _EXCEPTION_CHAIN_LIMIT = 8 @@ -55,6 +57,12 @@ def take(self) -> bool: def _is_retryable_transaction_error(exc: BaseException) -> bool: + """Recognize reviewed SQLSTATEs within the bounded explicit cause chain. + + A matching cause deeper than the review-frozen chain ceiling is + intentionally treated as non-retryable. That fails closed on availability + rather than allowing unbounded or attacker-shaped exception traversal. + """ current: BaseException | None = exc seen: set[int] = set() for _depth in range(_EXCEPTION_CHAIN_LIMIT): @@ -577,6 +585,7 @@ def _apply_record(self, cur: Any, record: Mapping[str, Any]) -> None: def _snapshot_database( self, *, retry_budget: _RetryBudget | None = None ) -> dict[str, Any]: + """Read one serializable snapshot under the bounded retry budget.""" # The limit historically means at most 8 total transaction attempts: # one initial attempt plus up to 7 retryable failures. Successful # snapshots must not consume the shared retry budget. @@ -619,6 +628,7 @@ def _snapshot_database( def _read_claim_row( self, approval_id: str ) -> tuple[str, str, str, str] | None: + """Read a claim with the same review-frozen transaction ceiling.""" for _attempt in range(_DB_TRANSACTION_RETRY_LIMIT): con = self._open() cur = None @@ -685,6 +695,7 @@ def _receipt_for_existing( ) def synchronize(self) -> dict[str, Any]: + """Reconcile PostgreSQL with the witness under bounded contention.""" # Share one retryable-failure budget across the whole logical # operation. Successful snapshot calls do not consume it. snapshot_retry_budget = _RetryBudget(_DB_TRANSACTION_RETRY_LIMIT - 1) diff --git a/tests/test_witnessed_approval_replay.py b/tests/test_witnessed_approval_replay.py index 70a4c28..72ee017 100644 --- a/tests/test_witnessed_approval_replay.py +++ b/tests/test_witnessed_approval_replay.py @@ -664,6 +664,22 @@ def test_retryable_transaction_error_detection_is_sqlstate_scoped(): RuntimeError("ordinary failure") ) + # The explicit cause-chain ceiling is a fail-closed safety boundary: + # a reviewed SQLSTATE inside the ceiling is retryable; one beyond it is not. + inside = RetryableTransactionError("inside bounded cause chain") + for _ in range(replay._EXCEPTION_CHAIN_LIMIT - 1): + wrapper = RuntimeError("wrapper") + wrapper.__cause__ = inside + inside = wrapper + assert replay._is_retryable_transaction_error(inside) + + outside = RetryableTransactionError("outside bounded cause chain") + for _ in range(replay._EXCEPTION_CHAIN_LIMIT): + wrapper = RuntimeError("wrapper") + wrapper.__cause__ = outside + outside = wrapper + assert not replay._is_retryable_transaction_error(outside) + def test_snapshot_retry_budget_counts_failures_not_successful_calls(): db = DbState()