Files
stack/v1/packages/mosaic/framework/tools/lease-broker/launch-runtime.py
T

223 lines
9.3 KiB
Python

#!/usr/bin/env python3
"""Register a runtime parent with the lease broker, then exec without changing PID."""
from __future__ import annotations
import argparse
import json
import os
import socket
import sys
import time
from collections.abc import Callable, Mapping, Sequence
from datetime import datetime, timezone
from pathlib import Path
from typing import Final
from activation_version_gate import (
EXPECTED_ACTIVATION_CAPABILITY,
ActivationCapability,
VersionCouplingError,
assert_activation_capability_matches,
default_probe_activation_capability,
)
from lease_generation import initialize_runtime_generation
MAX_FRAME: Final = 64 * 1024
BROKER_TIMEOUT_SECONDS: Final = 1.5
CLAUDE_DANGEROUS_FLAG: Final = "--dangerously-skip-permissions"
# Distinct, non-overlapping exit code for the C4 version-coupling gate (see
# `activation_version_gate.py`) — deliberately different from the `1`
# (broker registration failed closed) and `64` (usage error) codes already
# owned by this script, so a version-skew denial is unambiguous in caller
# logs/tests and is never confused with a broker-availability failure.
EXIT_VERSION_SKEW: Final = 65
def broker_request(socket_path: Path, request: dict[str, object]) -> dict[str, object]:
payload = (json.dumps(request, separators=(",", ":")) + "\n").encode()
response = bytearray()
with socket.socket(socket.AF_UNIX, socket.SOCK_STREAM) as connection:
connection.settimeout(BROKER_TIMEOUT_SECONDS)
connection.connect(str(socket_path))
connection.sendall(payload)
connection.shutdown(socket.SHUT_WR)
while len(response) <= MAX_FRAME:
chunk = connection.recv(min(4096, MAX_FRAME + 1 - len(response)))
if not chunk:
break
response.extend(chunk)
if len(response) > MAX_FRAME or not response.endswith(b"\n"):
raise ValueError("invalid broker reply")
value = json.loads(response)
if not isinstance(value, dict):
raise ValueError("invalid broker reply")
return value
def _self_starttime() -> str | None:
"""Field 22 of our own /proc stat — the anchor starttime the broker records.
Read past the comm field's parens, since a process name may contain them.
"""
try:
raw = Path(f"/proc/{os.getpid()}/stat").read_text()
return raw.rsplit(")", 1)[1].split()[19]
except (OSError, IndexError, ValueError):
return None
def _append_launch_record(environ: Mapping[str, str], record: dict[str, object]) -> None:
"""Append one NDJSON event to the #797 Runtime Session Ledger.
`fleet/run/sessions/` is operator-classified in framework-manifest.txt and is
already covered by test-upgrade-manifest-guard.sh, so an upgrade can neither
overwrite nor prune it. Files 0600 under a 0700 dir, matching what that guard
asserts.
Never raises: a launch must not be denied over bookkeeping. But it also never
fails silently — a missing record is exactly the kind of gap that made the
2026-08-06 MUTATOR_UNVERIFIED investigation cost a day.
"""
try:
mosaic_home = environ.get("MOSAIC_HOME") or str(Path.home() / ".config" / "mosaic")
directory = Path(mosaic_home) / "fleet" / "run" / "sessions"
directory.mkdir(parents=True, exist_ok=True)
os.chmod(directory, 0o700)
framed = {
"seq": time.time_ns() // 1_000_000,
"ts": datetime.now(timezone.utc).isoformat(),
**record,
}
path = directory / "events.ndjson"
descriptor = os.open(path, os.O_WRONLY | os.O_CREAT | os.O_APPEND, 0o600)
with os.fdopen(descriptor, "w") as handle:
handle.write(json.dumps(framed, separators=(",", ":")) + "\n")
except (OSError, ValueError, TypeError) as error:
print(f"[mosaic] WARNING: launch record not written: {error}", file=sys.stderr)
def main(
argv: Sequence[str] | None = None,
*,
environ: Mapping[str, str] | None = None,
request: Callable[[Path, dict[str, object]], dict[str, object]] = broker_request,
execute: Callable[[str, list[str], dict[str, str]], object] = os.execvpe,
initialize_generation: Callable[[Path, int], None] = initialize_runtime_generation,
probe_activation_capability: Callable[
[Mapping[str, str]], ActivationCapability | None
] = default_probe_activation_capability,
expected_activation_capability: ActivationCapability = EXPECTED_ACTIVATION_CAPABILITY,
) -> int:
parser = argparse.ArgumentParser()
parser.add_argument("--runtime", required=True, choices=("claude", "pi"))
parser.add_argument("--dangerous", action="store_true")
parser.add_argument("command", nargs=argparse.REMAINDER)
arguments = parser.parse_args(argv)
command = arguments.command
if command and command[0] == "--":
command = command[1:]
if not command:
print("lease-gated runtime command is required", file=sys.stderr)
return 64
if arguments.dangerous:
if arguments.runtime != "claude" or Path(command[0]).name != "claude":
print("dangerous mode is supported only for the Claude runtime", file=sys.stderr)
return 64
command = [command[0], CLAUDE_DANGEROUS_FLAG, *command[1:]]
source_environment = os.environ if environ is None else environ
# C4 version-coupling gate (#869 Point-1): before this ENFORCEMENT half
# chains into anything, assert that the ACTIVATION contract it is about
# to rely on (MOSAIC_LEASE_* injection, broker chaining) matches what
# this enforcement build expects. This is a build/deploy-defect check,
# not a broker-availability question, so it runs before — and
# independently of — broker registration below, and it FAILS LOUD: a
# clear stderr message plus a dedicated non-zero exit code, never a
# silent pass and never folded into the generic registration-failure
# branch.
try:
activation_capability = probe_activation_capability(source_environment)
assert_activation_capability_matches(
activation_capability,
expected_activation_capability,
)
except VersionCouplingError as version_error:
print(str(version_error), file=sys.stderr)
return EXIT_VERSION_SKEW
try:
socket_path = Path(source_environment["MOSAIC_LEASE_BROKER_SOCKET"])
generation = int(source_environment.get("MOSAIC_RUNTIME_GENERATION", "1"))
if generation < 0:
raise ValueError("invalid generation")
reply = request(
socket_path,
{
"action": "register_anchor",
"runtime_generation": generation,
},
)
session_id = reply.get("session_id")
if (
reply.get("ok") is not True
or not isinstance(session_id, str)
or len(session_id) != 64
or any(character not in "0123456789abcdef" for character in session_id)
):
raise ValueError("registration refused")
generation_file = socket_path.parent / f"generation-{session_id}.state"
initialize_generation(generation_file, generation)
except (KeyError, ValueError, OSError, json.JSONDecodeError):
print("Mosaic lease broker registration failed; runtime launch denied.", file=sys.stderr)
return 1
# Immutable launch record, half two. `mosaic` wrote `session.launch` with the
# config/provenance it knows; only this process knows the broker session id
# and the activation capability it just asserted. os.execvpe preserves the
# PID, so this PID is BOTH the anchor pid and the join key back to that
# record. Never fatal — bookkeeping must not deny a launch — but never
# silent either.
_append_launch_record(
source_environment,
{
"kind": "lease.register",
# Joins back to `mosaic`'s session.launch record. NOT pid: execRuntime()
# spawns rather than execs, so this process is a CHILD of mosaic with a
# different pid. This pid IS the broker anchor pid (os.execvpe below
# preserves it), which is a separate and still-useful fact.
"launch_id": source_environment.get("MOSAIC_LAUNCH_ID"),
"pid": os.getpid(),
"runtime": arguments.runtime,
"session_id": session_id,
"runtime_generation": generation,
"generation_file": str(generation_file),
"anchor_starttime": _self_starttime(),
"activation_capability": activation_capability,
"command": Path(command[0]).name,
},
)
environment = dict(source_environment)
environment["MOSAIC_LEASE_SESSION_ID"] = session_id
environment["MOSAIC_RUNTIME_GENERATION"] = str(generation)
environment["MOSAIC_LEASE_GENERATION_FILE"] = str(generation_file)
environment["MOSAIC_LEASE_RUNTIME"] = arguments.runtime
# Matches daemon.py's production default; deployments using a distinct
# observer socket may set this authenticated transport path explicitly.
environment.setdefault(
"MOSAIC_RECEIPT_OBSERVER_SOCKET",
str(socket_path.with_name("receipt-observer.sock")),
)
try:
execute(command[0], command, environment)
except OSError:
print("Mosaic lease-gated runtime exec failed.", file=sys.stderr)
return 1
return 0
if __name__ == "__main__":
raise SystemExit(main())