feat(mosaic): mechanically authorize lease promotion

This commit is contained in:
Jason Woltje
2026-08-11 20:51:03 -05:00
parent 239a2a93f1
commit 709a23d08c
6 changed files with 216 additions and 41 deletions
@@ -233,6 +233,14 @@ for runtime_file in \
copy_file_managed "$src" "$HOME/.claude/$runtime_file"
done
if [[ -d "$MOSAIC_HOME/runtime/claude/commands" ]]; then
mkdir -p "$HOME/.claude/commands"
for command_file in "$MOSAIC_HOME/runtime/claude/commands/"*; do
[[ -f "$command_file" ]] || continue
copy_file_managed "$command_file" "$HOME/.claude/commands/$(basename "$command_file")"
done
fi
# OpenCode runtime adapter (thin pointer to AGENTS.md)
opencode_adapter="$MOSAIC_HOME/runtime/opencode/AGENTS.md"
if [[ -f "$opencode_adapter" ]]; then
@@ -4,6 +4,7 @@
from __future__ import annotations
import fcntl
import importlib.util
import json
import os
import secrets
@@ -21,13 +22,26 @@ if _MODULE_DIRECTORY not in sys.path:
from receipt_challenge import receipt_for # noqa: E402
_observer_spec = importlib.util.spec_from_file_location(
"mosaic_receipt_observer_client", Path(__file__).resolve().with_name("receipt-observer-client.py")
)
if _observer_spec is None or _observer_spec.loader is None:
raise RuntimeError("unable to load receipt observer client")
_observer_module = importlib.util.module_from_spec(_observer_spec)
_observer_spec.loader.exec_module(_observer_module)
observer_request = _observer_module.observer_request
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"
AUTHORIZATION_DIRECTORY: Final = "authorizations"
AUTHORIZATION_TTL_SECONDS: Final = 60
LEASE_TTL_SECONDS: Final = 60 * 60
LOCK_FILE: Final = "promotion.lock"
RESULT_FILE: Final = "last-result.json"
EXPECTED_BEGIN_KEYS: Final = frozenset(
{"ok", "state", "receipt_challenge", "receipt", "binding"}
)
@@ -170,6 +184,57 @@ def sweep_stale_pending(directory_descriptor: int, current_time: float) -> None:
os.fsync(directory_descriptor)
def consume_authorization(directory_descriptor: int, session_id: str, wall_clock: float) -> str | None:
flags = os.O_RDONLY | getattr(os, "O_CLOEXEC", 0) | getattr(os, "O_DIRECTORY", 0) | getattr(os, "O_NOFOLLOW", 0)
try:
authorization_descriptor = os.open(AUTHORIZATION_DIRECTORY, flags, dir_fd=directory_descriptor)
except FileNotFoundError:
return None
try:
metadata = os.fstat(authorization_descriptor)
if not stat.S_ISDIR(metadata.st_mode) or metadata.st_uid != os.getuid() or stat.S_IMODE(metadata.st_mode) != 0o700:
raise ValueError("unsafe promotion authorization directory")
name = f"{session_id}.auth"
try:
descriptor = os.open(name, os.O_RDONLY | getattr(os, "O_CLOEXEC", 0) | getattr(os, "O_NOFOLLOW", 0), dir_fd=authorization_descriptor)
except FileNotFoundError:
return None
try:
token_metadata = os.fstat(descriptor)
if not stat.S_ISREG(token_metadata.st_mode) or token_metadata.st_uid != os.getuid() or stat.S_IMODE(token_metadata.st_mode) != 0o600 or token_metadata.st_size <= 0 or token_metadata.st_size > MAX_FRAME:
raise ValueError("unsafe promotion authorization")
raw = os.read(descriptor, MAX_FRAME + 1)
finally:
os.close(descriptor)
os.unlink(name, dir_fd=authorization_descriptor)
os.fsync(authorization_descriptor)
token = json.loads(raw, object_pairs_hook=reject_duplicate_json_keys)
if not isinstance(token, dict) or set(token) != {"nonce", "seat", "session_id", "expires_at", "ts"}:
return None
nonce = token.get("nonce")
expires_at = token.get("expires_at")
issued_at = token.get("ts")
if token.get("session_id") != session_id or not isinstance(token.get("seat"), str) or not isinstance(nonce, str) or len(nonce) != 64 or any(char not in "0123456789abcdef" for char in nonce) or type(expires_at) not in (int, float) or type(issued_at) not in (int, float) or expires_at <= wall_clock or expires_at > issued_at + AUTHORIZATION_TTL_SECONDS:
return None
return nonce
finally:
os.close(authorization_descriptor)
def write_result(directory_descriptor: int, attempt_id: str, verified: bool, reason: str | None, session_id: str, wall_clock: float) -> None:
temporary = f".{RESULT_FILE}.tmp-{secrets.token_hex(8)}"
descriptor = os.open(temporary, os.O_WRONLY | os.O_CREAT | os.O_EXCL | getattr(os, "O_CLOEXEC", 0) | getattr(os, "O_NOFOLLOW", 0), 0o600, dir_fd=directory_descriptor)
try:
os.fchmod(descriptor, 0o600)
with os.fdopen(descriptor, "w", encoding="utf-8", closefd=False) as stream:
json.dump({"attempt_id": attempt_id, "expires_at_wallclock": wall_clock + LEASE_TTL_SECONDS if verified else None, "reason": reason, "session_id": session_id, "ts": wall_clock, "verified": verified}, stream, separators=(",", ":"), sort_keys=True)
stream.flush(); os.fsync(stream.fileno())
os.replace(temporary, RESULT_FILE, src_dir_fd=directory_descriptor, dst_dir_fd=directory_descriptor)
os.fsync(directory_descriptor)
finally:
os.close(descriptor)
def write_pending(directory_descriptor: int, name: str, challenge: str) -> None:
temporary = f".{name}.tmp-{secrets.token_hex(8)}"
flags = (
@@ -285,42 +350,43 @@ def main(
lock_descriptor: int | None = None
try:
runtime_dir, pending_name = session_pending_name(source_environment)
session_id = source_environment["MOSAIC_LEASE_SESSION_ID"]
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}.")
wall_clock = now()
nonce = consume_authorization(directory_descriptor, session_id, wall_clock)
if nonce is None:
write_result(directory_descriptor, "0" * 64, False, "NOT_AUTHORIZED", session_id, wall_clock)
print("Mosaic promotion denied: NOT_AUTHORIZED.", file=error_stream)
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']}",
sweep_stale_pending(directory_descriptor, wall_clock)
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 or reply is None:
write_result(directory_descriptor, nonce, False, code or "PROMOTION_BEGIN_FAILED", session_id, now())
return 0
challenge = str(reply["receipt_challenge"])
observation = observer_request(
Path(source_environment["MOSAIC_RECEIPT_OBSERVER_SOCKET"]),
{"action": "record_runtime_observation", "session_id": session_id, "runtime_generation": int(source_environment["MOSAIC_RUNTIME_GENERATION"]), "runtime": "claude", "latest_assistant_message": reply["receipt"]},
)
if set(observation) != {"ok"} or observation.get("ok") is not True:
write_result(directory_descriptor, challenge, False, "OBSERVATION_REJECTED", session_id, now())
return 0
completion = run([sys.executable, "-I", "-S", "-B", str(PROMOTER), "--complete", challenge], check=False, capture_output=True, text=True, env=dict(source_environment), timeout=PROMOTER_TIMEOUT_SECONDS)
try:
outcome = json.loads(completion.stdout, object_pairs_hook=reject_duplicate_json_keys)
except (json.JSONDecodeError, ValueError):
outcome = None
if completion.returncode == 0 and isinstance(outcome, dict) and outcome.get("stage") == "promote_lease" and outcome.get("ok") is True and outcome.get("state") == "VERIFIED":
write_result(directory_descriptor, challenge, True, None, session_id, now())
else:
reason = outcome.get("code") if isinstance(outcome, dict) and isinstance(outcome.get("code"), str) else "PROMOTION_INCOMPLETE"
write_result(directory_descriptor, challenge, False, reason, session_id, now())
except PromotionAlreadyInProgress:
emit_context(
output_stream,
"Mosaic promotion did not begin: PROMOTION_ALREADY_IN_PROGRESS.",
)
print("Mosaic promotion denied: PROMOTION_ALREADY_IN_PROGRESS.", file=error_stream)
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)