From d15b4b83fc40c6f2d568a40e7597234bfdb8f0b0 Mon Sep 17 00:00:00 2001 From: Jason Woltje Date: Fri, 7 Aug 2026 19:31:44 -0500 Subject: [PATCH] feat(lease): add single-turn Claude promotion trigger --- .../runtime/claude/commands/mosaic-promote.md | 1 + .../framework/runtime/claude/settings.json | 16 +- .../tools/lease-broker/promote-begin.py | 333 +++++++++++ .../tools/lease-broker/promote-complete.py | 306 +++++++++++ .../lease-broker/receipt-observer-client.py | 7 +- .../tools/lease-broker/receipt_challenge.py | 2 +- packages/mosaic/package.json | 2 +- .../promotion_trigger_unittest.py | 519 ++++++++++++++++++ .../receipt_challenge_unittest.py | 18 + .../receipt_observer_client_unittest.py | 106 +++- 10 files changed, 1304 insertions(+), 6 deletions(-) create mode 100644 packages/mosaic/framework/runtime/claude/commands/mosaic-promote.md create mode 100644 packages/mosaic/framework/tools/lease-broker/promote-begin.py create mode 100644 packages/mosaic/framework/tools/lease-broker/promote-complete.py create mode 100644 packages/mosaic/src/lease-broker/promotion_trigger_unittest.py diff --git a/packages/mosaic/framework/runtime/claude/commands/mosaic-promote.md b/packages/mosaic/framework/runtime/claude/commands/mosaic-promote.md new file mode 100644 index 00000000..24c29784 --- /dev/null +++ b/packages/mosaic/framework/runtime/claude/commands/mosaic-promote.md @@ -0,0 +1 @@ +I invoked this registered command to authorize lease promotion; follow the local seat broker's injected receipt confirmation instruction exactly. diff --git a/packages/mosaic/framework/runtime/claude/settings.json b/packages/mosaic/framework/runtime/claude/settings.json index eada96fe..0e6dcec1 100644 --- a/packages/mosaic/framework/runtime/claude/settings.json +++ b/packages/mosaic/framework/runtime/claude/settings.json @@ -32,6 +32,18 @@ ] } ], + "UserPromptSubmit": [ + { + "matcher": "^/mosaic-promote$", + "hooks": [ + { + "type": "command", + "command": "python3 ~/.config/mosaic/tools/lease-broker/promote-begin.py", + "timeout": 15 + } + ] + } + ], "PreToolUse": [ { "matcher": ".*", @@ -81,8 +93,8 @@ "hooks": [ { "type": "command", - "command": "python3 ~/.config/mosaic/tools/lease-broker/receipt-observer-client.py --runtime claude --latest-entry", - "timeout": 3 + "command": "python3 ~/.config/mosaic/tools/lease-broker/receipt-observer-client.py --runtime claude --latest-entry; observer_status=$?; python3 ~/.config/mosaic/tools/lease-broker/promote-complete.py; exit $observer_status", + "timeout": 15 }, { "type": "command", diff --git a/packages/mosaic/framework/tools/lease-broker/promote-begin.py b/packages/mosaic/framework/tools/lease-broker/promote-begin.py new file mode 100644 index 00000000..dd93b13f --- /dev/null +++ b/packages/mosaic/framework/tools/lease-broker/promote-begin.py @@ -0,0 +1,333 @@ +#!/usr/bin/env python3 +"""Claude UserPromptSubmit hook for operator-triggered lease promotion.""" + +from __future__ import annotations + +import fcntl +import json +import os +import secrets +import stat +import subprocess +import sys +import time +from collections.abc import Callable, Mapping +from pathlib import Path +from typing import Final, TextIO + +_MODULE_DIRECTORY = str(Path(__file__).resolve().parent) +if _MODULE_DIRECTORY not in sys.path: + sys.path.insert(0, _MODULE_DIRECTORY) + +from receipt_challenge import receipt_for # noqa: E402 + +MAX_FRAME: Final = 64 * 1024 +PENDING_MAX_AGE_SECONDS: Final = 60 * 60 +PROMOTER_TIMEOUT_SECONDS: Final = 10.0 +PROMOTION_PROMPT: Final = "/mosaic-promote" +PROMOTER: Final = Path(__file__).resolve().with_name("lease_promote.py") +PENDING_DIRECTORY: Final = "mosaic-lease" +LOCK_FILE: Final = "promotion.lock" +EXPECTED_BEGIN_KEYS: Final = frozenset( + {"ok", "state", "receipt_challenge", "receipt", "binding"} +) +EXPECTED_BINDING_KEYS: Final = frozenset( + { + "compaction_epoch", + "request_epoch", + "h_source", + "h_payload", + "runtime_generation", + "schema_version", + } +) + + +class PromotionAlreadyInProgress(RuntimeError): + pass + + +def reject_duplicate_json_keys(pairs: list[tuple[str, object]]) -> dict[str, object]: + value: dict[str, object] = {} + for key, item in pairs: + if key in value: + raise ValueError("duplicate promoter JSON key") + value[key] = item + return value + + +def read_hook_input(stream: object) -> dict[str, object]: + raw = getattr(stream, "buffer", stream).read(MAX_FRAME + 1) + if not isinstance(raw, bytes) or len(raw) > MAX_FRAME: + raise ValueError("invalid UserPromptSubmit input") + value = json.loads(raw, object_pairs_hook=reject_duplicate_json_keys) + if not isinstance(value, dict): + raise ValueError("invalid UserPromptSubmit input") + return value + + +def emit_context(stream: TextIO, message: str) -> None: + json.dump( + { + "hookSpecificOutput": { + "hookEventName": "UserPromptSubmit", + "additionalContext": message, + } + }, + stream, + separators=(",", ":"), + ) + stream.write("\n") + + +def session_pending_name(environ: Mapping[str, str]) -> tuple[Path, str]: + runtime_dir = Path(environ["XDG_RUNTIME_DIR"]) + session_id = environ["MOSAIC_LEASE_SESSION_ID"] + if not runtime_dir.is_absolute(): + raise ValueError("XDG_RUNTIME_DIR must be absolute") + if len(session_id) != 64 or any(character not in "0123456789abcdef" for character in session_id): + raise ValueError("invalid lease session id") + return runtime_dir, f"pending-{session_id}" + + +def open_pending_directory(runtime_dir: Path) -> int: + directory_flags = ( + os.O_RDONLY + | getattr(os, "O_CLOEXEC", 0) + | getattr(os, "O_DIRECTORY", 0) + | getattr(os, "O_NOFOLLOW", 0) + ) + runtime_descriptor = os.open(runtime_dir, directory_flags) + try: + runtime_metadata = os.fstat(runtime_descriptor) + if ( + not stat.S_ISDIR(runtime_metadata.st_mode) + or runtime_metadata.st_uid != os.getuid() + or stat.S_IMODE(runtime_metadata.st_mode) != 0o700 + ): + raise ValueError("unsafe XDG runtime directory") + try: + os.mkdir(PENDING_DIRECTORY, mode=0o700, dir_fd=runtime_descriptor) + except FileExistsError: + pass + descriptor = os.open(PENDING_DIRECTORY, directory_flags, dir_fd=runtime_descriptor) + finally: + os.close(runtime_descriptor) + + metadata = os.fstat(descriptor) + if ( + not stat.S_ISDIR(metadata.st_mode) + or metadata.st_uid != os.getuid() + or stat.S_IMODE(metadata.st_mode) != 0o700 + ): + os.close(descriptor) + raise ValueError("unsafe promotion pending directory") + return descriptor + + +def acquire_lock(directory_descriptor: int) -> int: + flags = ( + os.O_RDWR + | os.O_CREAT + | getattr(os, "O_CLOEXEC", 0) + | getattr(os, "O_NOFOLLOW", 0) + ) + descriptor = os.open(LOCK_FILE, flags, 0o600, dir_fd=directory_descriptor) + metadata = os.fstat(descriptor) + if ( + not stat.S_ISREG(metadata.st_mode) + or metadata.st_uid != os.getuid() + or stat.S_IMODE(metadata.st_mode) != 0o600 + ): + os.close(descriptor) + raise ValueError("unsafe promotion lock file") + try: + fcntl.flock(descriptor, fcntl.LOCK_EX | fcntl.LOCK_NB) + except BlockingIOError as error: + os.close(descriptor) + raise PromotionAlreadyInProgress() from error + return descriptor + + +def sweep_stale_pending(directory_descriptor: int, current_time: float) -> None: + cutoff = current_time - PENDING_MAX_AGE_SECONDS + removed = False + with os.scandir(directory_descriptor) as entries: + for candidate in entries: + if not ( + candidate.name.startswith("pending-") + or candidate.name.startswith(".pending-") + ): + continue + try: + metadata = candidate.stat(follow_symlinks=False) + if metadata.st_mtime < cutoff and not stat.S_ISDIR(metadata.st_mode): + os.unlink(candidate.name, dir_fd=directory_descriptor) + removed = True + except FileNotFoundError: + continue + if removed: + os.fsync(directory_descriptor) + + +def write_pending(directory_descriptor: int, name: str, challenge: str) -> None: + temporary = f".{name}.tmp-{secrets.token_hex(8)}" + flags = ( + os.O_WRONLY + | os.O_CREAT + | os.O_EXCL + | getattr(os, "O_CLOEXEC", 0) + | getattr(os, "O_NOFOLLOW", 0) + ) + descriptor = os.open(temporary, flags, 0o600, dir_fd=directory_descriptor) + try: + os.fchmod(descriptor, 0o600) + with os.fdopen(descriptor, "w", encoding="utf-8", closefd=False) as stream: + stream.write(challenge) + stream.flush() + os.fsync(stream.fileno()) + os.replace( + temporary, + name, + src_dir_fd=directory_descriptor, + dst_dir_fd=directory_descriptor, + ) + os.fsync(directory_descriptor) + except Exception: + try: + os.unlink(temporary, dir_fd=directory_descriptor) + except FileNotFoundError: + pass + raise + finally: + os.close(descriptor) + + +def parse_begin_reply( + completed: subprocess.CompletedProcess[str], +) -> tuple[str, dict[str, object] | None]: + if completed.returncode != 0: + return f"PROMOTER_EXIT_{completed.returncode}", None + try: + value = json.loads( + completed.stdout, + object_pairs_hook=reject_duplicate_json_keys, + ) + except (json.JSONDecodeError, RecursionError, TypeError, ValueError): + return "INVALID_PROMOTER_REPLY", None + if not isinstance(value, dict): + return "INVALID_PROMOTER_REPLY", None + if value.get("ok") is False and set(value) == {"ok", "code"}: + code = value.get("code") + return code if isinstance(code, str) and code else "PROMOTION_BEGIN_REFUSED", value + if set(value) != EXPECTED_BEGIN_KEYS or value.get("ok") is not True: + return "INVALID_PROMOTER_REPLY", None + if value.get("state") != "PENDING_VERIFICATION": + return "INVALID_PROMOTER_REPLY", None + challenge = value.get("receipt_challenge") + receipt = value.get("receipt") + binding = value.get("binding") + if ( + not isinstance(challenge, str) + or len(challenge) != 64 + or any(character not in "0123456789abcdef" for character in challenge) + or not isinstance(receipt, str) + or not isinstance(binding, dict) + or set(binding) != EXPECTED_BINDING_KEYS + ): + return "INVALID_PROMOTER_REPLY", None + integer_fields = ( + "compaction_epoch", + "request_epoch", + "runtime_generation", + "schema_version", + ) + if any(type(binding.get(field)) is not int or binding[field] < 0 for field in integer_fields): + return "INVALID_PROMOTER_REPLY", None + if not all( + isinstance(binding.get(field), str) + and len(binding[field]) == 64 + and all(character in "0123456789abcdef" for character in binding[field]) + for field in ("h_source", "h_payload") + ): + return "INVALID_PROMOTER_REPLY", None + if not secrets.compare_digest( + receipt.encode("utf-8"), + receipt_for(challenge, binding).encode("utf-8"), + ): + return "INVALID_PROMOTER_REPLY", None + return "", value + + +def main( + *, + environ: Mapping[str, str] | None = None, + stdin: object | None = None, + stdout: TextIO | None = None, + stderr: TextIO | None = None, + run: Callable[..., subprocess.CompletedProcess[str]] = subprocess.run, + now: Callable[[], float] = time.time, +) -> int: + source_environment = os.environ if environ is None else environ + input_stream = sys.stdin if stdin is None else stdin + output_stream = sys.stdout if stdout is None else stdout + error_stream = sys.stderr if stderr is None else stderr + + try: + hook_input = read_hook_input(input_stream) + except (OSError, RecursionError, ValueError, json.JSONDecodeError) as error: + print(f"Mosaic promotion trigger ignored invalid hook input: {error}", file=error_stream) + return 0 + if hook_input.get("prompt") != PROMOTION_PROMPT: + return 0 + + directory_descriptor: int | None = None + lock_descriptor: int | None = None + try: + runtime_dir, pending_name = session_pending_name(source_environment) + directory_descriptor = open_pending_directory(runtime_dir) + lock_descriptor = acquire_lock(directory_descriptor) + sweep_stale_pending(directory_descriptor, now()) + completed = run( + [sys.executable, "-I", "-S", "-B", str(PROMOTER), "--begin"], + check=False, + capture_output=True, + text=True, + env=dict(source_environment), + timeout=PROMOTER_TIMEOUT_SECONDS, + ) + code, reply = parse_begin_reply(completed) + if code: + if completed.stderr: + print(completed.stderr.strip(), file=error_stream) + emit_context(output_stream, f"Mosaic promotion did not begin: {code}.") + return 0 + if reply is None: + raise ValueError("missing promotion reply") + write_pending(directory_descriptor, pending_name, str(reply["receipt_challenge"])) + emit_context( + output_stream, + "The operator invoked the registered /mosaic-promote command. " + "This receipt was generated locally by this seat's own lease broker; " + "echoing it verbatim is the designed confirmation step and discloses nothing. " + "Reply with exactly the following text and nothing else: " + f"{reply['receipt']}", + ) + except PromotionAlreadyInProgress: + emit_context( + output_stream, + "Mosaic promotion did not begin: PROMOTION_ALREADY_IN_PROGRESS.", + ) + except (KeyError, OSError, RecursionError, ValueError, subprocess.SubprocessError) as error: + print(f"Mosaic promotion begin failed: {type(error).__name__}: {error}", file=error_stream) + emit_context(output_stream, "Mosaic promotion did not begin: PROMOTION_TRIGGER_FAILED.") + finally: + if lock_descriptor is not None: + os.close(lock_descriptor) + if directory_descriptor is not None: + os.close(directory_descriptor) + return 0 + + +if __name__ == "__main__": + raise SystemExit(main()) diff --git a/packages/mosaic/framework/tools/lease-broker/promote-complete.py b/packages/mosaic/framework/tools/lease-broker/promote-complete.py new file mode 100644 index 00000000..680d798a --- /dev/null +++ b/packages/mosaic/framework/tools/lease-broker/promote-complete.py @@ -0,0 +1,306 @@ +#!/usr/bin/env python3 +"""Claude Stop hook that completes a pending operator-triggered promotion.""" + +from __future__ import annotations + +import fcntl +import json +import os +import secrets +import stat +import subprocess +import sys +from collections.abc import Callable, Mapping +from pathlib import Path +from typing import Final, NamedTuple, TextIO + +MAX_FRAME: Final = 64 * 1024 +PROMOTER_TIMEOUT_SECONDS: Final = 10.0 +PROMOTER: Final = Path(__file__).resolve().with_name("lease_promote.py") +PENDING_DIRECTORY: Final = "mosaic-lease" +LOCK_FILE: Final = "promotion.lock" +TERMINAL_FAILURE_CODES: Final = frozenset( + { + "RECEIPT_REPLAY", + "RECEIPT_MISMATCH", + "INVALID_LEASE_TRANSITION", + "PROMOTION_TOKEN_INVALID", + } +) + + +class PendingChallenge(NamedTuple): + value: str + device: int + inode: int + + +def reject_duplicate_json_keys(pairs: list[tuple[str, object]]) -> dict[str, object]: + value: dict[str, object] = {} + for key, item in pairs: + if key in value: + raise ValueError("duplicate promoter JSON key") + value[key] = item + return value + + +def session_pending_name(environ: Mapping[str, str]) -> tuple[Path, str]: + runtime_dir = Path(environ["XDG_RUNTIME_DIR"]) + session_id = environ["MOSAIC_LEASE_SESSION_ID"] + if not runtime_dir.is_absolute(): + raise ValueError("XDG_RUNTIME_DIR must be absolute") + if len(session_id) != 64 or any(character not in "0123456789abcdef" for character in session_id): + raise ValueError("invalid lease session id") + return runtime_dir, f"pending-{session_id}" + + +def open_pending_directory(runtime_dir: Path) -> int | None: + directory_flags = ( + os.O_RDONLY + | getattr(os, "O_CLOEXEC", 0) + | getattr(os, "O_DIRECTORY", 0) + | getattr(os, "O_NOFOLLOW", 0) + ) + try: + runtime_descriptor = os.open(runtime_dir, directory_flags) + except FileNotFoundError: + return None + try: + runtime_metadata = os.fstat(runtime_descriptor) + if ( + not stat.S_ISDIR(runtime_metadata.st_mode) + or runtime_metadata.st_uid != os.getuid() + or stat.S_IMODE(runtime_metadata.st_mode) != 0o700 + ): + raise ValueError("unsafe XDG runtime directory") + try: + descriptor = os.open(PENDING_DIRECTORY, directory_flags, dir_fd=runtime_descriptor) + except FileNotFoundError: + return None + finally: + os.close(runtime_descriptor) + + metadata = os.fstat(descriptor) + if ( + not stat.S_ISDIR(metadata.st_mode) + or metadata.st_uid != os.getuid() + or stat.S_IMODE(metadata.st_mode) != 0o700 + ): + os.close(descriptor) + raise ValueError("unsafe promotion pending directory") + return descriptor + + +def acquire_lock(directory_descriptor: int) -> int: + flags = ( + os.O_RDWR + | os.O_CREAT + | getattr(os, "O_CLOEXEC", 0) + | getattr(os, "O_NOFOLLOW", 0) + ) + descriptor = os.open(LOCK_FILE, flags, 0o600, dir_fd=directory_descriptor) + metadata = os.fstat(descriptor) + if ( + not stat.S_ISREG(metadata.st_mode) + or metadata.st_uid != os.getuid() + or stat.S_IMODE(metadata.st_mode) != 0o600 + ): + os.close(descriptor) + raise ValueError("unsafe promotion lock file") + try: + fcntl.flock(descriptor, fcntl.LOCK_EX | fcntl.LOCK_NB) + except BlockingIOError: + os.close(descriptor) + raise + return descriptor + + +def read_pending(directory_descriptor: int, name: str) -> PendingChallenge | None: + flags = os.O_RDONLY | getattr(os, "O_CLOEXEC", 0) | getattr(os, "O_NOFOLLOW", 0) + try: + descriptor = os.open(name, flags, dir_fd=directory_descriptor) + except FileNotFoundError: + return None + try: + metadata = os.fstat(descriptor) + if ( + not stat.S_ISREG(metadata.st_mode) + or metadata.st_uid != os.getuid() + or stat.S_IMODE(metadata.st_mode) != 0o600 + or metadata.st_size <= 0 + or metadata.st_size > MAX_FRAME + ): + raise ValueError("unsafe promotion pending file") + raw = os.read(descriptor, MAX_FRAME + 1) + finally: + os.close(descriptor) + if len(raw) > MAX_FRAME: + raise ValueError("oversized promotion challenge") + challenge = raw.decode("utf-8") + if ( + len(challenge) != 64 + or any(character not in "0123456789abcdef" for character in challenge) + ): + raise ValueError("invalid promotion challenge") + return PendingChallenge(challenge, metadata.st_dev, metadata.st_ino) + + +def delete_pending_if_unchanged( + directory_descriptor: int, + name: str, + pending: PendingChallenge, + error_stream: TextIO, +) -> None: + quarantine = f".{name}.delete-{secrets.token_hex(8)}" + try: + os.rename( + name, + quarantine, + src_dir_fd=directory_descriptor, + dst_dir_fd=directory_descriptor, + ) + except FileNotFoundError: + return + except OSError as error: + print(f"Mosaic promotion could not quarantine pending file: {error}", file=error_stream) + return + + try: + moved = os.stat( + quarantine, + dir_fd=directory_descriptor, + follow_symlinks=False, + ) + if (moved.st_dev, moved.st_ino) == (pending.device, pending.inode): + os.unlink(quarantine, dir_fd=directory_descriptor) + os.fsync(directory_descriptor) + return + + print("Mosaic promotion pending file changed; preserving replacement.", file=error_stream) + try: + os.link( + quarantine, + name, + src_dir_fd=directory_descriptor, + dst_dir_fd=directory_descriptor, + follow_symlinks=False, + ) + except FileExistsError: + print( + f"Mosaic promotion preserved replacement as {quarantine}.", + file=error_stream, + ) + else: + os.unlink(quarantine, dir_fd=directory_descriptor) + os.fsync(directory_descriptor) + except OSError as error: + print(f"Mosaic promotion could not resolve pending file: {error}", file=error_stream) + + +def parse_reply(completed: subprocess.CompletedProcess[str]) -> dict[str, object] | None: + if completed.returncode != 0: + return None + try: + value = json.loads( + completed.stdout, + object_pairs_hook=reject_duplicate_json_keys, + ) + except (json.JSONDecodeError, RecursionError, TypeError, ValueError): + return None + if not isinstance(value, dict): + return None + if set(value) == {"stage", "ok", "state"}: + if ( + value.get("stage") == "promote_lease" + and value.get("ok") is True + and value.get("state") == "VERIFIED" + ): + return value + return None + if set(value) == {"stage", "ok", "code"}: + if ( + value.get("stage") in {"observe_receipt", "promote_lease"} + and value.get("ok") is False + and isinstance(value.get("code"), str) + and value.get("code") + ): + return value + return None + + +def main( + *, + environ: Mapping[str, str] | None = None, + stderr: TextIO | None = None, + run: Callable[..., subprocess.CompletedProcess[str]] = subprocess.run, +) -> int: + source_environment = os.environ if environ is None else environ + error_stream = sys.stderr if stderr is None else stderr + directory_descriptor: int | None = None + lock_descriptor: int | None = None + + try: + runtime_dir, pending_name = session_pending_name(source_environment) + directory_descriptor = open_pending_directory(runtime_dir) + if directory_descriptor is None: + return 0 + try: + lock_descriptor = acquire_lock(directory_descriptor) + except (BlockingIOError, FileNotFoundError): + print("Mosaic promotion completion deferred: promotion is in progress.", file=error_stream) + return 0 + pending = read_pending(directory_descriptor, pending_name) + if pending is None: + return 0 + completed = run( + [ + sys.executable, + "-I", + "-S", + "-B", + str(PROMOTER), + "--complete", + pending.value, + ], + check=False, + capture_output=True, + text=True, + env=dict(source_environment), + timeout=PROMOTER_TIMEOUT_SECONDS, + ) + reply = parse_reply(completed) + if reply is not None and reply.get("ok") is True: + delete_pending_if_unchanged( + directory_descriptor, + pending_name, + pending, + error_stream, + ) + print("Mosaic lease promotion completed.", file=error_stream) + return 0 + + if reply is not None: + code = str(reply["code"]) + print(f"Mosaic promotion incomplete: {code}.", file=error_stream) + if code in TERMINAL_FAILURE_CODES: + delete_pending_if_unchanged( + directory_descriptor, + pending_name, + pending, + error_stream, + ) + else: + diagnostic = completed.stderr.strip() or f"promoter exit {completed.returncode}" + print(f"Mosaic promotion retryable failure: {diagnostic}.", file=error_stream) + except (KeyError, OSError, RecursionError, UnicodeError, ValueError, subprocess.SubprocessError) as error: + print(f"Mosaic promotion completion deferred: {type(error).__name__}: {error}", file=error_stream) + finally: + if lock_descriptor is not None: + os.close(lock_descriptor) + if directory_descriptor is not None: + os.close(directory_descriptor) + return 0 + + +if __name__ == "__main__": + raise SystemExit(main()) diff --git a/packages/mosaic/framework/tools/lease-broker/receipt-observer-client.py b/packages/mosaic/framework/tools/lease-broker/receipt-observer-client.py index 2289d606..4366d2e5 100644 --- a/packages/mosaic/framework/tools/lease-broker/receipt-observer-client.py +++ b/packages/mosaic/framework/tools/lease-broker/receipt-observer-client.py @@ -129,7 +129,12 @@ def main(argv: Sequence[str] | None = None, *, environ: Mapping[str, str] | None if arguments.runtime == "claude": if not arguments.latest_entry: raise ValueError("Claude observer requires --latest-entry") - message = claude_latest_entry(source) + if "last_assistant_message" in source: + message = source["last_assistant_message"] + if not isinstance(message, str): + raise ValueError("invalid Claude observer input") + else: + message = claude_latest_entry(source) else: if arguments.latest_entry: raise ValueError("Pi observer is message_end only") diff --git a/packages/mosaic/framework/tools/lease-broker/receipt_challenge.py b/packages/mosaic/framework/tools/lease-broker/receipt_challenge.py index 4c2f5c3d..abb642df 100644 --- a/packages/mosaic/framework/tools/lease-broker/receipt_challenge.py +++ b/packages/mosaic/framework/tools/lease-broker/receipt_challenge.py @@ -33,7 +33,7 @@ def is_verbatim_receipt(message: str, challenge: str, binding: dict[str, object] """Require the exact one current-cycle receipt, not a transcript substring.""" expected = receipt_for(challenge, binding) - return hmac.compare_digest(message, expected) + return hmac.compare_digest(message.encode("utf-8"), expected.encode("utf-8")) def latest_assistant_digest(message: str) -> str: diff --git a/packages/mosaic/package.json b/packages/mosaic/package.json index cc843cec..de24c717 100644 --- a/packages/mosaic/package.json +++ b/packages/mosaic/package.json @@ -25,7 +25,7 @@ "lint": "eslint src", "typecheck": "tsc --noEmit", "test": "vitest run --passWithNoTests && pnpm run test:framework-shell", - "test:framework-shell": "bash framework/tools/quality/scripts/check-test-enumeration.sh && bash framework/tools/quality/scripts/test-check-test-enumeration.sh && python3 src/lease-broker/daemon_deadline_unittest.py && python3 src/lease-broker/normative_fragments_unittest.py && python3 src/lease-broker/promotion_binding_unittest.py && python3 src/lease-broker/receipt_challenge_unittest.py && python3 src/lease-broker/context_recovery_unittest.py && python3 src/lease-broker/recovery_runtime_unittest.py && python3 src/lease-broker/recovery_b1_adversarial_unittest.py && python3 src/lease-broker/receipt_observer_client_unittest.py && python3 src/lease-broker/invariant_r_unittest.py && python3 src/lease-broker/framework_skill_portability_unittest.py && python3 src/mutator-gate/runtime_tools_unittest.py && python3 src/mutator-gate/runtime_launch_guard_unittest.py && python3 src/mutator-gate/version_coupling_unittest.py && python3 framework/tools/lease-broker/check-runtime-launches.py --root ../.. && bash framework/tools/codex/test-pr-diff-context.sh && bash framework/tools/qa/test-deps-preflight.sh && bash framework/tools/git/test-pr-review-gitea-comment.sh && bash framework/tools/git/test-pr-review-repo-host-override.sh && bash framework/tools/git/test-ci-queue-wait-branch-absent.sh && bash framework/tools/git/test-ci-queue-wait-tristate.sh && bash framework/tools/git/test-ci-queue-wait-github-checks.sh && bash framework/tools/git/test-pr-merge-queue-branch.sh && bash framework/tools/git/test-pr-merge-head-pin.sh && bash framework/tools/git/test-pr-merge-message-field.sh && bash framework/tools/git/test-git-credential-mosaic.sh && bash framework/tools/git/test-gitea-token-identity.sh && bash framework/tools/woodpecker/test-terminal-green-contract.sh && bash framework/tools/_scripts/test-install-ordering-guard.sh && bash framework/tools/tmux/agent-send.test.sh && bash framework/tools/wake/test-wake-store-ack.sh && bash framework/tools/wake/test-wake-store-enqueue-race.sh && bash framework/tools/wake/test-wake-digest-hmac.sh && bash framework/tools/wake/test-wake-digest-quarantine.sh && bash framework/tools/wake/test-wake-detector.sh && bash framework/tools/wake/test-wake-fn-oracle.sh && bash framework/tools/wake/test-wake-reconcile.sh && bash framework/tools/wake/test-wake-beacon.sh && bash framework/tools/wake/test-wake-preimage.sh && bash framework/tools/wake/test-wake-install.sh" + "test:framework-shell": "bash framework/tools/quality/scripts/check-test-enumeration.sh && bash framework/tools/quality/scripts/test-check-test-enumeration.sh && python3 src/lease-broker/daemon_deadline_unittest.py && python3 src/lease-broker/normative_fragments_unittest.py && python3 src/lease-broker/promotion_binding_unittest.py && python3 src/lease-broker/promotion_trigger_unittest.py && python3 src/lease-broker/receipt_challenge_unittest.py && python3 src/lease-broker/context_recovery_unittest.py && python3 src/lease-broker/recovery_runtime_unittest.py && python3 src/lease-broker/recovery_b1_adversarial_unittest.py && python3 src/lease-broker/receipt_observer_client_unittest.py && python3 src/lease-broker/invariant_r_unittest.py && python3 src/lease-broker/framework_skill_portability_unittest.py && python3 src/mutator-gate/runtime_tools_unittest.py && python3 src/mutator-gate/runtime_launch_guard_unittest.py && python3 src/mutator-gate/version_coupling_unittest.py && python3 framework/tools/lease-broker/check-runtime-launches.py --root ../.. && bash framework/tools/codex/test-pr-diff-context.sh && bash framework/tools/qa/test-deps-preflight.sh && bash framework/tools/git/test-pr-review-gitea-comment.sh && bash framework/tools/git/test-pr-review-repo-host-override.sh && bash framework/tools/git/test-ci-queue-wait-branch-absent.sh && bash framework/tools/git/test-ci-queue-wait-tristate.sh && bash framework/tools/git/test-ci-queue-wait-github-checks.sh && bash framework/tools/git/test-pr-merge-queue-branch.sh && bash framework/tools/git/test-pr-merge-head-pin.sh && bash framework/tools/git/test-pr-merge-message-field.sh && bash framework/tools/git/test-git-credential-mosaic.sh && bash framework/tools/git/test-gitea-token-identity.sh && bash framework/tools/woodpecker/test-terminal-green-contract.sh && bash framework/tools/_scripts/test-install-ordering-guard.sh && bash framework/tools/tmux/agent-send.test.sh && bash framework/tools/wake/test-wake-store-ack.sh && bash framework/tools/wake/test-wake-store-enqueue-race.sh && bash framework/tools/wake/test-wake-digest-hmac.sh && bash framework/tools/wake/test-wake-digest-quarantine.sh && bash framework/tools/wake/test-wake-detector.sh && bash framework/tools/wake/test-wake-fn-oracle.sh && bash framework/tools/wake/test-wake-reconcile.sh && bash framework/tools/wake/test-wake-beacon.sh && bash framework/tools/wake/test-wake-preimage.sh && bash framework/tools/wake/test-wake-install.sh" }, "dependencies": { "@mosaicstack/brain": "workspace:*", diff --git a/packages/mosaic/src/lease-broker/promotion_trigger_unittest.py b/packages/mosaic/src/lease-broker/promotion_trigger_unittest.py new file mode 100644 index 00000000..dea87f37 --- /dev/null +++ b/packages/mosaic/src/lease-broker/promotion_trigger_unittest.py @@ -0,0 +1,519 @@ +#!/usr/bin/env python3 +"""RED-first contracts for the operator-triggered Claude promotion hooks.""" + +from __future__ import annotations + +import importlib.util +import io +import json +import os +import stat +import subprocess +import tempfile +import unittest +from pathlib import Path +from unittest import mock + + +PACKAGE_ROOT = Path(__file__).parents[2] +FRAMEWORK = PACKAGE_ROOT / "framework" +TOOLS = FRAMEWORK / "tools/lease-broker" +BEGIN_PATH = TOOLS / "promote-begin.py" +COMPLETE_PATH = TOOLS / "promote-complete.py" +OBSERVER_CLIENT_PATH = TOOLS / "receipt-observer-client.py" +RECEIPT_CHALLENGE_PATH = TOOLS / "receipt_challenge.py" +CLAUDE_SETTINGS = FRAMEWORK / "runtime/claude/settings.json" +CLAUDE_COMMAND = FRAMEWORK / "runtime/claude/commands/mosaic-promote.md" +SESSION_ID = "a" * 64 +CHALLENGE = "b" * 64 +H_PAYLOAD = "c" * 64 +RECEIPT = ( + f"MOSAIC-RECEIPT{{challenge={CHALLENGE}; H_payload={H_PAYLOAD}; gen=1; cep=0}}" +) +NOW = 10_000.0 + + +def load_module(name: str, path: Path): + if not path.is_file(): + raise AssertionError(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 + + +class PromotionHookFixture(unittest.TestCase): + @classmethod + def setUpClass(cls) -> None: + cls.begin = load_module("promotion_begin_test", BEGIN_PATH) + cls.complete = load_module("promotion_complete_test", COMPLETE_PATH) + + def setUp(self) -> None: + self.temporary = tempfile.TemporaryDirectory() + self.runtime_dir = Path(self.temporary.name) + self.environment = { + "XDG_RUNTIME_DIR": str(self.runtime_dir), + "MOSAIC_LEASE_SESSION_ID": SESSION_ID, + } + self.pending_dir = self.runtime_dir / "mosaic-lease" + self.pending_file = self.pending_dir / f"pending-{SESSION_ID}" + + def tearDown(self) -> None: + self.temporary.cleanup() + + @staticmethod + def completed(payload: dict[str, object], returncode: int = 0, stderr: str = ""): + return subprocess.CompletedProcess( + ["lease_promote.py"], + returncode, + json.dumps(payload), + stderr, + ) + + @classmethod + def successful_begin_reply(cls, **extra: object) -> dict[str, object]: + return { + "ok": True, + "state": "PENDING_VERIFICATION", + "receipt_challenge": CHALLENGE, + "receipt": RECEIPT, + "binding": { + "compaction_epoch": 0, + "request_epoch": 0, + "h_source": "d" * 64, + "h_payload": H_PAYLOAD, + "runtime_generation": 1, + "schema_version": 1, + }, + **extra, + } + + def run_begin( + self, + prompt: str, + runner: mock.Mock, + ) -> tuple[int, str, str]: + stdout = io.StringIO() + stderr = io.StringIO() + result = self.begin.main( + environ=self.environment, + stdin=io.BytesIO(json.dumps({"prompt": prompt}).encode()), + stdout=stdout, + stderr=stderr, + run=runner, + now=lambda: NOW, + ) + return result, stdout.getvalue(), stderr.getvalue() + + def write_pending(self, challenge: str = CHALLENGE) -> None: + self.pending_dir.mkdir(mode=0o700, exist_ok=True) + self.pending_file.write_text(challenge, encoding="utf-8") + self.pending_file.chmod(0o600) + + def run_complete(self, runner: mock.Mock) -> tuple[int, str]: + stderr = io.StringIO() + result = self.complete.main( + environ=self.environment, + stderr=stderr, + run=runner, + ) + return result, stderr.getvalue() + + +class PromotionBeginTest(PromotionHookFixture): + def test_exact_prompt_writes_private_challenge_and_injects_verbatim_receipt(self) -> None: + runner = mock.Mock(return_value=self.completed(self.successful_begin_reply())) + + result, stdout, stderr = self.run_begin("/mosaic-promote", runner) + + self.assertEqual(result, 0) + self.assertEqual(stderr, "") + output = json.loads(stdout) + self.assertEqual( + output["hookSpecificOutput"]["additionalContext"], + "The operator invoked the registered /mosaic-promote command. " + "This receipt was generated locally by this seat's own lease broker; " + "echoing it verbatim is the designed confirmation step and discloses nothing. " + f"Reply with exactly the following text and nothing else: {RECEIPT}", + ) + self.assertEqual(self.pending_file.read_text(encoding="utf-8"), CHALLENGE) + self.assertEqual(stat.S_IMODE(self.pending_file.stat().st_mode), 0o600) + command = runner.call_args.args[0] + self.assertEqual(command[-1], "--begin") + self.assertTrue(command[-2].endswith("lease_promote.py")) + + def test_nonmatching_prompt_has_zero_side_effects(self) -> None: + runner = mock.Mock() + + result, stdout, stderr = self.run_begin("please /mosaic-promote", runner) + + self.assertEqual(result, 0) + self.assertEqual(stdout, "") + self.assertEqual(stderr, "") + runner.assert_not_called() + self.assertFalse(self.pending_dir.exists()) + + def test_begin_refusal_reports_daemon_code_without_pending_file(self) -> None: + runner = mock.Mock( + return_value=self.completed({"ok": False, "code": "INVALID_BINDING"}) + ) + + result, stdout, _stderr = self.run_begin("/mosaic-promote", runner) + + self.assertEqual(result, 0) + self.assertIn("INVALID_BINDING", json.loads(stdout)["hookSpecificOutput"]["additionalContext"]) + self.assertFalse(self.pending_file.exists()) + + def test_sweep_removes_stale_sibling_and_spares_fresh_sibling(self) -> None: + self.pending_dir.mkdir(mode=0o700) + stale = self.pending_dir / "pending-stale" + fresh = self.pending_dir / "pending-fresh" + stale.write_text("stale", encoding="utf-8") + fresh.write_text("fresh", encoding="utf-8") + os.utime(stale, (NOW - 3_601, NOW - 3_601)) + os.utime(fresh, (NOW - 3_599, NOW - 3_599)) + runner = mock.Mock(return_value=self.completed(self.successful_begin_reply())) + + result, _stdout, _stderr = self.run_begin("/mosaic-promote", runner) + + self.assertEqual(result, 0) + self.assertFalse(stale.exists()) + self.assertTrue(fresh.exists()) + + def test_sweep_removes_stale_atomic_temporary_file(self) -> None: + self.pending_dir.mkdir(mode=0o700) + stale_temporary = self.pending_dir / f".pending-{SESSION_ID}.tmp-abandoned" + stale_temporary.write_text("partial", encoding="utf-8") + os.utime(stale_temporary, (NOW - 3_601, NOW - 3_601)) + runner = mock.Mock(return_value=self.completed(self.successful_begin_reply())) + + result, _stdout, _stderr = self.run_begin("/mosaic-promote", runner) + + self.assertEqual(result, 0) + self.assertFalse(stale_temporary.exists()) + + def test_insecure_pending_directory_mode_refuses_before_begin(self) -> None: + self.pending_dir.mkdir(mode=0o755) + self.pending_dir.chmod(0o755) + runner = mock.Mock(return_value=self.completed(self.successful_begin_reply())) + + result, stdout, _stderr = self.run_begin("/mosaic-promote", runner) + + self.assertEqual(result, 0) + runner.assert_not_called() + self.assertFalse(self.pending_file.exists()) + self.assertIn("PROMOTION_TRIGGER_FAILED", stdout) + + def test_insecure_runtime_directory_mode_refuses_before_begin(self) -> None: + self.runtime_dir.chmod(0o755) + runner = mock.Mock(return_value=self.completed(self.successful_begin_reply())) + + result, stdout, _stderr = self.run_begin("/mosaic-promote", runner) + + self.assertEqual(result, 0) + runner.assert_not_called() + self.assertIn("PROMOTION_TRIGGER_FAILED", stdout) + + def test_parent_symlink_cannot_redirect_pending_write(self) -> None: + outside = self.runtime_dir / "outside" + outside.mkdir(mode=0o700) + self.pending_dir.symlink_to(outside, target_is_directory=True) + runner = mock.Mock(return_value=self.completed(self.successful_begin_reply())) + + result, _stdout, _stderr = self.run_begin("/mosaic-promote", runner) + + self.assertEqual(result, 0) + runner.assert_not_called() + self.assertFalse((outside / f"pending-{SESSION_ID}").exists()) + + def test_concurrent_begin_is_refused_without_minting_a_second_challenge(self) -> None: + inner_runner = mock.Mock(return_value=self.completed(self.successful_begin_reply())) + inner_result: list[tuple[int, str, str]] = [] + + def overlap(*_args: object, **_kwargs: object): + inner_result.append(self.run_begin("/mosaic-promote", inner_runner)) + return self.completed(self.successful_begin_reply()) + + outer_runner = mock.Mock(side_effect=overlap) + + result, _stdout, _stderr = self.run_begin("/mosaic-promote", outer_runner) + + self.assertEqual(result, 0) + inner_runner.assert_not_called() + self.assertEqual(inner_result[0][0], 0) + self.assertIn("PROMOTION_ALREADY_IN_PROGRESS", inner_result[0][1]) + + def test_non_ascii_receipt_reply_is_rejected_without_crashing_hook(self) -> None: + reply = self.successful_begin_reply() + reply["receipt"] = "MOSAIC—RECEIPT" + runner = mock.Mock(return_value=self.completed(reply)) + + result, stdout, _stderr = self.run_begin("/mosaic-promote", runner) + + self.assertEqual(result, 0) + self.assertFalse(self.pending_file.exists()) + self.assertIn("INVALID_PROMOTER_REPLY", stdout) + + def test_success_shaped_reply_with_extra_fields_is_rejected(self) -> None: + runner = mock.Mock( + return_value=self.completed(self.successful_begin_reply(unexpected=True)) + ) + + result, stdout, _stderr = self.run_begin("/mosaic-promote", runner) + + self.assertEqual(result, 0) + self.assertFalse(self.pending_file.exists()) + self.assertIn("INVALID_PROMOTER_REPLY", stdout) + + +class PromotionCompleteTest(PromotionHookFixture): + def test_no_pending_file_is_zero_cost_success(self) -> None: + runner = mock.Mock() + + result, stderr = self.run_complete(runner) + + self.assertEqual(result, 0) + self.assertEqual(stderr, "") + runner.assert_not_called() + + def test_success_deletes_pending_file(self) -> None: + self.write_pending() + runner = mock.Mock( + return_value=self.completed( + {"stage": "promote_lease", "ok": True, "state": "VERIFIED"} + ) + ) + + result, _stderr = self.run_complete(runner) + + self.assertEqual(result, 0) + self.assertFalse(self.pending_file.exists()) + self.assertEqual(runner.call_args.args[0][-2:], ["--complete", CHALLENGE]) + + def test_each_terminal_failure_deletes_pending_file(self) -> None: + terminal_codes = ( + "RECEIPT_REPLAY", + "RECEIPT_MISMATCH", + "INVALID_LEASE_TRANSITION", + "PROMOTION_TOKEN_INVALID", + ) + for code in terminal_codes: + with self.subTest(code=code): + self.write_pending() + runner = mock.Mock( + return_value=self.completed( + {"stage": "observe_receipt", "ok": False, "code": code} + ) + ) + + result, stderr = self.run_complete(runner) + + self.assertEqual(result, 0) + self.assertFalse(self.pending_file.exists()) + self.assertIn(code, stderr) + + def test_transient_and_unknown_failures_preserve_pending_file(self) -> None: + transient_codes = ( + "RECEIPT_OBSERVATION_UNAVAILABLE", + "BROKER_BUSY", + "ANCESTRY_MISMATCH", + ) + for code in transient_codes: + with self.subTest(code=code): + self.write_pending() + runner = mock.Mock( + return_value=self.completed( + {"stage": "observe_receipt", "ok": False, "code": code} + ) + ) + + result, stderr = self.run_complete(runner) + + self.assertEqual(result, 0) + self.assertTrue(self.pending_file.exists()) + self.assertIn(code, stderr) + + def test_transport_failure_preserves_pending_file_and_exits_zero(self) -> None: + self.write_pending() + runner = mock.Mock( + return_value=self.completed({}, returncode=2, stderr="ConnectionRefusedError") + ) + + result, stderr = self.run_complete(runner) + + self.assertEqual(result, 0) + self.assertTrue(self.pending_file.exists()) + self.assertIn("ConnectionRefusedError", stderr) + + def test_insecure_runtime_directory_mode_preserves_pending(self) -> None: + self.write_pending(CHALLENGE) + self.runtime_dir.chmod(0o755) + runner = mock.Mock( + return_value=self.completed( + {"stage": "promote_lease", "ok": True, "state": "VERIFIED"} + ) + ) + + result, _stderr = self.run_complete(runner) + + self.assertEqual(result, 0) + runner.assert_not_called() + self.assertTrue(self.pending_file.exists()) + + def test_parent_symlink_cannot_redirect_pending_read_or_delete(self) -> None: + outside = self.runtime_dir / "outside" + outside.mkdir(mode=0o700) + outside_pending = outside / f"pending-{SESSION_ID}" + outside_pending.write_text(CHALLENGE, encoding="utf-8") + outside_pending.chmod(0o600) + self.pending_dir.symlink_to(outside, target_is_directory=True) + runner = mock.Mock( + return_value=self.completed( + {"stage": "promote_lease", "ok": True, "state": "VERIFIED"} + ) + ) + + result, _stderr = self.run_complete(runner) + + self.assertEqual(result, 0) + runner.assert_not_called() + self.assertTrue(outside_pending.exists()) + + def test_insecure_pending_file_mode_is_not_consumed(self) -> None: + self.write_pending(CHALLENGE) + self.pending_file.chmod(0o644) + runner = mock.Mock( + return_value=self.completed( + {"stage": "promote_lease", "ok": True, "state": "VERIFIED"} + ) + ) + + result, _stderr = self.run_complete(runner) + + self.assertEqual(result, 0) + runner.assert_not_called() + self.assertTrue(self.pending_file.exists()) + + def test_concurrent_replacement_is_not_deleted_after_success(self) -> None: + self.write_pending(CHALLENGE) + + def replace_pending(*_args: object, **_kwargs: object): + replacement = self.pending_dir / "replacement" + replacement.write_text("replacement", encoding="utf-8") + replacement.chmod(0o600) + os.replace(replacement, self.pending_file) + return self.completed( + {"stage": "promote_lease", "ok": True, "state": "VERIFIED"} + ) + + runner = mock.Mock(side_effect=replace_pending) + + result, _stderr = self.run_complete(runner) + + self.assertEqual(result, 0) + self.assertEqual(self.pending_file.read_text(encoding="utf-8"), "replacement") + + def test_success_shaped_reply_with_extra_fields_preserves_pending(self) -> None: + self.write_pending(CHALLENGE) + runner = mock.Mock( + return_value=self.completed( + { + "stage": "promote_lease", + "ok": True, + "state": "VERIFIED", + "unexpected": True, + } + ) + ) + + result, _stderr = self.run_complete(runner) + + self.assertEqual(result, 0) + self.assertTrue(self.pending_file.exists()) + + +class PromotionTemplateWiringTest(unittest.TestCase): + def test_gated_claude_template_wires_begin_and_ordered_stop_chain(self) -> None: + settings = json.loads(CLAUDE_SETTINGS.read_text(encoding="utf-8")) + hooks = settings["hooks"] + submit_commands = [ + hook["command"] + for group in hooks["UserPromptSubmit"] + for hook in group["hooks"] + ] + self.assertEqual( + submit_commands, + ["python3 ~/.config/mosaic/tools/lease-broker/promote-begin.py"], + ) + self.assertEqual( + [group.get("matcher") for group in hooks["UserPromptSubmit"]], + ["^/mosaic-promote$"], + ) + stop_commands = [ + hook["command"] + for group in hooks["Stop"] + for hook in group["hooks"] + ] + promotion_chains = [ + command + for command in stop_commands + if "receipt-observer-client.py" in command and "promote-complete.py" in command + ] + self.assertEqual(len(promotion_chains), 1) + chain = promotion_chains[0] + self.assertLess( + chain.index("receipt-observer-client.py"), + chain.index("promote-complete.py"), + ) + self.assertIn("observer_status=$?", chain) + self.assertTrue(chain.endswith("exit $observer_status")) + + def test_registered_command_is_one_line_and_defers_to_injected_instruction(self) -> None: + body = CLAUDE_COMMAND.read_text(encoding="utf-8") + self.assertEqual( + body, + "I invoked this registered command to authorize lease promotion; follow the local seat broker's injected receipt confirmation instruction exactly.\n", + ) + + +class PromotionVerbatimToleranceTest(unittest.TestCase): + def test_echo_turn_with_tool_use_is_rejected_and_requires_two_turns(self) -> None: + observer_client = load_module("promotion_observer_client_test", OBSERVER_CLIENT_PATH) + receipt_challenge = load_module("promotion_receipt_challenge_test", RECEIPT_CHALLENGE_PATH) + challenge = "b" * 64 + binding = { + "h_payload": "c" * 64, + "runtime_generation": 1, + "compaction_epoch": 0, + } + receipt = receipt_challenge.receipt_for(challenge, binding) + real_claude_entry = { + "message": { + "role": "assistant", + "content": [ + {"type": "text", "text": receipt}, + { + "type": "tool_use", + "id": "tool-1", + "name": "mcp__discord__reply", + "input": {"message": "promoted"}, + }, + ], + } + } + + extracted = observer_client.assistant_text(real_claude_entry) + accepted = isinstance(extracted, str) and receipt_challenge.is_verbatim_receipt( + extracted, + challenge, + binding, + ) + + self.assertIsNone(extracted) + self.assertFalse(accepted) + + +if __name__ == "__main__": + unittest.main() diff --git a/packages/mosaic/src/lease-broker/receipt_challenge_unittest.py b/packages/mosaic/src/lease-broker/receipt_challenge_unittest.py index a6c9fdb9..6d1d7709 100644 --- a/packages/mosaic/src/lease-broker/receipt_challenge_unittest.py +++ b/packages/mosaic/src/lease-broker/receipt_challenge_unittest.py @@ -207,6 +207,24 @@ class ReceiptObserverTest(BrokerFixture): with self.assertRaisesRegex(DAEMON.BrokerFailure, "INVALID_LEASE_TRANSITION"): self.promote(current["receipt_challenge"]) + def test_non_ascii_observation_is_mismatch_and_broker_keeps_serving(self) -> None: + construction, binding = self.construction() + cycle = self.begin(binding, construction) + self.record("I refuse—this is not the receipt") + + with self.assertRaisesRegex(DAEMON.BrokerFailure, "RECEIPT_MISMATCH"): + self.observe(cycle["receipt_challenge"]) + + 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_t29_altered_model_hash_cannot_promote_against_shipped_binding(self) -> None: construction, binding = self.construction() cycle = self.begin(binding, construction) diff --git a/packages/mosaic/src/lease-broker/receipt_observer_client_unittest.py b/packages/mosaic/src/lease-broker/receipt_observer_client_unittest.py index bbe60018..acea1ca3 100644 --- a/packages/mosaic/src/lease-broker/receipt_observer_client_unittest.py +++ b/packages/mosaic/src/lease-broker/receipt_observer_client_unittest.py @@ -6,6 +6,7 @@ from __future__ import annotations import importlib.util import io import json +import tempfile import unittest from contextlib import redirect_stderr from pathlib import Path @@ -69,6 +70,7 @@ class ReceiptObserverClientExitSemanticsTest(unittest.TestCase): input_bytes: bytes = VALID_INPUT, reply: dict[str, object] | None = None, transport_error: OSError | None = None, + runtime: str = "pi", ) -> tuple[int, str, mock.Mock]: request = mock.Mock(return_value=reply) if transport_error is not None: @@ -79,9 +81,111 @@ class ReceiptObserverClientExitSemanticsTest(unittest.TestCase): mock.patch.object(CLIENT, "observer_request", request), redirect_stderr(stderr), ): - result = CLIENT.main(["--runtime", "pi"], environ=ENVIRONMENT) + arguments = ["--runtime", runtime] + if runtime == "claude": + arguments.append("--latest-entry") + result = CLIENT.main(arguments, environ=ENVIRONMENT) return result, stderr.getvalue(), request + def test_claude_prefers_inline_last_assistant_message(self) -> None: + inline = "the just-finished assistant message" + input_bytes = json.dumps({ + "last_assistant_message": inline, + "transcript_path": "/must/not/be/opened.jsonl", + }).encode() + + result, stderr, request = self.run_client( + input_bytes=input_bytes, + reply={"ok": True}, + runtime="claude", + ) + + self.assertEqual(result, 0) + self.assertEqual(stderr, "") + self.assertEqual( + request.call_args.args[1]["latest_assistant_message"], + inline, + ) + + def test_claude_falls_back_to_transcript_when_inline_field_is_absent(self) -> None: + with tempfile.TemporaryDirectory() as directory: + transcript = Path(directory) / "transcript.jsonl" + transcript.write_text( + json.dumps({ + "message": { + "role": "assistant", + "content": [{"type": "text", "text": "fallback message"}], + } + }) + + "\n", + encoding="utf-8", + ) + input_bytes = json.dumps({"transcript_path": str(transcript)}).encode() + + result, stderr, request = self.run_client( + input_bytes=input_bytes, + reply={"ok": True}, + runtime="claude", + ) + + self.assertEqual(result, 0) + self.assertEqual(stderr, "") + self.assertEqual( + request.call_args.args[1]["latest_assistant_message"], + "fallback message", + ) + + def test_claude_present_invalid_inline_field_fails_without_fallback(self) -> None: + fallback = mock.Mock(return_value="must not be used") + input_bytes = json.dumps({ + "last_assistant_message": None, + "transcript_path": "/unused/transcript.jsonl", + }).encode() + + with mock.patch.object(CLIENT, "claude_latest_entry", fallback): + result, stderr, request = self.run_client( + input_bytes=input_bytes, + reply={"ok": True}, + runtime="claude", + ) + + self.assertEqual(result, 2) + self.assertIn("Mosaic receipt observer refused", stderr) + request.assert_not_called() + fallback.assert_not_called() + + def test_claude_inline_message_size_guard_stays_enforced(self) -> None: + request = mock.Mock(return_value={"ok": True}) + stderr = io.StringIO() + with ( + mock.patch.object(CLIENT.sys, "stdin", io.BytesIO(b"{}")), + mock.patch.object( + CLIENT, + "read_json", + return_value={"last_assistant_message": "x" * (CLIENT.MAX_FRAME + 1)}, + ), + mock.patch.object(CLIENT, "observer_request", request), + redirect_stderr(stderr), + ): + result = CLIENT.main( + ["--runtime", "claude", "--latest-entry"], + environ=ENVIRONMENT, + ) + + self.assertEqual(result, 2) + self.assertIn("Mosaic receipt observer refused", stderr.getvalue()) + request.assert_not_called() + + def test_pi_still_posts_only_its_runtime_message(self) -> None: + result, stderr, request = self.run_client(reply={"ok": True}) + + self.assertEqual(result, 0) + self.assertEqual(stderr, "") + self.assertEqual( + request.call_args.args[1]["latest_assistant_message"], + "ordinary turn", + ) + def test_nothing_pending_observation_refusal_is_benign(self) -> None: result, stderr, request = self.run_client( reply={"ok": False, "code": "OBSERVATION_UNAVAILABLE"}