Restore supervising model recovery and retain failed response evidence

Preserve free live-executor waiting before ordinary managed model recovery.
Configured-session supervisors use the same upstream-observation path as
other managed tasks; external start and patch custody guards stay in place.
Negotiate private response evidence and preserve known terminal usage when
local message normalization fails, without introducing another retry owner.

Keep unknown accounting and exact source custody explicit (P1), repair the
shared recovery boundary (P2), and reuse existing transport/result owners
(P7). Update coupled tests and architecture/development descriptions.
This commit is contained in:
Ouroboros 2026-09-15 01:40:05 +03:00
parent 1e3d5335e9
commit 627ee94471
25 changed files with 622 additions and 133 deletions

File diff suppressed because one or more lines are too long

View file

@ -2206,7 +2206,14 @@ Focused regressions: `test_review_late_cas_recovery.py`, `test_delivery_control_
operation; never inject its credentials, run its tools, compact inside the
adapter or silently repeat a generation. Recover a local connection loss using
the same operation ID; record unknown outcomes as unknown. ACK only after the
existing private CAS owns the exact result. Optional host hints must be chosen
existing private CAS owns the exact result. Failed-response capture uses the
operation catalog's optional query, frozen before create and reused with the
same idempotency key; absence preserves the strict legacy result shape. Keep
full received bytes/exception chains private and compact diagnostics in the
ordinary problem context. A known terminal with unusable output is a settled
provider result plus local rejection, never unknown or not-dispatched: preserve
both existing stream-rejection markers across sync, async and process boundaries.
Local rejection must not rotate accounts. Optional host hints must be chosen
by their caller according to transport capability; explicit unsupported options
refuse, rather than being silently removed and retried. Record submitted model
options beside the engine's applied options on the usage row; an absent report
@ -2406,8 +2413,11 @@ by "Provider Independence" above. Call-site imperatives:
on the wire before a later reservation can be refused), so the record keeps
the attempt booked and the budget terminal, not the provider terminal, ends
the round; the bounded repeat rail belongs to interactive primary rounds. Ordinary managed
tasks and native API children use upstream-observed continuation, while exact
session nanny routes keep their independent hold. Every other caller — forced-final, fallback
tasks and native API children use upstream-observed continuation. Exact
session supervisors first try the live-leaf hold and otherwise use ordinary
managed recovery for their own model; no consumed/terminal/patch-disposition
predicate gates cognition. A successful hold closes and clears any prior
transport episode, so its acknowledged wake alone resumes the model. Every other caller — forced-final, fallback
candidates, review actors, safety, external-harness delegated runs — keeps
`transport_death_retries=0`. A round that holds a transport-death repeat
record sends nothing further except the typed-death repeats — a repeat that

View file

@ -18,10 +18,11 @@ Eligibility is deliberately narrow (owner Q8=A): configured-session/exact
actor routes with EXACTLY one open, non-terminal delegated run and no pending
invocations. The liveness probe is a READ-ONLY engine poll requiring a
positive non-terminal engine state — a refusal, fault, empty state, or open
containment fault is not evidence of a live leaf and takes today's terminal
path. A terminal-but-unsettled leaf also takes today's terminal path, whose
completion-wins reconciliation already preserves the leaf's output — holding
there would spin an instant-wake paid loop.
containment fault is not evidence of a live leaf. When no hold applies, the
supervising model uses ordinary managed transport recovery; external custody
remains unchanged. A terminal-but-unsettled leaf must not hold, since an
instant terminal wake would buy repeated model calls. A successful hold takes
over any existing transport episode, so its wake alone permits the next round.
"""
from __future__ import annotations
@ -151,12 +152,12 @@ def _unknown_attempt_id() -> str:
def latch_after_unknown(
tools: Any, *, error_kind: str, drive_logs: pathlib.Path, task_id: str,
emit_progress: Any,
emit_progress: Any, transport_episode: Any = None,
) -> bool:
"""Round-gate entry: latch a hold instead of terminalizing.
"""Round-gate entry: prefer a live leaf to a new model recovery episode.
Returns True when the caller should ``continue`` to the next round top,
where the durable latch parks the task in ``hold_step``.
clearing its transport episode; the durable latch parks in ``hold_step``.
"""
if str(error_kind or "") != "provider_outcome_unknown":
return False
@ -178,6 +179,15 @@ def latch_after_unknown(
**({"unknown_attempt_id": attempt_id} if attempt_id else {}),
}
write_unknown_hold(ctx, run_id, hold)
if transport_episode is not None:
from ouroboros.loop_transport import emit_network_wait_event
emit_network_wait_event(
drive_logs, task_id=task_id, phase="ended",
elapsed_sec=transport_episode.waited_sec, redials=transport_episode.redials,
model=str(getattr(ctx, "active_model", "") or ""), detail="hold_latched",
outcome_custody=transport_episode.outcome_custody,
)
_emit_hold_event(
drive_logs, task_id=task_id, phase="entered", run_id=run_id,
hold_cycles=hold["hold_cycles"], attempt_id=attempt_id,
@ -300,6 +310,7 @@ def hold_step(
return "terminal"
if controls.get("finalize_now"):
ctx._transport_repeat_control_reason = str(controls["finalize_now"]).splitlines()[0].strip()
return _exit_terminal("finalize_now")
if new_input:
# The round-top drain already appended owner/task dialogue: that IS the
@ -336,6 +347,9 @@ def hold_step(
pend_payload = None
if (_wake_is_control(payload) or _control_wakes(ctx)
or (isinstance(pend_payload, dict) and _wake_is_control(pend_payload))):
from ouroboros.loop_transport import transport_repeat_stop_requested
transport_repeat_stop_requested(ctx) # Preserve a current owner cause without delivering/ACKing it.
return _exit_terminal("control_wake")
if str(payload.get("status") or "") in _NON_WAKE_STATUSES:
# A refusal/fault is a daemon statement, not a leaf wake: no proof of a

View file

@ -286,17 +286,28 @@ def _is_loopback(host: str) -> bool:
return False
def account_catalog_supported(operations: list[dict], path: str) -> bool:
"""Opt in only when this exact operation declares the accounts query view."""
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") == "GET" and operation.get("path") == path
and any(parameter.get("name") == "view" and parameter.get("location") == "query"
and "accounts" in (parameter.get("enum") or [])
for parameter in operation.get("parameters", []) if isinstance(parameter, dict))
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")
class ClaudexorGateway:
"""Thin typed client over the Claudexor ``/v2`` control API."""
@ -568,11 +579,12 @@ class ClaudexorGateway:
return ref
def create_model_operation(self, request_ref: Dict[str, Any], *,
idempotency_key: str) -> 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", "/v2/model-operations", json_body={"request": _model_payload_ref(request_ref)},
"POST", path, json_body={"request": _model_payload_ref(request_ref)},
headers={"Idempotency-Key": key},
))

View file

@ -54,7 +54,9 @@ from ouroboros._usage_response import provider_cost_value
from ouroboros.anthropic_native_custody import scrub_native_custody
from ouroboros.claudexor_daemon import ensure_owned_gateway, owned_engine_version, read_owned_gateway
from ouroboros.deadline_utils import llm_transport_timeout_sec
from ouroboros.gateways.claudexor import ClaudexorUnavailable, engine_at_least, _READ_TIMEOUT_SEC
from ouroboros.gateways.claudexor import (
ClaudexorUnavailable, engine_at_least, model_failure_evidence_supported, _READ_TIMEOUT_SEC,
)
from ouroboros.llm_attempt import _attempt_request, _candidate_before_dispatch
from ouroboros.model_slots import MODEL_ACCOUNTS_KEY, model_role_option
from ouroboros.model_wait import ModelWaitInterrupted, current_model_wait, prepared_call_scope
@ -148,6 +150,11 @@ class ClaudexorModelError(RuntimeError):
self.status_code = 0 if unknown else int(context.get("httpStatus") or 0)
self.reset_at = str(context.get("resetsAt") or "")
self.retryable = False if unknown else problem.get("retryable") is True
if self.code == "response_rejected":
# Reconstructed process-boundary errors retain the same local
# rejection: no unknown outcome or same-request wire repair.
self.stream_rejected = self.stream_incomplete = True
self.retryable = False
self.model_role = model_role
self.operation_id = operation_id
self.route = copy.deepcopy(route or {})
@ -156,9 +163,12 @@ class ClaudexorModelError(RuntimeError):
def display_message(self) -> str:
"""Show typed provider details without changing exception classification text."""
context = self.problem.get("context") or {}
details = [] if self.code == "model_outcome_unknown" else [
fields = (("stage", "stage"), ("errorCode", "cause"))
if self.code != "model_outcome_unknown":
fields += (("vendorCode", "provider_code"), ("parameter", "parameter"))
details = [
f"{label}={value.strip()}"
for key, label in (("vendorCode", "provider_code"), ("parameter", "parameter"))
for key, label in fields
if isinstance(value := context.get(key), str) and value.strip()
]
# Details lead so the existing terminal preview can name the refusal.
@ -280,6 +290,8 @@ def adopt_turn_state(slot: ModelTurnState | None, payload: dict, result: dict) -
def _remember_failed_profile(target: dict, parameters: dict, error: ClaudexorModelError) -> None:
if getattr(error, "stream_rejected", False):
return # Local message normalization says nothing about account readiness.
route = error.route or {}
key = (parameters.get("cache_affinity"), route.get("source"), route.get("model"))
if (key == (parameters.get("cache_affinity"), target["source"], target["resolved_model"])
@ -391,6 +403,7 @@ class _ModelInvocation:
self.request_manifest_ref: dict = {}
self.interrupt_reason = ""
self.create_attempted = False
self.capture_failure_evidence = False
self.defer_close = False
self.io_active = False
self.io_lock = threading.Lock()
@ -429,11 +442,15 @@ class _ModelInvocation:
self.check_control()
try:
self.gateway = ensure_owned_gateway()
# Freeze once before create. A lost create reply or replaced gateway
# must reuse this same operation's diagnostic/idempotency contract.
self.capture_failure_evidence = model_failure_evidence_supported(self.gateway.operations())
self.request_ref = self.gateway.upload_model_request(self.payload, idempotency_key=self.invocation_id)
self.request_manifest_ref = persist_call(self.root, task_id=self.task_id, call_id=f"{self.invocation_id}_model_request",
call_type="llm_claudexor_request", payload=self.payload, keep_raw=True,
manifest={"invocation_id": self.invocation_id, "request_ref": self.request_ref,
"model_role": self.role})["manifest_ref"]
"model_role": self.role,
"capture_failure_evidence": self.capture_failure_evidence})["manifest_ref"]
except ClaudexorUnavailable as error:
raise ClaudexorModelError({"code": error.code, "message": str(error)}, model_role=self.role) from None
@ -447,7 +464,8 @@ class _ModelInvocation:
self.root, task_id=self.task_id, call_id=f"{self.invocation_id}_model_request",
call_type="llm_claudexor_request", payload=self.payload, keep_raw=True,
manifest={"invocation_id": self.invocation_id, "request_ref": self.request_ref,
"model_role": self.role, "operation_id": self.operation_id})["manifest_ref"]
"model_role": self.role, "operation_id": self.operation_id,
"capture_failure_evidence": self.capture_failure_evidence})["manifest_ref"]
except Exception as error:
self.request_manifest_ref = {}
log.warning("Model custody checkpoint unavailable: %s", type(error).__name__)
@ -475,7 +493,8 @@ class _ModelInvocation:
if not self.operation_id:
self.create_attempted = True
self.observe_operation()
detail = self.gateway.create_model_operation(self.request_ref, idempotency_key=self.invocation_id)
detail = self.gateway.create_model_operation(self.request_ref, idempotency_key=self.invocation_id,
**({"capture_failure_evidence": True} if self.capture_failure_evidence else {}))
self.operation_id = detail["id"]
self.observe_operation(accepted=True)
else:
@ -664,7 +683,10 @@ class _ModelInvocation:
raise error
message = result.get("message")
if not isinstance(message, dict):
error = self.error({"code": "malformed_response", "message": "The provider returned no model message."}, unknown=True)
error = ClaudexorModelError(result.get("problem") or {
"code": "response_rejected", "message": "The terminal provider response contained no usable model message."},
model_role=self.role, operation_id=self.operation_id, route=route)
error.stream_rejected = error.stream_incomplete = True
error.physical_attempt_capture = self.capture
error.usage = usage
raise error

View file

@ -553,6 +553,11 @@ def run_llm_loop(
tools._ctx._current_llm_call_meta = dict(accumulated_usage.get("_last_llm_call_meta") or {})
last_error_kind = str(accumulated_usage.get("_last_llm_error_kind") or "")
if msg is None and _delegate_hold_latch(
tools, error_kind=last_error_kind, drive_logs=drive_logs,
task_id=task_id, emit_progress=emit_progress, transport_episode=transport_wait):
transport_wait = None # The leaf wake owns resumption, without a provider probe.
continue
transport_wait = _reconcile_transport_wait(
transport_wait, ctx, msg_present=msg is not None, error_kind=last_error_kind,
drive_logs=drive_logs, task_id=task_id, model=active_model, emit_progress=emit_progress)
@ -589,10 +594,6 @@ def run_llm_loop(
emit_progress=emit_progress, incoming_messages=incoming_messages, owner_msg_seen=_owner_msg_seen):
free_redial = True
continue
if msg is None and _delegate_hold_latch(
tools, error_kind=last_error_kind, drive_logs=drive_logs,
task_id=task_id, emit_progress=emit_progress): # hold latched -> next round top parks
continue
if msg is None:
# Exact actor routes skip generic substitution and fail as infrastructure.
text, accumulated_usage, forced_trace = _handle_provider_unavailable(

View file

@ -143,10 +143,9 @@ def emit_network_wait_event(
def managed_transport_continuation(ctx: Any) -> bool:
"""Owner-selected continuation applies to ordinary managed cognition."""
"""Managed cognition may recover; a live delegated-leaf hold takes priority."""
return bool(ctx is not None and getattr(ctx, "task_id", "")
and not getattr(ctx, "is_direct_chat", False)
and getattr(ctx, "_configured_subagent_route_kind", "") != "agent_session")
and not getattr(ctx, "is_direct_chat", False))
def continue_unknown_transport(episode: TransportWaitEpisode, *, llm: Any, tools: Any,

View file

@ -58,7 +58,7 @@ def test_resource_refusal_does_not_suspend_cyber_execution(live_wait, tmp_path,
events = [json.loads(line) for line in (root / "logs/events.jsonl").read_text().splitlines()]
advisory = next(row for row in events if row["type"] == "safety_advisory")
assert advisory["assessment_allowed"] is False and code in advisory["assessment"]
assert len(transport.operations) == 1
assert len(transport.accepted_operations) == 1
assert [row["state"] for row in ledger(root)] == ["reserved", "dispatched", "released"]
assert current_model_wait() is controller and not controller.closed
assert "wait_for_resources" not in json.dumps(transport.uploads[0][0])

View file

@ -0,0 +1,198 @@
"""Complete failed model evidence and known-terminal rejection share existing custody."""
import asyncio
import base64
from copy import deepcopy
from dataclasses import asdict
import json
from types import SimpleNamespace
import pytest
from ouroboros import llm_claudexor as transport, usage_accounting as ua
from ouroboros.gateways.claudexor import ClaudexorUnavailable
from ouroboros.loop_llm_call import classify_llm_exception
from ouroboros.request_wire_recovery import plan_next_wire_retry
from ouroboros.tools import vision_process
from ouroboros.transport_custody import ProviderNotDispatched
from tests.test_llm_claudexor import Gateway, MODEL, ROUTE, ledger, result, retained, setup as setup
CAPTURE_OPERATION = {"method": "POST", "path": "/v2/model-operations", "parameters": [
{"name": "captureFailureEvidence", "location": "query", "enum": ["true", "false"]}]}
REJECTION = {"code": "response_rejected", "message": "The terminal response could not form a model message.",
"retryable": False, "context": {"stage": "message", "requestId": "request-one"}}
def failure_evidence(body=b'\xffprivate-wire-marker\r\ndata: invalid JSON\n\n'):
return {"bodyBase64": base64.b64encode(body).decode("ascii"), "receivedBytes": len(body),
"bodyComplete": False, "stage": "message", "causeCycle": False,
"errors": [{"name": "SyntaxError", "message": "private-error-marker", "stack": "private-stack-marker",
"code": None}]}
def call(client, asynchronous, **kwargs):
value = (client.chat_async if asynchronous else client.chat)([], MODEL, **kwargs)
return asyncio.run(value) if asynchronous else value
@pytest.mark.parametrize("asynchronous", [False, True])
@pytest.mark.parametrize("supported", [False, True])
def test_capture_is_negotiated_once_without_changing_provider_payload(setup, asynchronous, supported):
root, gateway, client = setup
gateway.operation_catalog = [deepcopy(CAPTURE_OPERATION)] if supported else []
call(client, asynchronous, model_role="main")
assert gateway.catalog_reads == 1
assert gateway.capture_requests == ([{"capture_failure_evidence": True}] if supported else [{}])
payload = gateway.uploads[0][0]
assert "captureFailureEvidence" not in json.dumps(payload)
assert "capture_failure_evidence" not in json.dumps(payload)
assert retained(root, "request") == payload
manifests = list((root / "observability/calls/task-one").glob("*_model_request.json"))
assert len(manifests) == 1
manifest = json.loads(manifests[0].read_text())
assert manifest["capture_failure_evidence"] is supported
assert manifest["operation_id"] == "op-0"
@pytest.mark.parametrize("asynchronous", [False, True])
def test_lost_create_and_gateway_replacement_reuse_frozen_capture(setup, monkeypatch, asynchronous):
root, gateway, client = setup
gateway.operation_catalog = [deepcopy(CAPTURE_OPERATION)]
gateway.lose_create = True
replacement = Gateway()
for name in ("accepted_operations", "creates", "capture_requests"):
setattr(replacement, name, getattr(gateway, name))
monkeypatch.setattr(transport, "read_owned_gateway", lambda: replacement)
monkeypatch.setattr(transport, "current_model_wait", lambda: SimpleNamespace(
tool_context=SimpleNamespace(task_id="task-one", is_direct_chat=False), control_reason=lambda: None))
monkeypatch.setattr(transport.config, "NETWORK_WAIT_BACKOFF_START_SEC", 0.001)
monkeypatch.setattr(transport.config, "NETWORK_WAIT_BACKOFF_MAX_SEC", 0.001)
answer, usage = call(client, asynchronous)
assert answer == result()["message"]
assert gateway.catalog_reads == 1 and replacement.catalog_reads == 0
assert gateway.capture_requests == [{"capture_failure_evidence": True}] * 2
assert len(gateway.creates) == 2 and len(set(gateway.creates)) == 1
assert len(gateway.accepted_operations) == len(usage["ledger_attempt_ids"]) == 1
assert gateway.closed == replacement.closed == 1
assert [row["state"] for row in ledger(root)] == ["reserved", "dispatched", "settled"]
@pytest.mark.parametrize("asynchronous", [False, True])
def test_catalog_read_failure_does_not_guess_unsupported_or_dispatch(setup, monkeypatch, asynchronous):
root, gateway, client = setup
def failed_catalog():
raise ClaudexorUnavailable("daemon_unreachable", "Metadata connection lost")
monkeypatch.setattr(gateway, "operations", failed_catalog)
with pytest.raises(transport.ClaudexorModelError) as caught:
call(client, asynchronous)
assert caught.value.code == "daemon_unreachable"
assert caught.value.physical_attempt_capture.state == "released"
assert not gateway.uploads and not gateway.creates and gateway.closed == 1
assert [row["state"] for row in ledger(root)] == ["reserved", "released"]
@pytest.mark.parametrize("asynchronous", [False, True])
@pytest.mark.parametrize("outcome", ["completed", "incomplete"])
@pytest.mark.parametrize("has_problem", [False, True])
def test_known_terminal_null_message_settles_then_rejects_without_private_projection(
setup, caplog, asynchronous, outcome, has_problem,
):
root, gateway, client = setup
gateway.operation_catalog = [deepcopy(CAPTURE_OPERATION)]
evidence = failure_evidence()
original_problem = deepcopy(REJECTION) if has_problem else None
gateway.results = [{**result(outcome=outcome, problem=original_problem),
"message": None, "failureEvidence": evidence}]
acknowledge = gateway.acknowledge_model_result
def retained_first(*args):
assert retained(root) == gateway.results[0]
assert ledger(root)[-1]["state"] == "settled"
return acknowledge(*args)
gateway.acknowledge_model_result = retained_first
with pytest.raises(transport.ClaudexorModelError) as caught:
call(client, asynchronous, model_role="vision")
error = caught.value
assert type(error) is transport.ClaudexorModelError
assert not isinstance(error, ProviderNotDispatched)
assert error.code == "response_rejected" and error.stream_rejected and error.stream_incomplete
assert error.operation_id == "op-0" and error.model_role == "vision" and error.route == ROUTE
if has_problem:
assert error.problem == original_problem
assert error.physical_attempt_capture.state == "settled"
assert error.usage["prompt_tokens"] == 20 and error.usage["completion_tokens"] == 7
assert error.usage["cost"] is None and error.usage["cost_final"] is False
assert error.usage["claudexor"]["outcome"] == outcome
assert error.usage["claudexor"]["result_custody"]["state"] == "acknowledged"
classified = classify_llm_exception(error)
assert classified.kind == "provider_error" and not classified.retry_same_request
assert plan_next_wire_retry({}, error=error) is None
assert [row["state"] for row in ledger(root)] == ["reserved", "dispatched", "settled"]
assert len(gateway.creates) == len(gateway.acks) == 1
assert not hasattr(error, "model_result")
public = json.dumps({"usage": error.usage, "ledger": ledger(root), "detail": gateway.detail(0)}) + caplog.text
public += "".join(path.read_text() for path in (root / "logs").glob("*.jsonl"))
for marker in (evidence["bodyBase64"], "private-error-marker", "private-stack-marker"):
assert marker not in public
assert retained(root)["failureEvidence"] == evidence
def test_response_rejection_survives_existing_vision_ipc_reconstruction():
capture = ua.PhysicalAttemptCapture("attempt-one", MODEL, "claudexor", "settled", "opaque")
receipt = {"receipt_id": "receipt-one", "custody": None, "capture": asdict(capture),
"kind": "model", "text": "", "usage": {"prompt_tokens": 20}, "ledger_attempt_ids": ["attempt-one"],
"error": "", "problem": deepcopy(REJECTION), "operation_id": "operation-one", "model_role": "vision",
"route": deepcopy(ROUTE), "unknown": False, "control_reason": "", "model_result": None}
with pytest.raises(transport.ClaudexorModelError) as caught:
vision_process._decode_terminal(json.loads(json.dumps(receipt)), "receipt-one")
error = caught.value
assert error.code == "response_rejected" and error.stream_rejected and error.stream_incomplete
assert error.problem == REJECTION and error.operation_id == "operation-one" and error.route == ROUTE
assert error.physical_attempt_capture.state == "settled" and error.usage == receipt["usage"]
assert classify_llm_exception(error).kind == "provider_error"
assert not isinstance(error, ProviderNotDispatched)
def test_unknown_diagnostic_display_keeps_only_compact_response_context():
error = transport.ClaudexorModelError({"code": "transport_unknown", "message": "Stream interrupted.",
"context": {"stage": "read", "errorCode": "UND_ERR_SOCKET", "requestId": "request-one",
"vendorCode": "not-a-terminal-provider-fact", "stack": "private-stack-marker"}}, unknown=True)
assert error.display_message == (
"stage=read, cause=UND_ERR_SOCKET; model_outcome_unknown: Stream interrupted.")
assert "private-stack-marker" not in error.display_message
assert "request-one" not in error.display_message
assert error.code == "model_outcome_unknown" and not error.retryable
def test_local_rejection_does_not_poison_the_next_account_preference(setup):
_, gateway, client = setup
gateway.results = [{**result(problem=deepcopy(REJECTION)), "message": None}, result()]
gateway.dispatch = ["response_received"] * 2
messages = [result()["message"]]
with pytest.raises(transport.ClaudexorModelError):
client.chat(messages, MODEL, cache_affinity="same-account-after-local-rejection")
client.chat(messages, MODEL, cache_affinity="same-account-after-local-rejection")
assert [payload["account"] for payload, _ in gateway.uploads] == [
{"mode": "auto", "preferredProfileId": "account-a"}] * 2
def test_large_unknown_result_retains_exact_private_bytes_and_existing_pending_ack(setup):
root, gateway, client = setup
gateway.operation_catalog = [deepcopy(CAPTURE_OPERATION)]
body = b'\xff\xc3\x28' + b"x" * (4 * 1024 * 1024) + b"unparsed-suffix\r\n"
evidence = failure_evidence(body)
gateway.results = [{**result(outcome="unknown"), "message": None, "failureEvidence": evidence}]
gateway.dispatch = ["unknown"]
with pytest.raises(transport.ClaudexorModelError) as caught:
client.chat([], MODEL)
assert caught.value.code == "model_outcome_unknown"
assert not getattr(caught.value, "stream_rejected", False)
stored = retained(root)["failureEvidence"]
assert base64.b64decode(stored["bodyBase64"]) == body
assert stored["receivedBytes"] == len(body)
assert ledger(root)[-1]["state"] == "unresolved"
assert len(gateway.creates) == 1 and not gateway.acks

View file

@ -427,3 +427,32 @@ def test_operation_identity_mismatch_is_not_an_accepted_response(gateway_factory
with pytest.raises(cx.ClaudexorUnavailable) as raised:
gateway.get_model_operation("op-one", timeout_sec=0.3)
assert raised.value.code == "malformed_response"
@pytest.mark.parametrize("capture", [False, True])
def test_failure_evidence_opt_in_preserves_create_body(gateway_factory, capture):
wire = ModelWire()
gateway = gateway_factory(wire)
ref = gateway.upload_model_request(_request(), idempotency_key="evidence")
gateway.create_model_operation(ref, idempotency_key="evidence", capture_failure_evidence=capture)
request = wire.calls[-1]
assert request.url.query == (b"captureFailureEvidence=true" if capture else b"")
assert json.loads(request.content) == {"request": ref}
assert wire.payload == _bytes(_request())
assert wire.generation_count == 1
def test_query_negotiation_uses_exact_wire_descriptor():
parameter = {"name": "captureFailureEvidence", "location": "query", "enum": ["true", "false"]}
operation = {"method": "POST", "path": "/v2/model-operations", "parameters": [parameter]}
assert cx.model_failure_evidence_supported([operation])
assert not cx.model_failure_evidence_supported([])
for change in ({"method": "GET"}, {"path": "/v2/model-operations/:id"}, {"parameters": []}):
assert not cx.model_failure_evidence_supported([{**operation, **change}])
for change in ({"name": "other"}, {"location": "header"}, {"enum": ["false"]}, {"enum": "true"}):
assert not cx.model_failure_evidence_supported([{**operation, "parameters": [{**parameter, **change}]}])
accounts = {"method": "GET", "path": "/v2/model-sources", "parameters": [
{"name": "view", "location": "query", "enum": ["accounts"]}]}
assert cx.account_catalog_supported([accounts], accounts["path"])
assert not cx.account_catalog_supported([operation], accounts["path"])
assert not cx.model_failure_evidence_supported([accounts])

View file

@ -0,0 +1,146 @@
"""Supervising cognition recovers without replaying external execution."""
import copy
import hashlib
import json
import httpx
import pytest
from ouroboros import delegate_custody as custody, delegate_hold, loop, loop_transport
from ouroboros import usage_accounting as ua
from ouroboros.delegate_start_claims import claimed_start_request
from tests.test_delegate_hold import _configured_registry, _loop_kwargs, _start_leaf
from tests.test_transport_death_retry import _LedgerLLM, _ledger
@pytest.mark.parametrize("external", ["inline", "consumed", "unread", "patch", "pending", "absent", "unreadable"])
def test_saved_external_work_does_not_gate_supervising_cognition(tmp_path, monkeypatch, external):
"""The incident's long consumed result and uncertain custody use one rail.
Only model I/O is scripted: physical accounting, custody replay, the main
loop and the fresh-start guard run as production code.
"""
task_id, run_id = "t-death", "run-completed"
registry = _configured_registry(tmp_path, task_id)
output = "External result, including Unicode: готово.\n" * (1000 if external == "consumed" else 1)
if external in {"inline", "consumed", "unread", "patch"}:
custody._CUSTODY.pop(run_id, None)
row = custody.RunCustody(run_id=run_id, task_id=task_id, route_id="claude", model="external")
if external == "patch":
row.snapshot_id = "snapshot-preserved"
assert custody.record_started(tmp_path, row)
assert custody.settle_run(tmp_path, None, row, {"summary": {
"state": "succeeded", "spendUsd": 0, "spendEstimated": False,
"inputTokens": 5, "outputTokens": 5,
}})["settled"]
if external in {"consumed", "unread"}:
data = output.encode("utf-8")
artifact = tmp_path / "delegated_runs" / f"{run_id}.json"
artifact.parent.mkdir()
artifact.write_bytes(data)
row.output_artifact, row.output_sha, row.output_complete = (
f"delegated_runs/{run_id}.json", hashlib.sha256(data).hexdigest(), True)
assert custody.emit(tmp_path, custody.OUTPUT_SPILLED, {
"run_id": run_id, "task_id": task_id, "artifact": row.output_artifact,
"sha256": row.output_sha, "bytes": len(data), "staged": True, "full_content": True,
})
if external == "consumed":
assert custody.record_output_consumed(tmp_path, row, artifact=row.output_artifact,
byte_length=len(data), sha256=row.output_sha, chars=len(output), lines=1001)
if external == "patch":
assert custody.record_patch_captured(tmp_path, row, patch_sha256="preserved-patch")
if external == "pending":
assert custody.record_start_requested(tmp_path, task_id=task_id,
invocation_id="pending-invocation", idempotency_key="pending-invocation", request={"prompt": "already sent"})
if external == "unreadable":
monkeypatch.setattr(custody, "custody_log_unreadable", lambda *_: True)
observations, messages_seen = [], []
class Model(_LedgerLLM):
def chat(self, **kwargs):
messages_seen.append(copy.deepcopy(kwargs["messages"]))
return super().chat(**kwargs)
model = Model(tmp_path, lambda: httpx.ReadError("host stream lost after external completion"))
monkeypatch.setattr(loop_transport, "upstream_transport_reachable",
lambda *a, **kw: observations.append(model.calls) or {"kind": "upstream_http", "status_code": 200})
monkeypatch.setattr(loop_transport, "interruptible_wait_sleep", lambda *a: False)
monkeypatch.setattr(delegate_hold, "_leaf_probe_live", lambda *a: pytest.fail("no holdable leaf"))
monkeypatch.setattr(custody, "release_task_runs", lambda *a: None)
monkeypatch.setenv("OUROBOROS_TASK_REVIEW_MODE", "off")
kwargs = _loop_kwargs(tmp_path, registry, [])
kwargs["llm"] = model
if external in {"inline", "consumed", "unread", "patch"}:
kwargs["messages"].extend([
{"role": "assistant", "content": None, "tool_calls": [{"id": "completed-wait", "type": "function",
"function": {"name": "delegate_wait", "arguments": json.dumps({"run_id": run_id})}}]},
{"role": "tool", "tool_call_id": "completed-wait", "content": output},
])
custody_before = custody.event_log_path(tmp_path).read_bytes() if custody.event_log_path(tmp_path).exists() else b""
with ua.usage_scope(ua.UsageScope(drive_root=tmp_path, task_id=task_id, root_task_id=task_id, global_limit_usd=100)):
text, usage, trace = loop.run_llm_loop(**kwargs)
assert text == "done" and model.calls == 2 and observations == [1]
assert not trace["tool_calls"] # No completed delegate/tool execution was replayed.
assert any("NEW physical model attempt" in str(row.get("content")) for row in messages_seen[-1])
if external in {"inline", "consumed", "unread", "patch"}:
assert [row["content"] for row in messages_seen[-1] if row.get("tool_call_id") == "completed-wait"] == [output]
ledger = _ledger(tmp_path)
assert [row["state"] for row in ledger] == ["reserved", "dispatched", "unresolved", "reserved", "dispatched", "settled"]
assert ledger[0]["attempt_id"] != ledger[3]["attempt_id"]
assert usage["transport_recovery"]["previous_attempt"]["physical_attempt_id"] == ledger[0]["attempt_id"]
assert ua.usage_projection(tmp_path)["unresolved_upper_bound_usd"] == 1.0
custody_after = custody.event_log_path(tmp_path).read_bytes() if custody.event_log_path(tmp_path).exists() else b""
rows = [json.loads(line) for line in custody_after.splitlines()]
assert [row for row in rows if row.get("type", "").startswith("delegate_")] == [
json.loads(line) for line in custody_before.splitlines() if json.loads(line).get("type", "").startswith("delegate_")]
if external in {"pending", "patch", "unreadable"}:
# Ability to think is not authority to duplicate the external work.
accepted, refusal = claimed_start_request(tmp_path, claim_target="", payload_busy=lambda *a: "",
actor_ctx=registry._ctx, enforce_actor_idle=True, task_id=task_id, invocation_id="duplicate")
assert not accepted
assert refusal["reason"] == ("replacement_custody_unknown" if external == "unreadable" else "replacement_requires_settlement")
def test_live_hold_takes_over_an_existing_transport_episode(tmp_path, monkeypatch):
registry = _configured_registry(tmp_path)
calls, probes, snapshots = [], [], []
def send(_llm, messages, *args, **kwargs):
usage = args[8]
calls.append(len(calls) + 1)
snapshots.append(copy.deepcopy(messages))
if len(calls) <= 2:
if len(calls) == 2:
# This attempt now owns a live external leaf; the first did not.
_start_leaf(tmp_path)
usage.update(_last_llm_error_kind="provider_outcome_unknown",
_pending_transport_outcome={"physical_attempt_id": f"old-{len(calls)}"})
return None, 0.0
usage.pop("_last_llm_error_kind", None)
return {"role": "assistant", "content": "integrated"}, 0.0
def observed(*args, **kwargs):
probes.append(len(calls))
assert len(calls) == 1, "A leaf wake must not require a new upstream observation"
return {"kind": "upstream_http", "status_code": 200}
monkeypatch.setattr(loop, "call_llm_with_retry", send)
monkeypatch.setattr(loop_transport, "upstream_transport_reachable", observed)
monkeypatch.setattr(loop_transport, "interruptible_wait_sleep", lambda *a: False)
monkeypatch.setattr(delegate_hold, "_leaf_probe_live", lambda *a: True)
monkeypatch.setattr(delegate_hold, "supervised_wait", lambda *a: json.dumps({
"status": "succeeded", "run_id": "run-leaf", "supervision_wake_id": "wake-after-hold"}))
monkeypatch.setattr(delegate_hold, "acknowledge_pending_wake", lambda *a, **kw: True)
monkeypatch.setattr(custody, "release_task_runs", lambda *a: None)
monkeypatch.setenv("OUROBOROS_TASK_REVIEW_MODE", "off")
result, usage, _trace = loop.run_llm_loop(**_loop_kwargs(tmp_path, registry, []))
assert result == "integrated" and len(calls) == 3 and probes == [1]
assert sum("NEW physical model attempt" in str(row.get("content")) for row in snapshots[-1]) == 1
assert any("[DELEGATED LEAF WAKE" in str(row.get("content")) for row in snapshots[-1])
events = [json.loads(line) for line in (tmp_path / "events.jsonl").read_text().splitlines()]
waits = [row for row in events if row.get("type") == "network_wait"]
assert [(row["phase"], row.get("detail", "")) for row in waits] == [
("entered", ""), ("waiting", ""), ("recovered", "new_attempt_after_unknown_outcome"), ("ended", "hold_latched")]
assert waits[-1]["outcome_custody"]["physical_attempt_id"] == "old-1"
assert usage["transport_recovery"]["old_outcome"] == "unknown"

View file

@ -97,7 +97,7 @@ def test_wide_route_switch_sends_once_and_reports_actual_model(packed):
assert text.endswith(REPORT)
assert any(model == MODEL_B for model, _, _ in packed.windows)
assert packed.engine.uploads[-1][0]["messages"][-1]["content"] == PACK
assert len(packed.engine.operations) == 2
assert len(packed.engine.accepted_operations) == 2
assert sum(row["state"] == "settled" for row in ledger(packed.root)) == 1
assert packed.records[-1][0].model == MODEL_B
assert packed.records[-1][0].session_profile == "account-b"
@ -110,7 +110,7 @@ def test_subfloor_switch_refuses_before_new_send_and_keeps_prior_custody(packed,
text, usage = packed.run()
assert usage["execution_status"] == "infra_failed"
assert "200,000" in text and "1,000,000" in text
assert len(packed.engine.operations) == 1
assert len(packed.engine.accepted_operations) == 1
assert sum(row["state"] == "settled" for row in ledger(packed.root)) == int(prior_dispatch == "response_received")
assert any(path.name.endswith("_model_response.json")
for path in (packed.root / "observability" / "calls").rglob("*.json"))
@ -125,7 +125,7 @@ def test_observed_only_subfloor_keeps_paid_report_with_actual_account(packed):
assert "incomplete=none" not in text and text.endswith(REPORT)
assert packed.windows[-1][2]["accountFingerprint"] == "identity-b"
assert packed.records[-1][2]["status"] == "error"
assert len(packed.engine.operations) == 2
assert len(packed.engine.accepted_operations) == 2
assert sum(row["state"] == "settled" for row in ledger(packed.root)) == 1
custody = usage["claudexor"]["result_custody"]
assert custody["state"] == "acknowledged"
@ -147,7 +147,7 @@ def test_auto_actual_account_is_revalidated_without_another_generation(packed, w
assert packed.engine.uploads[0][0]["account"] == {"mode": "auto"}
assert packed.windows[-1][1] == "account-b"
assert packed.windows[-1][2]["accountFingerprint"] == "identity-b"
assert len(packed.engine.operations) == 1
assert len(packed.engine.accepted_operations) == 1
assert sum(row["state"] == "settled" for row in ledger(packed.root)) == 1
@ -166,7 +166,7 @@ def test_changed_wide_route_still_checks_full_input_cap(packed, monkeypatch):
text, usage = packed.run()
assert usage["execution_status"] == "infra_failed"
assert "input" in text.lower() and MODEL_B in text
assert len(packed.engine.operations) == 1
assert len(packed.engine.accepted_operations) == 1
def test_resolve_packed_window_forwards_observed_account(monkeypatch):

View file

@ -3,8 +3,8 @@
A configured-session nanny whose metered round dies ``provider_outcome_unknown``
while EXACTLY one delegated leaf is alive must hold on the LEAF (zero provider
calls) and resume with a wake-bearing NEW round — the unknown request is never
resent. Control wakes and every ineligible shape keep today's no-resend
terminal, and the terminal cleanup (leaf cancellation) fires only on terminals.
resent. Control wakes keep their no-resend terminal. When no hold applies,
ordinary managed recovery can restore the supervising model.
"""
from __future__ import annotations
@ -18,6 +18,7 @@ import pytest
import ouroboros.delegate_hold as delegate_hold
import ouroboros.loop as loop_mod
import ouroboros.loop_transport as transport
from ouroboros import delegate_custody as custody
from ouroboros.delegate_supervision import read_unknown_hold, write_unknown_hold
from ouroboros.loop import run_llm_loop
@ -94,7 +95,21 @@ def _unknown_then_check_call(check):
return fake_call, calls
def _recover_model(monkeypatch):
monkeypatch.setattr(transport, "upstream_transport_reachable",
lambda *a, **kw: {"kind": "upstream_http", "status_code": 200})
monkeypatch.setattr(transport, "interruptible_wait_sleep", lambda *a: False)
def recovered(_messages, usage):
usage.pop("_last_llm_error_kind", None)
return {"role": "assistant", "content": "recovered"}, 0.0
return recovered
def test_unknown_with_live_leaf_holds_and_resumes_with_wake(tmp_path, monkeypatch, _quiet_probe):
monkeypatch.setattr(transport, "upstream_transport_reachable",
lambda *a, **kw: pytest.fail("live hold must precede provider recovery"))
wake_payload = {"status": "succeeded", "run_id": "run-leaf", "supervision_wake_id": "w1"}
monkeypatch.setattr(delegate_hold, "supervised_wait",
lambda _ctx, _run: json.dumps(wake_payload))
@ -127,6 +142,9 @@ def test_unknown_with_live_leaf_holds_and_resumes_with_wake(tmp_path, monkeypatc
assert not read_unknown_hold(registry._ctx).get("run_id") # inactive tombstone
assert _quiet_probe == ["t-hold"] # release only at the (successful) terminal
assert any("holding on the leaf" in note for note in notes)
assert not any(json.loads(line).get("type") == "network_wait"
for line in (tmp_path / "events.jsonl").read_text().splitlines())
assert "transport_recovery" not in usage
def test_terminal_leaf_never_enters_hold(tmp_path, monkeypatch, _quiet_probe):
@ -138,18 +156,17 @@ def test_terminal_leaf_never_enters_hold(tmp_path, monkeypatch, _quiet_probe):
)
monkeypatch.setattr(delegate_hold, "supervised_wait",
lambda *_a, **_k: pytest.fail("terminal leaf must not hold"))
fake_call, calls = _unknown_then_check_call(lambda *_: pytest.fail("no second dial"))
fake_call, calls = _unknown_then_check_call(_recover_model(monkeypatch))
monkeypatch.setattr(loop_mod, "call_llm_with_retry", fake_call)
monkeypatch.setenv("OUROBOROS_TASK_REVIEW_MODE", "off")
monkeypatch.delenv("USE_LOCAL_FALLBACK", raising=False)
registry = _configured_registry(tmp_path)
_start_leaf(tmp_path)
notes = []
_r, usage, trace = run_llm_loop(**_loop_kwargs(tmp_path, registry, notes))
result, usage, _trace = run_llm_loop(**_loop_kwargs(tmp_path, registry, notes))
assert calls["n"] == 1
assert usage.get("execution_status") == "infra_failed"
assert trace.get("forced_finalization", {}).get("source") == "provider_outcome_unknown_no_resend"
assert calls["n"] == 2 and result == "recovered"
assert usage["transport_recovery"]["old_outcome"] == "unknown"
assert _read_hold_events(tmp_path) == []
@ -219,31 +236,27 @@ def test_finalize_now_mid_hold_takes_no_call_terminal(tmp_path, monkeypatch, _qu
def test_generic_task_and_multi_run_never_hold(tmp_path, monkeypatch, _quiet_probe):
monkeypatch.setattr(delegate_hold, "supervised_wait",
lambda *_a, **_k: pytest.fail("ineligible shapes must not hold"))
fake_call, calls = _unknown_then_check_call(lambda *_: pytest.fail("no second dial"))
fake_call, calls = _unknown_then_check_call(_recover_model(monkeypatch))
monkeypatch.setattr(loop_mod, "call_llm_with_retry", fake_call)
monkeypatch.setenv("OUROBOROS_TASK_REVIEW_MODE", "off")
monkeypatch.delenv("USE_LOCAL_FALLBACK", raising=False)
from ouroboros import loop_transport
monkeypatch.setattr(loop_transport, "interruptible_wait_sleep", lambda *args: False)
def exhausted_network_wait(episode, *, tools, **kwargs):
tools._ctx.task_metadata = {"deadline_at": "2000-01-01T00:00:00Z"}
return False
monkeypatch.setattr(loop_mod, "_continue_unknown_transport", exhausted_network_wait)
# Generic tasks use the network owner, never the single-leaf nanny hold.
registry = ToolRegistry(repo_dir=tmp_path, drive_root=tmp_path)
registry._ctx.task_id = "t-generic"
_start_leaf(tmp_path, task_id="t-generic", run_id="run-g")
_r, usage, _t = run_llm_loop(**_loop_kwargs(tmp_path, registry, []))
assert calls["n"] == 1 and usage.get("execution_status") == "infra_failed"
result, usage, _t = run_llm_loop(**_loop_kwargs(tmp_path, registry, []))
assert calls["n"] == 2 and result == "recovered"
assert usage["transport_recovery"]["old_outcome"] == "unknown"
# Configured but TWO live leaves.
calls["n"] = 0
registry2 = _configured_registry(tmp_path, task_id="t-multi")
_start_leaf(tmp_path, task_id="t-multi", run_id="run-m1")
_start_leaf(tmp_path, task_id="t-multi", run_id="run-m2")
_r, usage2, _t = run_llm_loop(**_loop_kwargs(tmp_path, registry2, []))
assert calls["n"] == 1 and usage2.get("execution_status") == "infra_failed"
result, usage2, _t = run_llm_loop(**_loop_kwargs(tmp_path, registry2, []))
assert calls["n"] == 2 and result == "recovered"
assert usage2["transport_recovery"]["old_outcome"] == "unknown"
assert _read_hold_events(tmp_path) == []
@ -315,7 +328,7 @@ def test_repeated_unknown_reholds_with_backoff_floor(tmp_path, monkeypatch, _qui
def test_refused_probe_and_state_less_payload_never_hold(tmp_path, monkeypatch, _quiet_probe):
"""A daemon refusal or a state-less payload is not evidence of a live leaf
(grok #1/#2, fable F3): the probe fails closed to today's terminal."""
and does not prohibit recovery of the supervising model."""
import ouroboros.delegate_progress as progress_mod
monkeypatch.setattr(delegate_hold, "supervised_wait",
@ -327,19 +340,21 @@ def test_refused_probe_and_state_less_payload_never_hold(tmp_path, monkeypatch,
raise RuntimeError("daemon unreachable")
monkeypatch.setattr(progress_mod, "bounded_poll", raising_poll)
fake_call, calls = _unknown_then_check_call(lambda *_: pytest.fail("no second dial"))
fake_call, calls = _unknown_then_check_call(_recover_model(monkeypatch))
monkeypatch.setattr(loop_mod, "call_llm_with_retry", fake_call)
registry = _configured_registry(tmp_path, task_id="t-refused")
_start_leaf(tmp_path, task_id="t-refused", run_id="run-r1")
_r, usage, _t = run_llm_loop(**_loop_kwargs(tmp_path, registry, []))
assert calls["n"] == 1 and usage.get("execution_status") == "infra_failed"
result, usage, _t = run_llm_loop(**_loop_kwargs(tmp_path, registry, []))
assert calls["n"] == 2 and result == "recovered"
assert usage["transport_recovery"]["old_outcome"] == "unknown"
monkeypatch.setattr(progress_mod, "bounded_poll", lambda _gw, _run, _sec, **_k: {})
calls["n"] = 0
registry2 = _configured_registry(tmp_path, task_id="t-stateless")
_start_leaf(tmp_path, task_id="t-stateless", run_id="run-r2")
_r, usage2, _t = run_llm_loop(**_loop_kwargs(tmp_path, registry2, []))
assert calls["n"] == 1 and usage2.get("execution_status") == "infra_failed"
result, usage2, _t = run_llm_loop(**_loop_kwargs(tmp_path, registry2, []))
assert calls["n"] == 2 and result == "recovered"
assert usage2["transport_recovery"]["old_outcome"] == "unknown"
assert _read_hold_events(tmp_path) == []

View file

@ -127,7 +127,7 @@ class Gateway:
self.results = results or [result()]
self.dispatch = dispatch or ["response_received"] * len(self.results)
self.uploads = []
self.operations = {}
self.accepted_operations = {}
self.creates = []
self.reads = []
self.acks = []
@ -138,17 +138,25 @@ class Gateway:
self.pending = False
self.read_error = False
self.raw_result = None
self.operation_catalog = []
self.catalog_reads = 0
self.capture_requests = []
def operations(self):
self.catalog_reads += 1
return deepcopy(self.operation_catalog)
def upload_model_request(self, payload, *, idempotency_key):
self.uploads.append((deepcopy(payload), idempotency_key))
return REF
def create_model_operation(self, ref, *, idempotency_key):
def create_model_operation(self, ref, *, idempotency_key, **options):
assert ref == REF
self.capture_requests.append(deepcopy(options))
self.creates.append(idempotency_key)
if idempotency_key not in self.operations:
self.operations[idempotency_key] = len(self.operations)
index = self.operations[idempotency_key]
if idempotency_key not in self.accepted_operations:
self.accepted_operations[idempotency_key] = len(self.accepted_operations)
index = self.accepted_operations[idempotency_key]
if self.lose_create:
self.lose_create = False
raise ClaudexorUnavailable("daemon_unreachable", "lost create reply")
@ -162,7 +170,8 @@ class Gateway:
def detail(self, index):
value = self.results[index]
return {"id": f"op-{index}", "state": "running" if self.pending else "succeeded" if value["outcome"] == "completed" else "failed",
succeeded = value["outcome"] == "completed" and value["message"] is not None
return {"id": f"op-{index}", "state": "running" if self.pending else "succeeded" if succeeded else "failed",
"dispatch": {"state": "started" if self.pending else self.dispatch[index], "route": value["route"]},
"response": {"state": "absent"} if self.pending else {"state": "ready", "ref": REF},
"problem": value["problem"]}
@ -307,7 +316,7 @@ def test_lost_create_reply_rejoins_without_second_physical_attempt(setup):
answer, usage = client.chat([{"role": "user", "content": "hi"}], MODEL)
assert answer["content"] == "Ответ 🐍"
assert len(gateway.creates) == 2 and gateway.creates[0] == gateway.creates[1]
assert len(gateway.operations) == 1 and len(usage["ledger_attempt_ids"]) == 1
assert len(gateway.accepted_operations) == 1 and len(usage["ledger_attempt_ids"]) == 1
assert [row["state"] for row in ledger(root)] == ["reserved", "dispatched", "settled"]
@ -393,7 +402,7 @@ def test_control_connect_failure_after_acceptance_stays_unknown(setup):
error = raised.value
assert error.code == "model_outcome_unknown" and error.operation_id == "op-0"
assert not is_pre_dispatch_transport_failure(error) and not is_retryable_transport_death(error)
assert ledger(root)[-1]["state"] == "unresolved" and len(gateway.operations) == 1
assert ledger(root)[-1]["state"] == "unresolved" and len(gateway.accepted_operations) == 1
assert not gateway.acks
@ -408,7 +417,7 @@ def test_proven_never_started_quota_attempts_do_not_spend_generation_limit(setup
client.chat([], MODEL)
with pytest.raises(ua.PhysicalAttemptLimitExceeded):
client.chat([], MODEL)
assert len(gateway.operations) == 4
assert len(gateway.accepted_operations) == 4
assert [r['state'] for r in ledger(root)].count('settled') == 1
@ -420,7 +429,7 @@ def test_unknown_outcome_keeps_its_generation_limit_claim(setup):
client.chat([], MODEL)
with pytest.raises(ua.PhysicalAttemptLimitExceeded):
client.chat([], MODEL)
assert len(gateway.operations) == 1
assert len(gateway.accepted_operations) == 1
def test_confirmed_provider_failure_settles_real_usage_before_raising(setup):
@ -466,7 +475,7 @@ def test_field_refusal_retains_then_acknowledges_once_with_display(setup, asynch
assert f"provider_code={vendor}, parameter=input" in error.display_message[:220]
assert error.physical_attempt_capture.state == "settled"
assert error.usage["claudexor"]["result_custody"]["state"] == "acknowledged"
assert len(gateway.operations) == len(gateway.creates) == len(gateway.acks) == 1
assert len(gateway.accepted_operations) == len(gateway.creates) == len(gateway.acks) == 1
assert [row["state"] for row in ledger(root)] == ["reserved", "dispatched", "settled"]
@ -484,7 +493,7 @@ def test_proven_not_started_releases_and_never_fabricates_provider_usage(setup,
assert raised.value.physical_attempt_capture.provider_error_type == (vendor or "ClaudexorModelNotDispatched")
assert gateway.uploads[0][0]["options"]["temperature"] == 0.2
assert [row["state"] for row in ledger(root)] == ["reserved", "dispatched", "released"]
assert len(gateway.operations) == 1
assert len(gateway.accepted_operations) == 1
def test_typed_subject_refusal_suppresses_next_auto_preference(setup):
@ -519,7 +528,7 @@ def test_native_reset_requires_actual_account_change_and_keeps_canonical_tools(s
if not set(change) & {"credentialProfileId", "accountFingerprint"}:
with pytest.raises(transport.ClaudexorModelNotDispatched):
client.chat(messages, MODEL)
assert len(gateway.operations) == 1
assert len(gateway.accepted_operations) == 1
else:
_, usage = client.chat(messages, MODEL)
assert len(usage["ledger_attempt_ids"]) == 2
@ -540,7 +549,7 @@ def test_ack_failure_preserves_paid_result_and_does_not_repeat(setup, error):
answer, usage = client.chat([{"role": "user", "content": "hi"}], MODEL)
assert answer == result()["message"] and retained(root) == result()
assert usage["claudexor"]["result_custody"]["state"] == "pending"
assert len(gateway.operations) == 1 and ledger(root)[-1]["state"] == "settled"
assert len(gateway.accepted_operations) == 1 and ledger(root)[-1]["state"] == "settled"
def test_failed_local_result_retention_withholds_ack_but_keeps_answer(setup, monkeypatch):
@ -586,7 +595,7 @@ def test_async_tools_and_capture_remain_in_callers_context(setup):
assert answer == result()["message"]
asyncio.run(run())
assert len(gateway.operations) == 1 and ledger(root)[-1]["state"] == "settled"
assert len(gateway.accepted_operations) == 1 and ledger(root)[-1]["state"] == "settled"
def test_gigachat_async_tools_still_refuse_before_provider_io(setup, monkeypatch):
@ -685,7 +694,7 @@ def test_observer_failure_does_not_lose_response_or_repeat_generation(setup):
answer, _ = client.chat([], MODEL, model_operation_observer=failed)
assert answer == result()["message"] and retained(root) == result()
assert len(gateway.operations) == 1 and ledger(root)[-1]["state"] == "settled"
assert len(gateway.accepted_operations) == 1 and ledger(root)[-1]["state"] == "settled"
def test_async_local_switch_projects_its_new_capture_to_the_caller(setup, monkeypatch):
@ -758,7 +767,7 @@ def test_unknown_engine_outcome_retains_response_without_resend_or_false_zero(se
client.chat([{"role": "user", "content": "hi"}], MODEL)
assert raised.value.code == "model_outcome_unknown" and not raised.value.type
assert retained(root) == gateway.results[0] and not gateway.acks
assert len(gateway.operations) == 1 and ledger(root)[-1]["state"] == "unresolved"
assert len(gateway.accepted_operations) == 1 and ledger(root)[-1]["state"] == "unresolved"
assert ledger(root)[-1].get("cost_usd") is None
@ -770,7 +779,7 @@ def test_continuation_repair_is_bounded_to_one_unstarted_operation(setup):
gateway.dispatch = ["not_started", "not_started"]
with pytest.raises(transport.ClaudexorModelNotDispatched):
client.chat([result()["message"]], MODEL)
assert len(gateway.operations) == 2
assert len(gateway.accepted_operations) == 2
assert [row["state"] for row in ledger(root)] == ["reserved", "dispatched", "released"] * 2
@ -783,7 +792,7 @@ def test_unreleased_attempt_cannot_authorize_continuation_repair(setup, monkeypa
with pytest.raises(transport.ClaudexorModelNotDispatched) as raised:
client.chat([result()["message"]], MODEL)
assert raised.value.physical_attempt_capture.state == "unresolved"
assert len(gateway.operations) == 1 and ledger(root)[-1]["state"] == "unresolved"
assert len(gateway.accepted_operations) == 1 and ledger(root)[-1]["state"] == "unresolved"
def test_model_switch_strips_native_envelope_but_not_tool_results():
@ -802,7 +811,7 @@ def test_caller_control_interrupts_pending_operation_without_false_success(setup
gateway.pending = True
with pytest.raises(transport.ClaudexorModelError) as raised:
client.chat([{"role": "user", "content": "hi"}], MODEL,
model_poll_control=lambda: "deadline_exceeded" if gateway.operations else None)
model_poll_control=lambda: "deadline_exceeded" if gateway.accepted_operations else None)
assert raised.value.control_reason == "deadline_exceeded"
assert gateway.cancels == [("op-0", "host_cancelled")]
assert ledger(root)[-1]["state"] == "unresolved" and not gateway.acks
@ -873,7 +882,7 @@ def test_missing_old_account_does_not_authorize_native_reset(setup):
gateway.dispatch = ["not_started"]
with pytest.raises(transport.ClaudexorModelNotDispatched):
client.chat([message], MODEL)
assert len(gateway.operations) == 1
assert len(gateway.accepted_operations) == 1
def test_caller_control_before_create_proves_no_dispatch(setup):

View file

@ -460,7 +460,7 @@ def _install_model_operation_fake(stack: contextlib.ExitStack, recorder: _Record
class ModelGateway(Gateway):
def create_model_operation(self, ref, *, idempotency_key):
if idempotency_key not in self.operations:
if idempotency_key not in self.accepted_operations:
payload = self.uploads[-1][0]
step = recorder.record("claudexor.model_operation", payload=payload)
if step.get("kind") == "error" and step.get("code") != "provider_policy_refusal":
@ -469,7 +469,7 @@ def _install_model_operation_fake(stack: contextlib.ExitStack, recorder: _Record
"code": step.get("code"), "message": step.get("message")}, route={
"source": payload["source"], "model": payload["model"],
"credentialProfileId": "conformance-profile", "accountFingerprint": "conformance-identity"}))
index = len(self.operations)
index = len(self.accepted_operations)
if index == 0:
self.results[0], self.dispatch[0] = value, "response_received"
else:
@ -612,7 +612,7 @@ def _observe(spec: Dict[str, Any]) -> Dict[str, Any]:
client = LLMClient(**client_args)
call = copy.deepcopy(spec["call"])
if spec.get("cancel_model_after_create"):
call["kwargs"]["model_poll_control"] = lambda: "cancelled" if recorder.model_gateway.operations else None
call["kwargs"]["model_poll_control"] = lambda: "cancelled" if recorder.model_gateway.accepted_operations else None
try:
result = _call_route(client, call)
except BaseException as exc: # noqa: BLE001 - the raise IS the projection
@ -637,7 +637,7 @@ def _observe(spec: Dict[str, Any]) -> Dict[str, Any]:
observed["unused_script_steps"] = len(recorder.script)
if spec.get("model_operation"):
gateway = recorder.model_gateway
observed["model_control"] = {"create_posts": len(gateway.creates), "operations": len(gateway.operations),
observed["model_control"] = {"create_posts": len(gateway.creates), "operations": len(gateway.accepted_operations),
"unique_create_keys": len(set(gateway.creates)), "cancels": gateway.cancels}
return _jsonable(observed)

View file

@ -325,7 +325,7 @@ def test_async_cancellation_resolves_only_its_wait_and_keeps_shared_task(live_wa
assert terminal["resolution"] == "caller_cancelled"
asyncio.run(run())
assert not controller.closed and len(transport.operations) == 1
assert not controller.closed and len(transport.accepted_operations) == 1
assert all(row["state"] == "resolved" for row in load_task_result(root, "task-one")["model_waits"].values())
@ -623,7 +623,7 @@ def test_unproved_pool_cause_never_enters_resource_wait(live_wait, pool_context)
transport.dispatch = ["not_started"]
with pytest.raises(ClaudexorModelError, match="credential_pool_exhausted"):
client.chat([{"role": "user", "content": "Do not infer quota"}], MODEL, model_role="main")
assert not controller.waits and events.empty() and len(transport.operations) == 1
assert not controller.waits and events.empty() and len(transport.accepted_operations) == 1
def test_mixed_pool_cannot_bypass_unknown_physical_custody(live_wait):
@ -636,7 +636,7 @@ def test_mixed_pool_cannot_bypass_unknown_physical_custody(live_wait):
transport.dispatch = ["unknown"]
with pytest.raises(ClaudexorModelError, match="model_outcome_unknown"):
client.chat([{"role": "user", "content": "No duplicate generation"}], MODEL, model_role="main")
assert not controller.waits and events.empty() and len(transport.operations) == 1
assert not controller.waits and events.empty() and len(transport.accepted_operations) == 1
def test_mixed_pool_wait_keeps_calendar_deadline_and_existing_quota_union(live_wait, monkeypatch):
@ -691,14 +691,14 @@ def test_call_can_decline_resource_wait_without_losing_task_binding(live_wait, c
assert error.physical_attempt_capture.state == "released"
assert error.model_role_route == {"role": "light", "model": MODEL, "use_local": False,
"credential_profile_id": ""}
assert not controller.waits and events.empty() and len(transport.operations) == 1
assert not controller.waits and events.empty() and len(transport.accepted_operations) == 1
assert transport.uploads[0][0]["account"] == {"mode": "auto"}
assert "wait_for_resources" not in json.dumps(transport.uploads[0][0])
assert model_wait.current_model_wait() is controller and not controller.closed
assert [row["state"] for row in ledger(root)] == ["reserved", "dispatched", "released"]
answer, usage = call()
assert answer == result()["message"] and len(transport.operations) == 3
assert answer == result()["message"] and len(transport.accepted_operations) == 3
assert len(usage["ledger_attempt_ids"]) == 2
assert any(row["state"] == "waiting" for row in list(events.queue))
assert controller.overrides["light"]["model_account_override"] == ""
@ -716,4 +716,4 @@ def test_declining_resource_wait_still_honors_owner_control(live_wait, monkeypat
else:
client.chat(*args, **kwargs)
assert caught.value.control_reason == "finalize_requested"
assert not transport.operations and not controller.waits and events.empty()
assert not transport.accepted_operations and not controller.waits and events.empty()

View file

@ -121,14 +121,14 @@ def test_pinned_wait_rejects_other_catalog_account_before_new_generation(live_wa
polls = []
def catalog(source, profile=None, **kwargs):
assert profile == "account-a" and len(gateway.operations) == 1
assert profile == "account-a" and len(gateway.accepted_operations) == 1
polls.append(profile)
return {"source": source, "credentialProfileId": "account-b" if len(polls) == 1 else "account-a",
"models": [{"id": "exact-model"}]}
monkeypatch.setattr(client, "claudexor_model_catalog", catalog)
client.chat([], MODEL, model_role="light")
assert len(polls) == 2 and len(gateway.operations) == 2
assert len(polls) == 2 and len(gateway.accepted_operations) == 2
assert all(payload["account"] == {"mode": "pin", "profileId": "account-a"} for payload, _key in gateway.uploads)
@ -170,7 +170,7 @@ def test_main_wait_does_not_call_configured_api_fallback_before_owner_switch(mai
before_dispatch=_candidate_before_dispatch(request_body, request))
def catalog(*args, **kwargs):
assert api_calls == [] and len(gateway.operations) == 1
assert api_calls == [] and len(gateway.accepted_operations) == 1
row = next(event for event in reversed(list(events.queue)) if event.get("type") == "task_model_wait")
response = decide({"request_id": "switch-api", "decision_id": f"model_wait:task-one:{row['wait_id']}",
"revision": row["revision"], "action": "switch", "model": "openai::alternate",
@ -184,7 +184,7 @@ def test_main_wait_does_not_call_configured_api_fallback_before_owner_switch(mai
ctx.messages, tools, ctx.llm, ctx.drive_logs, lambda *_args, **_kwargs: None, queue.Queue(),
task_id="task-one", drive_root=ctx.drive_root, event_queue=events)
assert text == "Finished" and api_calls == ["openai"]
assert usage["_model_route"] == {} and len(gateway.operations) == 1
assert usage["_model_route"] == {} and len(gateway.accepted_operations) == 1
assert any(message.get("content") == "verified read A" for message in ctx.messages)
assert any(message.get("content") == "completed review B" for message in ctx.messages)
@ -437,7 +437,7 @@ def test_real_main_control_preserves_candidate_without_new_summary(main_call, mo
task_id="task-one", drive_root=ctx.drive_root, event_queue=events)
assert len(held) == 1 and text == completed["message"]["content"]
assert usage["reason_code"] == trace["forced_finalization"]["reason_code"] == expected_reason
assert len(gateway.operations) == 2 # Paid answer + interrupted call, never a summary retry.
assert len(gateway.accepted_operations) == 2 # Paid answer + interrupted call, never a summary retry.
assert trace["forced_finalization"]["source"].startswith("model_wait_retained_candidate")
if stop == "wrap_unknown":
assert usage["_last_llm_error_kind"] == "provider_outcome_unknown"
@ -510,7 +510,7 @@ def test_hard_cancel_returns_empty_events_to_real_worker_loop_and_keeps_queue_ow
return None # End the test's worker only after verifying retained ownership.
worker_process.worker_main(1, Input(), events, str(ctx.drive_root), str(ctx.drive_root))
assert len(reads) == 2 and len(gateway.operations) == 1 and not crashes
assert len(reads) == 2 and len(gateway.accepted_operations) == 1 and not crashes
assert ledger(ctx.drive_root)[-1]["state"] == "unresolved"

View file

@ -66,7 +66,7 @@ def test_response_received_with_proved_no_generation_reprepares_standard(setup,
_message, usage = asyncio.run(client.chat_async(**kwargs)) if asynchronous else client.chat(**kwargs)
first, second = [row[0] for row in gateway.uploads]
assert {**first, "options": {**first["options"], "processingPreference": "standard"}} == second
assert len(calls) == 1 and len(gateway.operations) == 2 and len(gateway.acks) == 2
assert len(calls) == 1 and len(gateway.accepted_operations) == 2 and len(gateway.acks) == 2
finals = list({row["attempt_id"]: row for row in ledger(root)}.values())
assert [row["state"] for row in finals] == ["released", "settled"]
assert finals[0]["candidate_raw_sha256"] != finals[1]["candidate_raw_sha256"]
@ -94,7 +94,7 @@ def test_only_explicit_no_start_advisory_proof_allows_retry(setup, axis):
cx.chat_claudexor(target, [], None, service_tier="flex")
else:
client.chat([], MODEL, **kwargs)
assert len(gateway.operations) == 1
assert len(gateway.accepted_operations) == 1
if axis.startswith("unknown"):
assert raised.value.code == "model_outcome_unknown" and ledger(root)[-1]["state"] == "unresolved"
if axis == "exact_native":
@ -143,7 +143,7 @@ def test_one_repair_on_each_axis_uses_the_same_bounded_preparation_loop(setup, o
gateway.results = [value for value, _dispatch in pair] + [{**result(), "processing": receipt()}]
gateway.dispatch = [dispatch for _value, dispatch in pair] + ["response_received"]
client.chat([result()["message"]], MODEL, processing_preference="economy")
assert len(gateway.operations) == 3
assert len(gateway.accepted_operations) == 3
assert [row["state"] for row in {r["attempt_id"]: r for r in ledger(root)}.values()] == ["released", "released", "settled"]
@ -167,7 +167,7 @@ def test_standard_must_be_supported_before_emitting_a_processing_retry(setup, mo
gateway.results = [refusal()]
with pytest.raises(cx.ClaudexorModelNotDispatched):
client.chat([], MODEL, processing_preference="economy")
assert len(gateway.operations) == 1, "Do not retry by omitting unsupported Standard and inheriting native premium"
assert len(gateway.accepted_operations) == 1, "Do not retry by omitting unsupported Standard and inheriting native premium"
def test_old_facade_does_not_send_new_query_parameter(monkeypatch):

View file

@ -93,11 +93,11 @@ def test_native_wait_switch_rechecks_bound_before_send_and_never_replays_read(li
with pytest.raises(ReviewRouteUnavailable) as raised:
executor.execute()
assert raised.value.code == "native_transcript_cap_exceeded"
assert len(gateway.operations) == 2
assert len(gateway.accepted_operations) == 2
assert executor.failure_custody()["native_transcript_bound"] == 0
with pytest.raises(ReviewRouteUnavailable):
executor.execute()
assert len(gateway.operations) == 2
assert len(gateway.accepted_operations) == 2
else:
answer = executor.execute()
assert answer.raw_text == final["message"]["content"]
@ -106,7 +106,7 @@ def test_native_wait_switch_rechecks_bound_before_send_and_never_replays_read(li
sent = gateway.uploads[-1][0]
assert sent["account"] == {"mode": "pin", "profileId": "account-b"}
assert any(message.get("role") == "tool" and "completed original read" in message["content"] for message in sent["messages"])
assert answer.usage["native_rounds"] == 2 and len(gateway.operations) == 3
assert answer.usage["native_rounds"] == 2 and len(gateway.accepted_operations) == 3
assert len(reads) == len(decisions) == 1
assert controller.overrides.keys() == {"reviewer:critic"}
assert [row["state"] for row in ledger(root)].count("settled") == (1 if narrow else 2)

View file

@ -63,7 +63,7 @@ def test_raw_default_hint_defers_but_explicit_temperature_stays_strict(setup, mo
assert raised.value.code == "unsupported_parameter"
assert gateway.uploads[0][0]["options"]["temperature"] == explicit
assert ledger(root)[-1]["state"] == "released"
assert len(gateway.operations) == 1
assert len(gateway.accepted_operations) == 1
assert "default_temperature" not in gateway.uploads[0][0]["options"]
assert kwargs["temperature"] is explicit and kwargs["default_temperature"] == 0.2
@ -132,7 +132,7 @@ def test_wait_switch_restores_api_hint_and_later_override_defers_again(live_wait
# A reused caller originally naming API must not carry its resolved 0.2 as explicit.
kwargs["model"] = "anthropic::other-model"
_message, usage = call()
assert usage["prompt_tokens"] == 20 and len(gateway.operations) == 2
assert usage["prompt_tokens"] == 20 and len(gateway.accepted_operations) == 2
assert all("temperature" not in payload["options"] for payload, _key in gateway.uploads)
assert [row["state"] for row in ledger(root)].count("released") == 1
@ -179,7 +179,7 @@ def test_actual_review_authors_reach_strict_raw_dispatch(setup, monkeypatch, sur
assert "temperature" not in payload["options"]
assert "default_temperature" not in payload["options"]
assert payload["account"] == {"mode": "pin", "profileId": "account-a"}
assert len(gateway.operations) == 1 and ledger(root)[-1]["state"] == "settled"
assert len(gateway.accepted_operations) == 1 and ledger(root)[-1]["state"] == "settled"
def test_review_custody_distinguishes_hint_from_explicit_without_changing_old_keys():
@ -212,4 +212,4 @@ def test_explicit_review_temperature_beats_both_host_hints(setup, monkeypatch, n
client.chat(**kwargs)
assert raised.value.code == "unsupported_parameter"
assert gateway.uploads[0][0]["options"]["temperature"] == expected
assert len(gateway.operations) == 1 and ledger(root)[-1]["state"] == "released"
assert len(gateway.accepted_operations) == 1 and ledger(root)[-1]["state"] == "released"

View file

@ -91,7 +91,7 @@ def test_native_account_repair_rebinds_real_physical_candidate_before_send(main_
assert dispatched[1]["physical_context"]["route_fp"] == "capacity-account-b"
assert dispatched[1]["physical_context"]["capacity_total_tokens"] == 240_000
assert observations[0]["model_route"] == ROUTE_B
assert len(gateway.operations) == 2 and gateway.creates[0] != gateway.creates[1]
assert len(gateway.accepted_operations) == 2 and gateway.creates[0] != gateway.creates[1]
resent = gateway.uploads[1][0]["messages"]
assert "nativeContinuation" not in resent[2]
assert resent[2]["tool_calls"] == original[2]["tool_calls"] and resent[3:] == original[3:]
@ -106,7 +106,7 @@ def test_quota_auto_wait_rejoins_same_round_then_repairs_changed_account(main_ca
gateway.dispatch = ["not_started", "not_started", "response_received"]
answer, cost, mode = _dispatch(ctx)
assert answer and cost is None and mode == "max"
assert len(gateway.operations) == 3
assert len(gateway.accepted_operations) == 3
assert ctx.accumulated_usage["rounds"] == 1
assert ctx.accumulated_usage["_model_route"] == ROUTE_B
assert ctx.accumulated_usage["_context_route_fp"] == "capacity-account-b"
@ -178,14 +178,14 @@ def test_manual_switch_updates_only_waiting_role_and_continues_current_main_call
def test_main_control_interrupt_is_typed_no_retry_and_keeps_operation_custody(main_call, monkeypatch):
ctx, gateway, controller, _events, _decide, _observations = main_call
gateway.pending = True
monkeypatch.setattr(controller, "control_reason", lambda: "cancelled" if gateway.operations else None)
monkeypatch.setattr(controller, "control_reason", lambda: "cancelled" if gateway.accepted_operations else None)
with pytest.raises(model_wait.ModelWaitInterrupted) as raised:
_dispatch(ctx)
error = raised.value
assert error.control_reason == "cancelled" and error.operation_id == "op-0"
assert error.physical_attempt_capture.state == "unresolved"
assert error.model_role_route["role"] == "main"
assert len(gateway.operations) == 1 and gateway.cancels == [("op-0", "host_cancelled")]
assert len(gateway.accepted_operations) == 1 and gateway.cancels == [("op-0", "host_cancelled")]
assert not controller.waits and not ctx.accumulated_usage.get("_last_llm_retry_same_request")
@ -198,7 +198,7 @@ def test_cancel_after_result_keeps_settled_usage_and_exact_result(main_call, mon
assert raised.value.usage["prompt_tokens"] == 20
assert raised.value.physical_attempt_capture.state == "settled"
assert raised.value.route == ROUTE
assert len(gateway.operations) == 1
assert len(gateway.accepted_operations) == 1
def test_native_repair_without_main_callback_refuses_stale_physical_fit(setup):
@ -212,7 +212,7 @@ def test_native_repair_without_main_callback_refuses_stale_physical_fit(setup):
with ua.bind_physical_attempt_context(physical):
with pytest.raises(model_wait.ModelWaitInterrupted, match="model_wait_reprepare_required"):
client.chat([result()["message"]], MODEL, model_role="main")
assert len(gateway.operations) == 1
assert len(gateway.accepted_operations) == 1
def test_unknown_main_outcome_never_retries_or_waits_for_quota(main_call):
@ -220,7 +220,7 @@ def test_unknown_main_outcome_never_retries_or_waits_for_quota(main_call):
gateway.results, gateway.dispatch = [result(outcome="unknown")], ["unknown"]
answer, _cost, _mode = _dispatch(ctx)
assert answer is None and ctx.accumulated_usage["_last_llm_error_kind"] == "provider_outcome_unknown"
assert len(gateway.operations) == 1 and not controller.waits
assert len(gateway.accepted_operations) == 1 and not controller.waits
assert ledger(ctx.drive_root)[-1]["state"] == "unresolved"
@ -312,7 +312,7 @@ def test_a_same_route_wait_and_reprepare_update_the_original_slot(main_call, tur
assert _dispatch(ctx)[0]
# The slot the loop still owns is the one the durable result must have replaced.
assert ctx.tools._ctx.model_turn_state is slot and slot.envelope == TURN
assert len(gateway.operations) == 2 and gateway.uploads[-1][0]["nativeContinuation"] is None
assert len(gateway.accepted_operations) == 2 and gateway.uploads[-1][0]["nativeContinuation"] is None
def test_a_helper_call_cannot_overwrite_the_running_loop_slot(setup, turn_engine):
@ -535,7 +535,7 @@ def test_live_owner_wait_reprojects_affinity(main_call, monkeypatch, destination
"execution_id": ctx.accumulated_usage["execution_id"],
"owner_switch_saved": decisions[0]["saved"],
"completed_tool_texts": [x["content"] for x in ctx.messages if x.get("role") == "tool"],
"gateway_operations": len(gateway.operations),
"gateway_operations": len(gateway.accepted_operations),
}
assert facts["completed_tool_texts"] == ["verified read A", "completed review B"]
assert ctx.accumulated_usage["execution_id"] == CACHE_REPREPARE_EXECUTION
@ -568,7 +568,7 @@ def test_recorded_wait_override_reprojects_affinity_before_send(main_call, initi
# API/local have no subscription quota wait of their own.
controller.overrides["main"] = {"model": MODEL, "use_local": False, "model_account_override": ""}
answer, _, _ = _dispatch(ctx)
assert answer and len(gateway.operations) == 1
assert answer and len(gateway.accepted_operations) == 1
payload = gateway.uploads[0][0]
facts = {
"initial_model": initial_model,
@ -576,7 +576,7 @@ def test_recorded_wait_override_reprojects_affinity_before_send(main_call, initi
"active_model": ctx.active_model,
"cache_key": payload["options"].get("cacheKey"),
"execution_id": ctx.accumulated_usage["execution_id"],
"gateway_operations": len(gateway.operations),
"gateway_operations": len(gateway.accepted_operations),
"completed_tool_texts": [x["content"] for x in ctx.messages if x.get("role") == "tool"],
}
assert facts["completed_tool_texts"] == ["verified read A", "completed review B"]

View file

@ -74,7 +74,7 @@ def test_preset_images_reach_real_main_transport(subscription_transport, catalog
assert gateway.uploads[0][0]["account"] == {"mode": "pin", "profileId": "explicit-main"}
assert calls == [("codex", "explicit-main", "exact-model")]
assert messages == original
assert len(gateway.creates) == 1 and len(gateway.operations) == 1
assert len(gateway.creates) == 1 and len(gateway.accepted_operations) == 1
assert [row["state"] for row in ledger(root)] == ["reserved", "dispatched", "settled"]

View file

@ -1,5 +1,6 @@
"""A cooperative stop interrupts only the unsent paid transport-repeat grant."""
import json
import threading
import time
from concurrent.futures import ThreadPoolExecutor
@ -168,9 +169,9 @@ def test_paid_repeat_empty_peek_reuses_existing_wait_proof(tmp_path, monkeypatch
assert ctx._loop_mailbox_seen_ids == {"old"}
@pytest.mark.parametrize("with_leaf", [False, True])
def test_wrapup_reason_survives_the_live_delegate_hold(tmp_path, monkeypatch, with_leaf):
from ouroboros import claudexor_daemon, delegate_custody, delegate_progress
@pytest.mark.parametrize("with_leaf, during_hold", [(False, False), (True, False), (True, True)])
def test_wrapup_reason_survives_the_live_delegate_hold(tmp_path, monkeypatch, with_leaf, during_hold):
from ouroboros import claudexor_daemon, delegate_custody, delegate_progress, delegate_hold
from tests.test_delegate_hold import _configured_registry, _start_leaf, _loop_kwargs as hold_kwargs
monkeypatch.setenv("OUROBOROS_TASK_REVIEW_MODE", "off")
@ -181,12 +182,23 @@ def test_wrapup_reason_survives_the_live_delegate_hold(tmp_path, monkeypatch, wi
if with_leaf:
_start_leaf(tmp_path, task_id="t-death", run_id="fixture-leaf")
def death():
def request_wrapup():
intent = cancel_intents.request_cancel(tmp_path, "t-death", requested_stop_policy=cancel_intents.STOP_POLICY_FINALIZE)
owner_mailbox.write_owner_message(tmp_path, REASON_OWNER_REQUESTED_FINALIZATION, "t-death",
msg_id=owner_stop_control_id(intent), kind=owner_mailbox.KIND_FINALIZE_NOW)
def death():
if not during_hold:
request_wrapup()
return httpx.ReadError("controlled post-dispatch failure")
def hold(*args):
request_wrapup()
return json.dumps({"status": "progress", "wake_events": [{"kind": "finalize_now"}]})
if during_hold:
monkeypatch.setattr(delegate_hold, "supervised_wait", hold)
llm = _LedgerLLM(tmp_path, death)
kwargs = hold_kwargs(tmp_path, registry, [])
kwargs["llm"] = llm

View file

@ -226,14 +226,15 @@ def test_subscription_requires_fresh_typed_upstream_observation(monkeypatch, fac
assert bool(transport.upstream_transport_reachable(None, "claudexor::codex=test", timeout=3)) is (fact == "upstream")
def test_claudexor_control_loss_keeps_same_operation_past_read_window(tmp_path, monkeypatch):
@pytest.mark.parametrize("route_kind", ["", "agent_session"])
def test_claudexor_control_loss_keeps_same_operation_past_read_window(tmp_path, monkeypatch, route_kind):
from ouroboros import llm_claudexor
from ouroboros.gateways.claudexor import ClaudexorUnavailable
inv = llm_claudexor._ModelInvocation({"usage_model": "claudexor::codex=test"}, {}, {"timeout": 1})
inv.operation_id, inv.invocation_id, inv.task_id, inv.root = "same-op", "same-attempt", "t", tmp_path
inv.create_attempted = True
now, reads = [0.0], []
ctx = SimpleNamespace(task_id="t")
ctx = SimpleNamespace(task_id="t", _configured_subagent_route_kind=route_kind)
waiter = SimpleNamespace(tool_context=ctx, control_reason=lambda: None)
monkeypatch.setattr(llm_claudexor, "current_model_wait", lambda: waiter)
monkeypatch.setattr(llm_claudexor, "time", SimpleNamespace(monotonic=lambda: now[0], sleep=lambda t: now.__setitem__(0, now[0]+t)))
@ -389,7 +390,10 @@ def test_upstream_head_uses_connection_window_in_every_socket_phase(monkeypatch,
assert requests[0].extensions["timeout"] == dict(connect=bound, read=bound, write=bound, pool=bound)
def test_unknown_policy_keeps_configured_session_nanny_out_of_managed_continuation(tmp_path):
def test_unknown_policy_admits_configured_session_model_to_managed_continuation(tmp_path):
ctx = SimpleNamespace(task_id="t", exact_model_route=True, _configured_subagent_route_kind="agent_session")
from ouroboros import loop_transport
assert loop_transport.reconcile_transport_wait(None, ctx, msg_present=False, error_kind="provider_outcome_unknown", drive_logs=tmp_path, task_id="t", model="m", emit_progress=lambda *a, **kw: pytest.fail("unexpected automatic continuation")) is None
episode = transport.reconcile_transport_wait(None, ctx, msg_present=False,
error_kind="provider_outcome_unknown", drive_logs=tmp_path, task_id="t", model="m",
emit_progress=lambda *a, **kw: None)
assert episode is not None and episode.wait_cause == "provider_outcome_unknown"
assert not episode.interactive

View file

@ -34,6 +34,9 @@ def _install_child_fixture(mode, events_path):
polls = 0
lost = False
def operations(self):
return [] # Legacy model-operation catalog.
def upload_model_request(self, payload, *, idempotency_key):
event("upload", key=idempotency_key, payload=payload)
return REF