WI-5 #832: Receipt-challenge protocol (compaction-refresh, milestone 188) (#845)
All checks were successful
ci/woodpecker/push/publish Pipeline was successful
ci/woodpecker/push/ci Pipeline was successful

This commit was merged in pull request #845.
This commit is contained in:
2026-07-19 23:18:56 +00:00
parent e522b22fa4
commit 07553ead33
10 changed files with 941 additions and 62 deletions

View File

@@ -0,0 +1,228 @@
#!/usr/bin/env python3
"""P5 Gate0 replay probe; BUILT ONLY, execution is Mos-gated.
Run only under fresh-executor authorization:
python3 -I -S -B docs/compaction-refresh/probes/p5_receipt_replay.py
Each of the default three isolated runs launches the shipped lease-broker daemon
in a distinct private temporary directory. This driver never changes broker
state directly and does not replace the promote gate: every transition is sent
over the daemon's real Unix socket. It proves the shipped order is
PENDING_DELIVERY -> observe/evidence commit -> consume -> VERIFIED and that a
consumed challenge cannot be replayed or reopen/renew its lease.
"""
from __future__ import annotations
import argparse
import base64
import importlib.util
import json
import os
import shutil
import socket
import subprocess
import sys
import tempfile
import time
from pathlib import Path
HERE = Path(__file__).resolve().parent
REPOSITORY = HERE.parents[2]
TOOLS = REPOSITORY / "packages/mosaic/framework/tools/lease-broker"
DAEMON = TOOLS / "daemon.py"
FRAGMENTS = TOOLS / "normative_fragments.py"
def load_shipped_fragments():
if not FRAGMENTS.is_file():
raise RuntimeError(f"shipped normative construction missing: {FRAGMENTS}")
spec = importlib.util.spec_from_file_location("p5_shipped_normative_fragments", FRAGMENTS)
if spec is None or spec.loader is None:
raise RuntimeError("unable to load shipped normative construction")
module = importlib.util.module_from_spec(spec)
spec.loader.exec_module(module)
return module
def request(socket_path: Path, value: dict[str, object]) -> dict[str, object]:
with socket.socket(socket.AF_UNIX, socket.SOCK_STREAM) as connection:
connection.settimeout(3.0)
connection.connect(str(socket_path))
connection.sendall((json.dumps(value, separators=(",", ":")) + "\n").encode())
connection.shutdown(socket.SHUT_WR)
response = bytearray()
while True:
chunk = connection.recv(4096)
if not chunk:
break
response.extend(chunk)
if not response.endswith(b"\n") or response.count(b"\n") != 1:
raise AssertionError(f"unframed broker reply: {bytes(response)!r}")
parsed = json.loads(response[:-1])
if not isinstance(parsed, dict):
raise AssertionError(f"non-object broker reply: {parsed!r}")
return parsed
def wait_ready(process: subprocess.Popen[str], socket_path: Path) -> None:
deadline = time.monotonic() + 5.0
while time.monotonic() < deadline:
if socket_path.exists():
return
if process.poll() is not None:
output = process.stdout.read() if process.stdout is not None else ""
raise RuntimeError(f"shipped daemon exited before READY: {output}")
time.sleep(0.02)
raise TimeoutError("shipped daemon did not create private probe socket")
def expect_refused(reply: dict[str, object], code: str) -> None:
if reply != {"ok": False, "code": code}:
raise AssertionError(f"expected refusal {code}, got {reply!r}")
def run_once(index: int) -> str:
fragments = load_shipped_fragments()
root = Path(tempfile.mkdtemp(prefix=f"mosaic-p5-replay-{index}-"))
os.chmod(root, 0o700)
socket_path = root / "broker.sock"
state_path = root / "state.json"
observer_path = root / "test-observer.json"
process = subprocess.Popen(
[
sys.executable, "-I", "-S", "-B", str(DAEMON), "--socket", str(socket_path),
"--state", str(state_path), "--test-observer-file", str(observer_path),
],
stdin=subprocess.DEVNULL,
stdout=subprocess.PIPE,
stderr=subprocess.STDOUT,
text=True,
)
try:
wait_ready(process, socket_path)
registered = request(socket_path, {"action": "register_anchor", "runtime_generation": 1})
if registered.get("ok") is not True or not isinstance(registered.get("session_id"), str):
raise AssertionError(f"registration failed: {registered!r}")
session_id = registered["session_id"]
construction = fragments.build_payload(
manifest_version=1,
generator_version="p5-replay-probe",
fragments=[
fragments.NormativeFragment(
"authority/probe",
b"P5 shipped transition driver\n",
"63537df1a6cb0d80195a96757ab11d629e5b5e1f23be167218b84cb195b1c1d6",
),
],
)
if construction.injectionDecision != "ACCEPTED" or not construction.promotion:
raise AssertionError("shipped normative construction refused P5 fixture")
binding = {
"compaction_epoch": index,
"request_epoch": index + 100,
"h_source": construction.h_source,
"h_payload": construction.h_payload,
"schema_version": 1,
}
construction_request = {
"manifest_version": 1,
"generator_version": "p5-replay-probe",
"fragments": [{
"source_id": "authority/probe",
"content_base64": base64.b64encode(b"P5 shipped transition driver\n").decode("ascii"),
"expected_sha256": "63537df1a6cb0d80195a96757ab11d629e5b5e1f23be167218b84cb195b1c1d6",
}],
}
pending = request(socket_path, {
"action": "begin_verification",
"session_id": session_id,
"runtime_generation": 1,
"runtime": "pi",
"binding": binding,
"construction": construction_request,
})
if pending.get("ok") is not True or pending.get("state") != "PENDING_VERIFICATION":
raise AssertionError(f"shipped pending-delivery transition failed: {pending!r}")
challenge = pending.get("receipt_challenge")
receipt = pending.get("receipt")
if not isinstance(challenge, str) or not isinstance(receipt, str):
raise AssertionError(f"shipped broker did not mint a receipt challenge: {pending!r}")
# Promotion before observation/evidence/consumption is forbidden.
expect_refused(request(socket_path, {
"action": "promote_lease",
"session_id": session_id,
"runtime_generation": 1,
"receipt_challenge": challenge,
}), "INVALID_LEASE_TRANSITION")
observer_path.write_text(json.dumps({
"session_id": session_id,
"runtime_generation": 1,
"latest_assistant_message": receipt,
}), encoding="utf-8")
os.chmod(observer_path, 0o600)
observed = request(socket_path, {
"action": "observe_receipt",
"session_id": session_id,
"runtime_generation": 1,
"receipt_challenge": challenge,
})
if observed.get("ok") is not True or observed.get("state") != "PENDING_PROMOTION":
raise AssertionError(f"shipped evidence transition failed: {observed!r}")
durable = json.loads(state_path.read_text(encoding="utf-8"))
evidence = durable["tokens"][challenge].get("evidence")
if not isinstance(evidence, dict) or not isinstance(evidence.get("h_latest_assistant"), str):
raise AssertionError("shipped receipt evidence was not committed before consume/promote")
promoted = request(socket_path, {
"action": "promote_lease",
"session_id": session_id,
"runtime_generation": 1,
"receipt_challenge": challenge,
})
if promoted.get("ok") is not True or promoted.get("state") != "VERIFIED":
raise AssertionError(f"shipped consume-before-promote transition failed: {promoted!r}")
# T25/T28: the actual consumed challenge, re-presented through the
# shipped daemon, can neither be observed again nor re-promote/reopen.
expect_refused(request(socket_path, {
"action": "observe_receipt",
"session_id": session_id,
"runtime_generation": 1,
"receipt_challenge": challenge,
}), "RECEIPT_REPLAY")
expect_refused(request(socket_path, {
"action": "promote_lease",
"session_id": session_id,
"runtime_generation": 1,
"receipt_challenge": challenge,
}), "RECEIPT_REPLAY")
return challenge
finally:
if process.poll() is None:
process.terminate()
try:
process.wait(timeout=3.0)
except subprocess.TimeoutExpired:
process.kill()
process.wait()
shutil.rmtree(root, ignore_errors=True)
def main() -> None:
parser = argparse.ArgumentParser()
parser.add_argument("--runs", type=int, default=3)
arguments = parser.parse_args()
if arguments.runs != 3:
raise SystemExit("P5 requires exactly three isolated runs")
challenges = [run_once(index) for index in range(arguments.runs)]
if len(set(challenges)) != arguments.runs:
raise AssertionError("separate shipped cycles did not mint unique challenges")
print("P5 receipt replay probe PASS: 3 isolated shipped-daemon runs")
if __name__ == "__main__":
main()

View File

@@ -0,0 +1,13 @@
# #832 Receipt-challenge protocol — build scratchpad
- **Objective:** Deliver WI-5 receipt-challenge protocol ACs T25, T26, T28, and T29 only.
- **Authority:** BUILD-BRIEF, SPEC-v5, ratification, and red-team hashes verified in STEP-0.
- **Base:** `e522b22fa4492861b0fcd4a956a8795c54eb9bfe` (`origin/main`).
- **Constraints:** Byte-build only: no live broker/socket/systemd/tmux mutation. No PR, self-review, or probe fire. T27/T30 are out of scope.
- **Plan:**
1. Add red-first deterministic T26/T29 in-build tests that call shipped normative construction and broker path.
2. Add an unexecuted, isolated P5 out-of-process replay harness that drives the shipped daemon and asserts consume-before-promote for T25/T28.
3. Implement the broker-minted receipt challenge and exact receipt observation/consume/promote path.
4. Run unit, framework-shell, compile, lint, and type checks; push after the required queue guard; report to `mosaic-100`.
- **Risks:** The standalone harness must drive the real daemon without a divergent fixture. If that is impossible, stop and flag Mos.
- **Evidence:** Initial RED recorded in `/home/hermes/agent-work/reviews/832-wi5-red-receipt-challenge.log`; initial green checks passed. Remediation RED recorded in `/home/hermes/agent-work/reviews/832-wi5-remediation-red.log` before observer/payload implementation; remediation green passed. Remediation-2 RED recorded in `/home/hermes/agent-work/reviews/832-wi5-remediation2-red.log`: each rejected begin restored prior VERIFIED authority. Remediation-2 GREEN: receipt unittest (5: all `INVALID_CONSTRUCTION`, `PAYLOAD_CONSTRUCTION_REFUSED`, and `PAYLOAD_BINDING_MISMATCH` cases preserve UNVERIFIED and deny the next mutator), normative-fragments unittest (5), state-store regression (10), full mutator-gate acceptance (20, including the real begin → observer → consume → promote path), `py_compile`, Mosaic package lint/typecheck, and targeted Prettier check. The P5 harness remains unfired. Coverage tooling remains unavailable (`python3 -m coverage`: module not installed). Push pending.

View File

@@ -6,6 +6,7 @@ from __future__ import annotations
import argparse import argparse
import copy import copy
import errno import errno
import hmac
from concurrent.futures import ThreadPoolExecutor from concurrent.futures import ThreadPoolExecutor
import json import json
import os import os
@@ -20,6 +21,15 @@ import time
from pathlib import Path from pathlib import Path
from typing import Final from typing import Final
# This script is loaded both as an executable and through isolated stdlib tests.
# Keep its co-located receipt implementation importable in both modes.
_MODULE_DIRECTORY = str(Path(__file__).resolve().parent)
if _MODULE_DIRECTORY not in sys.path:
sys.path.insert(0, _MODULE_DIRECTORY)
from normative_fragments import build_payload_from_wire
from receipt_challenge import is_verbatim_receipt, latest_assistant_digest, receipt_for
from receipt_observer import FileTestReceiptObserver, ReceiptObserver, UnavailableReceiptObserver
MAX_FRAME: Final = 64 * 1024 MAX_FRAME: Final = 64 * 1024
MAX_STATE: Final = 4 * 1024 * 1024 MAX_STATE: Final = 4 * 1024 * 1024
MAX_PENDING_TOKENS: Final = 256 MAX_PENDING_TOKENS: Final = 256
@@ -90,6 +100,14 @@ def valid_binding(binding: object) -> bool:
) )
def valid_receipt_evidence(evidence: object) -> bool:
return (
isinstance(evidence, dict)
and set(evidence) == {"h_latest_assistant"}
and is_hex_256(evidence["h_latest_assistant"])
)
def validate_state(value: object) -> dict[str, object]: def validate_state(value: object) -> dict[str, object]:
if not isinstance(value, dict) or set(value) != {"version", "sessions", "tokens"}: if not isinstance(value, dict) or set(value) != {"version", "sessions", "tokens"}:
raise BrokerFailure("STATE_INTEGRITY") raise BrokerFailure("STATE_INTEGRITY")
@@ -124,7 +142,11 @@ def validate_state(value: object) -> dict[str, object]:
for token_value, token in tokens.items(): for token_value, token in tokens.items():
if not is_hex_256(token_value) or not isinstance(token, dict): if not is_hex_256(token_value) or not isinstance(token, dict):
raise BrokerFailure("STATE_INTEGRITY") raise BrokerFailure("STATE_INTEGRITY")
if set(token) != {"session_id", "runtime_generation", "binding", "consumed"}: token_fields = set(token)
if token_fields not in (
{"session_id", "runtime_generation", "binding", "consumed"},
{"session_id", "runtime_generation", "binding", "consumed", "evidence"},
):
raise BrokerFailure("STATE_INTEGRITY") raise BrokerFailure("STATE_INTEGRITY")
session_id = token["session_id"] session_id = token["session_id"]
generation = token["runtime_generation"] generation = token["runtime_generation"]
@@ -137,6 +159,9 @@ def validate_state(value: object) -> dict[str, object]:
raise BrokerFailure("STATE_INTEGRITY") raise BrokerFailure("STATE_INTEGRITY")
if not valid_binding(token["binding"]) or token["consumed"] is not False: if not valid_binding(token["binding"]) or token["consumed"] is not False:
raise BrokerFailure("STATE_INTEGRITY") raise BrokerFailure("STATE_INTEGRITY")
evidence = token.get("evidence")
if evidence is not None and not valid_receipt_evidence(evidence):
raise BrokerFailure("STATE_INTEGRITY")
return value return value
@@ -275,8 +300,13 @@ class StateStore:
class Broker: class Broker:
def __init__(self, store: StateStore) -> None: def __init__(self, store: StateStore, observer: ReceiptObserver | None = None) -> None:
self.store = store self.store = store
self.observer: ReceiptObserver = observer if observer is not None else UnavailableReceiptObserver()
# Set only by begin_verification after its mandatory revoke-first fence.
# It is preserved if later cycle admission is refused; all other broker
# actions retain the normal snapshot rollback behavior.
self._rejected_cycle_fence: tuple[dict[str, object], dict[str, dict[str, object]]] | None = None
# VERIFIED authority is deliberately volatile: broker restart revokes all # VERIFIED authority is deliberately volatile: broker restart revokes all
# leases while preserving WI-1 identity and pending-token integrity. # leases while preserving WI-1 identity and pending-token integrity.
self.leases: dict[str, dict[str, object]] = {} self.leases: dict[str, dict[str, object]] = {}
@@ -348,6 +378,9 @@ class Broker:
"runtime_generation": generation, "runtime_generation": generation,
"binding": copy.deepcopy(binding), "binding": copy.deepcopy(binding),
"consumed": False, "consumed": False,
# Receipt evidence is durably committed before this challenge may
# be consumed and promotion made externally visible.
"evidence": None,
} }
return token return token
@@ -362,6 +395,7 @@ class Broker:
raise StateCommitUncertain() raise StateCommitUncertain()
previous = copy.deepcopy(self.store.value) previous = copy.deepcopy(self.store.value)
previous_leases = copy.deepcopy(self.leases) previous_leases = copy.deepcopy(self.leases)
self._rejected_cycle_fence = None
try: try:
response = self._handle(peer, request) response = self._handle(peer, request)
if self.store.value != previous: if self.store.value != previous:
@@ -376,9 +410,18 @@ class Broker:
except StateCommitUncertain: except StateCommitUncertain:
raise raise
except Exception: except Exception:
self.store.value = previous fence = self._rejected_cycle_fence
self.leases = previous_leases if fence is None:
self.store.value = previous
self.leases = previous_leases
else:
# A refused re-verification must never resurrect the preceding
# VERIFIED authority. Preserve only this post-revoke fence;
# every unrelated partial-write failure still rolls back.
self.store.value, self.leases = fence
raise raise
finally:
self._rejected_cycle_fence = None
def _handle(self, peer: tuple[int, int, int], request: dict[str, object]) -> dict[str, object]: def _handle(self, peer: tuple[int, int, int], request: dict[str, object]) -> dict[str, object]:
peer_pid, peer_uid, peer_gid = peer peer_pid, peer_uid, peer_gid = peer
@@ -430,6 +473,7 @@ class Broker:
session_id, _ = self.authenticate(peer_pid, request) session_id, _ = self.authenticate(peer_pid, request)
runtime = request.get("runtime") runtime = request.get("runtime")
binding = request.get("binding") binding = request.get("binding")
construction = request.get("construction")
ttl_seconds = request.get("ttl_seconds", MAX_LEASE_TTL_SECONDS) ttl_seconds = request.get("ttl_seconds", MAX_LEASE_TTL_SECONDS)
if runtime not in READ_ONLY_TOOLS: if runtime not in READ_ONLY_TOOLS:
raise BrokerFailure("INVALID_RUNTIME") raise BrokerFailure("INVALID_RUNTIME")
@@ -443,47 +487,117 @@ class Broker:
raise BrokerFailure("INVALID_LEASE_TTL") raise BrokerFailure("INVALID_LEASE_TTL")
# Revoke-first is a broker operation, not advisory adapter order. # Revoke-first is a broker operation, not advisory adapter order.
self.revoke_session_authority(session_id) self.revoke_session_authority(session_id)
token = self.mint_token(session_id, request["runtime_generation"], binding) # If subsequent construction admission rejects, handle() restores
# this fence rather than the pre-cycle VERIFIED snapshot.
self._rejected_cycle_fence = (
copy.deepcopy(self.store.value), copy.deepcopy(self.leases)
)
try:
constructed = build_payload_from_wire(construction)
except ValueError as exc:
raise BrokerFailure("INVALID_CONSTRUCTION") from exc
if (
constructed.injectionDecision != "ACCEPTED"
or not constructed.promotion
or not isinstance(constructed.h_source, str)
or not isinstance(constructed.h_payload, str)
):
raise BrokerFailure("PAYLOAD_CONSTRUCTION_REFUSED")
if (
not hmac.compare_digest(binding["h_source"], constructed.h_source)
or not hmac.compare_digest(binding["h_payload"], constructed.h_payload)
):
raise BrokerFailure("PAYLOAD_BINDING_MISMATCH")
cycle_binding = copy.deepcopy(binding)
cycle_binding["runtime_generation"] = request["runtime_generation"]
challenge = self.mint_token(session_id, request["runtime_generation"], binding)
self.leases[session_id] = { self.leases[session_id] = {
"state": LEASE_PENDING, "state": LEASE_PENDING,
"runtime": runtime, "runtime": runtime,
"runtime_generation": request["runtime_generation"], "runtime_generation": request["runtime_generation"],
"binding": copy.deepcopy(binding), "binding": copy.deepcopy(binding),
"promotion_token": token, "receipt_challenge": challenge,
"ttl_seconds": ttl_seconds, "ttl_seconds": ttl_seconds,
} }
# The broker constructs both the bound challenge and the delivery
# text. The model's role is an exact copy, never hashing its output.
return { return {
"ok": True, "ok": True,
"state": LEASE_PENDING, "state": LEASE_PENDING,
"promotion_token": token, "receipt_challenge": challenge,
"receipt": receipt_for(challenge, cycle_binding),
"binding": cycle_binding,
} }
if action == "promote_lease": if action == "observe_receipt":
session_id, _ = self.authenticate(peer_pid, request) session_id, _ = self.authenticate(peer_pid, request)
promotion_token = request.get("promotion_token") challenge = request.get("receipt_challenge")
lease = self.leases.get(session_id) if "latest_assistant_message" in request:
if not isinstance(lease, dict) or lease.get("state") != LEASE_PENDING: raise BrokerFailure("INVALID_RECEIPT")
raise BrokerFailure("INVALID_LEASE_TRANSITION") if not isinstance(challenge, str):
expected_token = lease.get("promotion_token") raise BrokerFailure("INVALID_RECEIPT")
if ( token = self.store.tokens().get(challenge)
not isinstance(promotion_token, str)
or not isinstance(expected_token, str)
or not secrets.compare_digest(promotion_token, expected_token)
):
raise BrokerFailure("PROMOTION_TOKEN_MISMATCH")
token = self.store.tokens().get(promotion_token)
if ( if (
not isinstance(token, dict) not isinstance(token, dict)
or token.get("session_id") != session_id or token.get("session_id") != session_id
or token.get("runtime_generation") != request.get("runtime_generation") or token.get("runtime_generation") != request.get("runtime_generation")
or token.get("binding") != lease.get("binding")
or token.get("consumed") is not False or token.get("consumed") is not False
): ):
raise BrokerFailure("PROMOTION_TOKEN_INVALID") raise BrokerFailure("RECEIPT_REPLAY")
del self.store.tokens()[promotion_token] lease = self.leases.get(session_id)
if not isinstance(lease, dict) or lease.get("state") != LEASE_PENDING:
raise BrokerFailure("RECEIPT_REPLAY")
expected_challenge = lease.get("receipt_challenge")
binding = lease.get("binding")
if (
not isinstance(expected_challenge, str)
or not secrets.compare_digest(challenge, expected_challenge)
or not isinstance(binding, dict)
or token.get("binding") != binding
):
raise BrokerFailure("RECEIPT_REPLAY")
receipt_binding = copy.deepcopy(binding)
receipt_binding["runtime_generation"] = request["runtime_generation"]
message = self.observer.observe_latest_assistant_message(
session_id, str(lease["runtime"]), request["runtime_generation"], receipt_binding
)
if not isinstance(message, str):
raise BrokerFailure("RECEIPT_OBSERVATION_UNAVAILABLE")
# This accepts one exact current-cycle assistant entry only. It is
# deliberately not a transcript search and rejects quoted/extra text.
if not is_verbatim_receipt(message, challenge, receipt_binding):
raise BrokerFailure("RECEIPT_MISMATCH")
evidence = {"h_latest_assistant": latest_assistant_digest(message)}
token["evidence"] = evidence
lease["state"] = LEASE_PENDING_PROMOTION lease["state"] = LEASE_PENDING_PROMOTION
lease["evidence"] = copy.deepcopy(evidence)
return {"ok": True, "state": LEASE_PENDING_PROMOTION}
if action == "promote_lease":
session_id, _ = self.authenticate(peer_pid, request)
challenge = request.get("receipt_challenge")
if not isinstance(challenge, str):
raise BrokerFailure("INVALID_RECEIPT")
lease = self.leases.get(session_id)
if not isinstance(lease, dict) or lease.get("state") == LEASE_UNVERIFIED:
raise BrokerFailure("INVALID_LEASE_TRANSITION")
token = self.store.tokens().get(challenge)
if (
not isinstance(token, dict)
or token.get("session_id") != session_id
or token.get("runtime_generation") != request.get("runtime_generation")
or token.get("consumed") is not False
):
raise BrokerFailure("RECEIPT_REPLAY")
if lease.get("state") != LEASE_PENDING_PROMOTION:
raise BrokerFailure("INVALID_LEASE_TRANSITION")
expected_challenge = lease.get("receipt_challenge")
if not isinstance(expected_challenge, str) or not secrets.compare_digest(challenge, expected_challenge):
raise BrokerFailure("RECEIPT_REPLAY")
if token.get("binding") != lease.get("binding") or not valid_receipt_evidence(token.get("evidence")):
raise BrokerFailure("PROMOTION_TOKEN_INVALID")
del self.store.tokens()[challenge]
lease["expires_at"] = time.monotonic() + int(lease["ttl_seconds"]) lease["expires_at"] = time.monotonic() + int(lease["ttl_seconds"])
# handle() commits token consumption before finish_promotion() makes # handle() commits the evidence-backed consumption before
# VERIFIED externally visible: promote-last by construction. # finish_promotion() makes VERIFIED externally visible: promote-last.
return {"ok": True, "state": LEASE_PENDING_PROMOTION} return {"ok": True, "state": LEASE_PENDING_PROMOTION}
if action == "revoke_lease": if action == "revoke_lease":
session_id, _ = self.authenticate(peer_pid, request) session_id, _ = self.authenticate(peer_pid, request)
@@ -604,12 +718,16 @@ def handle_connection(
return return
def serve(socket_path: Path, state_path: Path) -> None: def serve(
socket_path: Path,
state_path: Path,
observer: ReceiptObserver | None = None,
) -> None:
secure_parent(socket_path) secure_parent(socket_path)
if socket_path.exists() or socket_path.is_symlink(): if socket_path.exists() or socket_path.is_symlink():
raise BrokerFailure("SOCKET_ALREADY_EXISTS") raise BrokerFailure("SOCKET_ALREADY_EXISTS")
store = StateStore(state_path) store = StateStore(state_path)
broker = Broker(store) broker = Broker(store, observer)
broker_lock = threading.Lock() broker_lock = threading.Lock()
slots = threading.BoundedSemaphore(MAX_IN_FLIGHT_CONNECTIONS) slots = threading.BoundedSemaphore(MAX_IN_FLIGHT_CONNECTIONS)
fatal_lock = threading.Lock() fatal_lock = threading.Lock()
@@ -689,8 +807,14 @@ def main() -> None:
parser = argparse.ArgumentParser() parser = argparse.ArgumentParser()
parser.add_argument("--socket", required=True, type=Path) parser.add_argument("--socket", required=True, type=Path)
parser.add_argument("--state", required=True, type=Path) parser.add_argument("--state", required=True, type=Path)
parser.add_argument("--test-observer-file", type=Path)
arguments = parser.parse_args() arguments = parser.parse_args()
serve(arguments.socket, arguments.state) observer = (
FileTestReceiptObserver(arguments.test_observer_file)
if arguments.test_observer_file is not None
else None
)
serve(arguments.socket, arguments.state, observer)
if __name__ == "__main__": if __name__ == "__main__":

View File

@@ -9,6 +9,7 @@ as a payload input, preventing cryptographic self-reference.
from __future__ import annotations from __future__ import annotations
import base64
import hashlib import hashlib
import hmac import hmac
import struct import struct
@@ -164,6 +165,44 @@ def build_payload(
) )
def build_payload_from_wire(value: object) -> ConstructionResult:
"""Decode a bounded broker request then invoke the sole payload builder.
Wire inputs contain only source bytes and their claimed source digests. The
authoritative ``build_payload`` implementation remains the only code that
admits those bytes and derives ``h_source``/``h_payload``.
"""
if not isinstance(value, dict) or set(value) != {
"manifest_version", "generator_version", "fragments"
}:
raise ValueError("invalid construction")
raw_fragments = value["fragments"]
if not isinstance(raw_fragments, list):
raise ValueError("invalid construction")
fragments: list[NormativeFragment] = []
for raw_fragment in raw_fragments:
if not isinstance(raw_fragment, dict) or set(raw_fragment) != {
"source_id", "content_base64", "expected_sha256"
}:
raise ValueError("invalid construction")
source_id = raw_fragment["source_id"]
encoded = raw_fragment["content_base64"]
expected_sha256 = raw_fragment["expected_sha256"]
if not isinstance(source_id, str) or not isinstance(encoded, str) or not isinstance(expected_sha256, str):
raise ValueError("invalid construction")
try:
content = base64.b64decode(encoded.encode("ascii"), validate=True)
except (UnicodeEncodeError, ValueError) as exc:
raise ValueError("invalid construction") from exc
fragments.append(NormativeFragment(source_id, content, expected_sha256))
return build_payload(
manifest_version=value["manifest_version"],
generator_version=value["generator_version"],
fragments=fragments,
)
def build_for_claude(**kwargs: object) -> ConstructionResult: def build_for_claude(**kwargs: object) -> ConstructionResult:
"""Claude adapter entrypoint; delegates to the sole shared constructor.""" """Claude adapter entrypoint; delegates to the sole shared constructor."""

View File

@@ -0,0 +1,43 @@
#!/usr/bin/env python3
"""Broker-side construction and verification for one-time receipt challenges."""
from __future__ import annotations
import hashlib
import hmac
import struct
from typing import Final
LATEST_ASSISTANT_DOMAIN_SEPARATOR: Final = b"MOSAIC/H_LATEST_ASSISTANT/v1\x00"
def _length_frame(value: bytes) -> bytes:
return struct.pack(">Q", len(value)) + value
def receipt_for(challenge: str, binding: dict[str, object]) -> str:
"""Return the sole receipt text a model may copy for this broker cycle."""
return (
"MOSAIC-RECEIPT{"
f"challenge={challenge}; "
f"H_payload={binding['h_payload']}; "
f"gen={binding['runtime_generation']}; "
f"cep={binding['compaction_epoch']}"
"}"
)
def is_verbatim_receipt(message: str, challenge: str, binding: dict[str, object]) -> bool:
"""Require the exact one current-cycle receipt, not a transcript substring."""
expected = receipt_for(challenge, binding)
return hmac.compare_digest(message, expected)
def latest_assistant_digest(message: str) -> str:
"""Record the broker-computed digest of the exact observed assistant entry."""
encoded = message.encode("utf-8")
return hashlib.sha256(LATEST_ASSISTANT_DOMAIN_SEPARATOR + _length_frame(encoded)).hexdigest()

View File

@@ -0,0 +1,106 @@
#!/usr/bin/env python3
"""Trusted latest-assistant-message observer boundary for receipt promotion.
Runtime adapters must implement ``observe_latest_assistant_message`` directly:
Claude selects the exact latest assistant entry and Pi selects ``message_end``.
The broker accepts no observed message through its request protocol.
"""
from __future__ import annotations
import json
import os
import stat
from pathlib import Path
from typing import Protocol
class ReceiptObserver(Protocol):
def observe_latest_assistant_message(
self,
session_id: str,
runtime: str,
runtime_generation: int,
binding: dict[str, object],
) -> str | None: ...
class UnavailableReceiptObserver:
"""Production-safe default until a runtime adapter injects an observer."""
def observe_latest_assistant_message(
self,
_session_id: str,
_runtime: str,
_runtime_generation: int,
_binding: dict[str, object],
) -> str | None:
return None
class TestReceiptObserver:
"""Deterministic controlled observer used only by byte-build tests."""
def __init__(self) -> None:
self._messages: dict[tuple[str, int], str] = {}
def record_latest_assistant_message(
self, session_id: str, runtime_generation: int, message: str
) -> None:
self._messages[(session_id, runtime_generation)] = message
def observe_latest_assistant_message(
self,
session_id: str,
_runtime: str,
runtime_generation: int,
_binding: dict[str, object],
) -> str | None:
return self._messages.get((session_id, runtime_generation))
class FileTestReceiptObserver:
"""Private fixture-file observer for isolated out-of-process test drivers."""
def __init__(self, path: Path) -> None:
self.path = path
def observe_latest_assistant_message(
self,
session_id: str,
_runtime: str,
runtime_generation: int,
_binding: dict[str, object],
) -> str | None:
flags = os.O_RDONLY | getattr(os, "O_CLOEXEC", 0) | getattr(os, "O_NOFOLLOW", 0)
try:
descriptor = os.open(self.path, flags)
except OSError:
return None
try:
metadata = os.fstat(descriptor)
if (
not stat.S_ISREG(metadata.st_mode)
or stat.S_IMODE(metadata.st_mode) != 0o600
or metadata.st_uid != os.geteuid()
or metadata.st_size > 64 * 1024
):
return None
raw = os.read(descriptor, 64 * 1024 + 1)
finally:
os.close(descriptor)
if len(raw) > 64 * 1024:
return None
try:
value = json.loads(raw)
except (json.JSONDecodeError, UnicodeDecodeError):
return None
if (
not isinstance(value, dict)
or set(value) != {"session_id", "runtime_generation", "latest_assistant_message"}
or value["session_id"] != session_id
or value["runtime_generation"] != runtime_generation
or not isinstance(value["latest_assistant_message"], str)
):
return None
return value["latest_assistant_message"]

View File

@@ -25,7 +25,7 @@
"lint": "eslint src", "lint": "eslint src",
"typecheck": "tsc --noEmit", "typecheck": "tsc --noEmit",
"test": "vitest run --passWithNoTests && pnpm run test:framework-shell", "test": "vitest run --passWithNoTests && pnpm run test:framework-shell",
"test:framework-shell": "python3 src/lease-broker/daemon_deadline_unittest.py && python3 src/lease-broker/normative_fragments_unittest.py && python3 src/mutator-gate/runtime_tools_unittest.py && python3 src/mutator-gate/runtime_launch_guard_unittest.py && python3 framework/tools/lease-broker/check-runtime-launches.py --root ../.. && bash framework/tools/codex/test-pr-diff-context.sh" "test:framework-shell": "python3 src/lease-broker/daemon_deadline_unittest.py && python3 src/lease-broker/normative_fragments_unittest.py && python3 src/lease-broker/receipt_challenge_unittest.py && python3 src/mutator-gate/runtime_tools_unittest.py && python3 src/mutator-gate/runtime_launch_guard_unittest.py && python3 framework/tools/lease-broker/check-runtime-launches.py --root ../.. && bash framework/tools/codex/test-pr-diff-context.sh"
}, },
"dependencies": { "dependencies": {
"@mosaicstack/brain": "workspace:*", "@mosaicstack/brain": "workspace:*",

View File

@@ -1,4 +1,5 @@
import { createHash } from 'node:crypto'; import { createHash } from 'node:crypto';
import { chmod, writeFile } from 'node:fs/promises';
import { createConnection, type Socket } from 'node:net'; import { createConnection, type Socket } from 'node:net';
const DEFAULT_TIMEOUT_MS = 3_000; const DEFAULT_TIMEOUT_MS = 3_000;
@@ -185,3 +186,51 @@ export async function requestBrokerReply<T extends object>(
options, options,
); );
} }
export interface ReceiptChallengeCycle {
sessionId: string;
runtimeGeneration: number;
receiptChallenge: string;
receipt: string;
}
export interface ReceiptChallengeReply {
ok: boolean;
code?: string;
state?: 'UNVERIFIED' | 'PENDING_VERIFICATION' | 'PENDING_PROMOTION' | 'VERIFIED';
}
/**
* Complete the shipped begin -> trusted-observer -> consume -> promote path.
* The private fixture is read by the daemon's injected test observer; the
* observation request itself never carries assistant-message content.
*/
export async function observeAndPromoteReceiptChallenge(
socketPath: string,
observerFixturePath: string,
cycle: ReceiptChallengeCycle,
): Promise<ReceiptChallengeReply> {
await writeFile(
observerFixturePath,
`${JSON.stringify({
session_id: cycle.sessionId,
runtime_generation: cycle.runtimeGeneration,
latest_assistant_message: cycle.receipt,
})}\n`,
{ encoding: 'utf8', mode: 0o600 },
);
await chmod(observerFixturePath, 0o600);
const observed = await requestBrokerReply<ReceiptChallengeReply>(socketPath, {
action: 'observe_receipt',
session_id: cycle.sessionId,
runtime_generation: cycle.runtimeGeneration,
receipt_challenge: cycle.receiptChallenge,
});
if (observed.ok !== true || observed.state !== 'PENDING_PROMOTION') return observed;
return await requestBrokerReply<ReceiptChallengeReply>(socketPath, {
action: 'promote_lease',
session_id: cycle.sessionId,
runtime_generation: cycle.runtimeGeneration,
receipt_challenge: cycle.receiptChallenge,
});
}

View File

@@ -0,0 +1,227 @@
#!/usr/bin/env python3
"""RED-first contracts for the shipped receipt challenge and observer seam."""
from __future__ import annotations
import base64
import copy
import hashlib
import importlib.util
import os
import tempfile
import unittest
from pathlib import Path
TOOLS = Path(__file__).parents[2] / "framework/tools/lease-broker"
DAEMON_PATH = TOOLS / "daemon.py"
FRAGMENTS_PATH = TOOLS / "normative_fragments.py"
OBSERVER_PATH = TOOLS / "receipt_observer.py"
def load_module(name: str, path: Path):
assert path.is_file(), f"shipped module is missing: {path}"
spec = importlib.util.spec_from_file_location(name, path)
if spec is None or spec.loader is None:
raise RuntimeError(f"unable to load {name}")
module = importlib.util.module_from_spec(spec)
spec.loader.exec_module(module)
return module
DAEMON = load_module("lease_broker_receipt_daemon", DAEMON_PATH)
FRAGMENTS = load_module("lease_broker_normative_fragments", FRAGMENTS_PATH)
class BrokerFixture(unittest.TestCase):
def setUp(self) -> None:
self.temporary = tempfile.TemporaryDirectory()
root = Path(self.temporary.name)
os.chmod(root, 0o700)
self.peer = (os.getpid(), os.getuid(), os.getgid())
self.broker = DAEMON.Broker(DAEMON.StateStore(root / "state.json"))
registered = self.broker.handle(self.peer, {
"action": "register_anchor",
"runtime_generation": 7,
})
self.session_id = registered["session_id"]
self.assertIsInstance(self.session_id, str)
def tearDown(self) -> None:
self.temporary.cleanup()
def construction(self) -> tuple[dict[str, object], dict[str, object]]:
content = b"Constitution\n"
expected_sha256 = hashlib.sha256(content).hexdigest()
construction = {
"manifest_version": 1,
"generator_version": "wi5-receipt-test",
"fragments": [{
"source_id": "authority/constitution",
"content_base64": base64.b64encode(content).decode("ascii"),
"expected_sha256": expected_sha256,
}],
}
result = FRAGMENTS.build_payload(
manifest_version=construction["manifest_version"],
generator_version=construction["generator_version"],
fragments=[FRAGMENTS.NormativeFragment("authority/constitution", content, expected_sha256)],
)
self.assertEqual(result.injectionDecision, "ACCEPTED")
self.assertTrue(result.promotion)
return construction, {
"compaction_epoch": 3,
"request_epoch": 8,
"h_source": result.h_source,
"h_payload": result.h_payload,
"schema_version": 1,
}
def begin(self, binding: dict[str, object], construction: dict[str, object]) -> dict[str, object]:
response = self.broker.handle(self.peer, {
"action": "begin_verification",
"session_id": self.session_id,
"runtime_generation": 7,
"runtime": "pi",
"binding": binding,
"construction": construction,
})
self.assertEqual(response["state"], DAEMON.LEASE_PENDING)
self.assertIsInstance(response.get("receipt_challenge"), str)
self.assertIsInstance(response.get("receipt"), str)
return response
class BuildPayloadAdmissionTest(BrokerFixture):
def test_b3_forged_h_source_or_h_payload_is_refused_against_shipped_build_payload(self) -> None:
construction, trusted = self.construction()
for field in ("h_source", "h_payload"):
with self.subTest(field=field):
forged = dict(trusted)
forged[field] = "f" * 64
with self.assertRaisesRegex(DAEMON.BrokerFailure, "PAYLOAD_BINDING_MISMATCH"):
self.begin(forged, construction)
class ReceiptObserverTest(BrokerFixture):
def setUp(self) -> None:
super().setUp()
observers = load_module("lease_broker_test_observer", OBSERVER_PATH)
self.observer = observers.TestReceiptObserver()
self.broker = DAEMON.Broker(self.broker.store, observer=self.observer)
def record(self, message: str) -> None:
self.observer.record_latest_assistant_message(self.session_id, 7, message)
def observe(self, challenge: str, **untrusted: object) -> dict[str, object]:
return self.broker.handle(self.peer, {
"action": "observe_receipt",
"session_id": self.session_id,
"runtime_generation": 7,
"receipt_challenge": challenge,
**untrusted,
})
def promote(self, challenge: str) -> dict[str, object]:
return self.broker.handle(self.peer, {
"action": "promote_lease",
"session_id": self.session_id,
"runtime_generation": 7,
"receipt_challenge": challenge,
})
def test_b2_echoed_request_observation_is_refused_but_observer_source_promotes(self) -> None:
construction, binding = self.construction()
cycle = self.begin(binding, construction)
challenge = cycle["receipt_challenge"]
receipt = cycle["receipt"]
with self.assertRaisesRegex(DAEMON.BrokerFailure, "INVALID_RECEIPT"):
self.observe(challenge, latest_assistant_message=receipt)
with self.assertRaisesRegex(DAEMON.BrokerFailure, "RECEIPT_OBSERVATION_UNAVAILABLE"):
self.observe(challenge)
self.record(receipt)
self.assertEqual(self.observe(challenge)["state"], DAEMON.LEASE_PENDING_PROMOTION)
self.assertEqual(self.promote(challenge)["state"], DAEMON.LEASE_VERIFIED)
self.assertEqual(self.broker.handle(self.peer, {
"action": "authorize_tool",
"session_id": self.session_id,
"runtime_generation": 7,
"runtime": "pi",
"tool_name": "bash",
})["decision"], "allow")
def test_rejected_begin_keeps_revoke_first_fence_for_all_construction_refusals(self) -> None:
construction, binding = self.construction()
refusal_cases = {
"INVALID_CONSTRUCTION": {"bad": "construction"},
"PAYLOAD_CONSTRUCTION_REFUSED": {
**construction,
"fragments": [{
**construction["fragments"][0],
"expected_sha256": "0" * 64,
}],
},
"PAYLOAD_BINDING_MISMATCH": None,
}
for expected_code, rejected_construction in refusal_cases.items():
with self.subTest(expected_code=expected_code):
verified = self.begin(binding, construction)
self.record(verified["receipt"])
self.observe(verified["receipt_challenge"])
self.assertEqual(self.promote(verified["receipt_challenge"])["state"], DAEMON.LEASE_VERIFIED)
rejected_binding = copy.deepcopy(binding)
if expected_code == "PAYLOAD_BINDING_MISMATCH":
rejected_binding["h_payload"] = "f" * 64
rejected_construction = construction
with self.assertRaisesRegex(DAEMON.BrokerFailure, expected_code):
self.begin(rejected_binding, rejected_construction)
self.assertEqual(
self.broker.leases[self.session_id]["state"], DAEMON.LEASE_UNVERIFIED
)
denied = self.broker.handle(self.peer, {
"action": "authorize_tool",
"session_id": self.session_id,
"runtime_generation": 7,
"runtime": "pi",
"tool_name": "bash",
})
self.assertEqual(denied["decision"], "deny")
self.assertEqual(denied["state"], DAEMON.LEASE_UNVERIFIED)
def test_t26_stale_epoch_receipt_cannot_promote_against_shipped_binding(self) -> None:
construction, stale_binding = self.construction()
stale = self.begin(stale_binding, construction)
current_binding = dict(stale_binding)
current_binding["compaction_epoch"] = 4
current_binding["request_epoch"] = 9
current = self.begin(current_binding, construction)
self.assertNotEqual(stale["receipt_challenge"], current["receipt_challenge"])
self.record(stale["receipt"])
with self.assertRaisesRegex(DAEMON.BrokerFailure, "RECEIPT_MISMATCH"):
self.observe(current["receipt_challenge"])
with self.assertRaisesRegex(DAEMON.BrokerFailure, "INVALID_LEASE_TRANSITION"):
self.promote(current["receipt_challenge"])
def test_t29_altered_model_hash_cannot_promote_against_shipped_binding(self) -> None:
construction, binding = self.construction()
cycle = self.begin(binding, construction)
expected = cycle["receipt"]
altered_hash = "f" * 64
self.assertNotEqual(altered_hash, cycle["binding"]["h_payload"])
altered = expected.replace(cycle["binding"]["h_payload"], altered_hash, 1)
self.assertNotEqual(altered, expected)
self.record(altered)
with self.assertRaisesRegex(DAEMON.BrokerFailure, "RECEIPT_MISMATCH"):
self.observe(cycle["receipt_challenge"])
with self.assertRaisesRegex(DAEMON.BrokerFailure, "INVALID_LEASE_TRANSITION"):
self.promote(cycle["receipt_challenge"])
if __name__ == "__main__":
unittest.main()

View File

@@ -6,19 +6,31 @@ import { spawn, spawnSync, type ChildProcess } from 'node:child_process';
import { afterEach, describe, expect, test } from 'vitest'; import { afterEach, describe, expect, test } from 'vitest';
import { launchClaudex, type ClaudexHarnessAdapter } from '../commands/claudex.js'; import { launchClaudex, type ClaudexHarnessAdapter } from '../commands/claudex.js';
import { requestBrokerReply } from '../lease-broker/broker-test-client.js'; import {
observeAndPromoteReceiptChallenge,
requestBrokerReply,
} from '../lease-broker/broker-test-client.js';
interface BrokerReply { interface BrokerReply {
ok: boolean; ok: boolean;
code?: string; code?: string;
decision?: 'allow' | 'deny'; decision?: 'allow' | 'deny';
state?: 'UNVERIFIED' | 'PENDING_VERIFICATION' | 'VERIFIED'; state?: 'UNVERIFIED' | 'PENDING_VERIFICATION' | 'PENDING_PROMOTION' | 'VERIFIED';
session_id?: string; session_id?: string;
promotion_token?: string; receipt_challenge?: string;
receipt?: string;
}
interface PendingReceiptCycle {
sessionId: string;
runtimeGeneration: number;
receiptChallenge: string;
receipt: string;
} }
interface BrokerPaths { interface BrokerPaths {
socket: string; socket: string;
observerFixture: string;
} }
const frameworkRoot = new URL('../../framework/', import.meta.url).pathname; const frameworkRoot = new URL('../../framework/', import.meta.url).pathname;
@@ -38,14 +50,29 @@ const remediationHandlerPath = join(frameworkRoot, 'tools/qa/remediation-hook-ha
const children: ChildProcess[] = []; const children: ChildProcess[] = [];
const temporaryRoots: string[] = []; const temporaryRoots: string[] = [];
const construction = {
manifest_version: 1,
generator_version: 'mutator-gate-acceptance',
fragments: [
{
source_id: 'authority/mutator-gate-acceptance',
content_base64: 'bXV0YXRvci1nYXRlIGFjY2VwdGFuY2UK',
expected_sha256: 'cc3de191821d48037f60b4d006fce74b5dd394d39fb5c8bf681c28889ff0e623',
},
],
};
const binding = (compaction_epoch = 1) => ({ const binding = (compaction_epoch = 1) => ({
compaction_epoch, compaction_epoch,
request_epoch: 0, request_epoch: 0,
h_source: 'a'.repeat(64), h_source: '3c8fc6733d6a2bdc001ed9277d636d7cfac037ae48ed7d54203180a3839dc7a6',
h_payload: 'b'.repeat(64), h_payload: '42da889037c29c3a41397df0f86465ffd121bdb8cd18ad0d7c3df259ab502d3e',
schema_version: 1, schema_version: 1,
}); });
const pendingReceiptCycles = new Map<string, PendingReceiptCycle>();
const observerFixtures = new Map<string, string>();
async function request(socketPath: string, requestValue: object): Promise<BrokerReply> { async function request(socketPath: string, requestValue: object): Promise<BrokerReply> {
return await requestBrokerReply<BrokerReply>(socketPath, requestValue); return await requestBrokerReply<BrokerReply>(socketPath, requestValue);
} }
@@ -54,9 +81,18 @@ async function startBroker(): Promise<BrokerPaths> {
const root = await mkdtemp(join(tmpdir(), 'mosaic-mutator-gate-')); const root = await mkdtemp(join(tmpdir(), 'mosaic-mutator-gate-'));
await chmod(root, 0o700); await chmod(root, 0o700);
const socket = join(root, 'broker.sock'); const socket = join(root, 'broker.sock');
const observerFixture = join(root, 'test-observer.json');
const child = spawn( const child = spawn(
'python3', 'python3',
[daemonPath, '--socket', socket, '--state', join(root, 'state.json')], [
daemonPath,
'--socket',
socket,
'--state',
join(root, 'state.json'),
'--test-observer-file',
observerFixture,
],
{ {
stdio: ['ignore', 'pipe', 'pipe'], stdio: ['ignore', 'pipe', 'pipe'],
}, },
@@ -72,7 +108,8 @@ async function startBroker(): Promise<BrokerPaths> {
); );
child.stdout?.once('data', () => resolve()); child.stdout?.once('data', () => resolve());
}); });
return { socket }; observerFixtures.set(socket, observerFixture);
return { socket, observerFixture };
} }
interface RuntimeLaunchEntry { interface RuntimeLaunchEntry {
@@ -166,28 +203,43 @@ async function beginVerification(
ttl_seconds = 300, ttl_seconds = 300,
compactionEpoch = 1, compactionEpoch = 1,
): Promise<BrokerReply> { ): Promise<BrokerReply> {
return await request(socket, { const reply = await request(socket, {
action: 'begin_verification', action: 'begin_verification',
session_id, session_id,
runtime_generation, runtime_generation,
runtime, runtime,
ttl_seconds, ttl_seconds,
binding: binding(compactionEpoch), binding: binding(compactionEpoch),
construction,
}); });
if (typeof reply.receipt_challenge === 'string' && typeof reply.receipt === 'string') {
pendingReceiptCycles.set(reply.receipt_challenge, {
sessionId: session_id,
runtimeGeneration: runtime_generation,
receiptChallenge: reply.receipt_challenge,
receipt: reply.receipt,
});
}
return reply;
} }
async function promote( async function promote(
socket: string, socket: string,
session_id: string, session_id: string,
promotion_token: string, receipt_challenge: string,
runtime_generation = 1, runtime_generation = 1,
): Promise<BrokerReply> { ): Promise<BrokerReply> {
return await request(socket, { const cycle = pendingReceiptCycles.get(receipt_challenge);
action: 'promote_lease', const observerFixture = observerFixtures.get(socket);
session_id, if (cycle === undefined || observerFixture === undefined) {
runtime_generation, return await request(socket, {
promotion_token, action: 'promote_lease',
}); session_id,
runtime_generation,
receipt_challenge,
});
}
return await observeAndPromoteReceiptChallenge(socket, observerFixture, cycle);
} }
async function authorize( async function authorize(
@@ -249,13 +301,13 @@ describe('whole mutator-class lease gate', () => {
const pending = await beginVerification(socket, sessionId, 'claude'); const pending = await beginVerification(socket, sessionId, 'claude');
expect(pending).toMatchObject({ ok: true, state: 'PENDING_VERIFICATION' }); expect(pending).toMatchObject({ ok: true, state: 'PENDING_VERIFICATION' });
expect(pending.promotion_token).toMatch(/^[a-f0-9]{64}$/); expect(pending.receipt_challenge).toMatch(/^[a-f0-9]{64}$/);
expect(await authorize(socket, sessionId, 'claude', 'Write')).toMatchObject({ expect(await authorize(socket, sessionId, 'claude', 'Write')).toMatchObject({
ok: false, ok: false,
decision: 'deny', decision: 'deny',
}); });
expect(await promote(socket, sessionId, pending.promotion_token!)).toMatchObject({ expect(await promote(socket, sessionId, pending.receipt_challenge!)).toMatchObject({
ok: true, ok: true,
state: 'VERIFIED', state: 'VERIFIED',
}); });
@@ -271,9 +323,9 @@ describe('whole mutator-class lease gate', () => {
ok: false, ok: false,
decision: 'deny', decision: 'deny',
}); });
expect(await promote(socket, sessionId, pending.promotion_token!)).toMatchObject({ expect(await promote(socket, sessionId, pending.receipt_challenge!)).toMatchObject({
ok: false, ok: false,
code: 'PROMOTION_TOKEN_MISMATCH', code: 'RECEIPT_REPLAY',
}); });
}); });
@@ -401,7 +453,7 @@ describe('whole mutator-class lease gate', () => {
const { socket } = await startBroker(); const { socket } = await startBroker();
const sessionId = await register(socket); const sessionId = await register(socket);
const pending = await beginVerification(socket, sessionId, 'claude', 1, 1); const pending = await beginVerification(socket, sessionId, 'claude', 1, 1);
await promote(socket, sessionId, pending.promotion_token!); await promote(socket, sessionId, pending.receipt_challenge!);
// Intentionally invoke neither compaction observer: this is the amended // Intentionally invoke neither compaction observer: this is the amended
// D2-v5 bounded residual, not a fail-closed path. // D2-v5 bounded residual, not a fail-closed path.
@@ -431,7 +483,7 @@ describe('whole mutator-class lease gate', () => {
const { socket } = await startBroker(); const { socket } = await startBroker();
const sessionId = await register(socket); const sessionId = await register(socket);
const pending = await beginVerification(socket, sessionId, 'claude'); const pending = await beginVerification(socket, sessionId, 'claude');
await promote(socket, sessionId, pending.promotion_token!); await promote(socket, sessionId, pending.receipt_challenge!);
const revoked = spawnSync('python3', [revokerPath, '--runtime', 'claude', '--reason', reason], { const revoked = spawnSync('python3', [revokerPath, '--runtime', 'claude', '--reason', reason], {
encoding: 'utf8', encoding: 'utf8',
@@ -454,7 +506,7 @@ describe('whole mutator-class lease gate', () => {
const { socket } = await startBroker(); const { socket } = await startBroker();
const observerSessionId = await register(socket); const observerSessionId = await register(socket);
const observerPending = await beginVerification(socket, observerSessionId, 'claude'); const observerPending = await beginVerification(socket, observerSessionId, 'claude');
await promote(socket, observerSessionId, observerPending.promotion_token!); await promote(socket, observerSessionId, observerPending.receipt_challenge!);
expect(await authorize(socket, observerSessionId, 'claude', 'Bash')).toMatchObject({ expect(await authorize(socket, observerSessionId, 'claude', 'Bash')).toMatchObject({
ok: true, ok: true,
@@ -464,12 +516,10 @@ describe('whole mutator-class lease gate', () => {
const retriedPromotion = await promote( const retriedPromotion = await promote(
socket, socket,
observerSessionId, observerSessionId,
observerPending.promotion_token!, observerPending.receipt_challenge!,
); );
expect(retriedPromotion.ok).toBe(false); expect(retriedPromotion.ok).toBe(false);
expect(['PROMOTION_TOKEN_MISMATCH', 'INVALID_LEASE_TRANSITION']).toContain( expect(retriedPromotion.code).toBe('RECEIPT_REPLAY');
retriedPromotion.code,
);
const revoked = spawnSync( const revoked = spawnSync(
'python3', 'python3',
@@ -495,7 +545,7 @@ describe('whole mutator-class lease gate', () => {
const expirySessionId = await register(expirySocket); const expirySessionId = await register(expirySocket);
expect(expirySessionId).not.toBe(observerSessionId); expect(expirySessionId).not.toBe(observerSessionId);
const expiryPending = await beginVerification(expirySocket, expirySessionId, 'claude', 1, 1); const expiryPending = await beginVerification(expirySocket, expirySessionId, 'claude', 1, 1);
await promote(expirySocket, expirySessionId, expiryPending.promotion_token!); await promote(expirySocket, expirySessionId, expiryPending.receipt_challenge!);
expect(await authorize(expirySocket, expirySessionId, 'claude', 'Bash')).toMatchObject({ expect(await authorize(expirySocket, expirySessionId, 'claude', 'Bash')).toMatchObject({
ok: true, ok: true,
decision: 'allow', decision: 'allow',
@@ -514,7 +564,7 @@ describe('whole mutator-class lease gate', () => {
const { socket } = await startBroker(); const { socket } = await startBroker();
const sessionId = await register(socket); const sessionId = await register(socket);
const pending = await beginVerification(socket, sessionId, 'claude'); const pending = await beginVerification(socket, sessionId, 'claude');
await promote(socket, sessionId, pending.promotion_token!); await promote(socket, sessionId, pending.receipt_challenge!);
const root = await mkdtemp(join(tmpdir(), 'mosaic-observer-fence-')); const root = await mkdtemp(join(tmpdir(), 'mosaic-observer-fence-'));
temporaryRoots.push(root); temporaryRoots.push(root);
@@ -547,7 +597,7 @@ describe('whole mutator-class lease gate', () => {
const { socket } = await startBroker(); const { socket } = await startBroker();
const sessionId = await register(socket); const sessionId = await register(socket);
const pending = await beginVerification(socket, sessionId, 'pi'); const pending = await beginVerification(socket, sessionId, 'pi');
await promote(socket, sessionId, pending.promotion_token!); await promote(socket, sessionId, pending.receipt_challenge!);
const anchorPid = process.pid; const anchorPid = process.pid;
const root = await mkdtemp(join(tmpdir(), 'mosaic-generation-bump-')); const root = await mkdtemp(join(tmpdir(), 'mosaic-generation-bump-'));
@@ -612,7 +662,7 @@ describe('whole mutator-class lease gate', () => {
const { socket } = await startBroker(); const { socket } = await startBroker();
const sessionId = await register(socket); const sessionId = await register(socket);
const pending = await beginVerification(socket, sessionId, 'claude', 1, 1); const pending = await beginVerification(socket, sessionId, 'claude', 1, 1);
await promote(socket, sessionId, pending.promotion_token!); await promote(socket, sessionId, pending.receipt_challenge!);
expect(await authorize(socket, sessionId, 'claude', 'Bash')).toMatchObject({ expect(await authorize(socket, sessionId, 'claude', 'Bash')).toMatchObject({
ok: true, ok: true,
@@ -626,7 +676,7 @@ describe('whole mutator-class lease gate', () => {
}); });
const refreshed = await beginVerification(socket, sessionId, 'claude', 1, 300, 2); const refreshed = await beginVerification(socket, sessionId, 'claude', 1, 300, 2);
await promote(socket, sessionId, refreshed.promotion_token!); await promote(socket, sessionId, refreshed.receipt_challenge!);
expect( expect(
await request(socket, { await request(socket, {
action: 'revoke_lease', action: 'revoke_lease',
@@ -646,7 +696,7 @@ describe('whole mutator-class lease gate', () => {
const { socket } = await startBroker(); const { socket } = await startBroker();
const sessionId = await register(socket); const sessionId = await register(socket);
const pending = await beginVerification(socket, sessionId, 'pi'); const pending = await beginVerification(socket, sessionId, 'pi');
await promote(socket, sessionId, pending.promotion_token!); await promote(socket, sessionId, pending.receipt_challenge!);
expect(await authorize(socket, sessionId, 'pi', 'bash', 2)).toMatchObject({ expect(await authorize(socket, sessionId, 'pi', 'bash', 2)).toMatchObject({
ok: false, ok: false,
@@ -859,7 +909,7 @@ raise SystemExit(0 if len(session_id) == 64 and hook_present and observers_prese
expect(runRuntimeGate(socket, sessionId, 'pi', 'unknown_custom_tool').status).toBe(2); expect(runRuntimeGate(socket, sessionId, 'pi', 'unknown_custom_tool').status).toBe(2);
const pending = await beginVerification(socket, sessionId, 'claude'); const pending = await beginVerification(socket, sessionId, 'claude');
await promote(socket, sessionId, pending.promotion_token!); await promote(socket, sessionId, pending.receipt_challenge!);
expect(runRuntimeGate(socket, sessionId, 'claude', 'Bash').status).toBe(0); expect(runRuntimeGate(socket, sessionId, 'claude', 'Bash').status).toBe(0);
const settings = JSON.parse(await readFile(claudeSettingsPath, 'utf8')) as { const settings = JSON.parse(await readFile(claudeSettingsPath, 'utf8')) as {