fix(#832): bind receipt promotion to trusted observation

This commit is contained in:
ms-lead-reviewer
2026-07-19 17:15:10 -05:00
parent 1713dfe5d0
commit e55262708e
7 changed files with 347 additions and 45 deletions

View File

@@ -15,6 +15,7 @@ consumed challenge cannot be replayed or reopen/renew its lease.
from __future__ import annotations from __future__ import annotations
import argparse import argparse
import base64
import importlib.util import importlib.util
import json import json
import os import os
@@ -88,8 +89,12 @@ def run_once(index: int) -> str:
os.chmod(root, 0o700) os.chmod(root, 0o700)
socket_path = root / "broker.sock" socket_path = root / "broker.sock"
state_path = root / "state.json" state_path = root / "state.json"
observer_path = root / "test-observer.json"
process = subprocess.Popen( process = subprocess.Popen(
[sys.executable, "-I", "-S", "-B", str(DAEMON), "--socket", str(socket_path), "--state", str(state_path)], [
sys.executable, "-I", "-S", "-B", str(DAEMON), "--socket", str(socket_path),
"--state", str(state_path), "--test-observer-file", str(observer_path),
],
stdin=subprocess.DEVNULL, stdin=subprocess.DEVNULL,
stdout=subprocess.PIPE, stdout=subprocess.PIPE,
stderr=subprocess.STDOUT, stderr=subprocess.STDOUT,
@@ -121,12 +126,22 @@ def run_once(index: int) -> str:
"h_payload": construction.h_payload, "h_payload": construction.h_payload,
"schema_version": 1, "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, { pending = request(socket_path, {
"action": "begin_verification", "action": "begin_verification",
"session_id": session_id, "session_id": session_id,
"runtime_generation": 1, "runtime_generation": 1,
"runtime": "pi", "runtime": "pi",
"binding": binding, "binding": binding,
"construction": construction_request,
}) })
if pending.get("ok") is not True or pending.get("state") != "PENDING_VERIFICATION": if pending.get("ok") is not True or pending.get("state") != "PENDING_VERIFICATION":
raise AssertionError(f"shipped pending-delivery transition failed: {pending!r}") raise AssertionError(f"shipped pending-delivery transition failed: {pending!r}")
@@ -143,12 +158,17 @@ def run_once(index: int) -> str:
"receipt_challenge": challenge, "receipt_challenge": challenge,
}), "INVALID_LEASE_TRANSITION") }), "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, { observed = request(socket_path, {
"action": "observe_receipt", "action": "observe_receipt",
"session_id": session_id, "session_id": session_id,
"runtime_generation": 1, "runtime_generation": 1,
"receipt_challenge": challenge, "receipt_challenge": challenge,
"latest_assistant_message": receipt,
}) })
if observed.get("ok") is not True or observed.get("state") != "PENDING_PROMOTION": if observed.get("ok") is not True or observed.get("state") != "PENDING_PROMOTION":
raise AssertionError(f"shipped evidence transition failed: {observed!r}") raise AssertionError(f"shipped evidence transition failed: {observed!r}")
@@ -173,7 +193,6 @@ def run_once(index: int) -> str:
"session_id": session_id, "session_id": session_id,
"runtime_generation": 1, "runtime_generation": 1,
"receipt_challenge": challenge, "receipt_challenge": challenge,
"latest_assistant_message": receipt,
}), "RECEIPT_REPLAY") }), "RECEIPT_REPLAY")
expect_refused(request(socket_path, { expect_refused(request(socket_path, {
"action": "promote_lease", "action": "promote_lease",

View File

@@ -10,4 +10,4 @@
3. Implement the broker-minted receipt challenge and exact receipt observation/consume/promote path. 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`. 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. - **Risks:** The standalone harness must drive the real daemon without a divergent fixture. If that is impossible, stop and flag Mos.
- **Evidence:** RED recorded in `/home/hermes/agent-work/reviews/832-wi5-red-receipt-challenge.log`: receipt module absent before implementation. GREEN: receipt T26/T29 unittest (2 tests), normative-fragments unittest (5 tests), state-store regression (10 tests), `py_compile`, Mosaic package lint, and Mosaic package typecheck. Coverage tooling is unavailable (`python3 -m coverage`: module not installed); no probe was fired. Push pending. - **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: forged `h_source`/`h_payload` was accepted and the observer module was absent. Remediation GREEN: observer/payload/T26/T29 unittest (4 tests), 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 was updated for the test-observer seam but was not fired. 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
@@ -25,7 +26,9 @@ from typing import Final
_MODULE_DIRECTORY = str(Path(__file__).resolve().parent) _MODULE_DIRECTORY = str(Path(__file__).resolve().parent)
if _MODULE_DIRECTORY not in sys.path: if _MODULE_DIRECTORY not in sys.path:
sys.path.insert(0, _MODULE_DIRECTORY) 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_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
@@ -297,8 +300,9 @@ 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()
# 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]] = {}
@@ -455,6 +459,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")
@@ -468,6 +473,22 @@ 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)
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 = copy.deepcopy(binding)
cycle_binding["runtime_generation"] = request["runtime_generation"] cycle_binding["runtime_generation"] = request["runtime_generation"]
challenge = self.mint_token(session_id, request["runtime_generation"], binding) challenge = self.mint_token(session_id, request["runtime_generation"], binding)
@@ -491,8 +512,9 @@ class Broker:
if action == "observe_receipt": if action == "observe_receipt":
session_id, _ = self.authenticate(peer_pid, request) session_id, _ = self.authenticate(peer_pid, request)
challenge = request.get("receipt_challenge") challenge = request.get("receipt_challenge")
message = request.get("latest_assistant_message") if "latest_assistant_message" in request:
if not isinstance(challenge, str) or not isinstance(message, str): raise BrokerFailure("INVALID_RECEIPT")
if not isinstance(challenge, str):
raise BrokerFailure("INVALID_RECEIPT") raise BrokerFailure("INVALID_RECEIPT")
token = self.store.tokens().get(challenge) token = self.store.tokens().get(challenge)
if ( if (
@@ -516,6 +538,11 @@ class Broker:
raise BrokerFailure("RECEIPT_REPLAY") raise BrokerFailure("RECEIPT_REPLAY")
receipt_binding = copy.deepcopy(binding) receipt_binding = copy.deepcopy(binding)
receipt_binding["runtime_generation"] = request["runtime_generation"] 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 # This accepts one exact current-cycle assistant entry only. It is
# deliberately not a transcript search and rejects quoted/extra text. # deliberately not a transcript search and rejects quoted/extra text.
if not is_verbatim_receipt(message, challenge, receipt_binding): if not is_verbatim_receipt(message, challenge, receipt_binding):
@@ -530,6 +557,9 @@ class Broker:
challenge = request.get("receipt_challenge") challenge = request.get("receipt_challenge")
if not isinstance(challenge, str): if not isinstance(challenge, str):
raise BrokerFailure("INVALID_RECEIPT") 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) token = self.store.tokens().get(challenge)
if ( if (
not isinstance(token, dict) not isinstance(token, dict)
@@ -538,8 +568,7 @@ class Broker:
or token.get("consumed") is not False or token.get("consumed") is not False
): ):
raise BrokerFailure("RECEIPT_REPLAY") raise BrokerFailure("RECEIPT_REPLAY")
lease = self.leases.get(session_id) if lease.get("state") != LEASE_PENDING_PROMOTION:
if not isinstance(lease, dict) or lease.get("state") != LEASE_PENDING_PROMOTION:
raise BrokerFailure("INVALID_LEASE_TRANSITION") raise BrokerFailure("INVALID_LEASE_TRANSITION")
expected_challenge = lease.get("receipt_challenge") expected_challenge = lease.get("receipt_challenge")
if not isinstance(expected_challenge, str) or not secrets.compare_digest(challenge, expected_challenge): if not isinstance(expected_challenge, str) or not secrets.compare_digest(challenge, expected_challenge):
@@ -670,12 +699,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()
@@ -755,8 +788,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,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

@@ -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

@@ -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> {
const cycle = pendingReceiptCycles.get(receipt_challenge);
const observerFixture = observerFixtures.get(socket);
if (cycle === undefined || observerFixture === undefined) {
return await request(socket, { return await request(socket, {
action: 'promote_lease', action: 'promote_lease',
session_id, session_id,
runtime_generation, runtime_generation,
promotion_token, 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 {