feat(lease): add single-turn Claude promotion trigger
This commit is contained in:
@@ -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 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")
|
||||
|
||||
@@ -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:
|
||||
|
||||
Reference in New Issue
Block a user