feat(lease): add single-turn Claude promotion trigger
This commit is contained in:
@@ -0,0 +1 @@
|
|||||||
|
I invoked this registered command to authorize lease promotion; follow the local seat broker's injected receipt confirmation instruction exactly.
|
||||||
@@ -32,6 +32,18 @@
|
|||||||
]
|
]
|
||||||
}
|
}
|
||||||
],
|
],
|
||||||
|
"UserPromptSubmit": [
|
||||||
|
{
|
||||||
|
"matcher": "^/mosaic-promote$",
|
||||||
|
"hooks": [
|
||||||
|
{
|
||||||
|
"type": "command",
|
||||||
|
"command": "python3 ~/.config/mosaic/tools/lease-broker/promote-begin.py",
|
||||||
|
"timeout": 15
|
||||||
|
}
|
||||||
|
]
|
||||||
|
}
|
||||||
|
],
|
||||||
"PreToolUse": [
|
"PreToolUse": [
|
||||||
{
|
{
|
||||||
"matcher": ".*",
|
"matcher": ".*",
|
||||||
@@ -81,8 +93,8 @@
|
|||||||
"hooks": [
|
"hooks": [
|
||||||
{
|
{
|
||||||
"type": "command",
|
"type": "command",
|
||||||
"command": "python3 ~/.config/mosaic/tools/lease-broker/receipt-observer-client.py --runtime claude --latest-entry",
|
"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": 3
|
"timeout": 15
|
||||||
},
|
},
|
||||||
{
|
{
|
||||||
"type": "command",
|
"type": "command",
|
||||||
|
|||||||
@@ -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())
|
||||||
@@ -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())
|
||||||
@@ -129,7 +129,12 @@ def main(argv: Sequence[str] | None = None, *, environ: Mapping[str, str] | None
|
|||||||
if arguments.runtime == "claude":
|
if arguments.runtime == "claude":
|
||||||
if not arguments.latest_entry:
|
if not arguments.latest_entry:
|
||||||
raise ValueError("Claude observer requires --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:
|
else:
|
||||||
if arguments.latest_entry:
|
if arguments.latest_entry:
|
||||||
raise ValueError("Pi observer is message_end only")
|
raise ValueError("Pi observer is message_end only")
|
||||||
|
|||||||
@@ -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."""
|
"""Require the exact one current-cycle receipt, not a transcript substring."""
|
||||||
|
|
||||||
expected = receipt_for(challenge, binding)
|
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:
|
def latest_assistant_digest(message: str) -> str:
|
||||||
|
|||||||
@@ -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": "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": {
|
"dependencies": {
|
||||||
"@mosaicstack/brain": "workspace:*",
|
"@mosaicstack/brain": "workspace:*",
|
||||||
|
|||||||
@@ -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()
|
||||||
@@ -207,6 +207,24 @@ class ReceiptObserverTest(BrokerFixture):
|
|||||||
with self.assertRaisesRegex(DAEMON.BrokerFailure, "INVALID_LEASE_TRANSITION"):
|
with self.assertRaisesRegex(DAEMON.BrokerFailure, "INVALID_LEASE_TRANSITION"):
|
||||||
self.promote(current["receipt_challenge"])
|
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:
|
def test_t29_altered_model_hash_cannot_promote_against_shipped_binding(self) -> None:
|
||||||
construction, binding = self.construction()
|
construction, binding = self.construction()
|
||||||
cycle = self.begin(binding, construction)
|
cycle = self.begin(binding, construction)
|
||||||
|
|||||||
@@ -6,6 +6,7 @@ from __future__ import annotations
|
|||||||
import importlib.util
|
import importlib.util
|
||||||
import io
|
import io
|
||||||
import json
|
import json
|
||||||
|
import tempfile
|
||||||
import unittest
|
import unittest
|
||||||
from contextlib import redirect_stderr
|
from contextlib import redirect_stderr
|
||||||
from pathlib import Path
|
from pathlib import Path
|
||||||
@@ -69,6 +70,7 @@ class ReceiptObserverClientExitSemanticsTest(unittest.TestCase):
|
|||||||
input_bytes: bytes = VALID_INPUT,
|
input_bytes: bytes = VALID_INPUT,
|
||||||
reply: dict[str, object] | None = None,
|
reply: dict[str, object] | None = None,
|
||||||
transport_error: OSError | None = None,
|
transport_error: OSError | None = None,
|
||||||
|
runtime: str = "pi",
|
||||||
) -> tuple[int, str, mock.Mock]:
|
) -> tuple[int, str, mock.Mock]:
|
||||||
request = mock.Mock(return_value=reply)
|
request = mock.Mock(return_value=reply)
|
||||||
if transport_error is not None:
|
if transport_error is not None:
|
||||||
@@ -79,9 +81,111 @@ class ReceiptObserverClientExitSemanticsTest(unittest.TestCase):
|
|||||||
mock.patch.object(CLIENT, "observer_request", request),
|
mock.patch.object(CLIENT, "observer_request", request),
|
||||||
redirect_stderr(stderr),
|
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
|
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:
|
def test_nothing_pending_observation_refusal_is_benign(self) -> None:
|
||||||
result, stderr, request = self.run_client(
|
result, stderr, request = self.run_client(
|
||||||
reply={"ok": False, "code": "OBSERVATION_UNAVAILABLE"}
|
reply={"ok": False, "code": "OBSERVATION_UNAVAILABLE"}
|
||||||
|
|||||||
Reference in New Issue
Block a user