mirror of
https://github.com/razzant/ouroboros.git
synced 2026-10-03 04:07:04 +00:00
Add the transport half of the fifth nanny verb to the Claudexor gateway: POST /v2/runs/:id/messages with a REQUIRED Idempotency-Key (the caller's message identity), an optional expectedAttemptId and an optional per-call timeout. Like answer_interaction it never goes through _request: any body carrying a typed outcome (LIVE_MESSAGE_OUTCOMES) is returned as the answer whatever the HTTP status; every untyped refusal (the daemon's 404, the 409 idempotency problems, 400, 501, 5xx) is typed through _problem with its status code so the verb classifies by code AND status; transport death is daemon_unreachable; a 2xx without a typed outcome is malformed_response. run_message_supported negotiates the route structurally from the engine's /v2/operations catalog (Express-style template, like the control and interaction-answer rows), so an engine older than the live-message release never receives a POST whose 404 would read as "no such run". Co-authored-by: Ouroboros <311266734+ouroboros-agent@users.noreply.github.com>
1495 lines
74 KiB
Python
1495 lines
74 KiB
Python
"""Claudexor v3 control-plane transport.
|
||
|
||
Pure I/O, like every gateway (docs/DEVELOPMENT.md "Gateway Boundary Pattern"): discover the
|
||
loopback daemon, negotiate the protocol, and translate the typed ``/v2`` surface
|
||
into Python primitives. No routing, no policy, no harness identity branches — the
|
||
caller asks for a CAPABILITY (an access profile, an opaque route id) and reads the
|
||
manifest Claudexor publishes.
|
||
|
||
Two Claudexor I/O surfaces live here, not one: the HTTP control plane, and the run
|
||
tree Claudexor writes on disk. The engine deliberately records some APPLIED facts
|
||
only as run artifacts (an attempt's harness HOME is one), so a caller that has to
|
||
verify what was enforced has nowhere else to read them — and the path layout is
|
||
Claudexor's, which makes it this module's business rather than a policy caller's.
|
||
|
||
Token custody: the daemon bearer token grants the ENTIRE ``/v2`` surface. It is
|
||
read here, held in this process only, and never returned to a caller, written into
|
||
a ``ToolContext``, put in a child's environment, or handed to a harness sandbox.
|
||
"""
|
||
|
||
from __future__ import annotations
|
||
|
||
import hashlib
|
||
import json
|
||
import logging
|
||
import os
|
||
import pathlib
|
||
import re
|
||
import uuid
|
||
from dataclasses import dataclass
|
||
from typing import Any, Dict, List, Optional
|
||
|
||
import httpx
|
||
|
||
from ouroboros.config import (
|
||
CLAUDEXOR_MIN_VERSION,
|
||
CLAUDEXOR_PROTOCOL_MAJOR,
|
||
get_claudexor_quota_refresh_timeout_sec,
|
||
)
|
||
|
||
log = logging.getLogger(__name__)
|
||
|
||
CONTROL_API_REL = ".claudexor/v3/daemon/control-api.json"
|
||
PROTOCOL_HEADER = "X-Claudexor-Protocol-Major"
|
||
_CONNECT_TIMEOUT_SEC = 5.0
|
||
# The client-wide read default, and the CEILING on any polling/self-bounding caller's ask
|
||
# (`delegate_progress.poll_bound` reads it): a per-request value above it is not a bound
|
||
# at all, it is a hung read granted more rope than it would have had. Generous on
|
||
# purpose: most calls here would rather wait than fail, and a run start can take a while
|
||
# to answer.
|
||
_READ_TIMEOUT_SEC = 60.0
|
||
# The FLOOR under non-strict bounded wait/admission asks. Owned-daemon startup also
|
||
# imports it as the CEILING for each fast loopback liveness probe. Every ordinary
|
||
# `delegate_wait` poll asks for what its window has left; this is where that narrowing
|
||
# stops, because a nearly spent window asking for its own 0.2s turns a healthy daemon
|
||
# into a timeout, while the 60s default would outrun the very deadline the wait clamps
|
||
# itself to (measured, 4.51s of wall against a 4.2s window, and that was a fast daemon).
|
||
# Five seconds is a real answer from a healthy loopback daemon and a rounding error
|
||
# against the finalization grace the clamp reserves. It bounds the READ phase only:
|
||
# `_request` applies the caller bound to every HTTPX phase. Strict review polls add
|
||
# their total wall-clock bound outside this phase-local adapter.
|
||
# Passed per request via ``_request(timeout_sec=...)``; it never changes the default.
|
||
SHORT_POLL_TIMEOUT_SEC = 5.0
|
||
# The httpx failures a READ-ONLY observer may retry as the same unresolved read: the
|
||
# socket delivered no daemon answer, so nothing is known about the run and nothing
|
||
# was claimed. Classified by exception TYPE plus received status, never by prose; a
|
||
# received 4xx/5xx still wins (see ``_request``).
|
||
_OBSERVATION_RETRYABLE_ERRORS = (
|
||
httpx.ReadTimeout, httpx.ConnectError, httpx.ConnectTimeout, httpx.PoolTimeout,
|
||
httpx.ReadError, httpx.WriteError, httpx.RemoteProtocolError,
|
||
)
|
||
# ...and WHICH hole it was, for the observer that has to tell the owner apart a
|
||
# daemon that is merely slow from one that is not there. Our own read bound
|
||
# expiring says nothing about the daemon; a socket that could not be opened or
|
||
# that broke mid-exchange says it did not answer. Classified by exception TYPE,
|
||
# never by prose, and carried beside ``code`` rather than replacing it: the
|
||
# model-control retry loop, the daemon liveness probe and the definite-unrun set
|
||
# all key on ``daemon_unreachable`` for reasons that have nothing to do with an
|
||
# observation.
|
||
_OBSERVATION_READ_TIMEOUT = "observation_read_timeout"
|
||
|
||
|
||
def _observation_reason(exc: BaseException) -> str:
|
||
"""The typed reason for a read-only observation hole."""
|
||
return (_OBSERVATION_READ_TIMEOUT if isinstance(exc, httpx.ReadTimeout)
|
||
else "daemon_unreachable")
|
||
_ATTEMPTS_REL = "attempts"
|
||
_ATTEMPT_RECORD = "attempt.yaml"
|
||
|
||
|
||
class ClaudexorUnavailable(RuntimeError):
|
||
"""Typed lane refusal: the delegated route cannot run right now.
|
||
|
||
Carries the machine-readable ``code`` so callers classify instead of matching
|
||
prose. Never raised for an ordinary in-run failure — only for "this transport
|
||
is not usable".
|
||
|
||
``required_actions`` retains the daemon's TOP-LEVEL ``ControlProblem.requiredActions``
|
||
string list when the refusal carried one (e.g. the reconcile 409's
|
||
``retry_setup_reconciliation``), bounded to the daemon's own wire limit. It is
|
||
a preserved fact for the typed error seam, not a client action framework.
|
||
"""
|
||
|
||
# What the engine REPORTED about a failed run ("" = nothing reported); set only by
|
||
# ``run_failure_error``. An opaque fact: carried and shown, never branched on.
|
||
reported_cause = ""
|
||
|
||
def __init__(self, code: str, message: str, *, status_code: int = 0,
|
||
required_actions: tuple[str, ...] = (), observation_timeout: bool = False,
|
||
observation_reason: str = "") -> None:
|
||
super().__init__(message)
|
||
self.code = str(code or "claudexor_unavailable")
|
||
self.status_code = int(status_code or 0)
|
||
self.required_actions = tuple(required_actions or ())
|
||
# Read-only observers may retry this exact HTTP read without claiming
|
||
# anything about the worker. A received HTTP refusal still wins.
|
||
self.observation_timeout = bool(observation_timeout)
|
||
# Which hole it was, for the observer only (``_observation_reason``):
|
||
# a read bound that expired against a live daemon is not the same fact
|
||
# as a socket that never carried an answer, and only the second is
|
||
# worth an owner-facing outage line.
|
||
self.observation_reason = str(observation_reason or "")
|
||
|
||
|
||
# Cross-repo contract (B1): the engine's window-exhausted RunFailure codes. A
|
||
# newer engine (rotation PR-A) reports a spent credential POOL under its own
|
||
# code. A spent subscription window always heals on a timer (the engine always
|
||
# dates it); a pool heals on a timer only when the engine DATED it (`resetsAt`),
|
||
# an undated pool (every enabled account refused the model, every account
|
||
# disabled) is structural. `_window_exhausted_refusal` is the one reader; the
|
||
# ORIGINAL code is preserved either way. Any other code stays a generic
|
||
# ClaudexorUnavailable: fail-open, old engines emitting code:null included.
|
||
WINDOW_EXHAUSTED_CODES = ("subscription_window_exhausted", "credential_pool_exhausted")
|
||
|
||
|
||
class ClaudexorSubscriptionWindowExhausted(ClaudexorUnavailable):
|
||
"""The subscription window is spent and heals on a timer, not on payment."""
|
||
|
||
def __init__(self, message: str, *, reset_at: str = "", status_code: int = 0,
|
||
code: str = "subscription_window_exhausted") -> None:
|
||
super().__init__(code, message, status_code=status_code)
|
||
self.reset_at = str(reset_at or "")
|
||
|
||
|
||
def _window_exhausted_refusal(code: str, message: str, resets_at: Any, *,
|
||
status_code: int = 0) -> Optional[ClaudexorSubscriptionWindowExhausted]:
|
||
"""The timer-healing class for a window-exhausted code, or None for the caller's
|
||
plain refusal. Keyed on the structured ``resetsAt`` field only, never on prose: a
|
||
pool maps only when the engine dated it (a non-empty string), so an undated pool
|
||
never masquerades as a quota timer."""
|
||
dated = isinstance(resets_at, str) and bool(resets_at.strip())
|
||
if code not in WINDOW_EXHAUSTED_CODES or (
|
||
code != "subscription_window_exhausted" and not dated):
|
||
return None
|
||
return ClaudexorSubscriptionWindowExhausted(
|
||
message, reset_at=str(resets_at or ""), status_code=status_code, code=code)
|
||
|
||
|
||
REPORTED_CAUSE_CHARS = 512 # strict bound of the record field, omission marker included
|
||
|
||
|
||
def run_failure_cause(failure: Any) -> str:
|
||
"""What the engine REPORTED about a failed run (``failure.safeMessage``), whitespace-
|
||
collapsed, secret-redacted and strictly bounded; "" when it reported nothing. An OPAQUE
|
||
fact: stored and displayed, never parsed or branched on (BIBLE P5) — presence is the
|
||
only test a caller may make."""
|
||
from ouroboros.utils import sanitize_tool_result_for_log, truncate_within_limit
|
||
|
||
words = (failure if isinstance(failure, dict) else {}).get("safeMessage")
|
||
return truncate_within_limit(
|
||
sanitize_tool_result_for_log(" ".join(str(words or "").split())), REPORTED_CAUSE_CHARS)
|
||
|
||
|
||
def run_failure_error(run_id: str, run_state: str, failure: Any) -> ClaudexorUnavailable:
|
||
"""The typed refusal for a delegated review run that did not succeed. A null engine code
|
||
keeps the host-derived ``run_<state>`` code: a STATE fact, never a cause — the cause the
|
||
engine reported rides beside it as ``reported_cause``."""
|
||
failure = failure if isinstance(failure, dict) else {}
|
||
message = (f"delegated review session {run_id} ended {run_state or 'unknown'}"
|
||
+ (f": {json.dumps(failure, ensure_ascii=False)}" if failure else ""))
|
||
code = str(failure.get("code") or "")
|
||
exc = (_window_exhausted_refusal(code, message, failure.get("resetsAt"))
|
||
or ClaudexorUnavailable(code or f"run_{run_state or 'unknown'}", message))
|
||
exc.reported_cause = run_failure_cause(failure)
|
||
return exc
|
||
|
||
|
||
@dataclass(frozen=True)
|
||
class DaemonEndpoint:
|
||
"""Loopback control-plane address. ``token`` stays host-side (see module doc)."""
|
||
|
||
host: str
|
||
port: int
|
||
token: str
|
||
|
||
|
||
def engine_at_least(version: str, minimum: str) -> bool:
|
||
"""Is the reported engine at or past ``minimum``? THE version-floor predicate.
|
||
|
||
One reader, so the handshake's transport floor and a lane's own feature floor
|
||
cannot disagree about what "old enough to refuse" means. An unparsable or absent
|
||
version compares as ``(0,)`` — below every floor, so it fails CLOSED.
|
||
"""
|
||
pair: List[tuple] = []
|
||
for value in (version, minimum):
|
||
parts: List[int] = []
|
||
for chunk in str(value or "").split("."):
|
||
digits = "".join(c for c in chunk if c.isdigit())
|
||
parts.append(int(digits) if digits else 0)
|
||
pair.append(tuple(parts or [0]))
|
||
return pair[0] >= pair[1]
|
||
|
||
|
||
def operator_home() -> pathlib.Path:
|
||
"""The home ``discover_daemon`` reads the control token from.
|
||
|
||
Named once because it is the exact directory a delegated harness must never
|
||
inherit: it holds ``~/.claudexor/v3/daemon/token``, which grants the whole ``/v2``
|
||
surface. Both the discovery below and the applied-isolation check read it here.
|
||
"""
|
||
return pathlib.Path(os.path.expanduser("~"))
|
||
|
||
|
||
def discover_daemon(home: Optional[pathlib.Path] = None) -> DaemonEndpoint:
|
||
"""Read the daemon descriptor plus the referenced token.
|
||
|
||
With an explicit ``home`` this reads that home's ``~/.claudexor/v3`` layout
|
||
verbatim. With none, the OWNED daemon is preferred WHEN PROVISIONED (D30):
|
||
once Ouroboros has spawned its own daemon under the data-plane config dir,
|
||
every default discovery — delegated subagents, review sessions, the account
|
||
surfaces — talks to that one, and the operator's personal daemon is left
|
||
alone. An unprovisioned owned home falls through to the operator layout,
|
||
which is the entire pre-D30 behavior; the cutover is the owner's own
|
||
provisioning action, never a silent boot-time switch.
|
||
"""
|
||
if home is None:
|
||
from ouroboros.claudexor_daemon import owned_daemon_provisioned, owned_descriptor_path
|
||
|
||
if owned_daemon_provisioned():
|
||
return _endpoint_from_descriptor(owned_descriptor_path())
|
||
root = pathlib.Path(home) if home is not None else operator_home()
|
||
return _endpoint_from_descriptor(root / CONTROL_API_REL)
|
||
|
||
|
||
def discover_daemon_at(config_dir: pathlib.Path) -> DaemonEndpoint:
|
||
"""Discovery for an explicit ``CLAUDEXOR_CONFIG_DIR`` root.
|
||
|
||
Under an override the override IS the complete relocatable root, so the
|
||
descriptor lives at ``<config_dir>/daemon/control-api.json`` — a different
|
||
shape from the default ``~/.claudexor/v3`` layout ``discover_daemon`` reads.
|
||
"""
|
||
return _endpoint_from_descriptor(pathlib.Path(config_dir) / "daemon" / "control-api.json")
|
||
|
||
|
||
def _endpoint_from_descriptor(control_path: pathlib.Path) -> DaemonEndpoint:
|
||
try:
|
||
raw = json.loads(control_path.read_text(encoding="utf-8"))
|
||
except FileNotFoundError as exc:
|
||
raise ClaudexorUnavailable(
|
||
"daemon_not_discovered",
|
||
f"Claudexor control-api descriptor not found at {control_path}",
|
||
) from exc
|
||
except (OSError, ValueError) as exc:
|
||
raise ClaudexorUnavailable(
|
||
"daemon_descriptor_unreadable",
|
||
f"Claudexor control-api descriptor unreadable: {type(exc).__name__}: {exc}",
|
||
) from exc
|
||
if not isinstance(raw, dict):
|
||
raise ClaudexorUnavailable("daemon_descriptor_unreadable", "control-api.json is not an object")
|
||
host = str(raw.get("host") or "").strip()
|
||
token_path = str(raw.get("tokenPath") or "").strip()
|
||
try:
|
||
port = int(raw.get("port") or 0)
|
||
except (TypeError, ValueError):
|
||
port = 0
|
||
if not host or port <= 0 or not token_path:
|
||
raise ClaudexorUnavailable(
|
||
"daemon_descriptor_incomplete",
|
||
"control-api.json is missing host/port/tokenPath",
|
||
)
|
||
try:
|
||
# The SAME set as the descriptor read four lines up, and for the same reason: a
|
||
# path out of a JSON descriptor can carry an embedded null and a token file can
|
||
# hold bytes that are not UTF-8, both of which `read_text` raises as `ValueError`
|
||
# (`UnicodeDecodeError` is one), and either escaping here is a traceback where a
|
||
# typed refusal belongs. `RuntimeError` is deliberately NOT in this set, and the
|
||
# v6.87.44 comment claiming it covered a symlink loop was wrong: `read_text` on a
|
||
# loop raises `OSError` (ELOOP) — `resolve()` is what raises `RuntimeError`, which
|
||
# is why `_resolved` in delegate.py catches it and this does not. It was also a
|
||
# hazard, since `ClaudexorUnavailable` IS a `RuntimeError`: the moment this block
|
||
# grew a call that refuses typed, the catch would re-wrap it under this code.
|
||
token = pathlib.Path(token_path).read_text(encoding="utf-8").strip()
|
||
except (OSError, ValueError) as exc:
|
||
raise ClaudexorUnavailable(
|
||
"daemon_token_unreadable",
|
||
f"Claudexor daemon token unreadable: {type(exc).__name__}",
|
||
) from exc
|
||
if not token:
|
||
raise ClaudexorUnavailable("daemon_token_unreadable", "Claudexor daemon token file is empty")
|
||
if not _is_loopback(host):
|
||
# The token this descriptor points at grants the ENTIRE /v2 control API — start
|
||
# runs, read every artifact, cancel anything. The loopback boundary was
|
||
# documented and never enforced, so anything able to write one file under
|
||
# ~/.claudexor could redirect the bearer to a host it controls: token
|
||
# exfiltration plus authenticated SSRF, from a file write. Refused BEFORE the
|
||
# client exists, so no request can be built against a non-loopback endpoint.
|
||
raise ClaudexorUnavailable(
|
||
"daemon_endpoint_not_loopback",
|
||
f"Claudexor control-api descriptor names the non-loopback host {host!r}. The "
|
||
f"daemon control token grants the whole /v2 surface and is only ever sent to "
|
||
f"the local daemon; refusing rather than shipping it off-host.",
|
||
)
|
||
return DaemonEndpoint(host=host, port=port, token=token)
|
||
|
||
|
||
def _is_loopback(host: str) -> bool:
|
||
"""True only for a literal loopback ADDRESS, or the exact name ``localhost``.
|
||
|
||
An IP literal is decided by the stdlib (``127.0.0.0/8``, ``::1``, and their
|
||
zone/bracket spellings), never by string prefixes: ``127.0.0.1.evil.com`` is a
|
||
NAME, and ``0x7f.1`` is not one this reader will guess at. Resolution is
|
||
deliberately NOT attempted — a name that resolves to loopback today can resolve
|
||
elsewhere on the next lookup, so only ``localhost`` (which every platform pins to
|
||
loopback) is accepted by name. Everything else is refused.
|
||
"""
|
||
import ipaddress
|
||
|
||
candidate = str(host or "").strip().strip("[]")
|
||
if not candidate:
|
||
return False
|
||
if candidate.lower() == "localhost":
|
||
return True
|
||
candidate = candidate.split("%", 1)[0] # drop an IPv6 zone id
|
||
try:
|
||
return ipaddress.ip_address(candidate).is_loopback
|
||
except ValueError:
|
||
return False
|
||
|
||
|
||
def operation_query_supported(operations: list[dict], *, method: str, path: str,
|
||
name: str, value: str) -> bool:
|
||
"""Negotiate an exact query value from the serving operation's descriptor."""
|
||
return any(
|
||
operation.get("method") == method and operation.get("path") == path
|
||
and any(parameter.get("name") == name and parameter.get("location") == "query"
|
||
and isinstance(parameter.get("enum"), list) and value in parameter["enum"]
|
||
for parameter in (operation.get("parameters") or []) if isinstance(parameter, dict))
|
||
for operation in operations if isinstance(operation, dict)
|
||
)
|
||
|
||
|
||
def account_catalog_supported(operations: list[dict], path: str) -> bool:
|
||
"""Opt in only when this exact operation declares the accounts query view."""
|
||
return operation_query_supported(operations, method="GET", path=path, name="view", value="accounts")
|
||
|
||
|
||
def model_failure_evidence_supported(operations: list[dict]) -> bool:
|
||
return operation_query_supported(operations, method="POST", path="/v2/model-operations",
|
||
name="captureFailureEvidence", value="true")
|
||
|
||
|
||
# The engine's typed live-message outcomes (``LiveMessageOutcome``): the two
|
||
# positive boundaries plus the four typed non-deliveries. Mirrored 1:1 by
|
||
# ``delegate_message``; the host adds only its own ``not_found`` vocabulary.
|
||
LIVE_MESSAGE_OUTCOMES = frozenset({
|
||
"delivered", "accepted", "rejected", "not_active", "unsupported", "delivery_unknown",
|
||
})
|
||
# The catalog spelling of the live-message route (Express-style template, the
|
||
# same shape as ``/v2/runs/:id/control`` and the interaction-answer row).
|
||
RUN_MESSAGE_OPERATION = ("POST", "/v2/runs/:id/messages")
|
||
|
||
|
||
def run_message_supported(operations: list[dict]) -> bool:
|
||
"""Does this engine's own route catalog list ``POST /v2/runs/:id/messages``?
|
||
|
||
Presence is negotiated structurally, like every other route here: an engine
|
||
older than the live-message release answers a route 404, which the verb must
|
||
never reach (a 404 is otherwise the daemon's "no such run").
|
||
"""
|
||
method, path = RUN_MESSAGE_OPERATION
|
||
return any(
|
||
operation.get("method") == method and operation.get("path") == path
|
||
for operation in operations if isinstance(operation, dict)
|
||
)
|
||
|
||
|
||
class ClaudexorGateway:
|
||
"""Thin typed client over the Claudexor ``/v2`` control API."""
|
||
|
||
def __init__(self, endpoint: Optional[DaemonEndpoint] = None, *, home: Optional[pathlib.Path] = None):
|
||
self._endpoint = endpoint if endpoint is not None else discover_daemon(home)
|
||
self._engine_version = ""
|
||
self._engine_build_sha = ""
|
||
# trust_env=False: a shell HTTP(S)_PROXY must never be able to intercept the
|
||
# loopback control plane (the bearer token rides these requests).
|
||
self._client = httpx.Client(
|
||
base_url=f"http://{self._endpoint.host}:{self._endpoint.port}",
|
||
timeout=httpx.Timeout(_READ_TIMEOUT_SEC, connect=_CONNECT_TIMEOUT_SEC),
|
||
trust_env=False,
|
||
headers={
|
||
"Authorization": f"Bearer {self._endpoint.token}",
|
||
PROTOCOL_HEADER: str(CLAUDEXOR_PROTOCOL_MAJOR),
|
||
"Content-Type": "application/json",
|
||
},
|
||
)
|
||
|
||
# -- lifecycle -------------------------------------------------------------
|
||
|
||
def close(self) -> None:
|
||
try:
|
||
self._client.close()
|
||
except Exception:
|
||
log.debug("Claudexor client close failed", exc_info=True)
|
||
|
||
def __enter__(self) -> "ClaudexorGateway":
|
||
return self
|
||
|
||
def __exit__(self, *_exc: Any) -> None:
|
||
self.close()
|
||
|
||
@property
|
||
def engine_version(self) -> str:
|
||
return self._engine_version
|
||
|
||
@property
|
||
def engine_build_sha(self) -> str:
|
||
return self._engine_build_sha
|
||
|
||
# -- transport -------------------------------------------------------------
|
||
|
||
def _request(self, method: str, path: str, *, json_body: Any = None,
|
||
headers: Optional[Dict[str, str]] = None,
|
||
timeout_sec: Optional[float] = None,
|
||
content_body: Optional[bytes] = None, raw_bytes: bool = False) -> Any:
|
||
# ``timeout_sec`` replaces the client's read default for THIS call only, for a
|
||
# caller that is bounding ITSELF — today every `delegate_wait` poll, each asking
|
||
# for what its window has left, floored at ``SHORT_POLL_TIMEOUT_SEC`` and never
|
||
# raised above the default (``delegate_progress.poll_bound``). The caller's
|
||
# bound applies to connect as well as read/write/pool phases. HTTPX treats
|
||
# these as phase-local, so strict total budgeting happens in bounded_poll.
|
||
# Absent, the call is not passed at all rather than passed as None — httpx reads
|
||
# an explicit ``timeout=None`` as "no timeout", which is the opposite of the
|
||
# default it would otherwise inherit.
|
||
if timeout_sec is None:
|
||
bound: Dict[str, Any] = {}
|
||
else:
|
||
bounded = max(0.000001, float(timeout_sec))
|
||
bound = {
|
||
"timeout": httpx.Timeout(
|
||
bounded,
|
||
connect=min(_CONNECT_TIMEOUT_SEC, bounded),
|
||
)
|
||
}
|
||
response = None
|
||
try:
|
||
# Preserve received headers even when decoding or reading the body
|
||
# fails. In particular, a 401/403 must survive a later read timeout.
|
||
payload = {"content": content_body} if content_body is not None else {"json": json_body}
|
||
with self._client.stream(method, path, **payload,
|
||
headers=headers or None, **bound) as response:
|
||
response.read()
|
||
except httpx.HTTPError as exc:
|
||
retryable = (isinstance(exc, _OBSERVATION_RETRYABLE_ERRORS)
|
||
and (response is None or response.status_code < 400))
|
||
raise ClaudexorUnavailable(
|
||
"daemon_unreachable",
|
||
f"Claudexor daemon unreachable: {type(exc).__name__}: {exc}",
|
||
status_code=response.status_code if response is not None else 0,
|
||
observation_timeout=retryable,
|
||
observation_reason=_observation_reason(exc) if retryable else "",
|
||
) from exc
|
||
if response.status_code >= 400:
|
||
raise self._problem(response)
|
||
if raw_bytes:
|
||
return response.content
|
||
if not response.content:
|
||
return None
|
||
try:
|
||
return response.json()
|
||
except ValueError as exc:
|
||
raise ClaudexorUnavailable(
|
||
"malformed_response",
|
||
f"Claudexor returned a non-JSON body for {method} {path}: {exc}",
|
||
) from exc
|
||
|
||
def _problem(self, response: httpx.Response) -> ClaudexorUnavailable:
|
||
"""Translate a ControlProblem body into a typed refusal."""
|
||
code = f"http_{response.status_code}"
|
||
message = response.text[:500]
|
||
context: Dict[str, Any] = {}
|
||
required_actions: tuple[str, ...] = ()
|
||
try:
|
||
body = response.json()
|
||
except ValueError:
|
||
body = None
|
||
if isinstance(body, dict):
|
||
code = str(body.get("code") or code)
|
||
message = str(body.get("message") or message)
|
||
raw_context = body.get("context")
|
||
context = raw_context if isinstance(raw_context, dict) else {}
|
||
# The daemon serializes `requiredActions` at the ControlProblem TOP LEVEL
|
||
# (`daemon-server` projects the field beside code/message; `problem-safety`
|
||
# bounds the wire list to at most 16 redacted strings of at most 512 chars).
|
||
# It is deliberately NOT read from `context`: no producer puts it there, and
|
||
# a context sniff would resurrect the exact had-it-both-ways bug the reset
|
||
# classification below documents. The bound is mirrored so a foreign body
|
||
# cannot balloon the retained tuple.
|
||
raw_actions = body.get("requiredActions")
|
||
if isinstance(raw_actions, list):
|
||
required_actions = tuple(
|
||
str(item)[:512] for item in raw_actions[:16]
|
||
if isinstance(item, str) and item
|
||
)
|
||
# The CODE decides, exactly as every other classification on this seam does.
|
||
# Sniffing `context` for a reset key on ANY code had it both ways: an unrelated
|
||
# refusal that happened to carry one, an `idempotency_conflict` say, was announced
|
||
# as a spent subscription window and retried on a timer. So a reset key is read
|
||
# only inside a window code, and there only the pool's own `resetsAt` says whether
|
||
# it is a timer at all. At engine 3.14.0 the daemon serializes no `resetsAt` into a
|
||
# pool ControlProblem context (the dated producer is the run-detail RunFailure, and
|
||
# `cooldown_until` lives in a quota snapshot), so this seam yields the plain class.
|
||
return (_window_exhausted_refusal(code, message, context.get("resetsAt"),
|
||
status_code=response.status_code)
|
||
or ClaudexorUnavailable(code, message, status_code=response.status_code,
|
||
required_actions=required_actions))
|
||
|
||
# -- operations ------------------------------------------------------------
|
||
|
||
def handshake(self, *, timeout_sec: Optional[float] = None) -> Dict[str, Any]:
|
||
"""Negotiate protocol major and enforce the minimum engine version.
|
||
|
||
``timeout_sec`` bounds this call the way ``get_run`` is bounded: a caller
|
||
holding a clamped window pays for the OPENING round trip out of that window,
|
||
so an unbounded handshake could spend it before the window began.
|
||
"""
|
||
body = self._request(
|
||
"POST", "/v2/handshake", timeout_sec=timeout_sec,
|
||
json_body={"protocolMajor": CLAUDEXOR_PROTOCOL_MAJOR, "client": "ouroboros"},
|
||
)
|
||
if not isinstance(body, dict):
|
||
raise ClaudexorUnavailable("malformed_response", "handshake did not return an object")
|
||
if not body.get("compatible") or int(body.get("protocolMajor") or 0) != CLAUDEXOR_PROTOCOL_MAJOR:
|
||
raise ClaudexorUnavailable(
|
||
"protocol_incompatible",
|
||
f"Claudexor refused protocol major {CLAUDEXOR_PROTOCOL_MAJOR}: {body!r}",
|
||
)
|
||
engine = body.get("engine") if isinstance(body.get("engine"), dict) else {}
|
||
version = str(engine.get("version") or "")
|
||
if not engine_at_least(version, CLAUDEXOR_MIN_VERSION):
|
||
raise ClaudexorUnavailable(
|
||
"engine_too_old",
|
||
f"Claudexor {version or 'unknown'} is older than the required {CLAUDEXOR_MIN_VERSION}",
|
||
)
|
||
self._engine_version = version
|
||
self._engine_build_sha = str(engine.get("sha") or "")
|
||
return body
|
||
|
||
def agent_capabilities(self) -> Dict[str, Any]:
|
||
body = self._request("GET", "/v2/agent-capabilities")
|
||
return body if isinstance(body, dict) else {}
|
||
|
||
# Model operations use the same private control transport, not Agent runs
|
||
# or the redacted/size-capped artifact surface. Callers own all admission,
|
||
# waiting, billing and explicit acknowledgement. The engine owns account
|
||
# choice; this client preserves the caller's Auto/pin request verbatim.
|
||
|
||
def list_model_sources(self, *, view: Optional[str] = None) -> Dict[str, Any]:
|
||
"""Return the engine's opaque source ids and credential-harness bindings."""
|
||
from urllib.parse import urlencode
|
||
|
||
path = "/v2/model-sources"
|
||
if view is not None:
|
||
path += "?" + urlencode({"view": view})
|
||
return _model_object(self._request("GET", path))
|
||
|
||
def list_source_models(self, source: str,
|
||
credential_profile_id: Optional[str] = None, *,
|
||
requested_model: Optional[str] = None,
|
||
timeout_sec: Optional[float] = None,
|
||
view: Optional[str] = None) -> Dict[str, Any]:
|
||
"""Preserve the exact-profile catalog envelope; an omitted pin means engine Auto."""
|
||
from urllib.parse import quote, urlencode
|
||
|
||
path = f"/v2/model-sources/{quote(str(source), safe='')}/models"
|
||
query = {}
|
||
if view is not None:
|
||
query["view"] = view
|
||
if credential_profile_id is not None:
|
||
query["credentialProfileId"] = credential_profile_id
|
||
if requested_model is not None:
|
||
query["requestedModel"] = requested_model
|
||
if query:
|
||
path += "?" + urlencode(query)
|
||
return _model_object(self._request("GET", path, **({"timeout_sec": timeout_sec} if timeout_sec is not None else {})))
|
||
|
||
def upload_model_request(self, request: Dict[str, Any], *,
|
||
idempotency_key: str) -> Dict[str, Any]:
|
||
"""Upload exact request JSON without spending a generation or logging content.
|
||
|
||
Stage keys derive from the caller's invocation identity, never the prompt.
|
||
An explicit retry can finish an uploaded/finalized handle after a lost
|
||
reply. A cancelled or still-writing upload remains a typed refusal; this
|
||
method never creates another upload or inference to hide that outcome.
|
||
"""
|
||
from urllib.parse import quote
|
||
|
||
key = _model_idempotency_key(idempotency_key)
|
||
if not isinstance(request, dict):
|
||
raise ClaudexorUnavailable("invalid_model_payload", "Model request must be a JSON object")
|
||
try:
|
||
data = json.dumps(request, ensure_ascii=False, allow_nan=False,
|
||
sort_keys=True, separators=(",", ":")).encode("utf-8")
|
||
except (TypeError, ValueError, UnicodeError) as exc:
|
||
raise ClaudexorUnavailable("invalid_model_payload", "Model request is not valid JSON") from exc
|
||
digest = "sha256:" + hashlib.sha256(data).hexdigest()
|
||
stage_key = hashlib.sha256(key.encode("utf-8")).hexdigest()
|
||
created = _model_object(self._request(
|
||
"POST", "/v2/uploads",
|
||
json_body={"purpose": "model", "kind": "file", "mime": "application/json",
|
||
# Bind create's existing metadata identity to the bytes too:
|
||
# equal length alone would let an open upload change payload.
|
||
"name": f"model-request-{digest.removeprefix('sha256:')}.json",
|
||
"sizeBytes": len(data)},
|
||
headers={"Idempotency-Key": f"model-upload-{stage_key}"},
|
||
))
|
||
upload_id = created.get("uploadId")
|
||
if not isinstance(upload_id, str) or not upload_id:
|
||
raise ClaudexorUnavailable("malformed_response", "Model upload returned no upload id")
|
||
path = f"/v2/uploads/{quote(upload_id, safe='')}"
|
||
try:
|
||
status = _model_object(self._request("GET", path))
|
||
except ClaudexorUnavailable as exc:
|
||
if exc.status_code != 404 or exc.code != "upload_not_found":
|
||
raise
|
||
# Finalize removes the upload handle but retains its idempotent receipt.
|
||
else:
|
||
if status.get("uploadId") != upload_id:
|
||
raise ClaudexorUnavailable("malformed_response", "Model upload status identity mismatch")
|
||
if status.get("state") == "open":
|
||
self._request("PUT", f"{path}/bytes", content_body=data,
|
||
headers={"Content-Type": "application/octet-stream"})
|
||
elif status.get("state") != "uploaded":
|
||
raise ClaudexorUnavailable("model_upload_unavailable", "Model upload is not ready for finalization")
|
||
resource = _model_object(self._request(
|
||
"POST", f"{path}/finalize", json_body={"expectedSha256": digest},
|
||
headers={"Idempotency-Key": f"model-finalize-{stage_key}"},
|
||
))
|
||
if resource.get("purpose") != "model":
|
||
raise ClaudexorUnavailable("resource_purpose_mismatch", "Model upload returned an Agent attachment")
|
||
ref = _model_payload_ref({name: resource.get(name) for name in ("resourceId", "sha256", "sizeBytes")})
|
||
if ref["sha256"] != digest or ref["sizeBytes"] != len(data):
|
||
raise ClaudexorUnavailable("model_payload_integrity_error", "Finalized model request differs from uploaded bytes")
|
||
return ref
|
||
|
||
def create_model_operation(self, request_ref: Dict[str, Any], *,
|
||
idempotency_key: str, capture_failure_evidence: bool = False) -> Dict[str, Any]:
|
||
"""Create or rejoin exactly one caller-identified generation; never mint a retry key."""
|
||
key = _model_idempotency_key(idempotency_key)
|
||
path = "/v2/model-operations" + ("?captureFailureEvidence=true" if capture_failure_evidence else "")
|
||
return _model_operation(self._request(
|
||
"POST", path, json_body={"request": _model_payload_ref(request_ref)},
|
||
headers={"Idempotency-Key": key},
|
||
))
|
||
|
||
def get_model_operation(self, operation_id: str, *,
|
||
timeout_sec: Optional[float] = None) -> Dict[str, Any]:
|
||
from urllib.parse import quote
|
||
|
||
return _model_operation(self._request(
|
||
"GET", f"/v2/model-operations/{quote(str(operation_id), safe='')}",
|
||
timeout_sec=timeout_sec,
|
||
), operation_id)
|
||
|
||
def get_model_result(self, operation_id: str, *, expected_ref: Dict[str, Any],
|
||
timeout_sec: Optional[float] = None,
|
||
raw_bytes: bool = False) -> Dict[str, Any] | bytes:
|
||
"""Read and verify the complete result, without ACK, redaction or artifact caps.
|
||
|
||
The expected reference comes from this operation's ready custody record.
|
||
Failure preserves that handle: only another read of the same operation is
|
||
appropriate here, never a new generation. The caller acknowledges after
|
||
it has retained the returned result under its own custody contract.
|
||
``raw_bytes`` retains the exact verified JSON encoding for that custody;
|
||
it never skips the size, digest, UTF-8 or object validation below.
|
||
"""
|
||
from urllib.parse import quote
|
||
|
||
ref = _model_payload_ref(expected_ref)
|
||
data = self._request(
|
||
"GET", f"/v2/model-operations/{quote(str(operation_id), safe='')}/result",
|
||
timeout_sec=timeout_sec, raw_bytes=True,
|
||
)
|
||
if len(data) != ref["sizeBytes"] or "sha256:" + hashlib.sha256(data).hexdigest() != ref["sha256"]:
|
||
raise ClaudexorUnavailable("model_payload_integrity_error", "Model result does not match its size and SHA-256")
|
||
try:
|
||
body = json.loads(data.decode("utf-8", errors="strict"), parse_constant=_reject_json_constant)
|
||
except (ValueError, UnicodeError) as exc:
|
||
raise ClaudexorUnavailable("malformed_response", "Model result is not valid UTF-8 JSON") from exc
|
||
body = _model_object(body)
|
||
return data if raw_bytes else body
|
||
|
||
def acknowledge_model_result(self, operation_id: str, sha256: str) -> Dict[str, Any]:
|
||
"""Acknowledge only the exact result the caller has retained; no implicit ACK."""
|
||
from urllib.parse import quote
|
||
|
||
return _model_operation(self._request(
|
||
"POST", f"/v2/model-operations/{quote(str(operation_id), safe='')}/ack",
|
||
json_body={"sha256": sha256},
|
||
), operation_id)
|
||
|
||
def cancel_model_operation(self, operation_id: str, *, reason_code: str = "") -> Dict[str, Any]:
|
||
"""Request cancellation; the returned engine lifecycle, not this POST, proves settlement."""
|
||
from urllib.parse import quote
|
||
|
||
control = {"action": "cancel"}
|
||
if reason_code:
|
||
control["reasonCode"] = reason_code
|
||
return _model_operation(self._request(
|
||
"POST", f"/v2/model-operations/{quote(str(operation_id), safe='')}/control",
|
||
json_body=control,
|
||
), operation_id)
|
||
|
||
def harnesses(self) -> List[Dict[str, Any]]:
|
||
"""GET /v2/harnesses — per-harness status rows WITH the full manifest.
|
||
|
||
The agent-capability catalog is a derived projection that deliberately
|
||
drops the manifest's transport flags (``json_schema_output``,
|
||
``interactive``); this is the surface that still carries them, so
|
||
transport-capability questions are asked here, not of the catalog.
|
||
"""
|
||
body = self._request("GET", "/v2/harnesses")
|
||
rows = body.get("harnesses") if isinstance(body, dict) else None
|
||
return [row for row in (rows or []) if isinstance(row, dict)]
|
||
|
||
def quota_state(self) -> Dict[str, Any]:
|
||
"""GET /v2/quota once, retaining its one-epoch evidence envelope."""
|
||
body = self._request("GET", "/v2/quota")
|
||
return body if isinstance(body, dict) else {}
|
||
|
||
def refresh_quota(self) -> Dict[str, Any]:
|
||
"""POST /v2/quota once, returning the foreground evidence envelope."""
|
||
return self._request(
|
||
"POST",
|
||
"/v2/quota",
|
||
json_body={},
|
||
timeout_sec=get_claudexor_quota_refresh_timeout_sec(),
|
||
)
|
||
|
||
def quota_snapshots(self) -> List[Dict[str, Any]]:
|
||
body = self.quota_state()
|
||
snapshots = body.get("snapshots") if isinstance(body, dict) else None
|
||
return [row for row in (snapshots or []) if isinstance(row, dict)]
|
||
|
||
def quota_absences(self) -> List[Dict[str, Any]]:
|
||
"""Profiles whose quota could NOT be read (a 429/failed refresh, no login).
|
||
|
||
Legacy compatibility projection of the same one-epoch quota envelope.
|
||
An absence is typed evidence that quota could not be read for a profile;
|
||
route health treats that state as unknown and therefore fail-open.
|
||
"""
|
||
body = self.quota_state()
|
||
absences = body.get("absences") if isinstance(body, dict) else None
|
||
return [row for row in (absences or []) if isinstance(row, dict)]
|
||
|
||
def register_project(self, root: str) -> str:
|
||
"""Register a run root and return its project id (idempotent per root).
|
||
|
||
Claudexor answers the FIRST run against an unregistered root with
|
||
404 ``project_not_registered``, so registration is a required step, not an
|
||
optimization. Re-registering an existing root returns the existing id.
|
||
"""
|
||
body = self._request(
|
||
"POST", "/v2/projects",
|
||
json_body={"root": str(root)},
|
||
headers={"Idempotency-Key": uuid.uuid4().hex},
|
||
)
|
||
project_id = str((body or {}).get("id") or "") if isinstance(body, dict) else ""
|
||
if not project_id:
|
||
raise ClaudexorUnavailable("malformed_response", "project registration returned no id")
|
||
return project_id
|
||
|
||
def remove_project(self, project_id: str) -> Dict[str, Any]:
|
||
"""Retire a project registration. Non-destructive: artifacts are retained."""
|
||
body = self._request("DELETE", f"/v2/projects/{project_id}")
|
||
return body if isinstance(body, dict) else {}
|
||
|
||
def find_project_id(self, root: str) -> str:
|
||
body = self._request("GET", "/v2/projects")
|
||
target = str(root)
|
||
for row in (body or {}).get("projects") or [] if isinstance(body, dict) else []:
|
||
if isinstance(row, dict) and str(row.get("root") or "") == target:
|
||
return str(row.get("id") or "")
|
||
return ""
|
||
|
||
def ensure_full_access(self, root: str) -> Dict[str, Any]:
|
||
"""Create a missing scoped full grant; the caller owns when to request it.
|
||
|
||
Scoped GET also returns default false for an absent trust file. The
|
||
actual-file list distinguishes absence from an explicit refusal; its
|
||
path remains authoritative when a legacy record has no repoRoot.
|
||
"""
|
||
from urllib.parse import quote
|
||
|
||
target = str(root)
|
||
body = self._request("GET", f"/v2/trust?repoRoot={quote(target, safe='')}")
|
||
entries = body.get("entries") if isinstance(body, dict) else None
|
||
if not isinstance(entries, list) or len(entries) != 1:
|
||
raise ClaudexorUnavailable("malformed_response", "Scoped trust lookup returned no unique state")
|
||
state = _trust_state(entries[0])
|
||
if state["repoRoot"] != target:
|
||
raise ClaudexorUnavailable("malformed_response", "Scoped trust lookup returned a different root")
|
||
if state["allowFullAccess"]:
|
||
return state
|
||
body = self._request("GET", "/v2/trust")
|
||
entries = body.get("entries") if isinstance(body, dict) else None
|
||
if not isinstance(entries, list):
|
||
raise ClaudexorUnavailable("malformed_response", "Trust listing returned no entries")
|
||
for entry in entries:
|
||
entry = _trust_state(entry)
|
||
if entry["path"] == state["path"]:
|
||
if entry["allowFullAccess"]:
|
||
return entry
|
||
raise ClaudexorUnavailable(
|
||
"trust_full_access_required", "Full access is disabled for this scope; its existing choice was preserved. "
|
||
"Lower this call with access=workspace_write or select Working files on the actor row.",
|
||
status_code=403,
|
||
)
|
||
granted = _trust_state(self._request(
|
||
"POST", "/v2/trust", json_body={"repoRoot": target, "allowFullAccess": True},
|
||
))
|
||
if (granted["repoRoot"] != target or granted["path"] != state["path"]
|
||
or not granted["allowFullAccess"]):
|
||
raise ClaudexorUnavailable("malformed_response", "Trust update did not confirm full access for this scope")
|
||
return granted
|
||
|
||
def start_run(self, request: Dict[str, Any], *, idempotency_key: str = "") -> Dict[str, Any]:
|
||
"""POST /v2/runs with a caller-built, schema-valid request body.
|
||
|
||
``idempotency_key`` is the caller's LOGICAL INVOCATION ID — minted once per
|
||
intended invocation, reused verbatim on a transport retry. The engine's replay
|
||
check (control-api ``handleRunCreate`` → daemon command store) runs BEFORE any
|
||
preflight: a replayed key with the byte-identical request returns the ORIGINAL
|
||
accepted job's handle, and a replayed key with a different request digest is a
|
||
409 ``idempotency_conflict`` — so a retry must reuse the id AND reproduce the
|
||
body exactly. A fresh random key per POST makes an accepted start whose
|
||
response was lost come back as a SECOND live run that nothing knows about; a
|
||
stale content-stable key makes a deliberate re-run come back as the finished
|
||
OLD run. Callers with no invocation identity keep the random default.
|
||
"""
|
||
body = self._request(
|
||
"POST", "/v2/runs",
|
||
json_body=dict(request),
|
||
headers={"Idempotency-Key": str(idempotency_key or "") or uuid.uuid4().hex},
|
||
)
|
||
if not isinstance(body, dict):
|
||
raise ClaudexorUnavailable("malformed_response", "run start returned no handle")
|
||
return body
|
||
|
||
def create_thread(self, request: Dict[str, Any], *, idempotency_key: str) -> Dict[str, Any]:
|
||
"""Create one durable v3 conversation thread.
|
||
|
||
Claudexor owns continuity and profile routing. This client only carries
|
||
the strict request and the caller's stable idempotency identity.
|
||
"""
|
||
body = self._request(
|
||
"POST", "/v2/threads", json_body=dict(request),
|
||
headers={"Idempotency-Key": str(idempotency_key)},
|
||
)
|
||
if not isinstance(body, dict) or not str(body.get("id") or ""):
|
||
raise ClaudexorUnavailable("malformed_response", "thread create returned no id")
|
||
return body
|
||
|
||
def start_thread_turn(
|
||
self, thread_id: str, request: Dict[str, Any], *, idempotency_key: str,
|
||
) -> Dict[str, Any]:
|
||
"""Append one turn through the public v3 thread pipeline."""
|
||
from urllib.parse import quote
|
||
|
||
body = self._request(
|
||
"POST", f"/v2/threads/{quote(str(thread_id), safe='')}/turns",
|
||
json_body=dict(request), headers={"Idempotency-Key": str(idempotency_key)},
|
||
)
|
||
if not isinstance(body, dict):
|
||
raise ClaudexorUnavailable("malformed_response", "thread turn returned no handle")
|
||
return body
|
||
|
||
def get_thread(self, thread_id: str) -> Dict[str, Any]:
|
||
"""Read turns, native-session bindings, and continuity receipts."""
|
||
from urllib.parse import quote
|
||
|
||
body = self._request("GET", f"/v2/threads/{quote(str(thread_id), safe='')}")
|
||
return body if isinstance(body, dict) else {}
|
||
|
||
def get_run(self, run_id: str, *, timeout_sec: Optional[float] = None) -> Dict[str, Any]:
|
||
body = self._request("GET", f"/v2/runs/{run_id}", timeout_sec=timeout_sec)
|
||
return body if isinstance(body, dict) else {}
|
||
|
||
def get_run_artifact(self, run_id: str, path: str) -> bytes:
|
||
"""GET /v2/runs/:id/artifacts/<path> — the FULL artifact body, raw bytes.
|
||
|
||
The run detail's ``primaryOutput.text`` is a bounded 256 KiB PREVIEW
|
||
(control-api ``PRIMARY_OUTPUT_PREVIEW_BYTES``) with ``bytes``/``truncated``
|
||
beside it; this route is where the full file actually lives. The engine serves
|
||
text artifacts through ``redactSecrets`` (so the served length may differ from
|
||
the on-disk ``bytes`` when a secret was rewritten), refuses credential-shaped
|
||
files with a typed 409, caps text at 4 MiB with a 413, and answers a
|
||
retention-reclaimed run with a 410 tombstone — all of which surface here as
|
||
typed ``ClaudexorUnavailable`` refusals, never as a silent empty body.
|
||
"""
|
||
from urllib.parse import quote
|
||
|
||
try:
|
||
response = self._client.request(
|
||
"GET", f"/v2/runs/{quote(str(run_id), safe='')}/artifacts/{quote(str(path), safe='/')}")
|
||
except httpx.HTTPError as exc:
|
||
raise ClaudexorUnavailable(
|
||
"daemon_unreachable",
|
||
f"Claudexor daemon unreachable: {type(exc).__name__}: {exc}",
|
||
) from exc
|
||
if response.status_code >= 400:
|
||
raise self._problem(response)
|
||
return response.content
|
||
|
||
def stream_run_artifact(self, run_id: str, path: str, sink: Any,
|
||
*, expected: Optional[Dict[str, Any]] = None) -> Dict[str, Any]:
|
||
"""Stream exact run bytes into a caller-owned temporary file.
|
||
|
||
The caller publishes only after this read verifies the manifest identity.
|
||
HTTP refusals and partial streams never become an empty successful file;
|
||
the existing small diagnostic-artifact reader keeps its bytes contract.
|
||
"""
|
||
from urllib.parse import quote
|
||
|
||
digest, size = hashlib.sha256(), 0
|
||
response = None
|
||
try:
|
||
with self._client.stream(
|
||
"GET", f"/v2/runs/{quote(str(run_id), safe='')}/artifacts/{quote(str(path), safe='/')}",
|
||
) as response:
|
||
if response.status_code >= 400:
|
||
response.read()
|
||
raise self._problem(response)
|
||
for chunk in response.iter_bytes(chunk_size=1024 * 1024):
|
||
sink.write(chunk)
|
||
digest.update(chunk)
|
||
size += len(chunk)
|
||
except httpx.HTTPError as exc:
|
||
retryable = (isinstance(exc, _OBSERVATION_RETRYABLE_ERRORS)
|
||
and (response is None or response.status_code < 400))
|
||
raise ClaudexorUnavailable(
|
||
"daemon_unreachable", f"Claudexor artifact read failed: {type(exc).__name__}",
|
||
status_code=response.status_code if response is not None else 0,
|
||
observation_timeout=retryable,
|
||
observation_reason=_observation_reason(exc) if retryable else "",
|
||
) from exc
|
||
measured = {"size": size, "sha256": digest.hexdigest()}
|
||
if expected is not None:
|
||
want_size = expected.get("sizeBytes", expected.get("size"))
|
||
want_hash = str(expected.get("sha256") or "").removeprefix("sha256:")
|
||
if (want_size is not None and want_size != size) or want_hash != measured["sha256"]:
|
||
raise ClaudexorUnavailable("artifact_integrity_error", "Run artifact differs from its captured manifest")
|
||
return measured
|
||
|
||
def apply_run(self, run_id: str, request: Dict[str, Any], *, idempotency_key: str) -> Dict[str, Any]:
|
||
"""Apply the existing run product; the caller retains intent and custody."""
|
||
from urllib.parse import quote
|
||
|
||
return _model_object(self._request(
|
||
"POST", f"/v2/runs/{quote(str(run_id), safe='')}/apply", json_body=request,
|
||
headers={"Idempotency-Key": idempotency_key},
|
||
))
|
||
|
||
def decide_run(self, run_id: str, request: Dict[str, Any], *, idempotency_key: str) -> Dict[str, Any]:
|
||
"""Submit an explicit disposition through the existing engine decision route."""
|
||
from urllib.parse import quote
|
||
|
||
body = self._request("POST", f"/v2/runs/{quote(str(run_id), safe='')}/decision", json_body=request,
|
||
headers={"Idempotency-Key": idempotency_key})
|
||
if not isinstance(body, dict):
|
||
raise ClaudexorUnavailable("malformed_response", "Run decision returned a non-object response")
|
||
return body
|
||
|
||
def answer_interaction(self, run_id: str, interaction_id: str,
|
||
answers: List[Dict[str, Any]]) -> Dict[str, Any]:
|
||
"""POST /v2/runs/:id/interactions/:iid/answer — deliver one answer set.
|
||
|
||
``answers`` rows are already in the wire shape (``questionId`` /
|
||
``selectedLabels`` / ``freeText`` — the strict ``ControlInteractionAnswerRequest``);
|
||
this method is transport, not translation.
|
||
|
||
The engine's reply is TYPED at every HTTP status it owns: 200 carries
|
||
``{accepted, status: "delivered"}``, and a 404/409 refusal carries the SAME
|
||
``ControlInteractionAnswerResponse`` shape with ``status`` ``not_found`` /
|
||
``already_resolved`` / ``rejected`` (daemon-server answers the route with the
|
||
parsed response at 200/404/409). Any body carrying one of those statuses is
|
||
returned as the ANSWER it is — an engine's ``already_resolved`` is a fact,
|
||
not an outage. What still raises ``ClaudexorUnavailable``: transport
|
||
failures, a bodyless 404 (``no such run``), the 501 of an engine build with
|
||
no answer service, and any other refusal without a typed status.
|
||
"""
|
||
from urllib.parse import quote
|
||
|
||
path = (f"/v2/runs/{quote(str(run_id), safe='')}"
|
||
f"/interactions/{quote(str(interaction_id), safe='')}/answer")
|
||
try:
|
||
response = self._client.request("POST", path,
|
||
json={"answers": list(answers or [])})
|
||
except httpx.HTTPError as exc:
|
||
raise ClaudexorUnavailable(
|
||
"daemon_unreachable",
|
||
f"Claudexor daemon unreachable: {type(exc).__name__}: {exc}",
|
||
) from exc
|
||
body: Any = None
|
||
if response.content:
|
||
try:
|
||
body = response.json()
|
||
except ValueError:
|
||
body = None
|
||
if isinstance(body, dict) and str(body.get("status") or "") in (
|
||
"delivered", "not_found", "already_resolved", "rejected"):
|
||
return body
|
||
if response.status_code >= 400:
|
||
raise self._problem(response)
|
||
raise ClaudexorUnavailable(
|
||
"malformed_response",
|
||
f"interaction answer returned no typed status (HTTP {response.status_code})",
|
||
)
|
||
|
||
def send_run_message(self, run_id: str, text: str, *, idempotency_key: str,
|
||
expected_attempt_id: str = "",
|
||
timeout_sec: Optional[float] = None) -> Dict[str, Any]:
|
||
"""POST /v2/runs/:id/messages — one live message into a running run.
|
||
|
||
``idempotency_key`` is REQUIRED and is the caller's message identity: the
|
||
engine serves the route through its idempotent-delivery store, so a replay
|
||
under the same key returns the stored receipt instead of delivering twice
|
||
(``delegate_message`` mints it as ``message_id`` and hands it back).
|
||
``expected_attempt_id`` pins the live attempt when a caller holds one.
|
||
|
||
Like ``answer_interaction`` this is transport, not translation, and it never
|
||
goes through ``_request`` (which raises on every status >= 400): any body
|
||
carrying a typed ``outcome`` (``LIVE_MESSAGE_OUTCOMES``) is returned as the
|
||
ANSWER it is, whatever the HTTP status. What raises ``ClaudexorUnavailable``:
|
||
transport failures (``daemon_unreachable``), and every refusal without a
|
||
typed outcome — the daemon's 404 ``no such run``, the 409 idempotency
|
||
problems (``idempotency_conflict`` / ``delivery_in_progress`` /
|
||
``delivery_interrupted``), 400 (malformed, secret, too long), 501 (no
|
||
service), 5xx — each typed through ``_problem`` with its status code, so the
|
||
verb classifies by code AND status. A 2xx without a typed outcome is
|
||
``malformed_response``.
|
||
"""
|
||
from urllib.parse import quote
|
||
|
||
path = f"/v2/runs/{quote(str(run_id), safe='')}/messages"
|
||
payload: Dict[str, Any] = {"text": str(text)}
|
||
if expected_attempt_id:
|
||
payload["expectedAttemptId"] = str(expected_attempt_id)
|
||
bound: Dict[str, Any] = {}
|
||
if timeout_sec is not None:
|
||
bounded = max(0.000001, float(timeout_sec))
|
||
bound = {"timeout": httpx.Timeout(bounded, connect=min(_CONNECT_TIMEOUT_SEC, bounded))}
|
||
try:
|
||
response = self._client.request(
|
||
"POST", path, json=payload,
|
||
headers={"Idempotency-Key": str(idempotency_key)}, **bound)
|
||
except httpx.HTTPError as exc:
|
||
raise ClaudexorUnavailable(
|
||
"daemon_unreachable",
|
||
f"Claudexor daemon unreachable: {type(exc).__name__}: {exc}",
|
||
) from exc
|
||
body: Any = None
|
||
if response.content:
|
||
try:
|
||
body = response.json()
|
||
except ValueError:
|
||
body = None
|
||
if isinstance(body, dict) and str(body.get("outcome") or "") in LIVE_MESSAGE_OUTCOMES:
|
||
return body
|
||
if response.status_code >= 400:
|
||
raise self._problem(response)
|
||
raise ClaudexorUnavailable(
|
||
"malformed_response",
|
||
f"run message returned no typed outcome (HTTP {response.status_code})",
|
||
)
|
||
|
||
def cancel_run(self, run_id: str, *, reason: str = "") -> Dict[str, Any]:
|
||
control: Dict[str, Any] = {"kind": "cancel"}
|
||
if reason:
|
||
control["reason"] = str(reason)
|
||
body = self._request("POST", f"/v2/runs/{run_id}/control", json_body={"control": control})
|
||
return body if isinstance(body, dict) else {}
|
||
|
||
# -- account surfaces (D30: read/translate only; the daemon owns ALL auth
|
||
# logic — profiles, login jobs, device-code custody, verification) ---------
|
||
|
||
def credential_profiles(self) -> Dict[str, Any]:
|
||
"""GET /v2/credential-profiles — profiles + per-harness native rows.
|
||
|
||
The daemon's payload already distinguishes the two verification
|
||
truths the UI must show honestly (Q2-а): ``verification_source``
|
||
``local_store`` (material present, liveness UNPROVEN) vs ``vendor``
|
||
(the vendor answered a request with this credential)."""
|
||
body = self._request("GET", "/v2/credential-profiles")
|
||
return body if isinstance(body, dict) else {}
|
||
|
||
def create_credential_profile(self, harness_id: str, profile_id: str,
|
||
display_name: str = "") -> Dict[str, Any]:
|
||
request: Dict[str, Any] = {"harnessId": str(harness_id), "profileId": str(profile_id)}
|
||
if display_name:
|
||
request["displayName"] = str(display_name)
|
||
body = self._request(
|
||
"POST", "/v2/credential-profiles", json_body=request,
|
||
headers={"Idempotency-Key": uuid.uuid4().hex},
|
||
)
|
||
return body if isinstance(body, dict) else {}
|
||
|
||
def update_credential_profile(self, harness_id: str, profile_id: str,
|
||
*, enabled: bool) -> Dict[str, Any]:
|
||
"""PATCH /v2/credential-profiles/:harness/:profileId — the engine's own
|
||
Enabled toggle for a NAMED account (``{enabled}`` is the one
|
||
user-settable routing control the profile row carries).
|
||
|
||
Translate-only, like every account surface here: the daemon owns the
|
||
registry row and rotation policy, and its refusal is the answer. The
|
||
route exists on 3.5.0 engines already; unified-model engines serve the
|
||
migrated default logins through it too, because those are ordinary
|
||
registry rows there."""
|
||
from urllib.parse import quote
|
||
|
||
body = self._request(
|
||
"PATCH",
|
||
f"/v2/credential-profiles/{quote(str(harness_id), safe='')}"
|
||
f"/{quote(str(profile_id), safe='')}",
|
||
json_body={"enabled": bool(enabled)},
|
||
)
|
||
return body if isinstance(body, dict) else {}
|
||
|
||
def delete_credential_profile(self, harness_id: str, profile_id: str) -> Dict[str, Any]:
|
||
"""DELETE /v2/credential-profiles/:harness/:profileId — the engine's own
|
||
removal contract for a NAMED account.
|
||
|
||
Ouroboros never deletes vendor credential material itself: the daemon
|
||
owns the profile record and whatever it stored for it, so removal is a
|
||
request to that owner and its refusal is the answer. There is no
|
||
counterpart for a native CLI login — that account belongs to the
|
||
vendor's own CLI, and simulating a sign-out here would claim an effect
|
||
this process cannot have."""
|
||
from urllib.parse import quote
|
||
|
||
body = self._request(
|
||
"DELETE",
|
||
f"/v2/credential-profiles/{quote(str(harness_id), safe='')}"
|
||
f"/{quote(str(profile_id), safe='')}",
|
||
)
|
||
return body if isinstance(body, dict) else {}
|
||
|
||
def harness_models(self, harness_id: str) -> List[Dict[str, Any]]:
|
||
"""GET /v2/harnesses/:id/models — the discovered model list (owner
|
||
directive: models are a dropdown fed by discovery, never free input)."""
|
||
from urllib.parse import quote
|
||
|
||
body = self._request("GET", f"/v2/harnesses/{quote(str(harness_id), safe='')}/models")
|
||
models = body.get("models") if isinstance(body, dict) else None
|
||
return [row for row in (models or []) if isinstance(row, dict)]
|
||
|
||
def harness_model_catalog(self, harness_id: str, *,
|
||
credential_profile_id: Optional[str] = None,
|
||
view: Optional[str] = None) -> Dict[str, Any]:
|
||
"""Retain the catalog envelope for an explicitly negotiated view."""
|
||
from urllib.parse import quote, urlencode
|
||
|
||
path = f"/v2/harnesses/{quote(str(harness_id), safe='')}/models"
|
||
query = {}
|
||
if view is not None:
|
||
query["view"] = view
|
||
if credential_profile_id is not None:
|
||
query["credentialProfileId"] = credential_profile_id
|
||
if query:
|
||
path += "?" + urlencode(query)
|
||
return _model_object(self._request("GET", path))
|
||
|
||
def setup_job_create(self, request: Dict[str, Any], *, idempotency_key: str = "") -> Dict[str, Any]:
|
||
body = self._request(
|
||
"POST", "/v2/setup/jobs", json_body=dict(request),
|
||
headers={"Idempotency-Key": str(idempotency_key or "") or uuid.uuid4().hex},
|
||
)
|
||
return body if isinstance(body, dict) else {}
|
||
|
||
def setup_job_call(self, job_id: str, op: str, *, value: str = "") -> Dict[str, Any]:
|
||
"""One job-scoped setup call: ``snapshot`` (GET; the transient
|
||
device-code/oauth_url disclosure rides it and is never journaled),
|
||
``cancel`` (POST), ``input`` (POST /v2/setup/jobs/{id}/input —
|
||
deliver ONE line of user input, the claude OAuth paste-code, to a
|
||
login job that awaits it; engine 3.3.7+), or ``reconcile``
|
||
(POST /v2/setup/jobs/{id}/reconcile — ask the daemon to prove an
|
||
unconfirmed termination's process group empty; supported floor 3.2.0).
|
||
|
||
The input value is live login material, the same custody rule as the
|
||
device code: it rides this loopback request once and is never logged,
|
||
stored, or echoed by this client. An engine that predates the input
|
||
route answers 404, which the caller must treat as a typed capability
|
||
gap, not a bug.
|
||
"""
|
||
from urllib.parse import quote
|
||
|
||
base = f"/v2/setup/jobs/{quote(str(job_id), safe='')}"
|
||
if op == "snapshot":
|
||
body = self._request("GET", f"{base}/snapshot")
|
||
elif op == "cancel":
|
||
body = self._request("POST", f"{base}/cancel")
|
||
elif op == "input":
|
||
body = self._request("POST", f"{base}/input", json_body={"value": str(value)})
|
||
elif op == "reconcile":
|
||
body = self._request("POST", f"{base}/reconcile")
|
||
else:
|
||
raise ValueError(f"unknown setup job op: {op!r}")
|
||
return body if isinstance(body, dict) else {}
|
||
|
||
def operations(self) -> List[Dict[str, Any]]:
|
||
"""GET /v2/operations — the engine's own implemented-route catalog.
|
||
|
||
The handshake advertises this path (``operationsPath``); the catalog is
|
||
how a caller discovers a capability structurally instead of branching
|
||
on version folklore (verified live: ``{protocolMajor, operations:[{id,
|
||
method, path, ...}]}``).
|
||
"""
|
||
body = self._request("GET", "/v2/operations")
|
||
ops = body.get("operations") if isinstance(body, dict) else None
|
||
return [row for row in (ops or []) if isinstance(row, dict)]
|
||
|
||
def get_settings(self) -> Dict[str, Any]:
|
||
"""GET /v2/settings — the daemon's effective settings snapshot.
|
||
|
||
The read half of the rotation reconcile (B3): the snapshot's
|
||
``harnesses`` map carries each configured harness's
|
||
``profileLimitAction``, so provisioning patches only what is actually
|
||
missing instead of blind-writing every discovered harness. The route
|
||
has served since the v2 boundary existed, so every engine past
|
||
``CLAUDEXOR_MIN_VERSION`` answers it.
|
||
"""
|
||
body = self._request("GET", "/v2/settings")
|
||
return body if isinstance(body, dict) else {}
|
||
|
||
def patch_settings(self, request: Dict[str, Any]) -> Dict[str, Any]:
|
||
"""POST /v2/settings — the daemon's own live settings patch (the
|
||
write half of the rotation reconcile, D28/B3)."""
|
||
body = self._request("POST", "/v2/settings", json_body=dict(request))
|
||
return body if isinstance(body, dict) else {}
|
||
|
||
def set_secret(self, name: str, value: str) -> Dict[str, Any]:
|
||
"""Store one managed secret through the daemon's non-journaled route.
|
||
|
||
The value stays in this loopback request and is never returned or logged.
|
||
This is transport only; callers own the choice of managed slot.
|
||
"""
|
||
body = self._request(
|
||
"POST", "/v2/secrets",
|
||
json_body={"name": str(name), "value": str(value)},
|
||
)
|
||
return body if isinstance(body, dict) else {}
|
||
|
||
|
||
def _trust_state(body: Any) -> Dict[str, Any]:
|
||
if (not isinstance(body, dict) or "repoRoot" not in body
|
||
or (body["repoRoot"] is not None and not isinstance(body["repoRoot"], str))
|
||
or not isinstance(body.get("path"), str) or not body["path"]
|
||
or type(body.get("allowFullAccess")) is not bool
|
||
or not isinstance(body.get("accessDefault"), str) or not body["accessDefault"]):
|
||
raise ClaudexorUnavailable("malformed_response", "Trust lookup returned an invalid state")
|
||
return body
|
||
|
||
|
||
def _model_object(body: Any) -> Dict[str, Any]:
|
||
if not isinstance(body, dict):
|
||
raise ClaudexorUnavailable("malformed_response", "Model transport returned a non-object response")
|
||
return body
|
||
|
||
|
||
def _model_operation(body: Any, operation_id: Optional[str] = None) -> Dict[str, Any]:
|
||
body = _model_object(body)
|
||
if (not isinstance(body.get("id"), str) or not body["id"]
|
||
or (operation_id is not None and body["id"] != operation_id)):
|
||
raise ClaudexorUnavailable("malformed_response", "Model operation response identity mismatch")
|
||
return body
|
||
|
||
|
||
def _model_idempotency_key(value: str) -> str:
|
||
if not isinstance(value, str) or not value.strip() or len(value.strip()) > 256:
|
||
raise ClaudexorUnavailable("invalid_idempotency_key", "A stable caller-supplied model invocation key is required")
|
||
return value.strip()
|
||
|
||
|
||
def _model_payload_ref(value: Dict[str, Any]) -> Dict[str, Any]:
|
||
if (not isinstance(value, dict) or set(value) != {"resourceId", "sha256", "sizeBytes"}
|
||
or not isinstance(value.get("resourceId"), str) or not value["resourceId"]
|
||
or not isinstance(value.get("sha256"), str)
|
||
or re.fullmatch(r"sha256:[a-f0-9]{64}", value["sha256"]) is None
|
||
or type(value.get("sizeBytes")) is not int or value["sizeBytes"] < 0):
|
||
raise ClaudexorUnavailable("model_payload_ref_invalid", "Invalid model payload reference")
|
||
return dict(value)
|
||
|
||
|
||
def _reject_json_constant(_value: str) -> None:
|
||
raise ValueError("Non-finite JSON number")
|
||
|
||
|
||
def pending_interactions(detail: Dict[str, Any]) -> List[Dict[str, Any]]:
|
||
"""The run detail's live interactive questions, normalized and complete.
|
||
|
||
``GET /v2/runs/:id`` carries ``pendingInteractions`` — full
|
||
``ControlPendingInteraction`` rows with the question TEXT, header, options and
|
||
``multi_select``, not just the ``summary.waitingOnUser`` boolean the old wait
|
||
kept. This is the ONE reader of that wire shape: snake_case keys out, absent
|
||
strings normalized to ``None``/empty, rows without an interaction id dropped
|
||
(an unanswerable row is noise, not a question). Purely shape translation — no
|
||
truncation here; bounding belongs to the delivery layer that knows its budget.
|
||
"""
|
||
rows = detail.get("pendingInteractions") if isinstance(detail, dict) else None
|
||
out: List[Dict[str, Any]] = []
|
||
for row in rows or []:
|
||
if not isinstance(row, dict):
|
||
continue
|
||
questions: List[Dict[str, Any]] = []
|
||
for question in row.get("questions") or []:
|
||
if not isinstance(question, dict):
|
||
continue
|
||
questions.append({
|
||
"question_id": str(question.get("id") or ""),
|
||
"question": str(question.get("question") or ""),
|
||
"header": str(question.get("header") or "") or None,
|
||
"options": [
|
||
{"label": str(option.get("label") or ""),
|
||
"description": str(option.get("description") or "") or None}
|
||
for option in (question.get("options") or [])
|
||
if isinstance(option, dict)
|
||
],
|
||
"multi_select": bool(question.get("multi_select")),
|
||
})
|
||
interaction_id = str(row.get("interactionId") or "")
|
||
if not interaction_id:
|
||
continue
|
||
out.append({
|
||
"interaction_id": interaction_id,
|
||
"source_tool": str(row.get("sourceTool") or "") or None,
|
||
"requested_at": str(row.get("requestedAt") or ""),
|
||
"timeout_at": str(row.get("timeoutAt") or "") or None,
|
||
"questions": questions,
|
||
})
|
||
return out
|
||
|
||
|
||
# -- applied-fact artifacts ----------------------------------------------------
|
||
|
||
|
||
@dataclass(frozen=True)
|
||
class AttemptContainment:
|
||
"""What one attempt's harness ACTUALLY ran under, as the engine recorded it.
|
||
|
||
TWO axes, one record, because they are two different guarantees and reading only
|
||
the first is how a `~`-redirect gets reported as a sandbox:
|
||
|
||
``home_isolated`` — the scoped ``HOME``. ``None`` when the attempt recorded no such
|
||
fact at all: "unverified", which is not "verified false" and must not be collapsed
|
||
into it by a caller reading a bare boolean. The consequence of a recorded false is a
|
||
CANCELLATION, so absence must stay absence.
|
||
|
||
``boundary_mechanism`` — the OS-enforced filesystem boundary the engine APPLIED
|
||
(``""`` when the attempt named none). Silence collapses to "no boundary" here, and
|
||
deliberately not for the home above: an engine that reports nothing is
|
||
indistinguishable from an engine that applied nothing, and the consequence of
|
||
reading it that way is a DISCLOSURE rather than a refusal. That direction is safe;
|
||
the opposite one would let an unconfined run pass as confined.
|
||
|
||
``confinement_unavailable_reason`` — the engine's own typed explanation for a
|
||
missing boundary (e.g. no mechanism exists for this host), read from the SAME
|
||
attempt artifact. Telemetry that AMPLIFIES the unconfined disclosure — never
|
||
an admission token: an old engine that writes nothing here changes no
|
||
decision, and a reason's presence never excuses a recorded FALSE.
|
||
"""
|
||
|
||
attempt_id: str
|
||
home_isolated: Optional[bool]
|
||
home_dir: str
|
||
boundary_mechanism: str = ""
|
||
confinement_unavailable_reason: str = ""
|
||
|
||
|
||
def attempt_containment(run_dir: str) -> List[AttemptContainment]:
|
||
"""Read every attempt's APPLIED containment facts from the run's own artifacts.
|
||
|
||
Claudexor writes ``harness_home_isolated`` / ``harness_home_dir`` and the applied
|
||
boundary (``confinement_mechanism``, with ``confinement_verified_denied_path`` as
|
||
its proof) onto ``<runDir>/attempts/<id>/attempt.yaml``. The HOME pair is projected
|
||
onto no ``/v2`` response, so for that half this file is the only evidence a caller
|
||
has that the confinement it asked for was applied rather than merely requested; the
|
||
BOUNDARY half is also served on the run detail as ``candidates[].confinement`` (since
|
||
3.3.6). The artifact is read for both, because it is where the two meet per attempt.
|
||
|
||
The mechanism is read as an OPAQUE string. Which mechanisms exist, and on which
|
||
hosts, is the engine's business — Ouroboros asks what was applied and never which
|
||
OS it is sitting on, so the day a second mechanism ships this reader is unchanged.
|
||
|
||
Empty means NO EVIDENCE, never "not isolated": the record is written when an
|
||
attempt finishes, so a young run legitimately has none. Unreadable artifacts are
|
||
skipped for the same reason — a caller distinguishes absence from a recorded false.
|
||
"""
|
||
root = pathlib.Path(str(run_dir or "")) / _ATTEMPTS_REL
|
||
try:
|
||
attempt_dirs = sorted(entry for entry in root.iterdir() if entry.is_dir())
|
||
except (OSError, ValueError, RuntimeError):
|
||
return []
|
||
import yaml # type: ignore
|
||
|
||
applied: List[AttemptContainment] = []
|
||
for attempt_dir in attempt_dirs:
|
||
try:
|
||
record = yaml.safe_load((attempt_dir / _ATTEMPT_RECORD).read_text(encoding="utf-8"))
|
||
except (OSError, ValueError, RuntimeError, yaml.YAMLError):
|
||
continue
|
||
if not isinstance(record, dict):
|
||
continue
|
||
raw = record.get("harness_home_isolated")
|
||
# A mechanism is only evidence WITH its proof: 3.3.2 records the mechanism
|
||
# beside the path the policy was executed against and refused, on this host,
|
||
# before the harness spawned. A name on its own is the promise the applied-fact
|
||
# block exists to replace, so it is read as no boundary at all.
|
||
mechanism = str(record.get("confinement_mechanism") or "").strip()
|
||
proven = str(record.get("confinement_verified_denied_path") or "").strip()
|
||
applied.append(AttemptContainment(
|
||
attempt_id=str(record.get("attempt_id") or attempt_dir.name),
|
||
home_isolated=raw if isinstance(raw, bool) else None,
|
||
home_dir=str(record.get("harness_home_dir") or ""),
|
||
boundary_mechanism=mechanism if (mechanism and proven) else "",
|
||
confinement_unavailable_reason=str(
|
||
record.get("confinement_unavailable_reason") or ""
|
||
).strip(),
|
||
))
|
||
return applied
|
||
|
||
|
||
def final_attempt_facts(detail: Dict[str, Any], run_id: str) -> Dict[str, str]:
|
||
"""Read the final attempt's route facts from engine-owned telemetry.
|
||
|
||
The summary's model and harnesses echo the request; its route/authRoute
|
||
projections may borrow facts from earlier attempts. Only the unique row
|
||
named by final_attempt_id belongs to the delivered result. Missing facts
|
||
stay unknown, never filled from another attempt or the request.
|
||
"""
|
||
summary = detail.get("summary") if isinstance(detail, dict) else None
|
||
if not isinstance(summary, dict) or not isinstance(run_id, str) or not run_id.strip():
|
||
return {}
|
||
run_dir = summary.get("runDir")
|
||
if not isinstance(run_dir, str) or not run_dir.strip():
|
||
return {}
|
||
import yaml # type: ignore
|
||
|
||
try:
|
||
record = yaml.safe_load(
|
||
(pathlib.Path(run_dir) / "final" / "telemetry.yaml").read_text(encoding="utf-8"))
|
||
except (OSError, ValueError, RuntimeError, yaml.YAMLError):
|
||
return {}
|
||
if not isinstance(record, dict) or record.get("run_id") != run_id:
|
||
return {}
|
||
final_id, attempts = record.get("final_attempt_id"), record.get("attempts")
|
||
if not isinstance(final_id, str) or not final_id.strip() or not isinstance(attempts, list):
|
||
return {}
|
||
matching = [row for row in attempts if isinstance(row, dict) and row.get("attempt_id") == final_id]
|
||
if len(matching) != 1:
|
||
return {}
|
||
row = matching[0]
|
||
return {
|
||
target: row.get(source) if isinstance(row.get(source), str) else ""
|
||
for target, source in (
|
||
("attempt_id", "attempt_id"), ("harness_id", "harness_id"),
|
||
("model", "observed_model"), ("profile_id", "profile_id"),
|
||
)
|
||
}
|
||
|
||
|
||
__all__ = [
|
||
"AttemptContainment",
|
||
"ClaudexorGateway",
|
||
"ClaudexorSubscriptionWindowExhausted",
|
||
"ClaudexorUnavailable",
|
||
"DaemonEndpoint",
|
||
"WINDOW_EXHAUSTED_CODES",
|
||
"attempt_containment",
|
||
"discover_daemon",
|
||
"discover_daemon_at",
|
||
"engine_at_least",
|
||
"final_attempt_facts",
|
||
"operator_home",
|
||
"pending_interactions",
|
||
"run_failure_cause",
|
||
"run_failure_error",
|
||
]
|