diff --git a/devtools/benchmarks/cybergym/cybergym_protocol.py b/devtools/benchmarks/cybergym/cybergym_protocol.py index 9daa4e915..b18ba0f2e 100644 --- a/devtools/benchmarks/cybergym/cybergym_protocol.py +++ b/devtools/benchmarks/cybergym/cybergym_protocol.py @@ -41,6 +41,7 @@ DEFAULT_DISABLED_TOOLS = ( "delegate_wait", "delegate_cancel", "delegate_answer", + "delegate_message", "claude_code_edit", "analyze_screenshot", "vlm_query", @@ -198,7 +199,7 @@ def derive_disabled_tools(extra: Iterable[str] = ()) -> tuple[str, ...]: "analyze_screenshot", "vlm_query", "view_image", "ocr_pdf", "extract_video_frames", "send_photo", "send_video", "switch_model", "schedule_subagent", "delegate_start", "delegate_wait", "delegate_cancel", - "delegate_answer", "claude_code_edit", "wait_task", "wait_tasks", + "delegate_answer", "delegate_message", "claude_code_edit", "wait_task", "wait_tasks", "get_task_result", "peek_task", "cancel_task", "discard_child_result", "task_acceptance_review", "request_deep_self_review", } diff --git a/ouroboros/delegate_interactions.py b/ouroboros/delegate_interactions.py index 8ec113126..a965b608f 100644 --- a/ouroboros/delegate_interactions.py +++ b/ouroboros/delegate_interactions.py @@ -1,4 +1,5 @@ -"""A delegated run's interactive questions, and the nanny's answer to them. +"""Nanny input into a delegated run's LIVE session: the answer to its question, +and a live message into its running turn. Extracted from ``ouroboros/tools/delegate.py`` when that module crossed its size gate — the same split as ``delegate_output`` / ``delegate_progress`` / @@ -7,6 +8,9 @@ AskUserQuestion-style interaction (full question text/options on the run detail; a row carrying ``timeout_at`` benign-declines at the engine timeout, a null one waits until answered), and the nanny — the task that owns the run — is the party that answers (owner decision 7=A, poltergeist phase B). +``_delegate_message`` is the same concern one step earlier: a message placed +into the run's live turn through the engine's capability-declared channel +(``liveInput`` on the route's catalog row), typed end to end and never guessed. ``tools.delegate`` re-exports these names, so every existing reference (and the tests) still finds them there, and ``_REPORTED_INTERACTIONS`` stays one object. """ @@ -17,10 +21,19 @@ import hashlib import json import logging import time +import uuid from typing import Any, Dict, List, Optional, Tuple from ouroboros.delegate_output import _PAYLOAD_ENVELOPE_HEADROOM, _stage_full_output -from ouroboros.delegate_shared import AGENT_FAULT_CODE, SUBSTRATE_REFUSAL_CODE, _emit, _fail, _owned_run, delegate_result +from ouroboros.delegate_shared import ( + AGENT_FAULT_CODE, + SUBSTRATE_REFUSAL_CODE, + _emit, + _fail, + _owned_run, + delegate_result, + refusal_host_code, +) from ouroboros.tool_capabilities import tool_result_limit from ouroboros.tools.registry import ToolContext from ouroboros.tools.tool_result import ToolResult @@ -637,3 +650,235 @@ def _delegate_answer( return _answer_delivery_unknown(gateway, rid, iid, exc, seconds_left=_left()) finally: gateway.close() + + +# -- delegate_message: a live message into the run's running turn ---------------- +# +# What each typed live-message outcome MEANS for the nanny, relayed verbatim like +# ``_ANSWER_NOTES``. Custody of ``message_id`` (the wire Idempotency-Key, minted +# here, returned in EVERY result): the engine REPLAYS a succeeded command's stored +# receipt under the same key, so every FINAL verdict (delivered / accepted / +# rejected / not_active / unsupported) is that invocation's answer forever — the +# SAME message_id is re-sent ONLY after ``delivery_unknown``; a new message after +# any final verdict needs a NEW id (omit message_id). Never a content-stable key. +_MESSAGE_NOTES = { + "delivered": "The harness CONSUMED this message inside the live turn (a correlated " + "native echo); obedience is unproved. Keep watching with delegate_wait " + "(the timeline's message.* rows carry messageId and outcome). This " + "message_id is spent: a further message needs a NEW one (omit message_id).", + "accepted": "The harness's acceptance boundary was observed; CONSUMPTION is unproved " + "until a message.delivered timeline row appears. Keep watching with " + "delegate_wait. This message_id is spent: a further message needs a NEW one.", + "rejected": "An explicit refusal of THIS submission (see reason/detail): the vendor " + "refused the steer on a still-active turn, the call itself was in error, " + "or the daemon could not persist admission. The message did NOT land. A " + "retry needs a NEW message_id — a replay of this one returns this verdict.", + "not_active": "No eligible live target existed before dispatch (no live attempt, a " + "terminal or settled run, a turn gap, an attempt mismatch, or a PENDING " + "question — answer that with delegate_answer). Nothing was written to the " + "harness. A later message needs a NEW message_id.", + "unsupported": "This route/run has no live-input channel (engine without the operation, " + "harness liveInput none, or a thread-bound run). Nothing was written. " + "Steer by cancel + a new delegate_start, or wait for the terminal.", + "delivery_unknown": "The message MAY have landed (transport loss, timeout, malformed " + "reply, receipt-save failure). Do NOT send a different message. " + "Re-check the timeline with delegate_wait (message.* rows carry " + "messageId and outcome); to retry, call delegate_message again " + "with the SAME message_id and the SAME text — the engine replays " + "the stored receipt instead of delivering twice.", + "not_found": "The daemon answered 404 for this run after advertising the operation " + "and the route's live-input capability: the run is unknown to it. " + "Custody is untouched; re-read the run with delegate_wait first.", +} + +# The 4xx problem bodies that are a PAYLOAD verdict about these message bytes +# (malformed, secret-bearing, too long): ``rejected`` as the agent fault it is. +# 409 is deliberately absent — on this route it is the idempotency store +# speaking (``idempotency_conflict`` = the same message_id with different text, +# a rejection; ``delivery_in_progress`` / ``delivery_interrupted`` = a delivery +# whose fate is unknown), read from the typed ``code`` FIRST. +_MESSAGE_PAYLOAD_VERDICT_CODES = frozenset({400, 413, 422}) +# The internal wall-clock budget for ONE delegate_message call, strictly below +# its ToolEntry timeout (120s): handshake, the two capability reads and the +# POST are budgeted against what remains (the codex steer itself is bounded at +# 30s inside the engine), and exhaustion returns a typed outcome instead of an +# executor thread-kill mid-wire. +_MESSAGE_DEADLINE_SEC = 100.0 + + +def _live_input_unsupported(gateway: Any, route_id: str) -> Tuple[str, str, str]: + """``(reason, detail, live_input)`` — reason empty when the route CAN take a + live message. Discovery is structural (A18): the engine's own route catalog + must list the operation AND the route's catalog row must declare a + ``liveInput`` other than ``none``; a read that fails is ``unsupported`` too + (no POST on a guess), never a refusal that spends a model round.""" + from ouroboros.gateways.claudexor import RUN_MESSAGE_OPERATION, run_message_supported + + try: + if not run_message_supported(gateway.operations()): + return ("engine_lacks_operation", + f"the engine's /v2/operations catalog does not list " + f"{' '.join(RUN_MESSAGE_OPERATION)}", "") + catalog = gateway.agent_capabilities() + row = next((item for item in (catalog.get("harnesses") or []) + if isinstance(item, dict) and str(item.get("id") or "") == route_id), None) + except Exception as exc: # noqa: BLE001 — a failed read is "unknown", answered typed + return ("capability_read_failed", f"{type(exc).__name__}: {exc}", "unknown") + if row is None: + return ("route_not_in_capability_catalog", + f"route {route_id!r} has no row in /v2/agent-capabilities", "") + live_input = str(row.get("liveInput") or "none") + if live_input == "none": + return ("route_live_input_none", + f"route {route_id!r} declares liveInput={live_input!r}", live_input) + return ("", "", live_input) + + +def _message_problem_outcome(exc: Exception) -> Tuple[str, str, Optional[str]]: + """``(outcome, reason, host_code)`` for an untyped refusal of the message POST. + + The typed ``code`` is read FIRST (A16): a 409 ``idempotency_conflict`` is a + definite rejection of this submission (same message_id, different text); the + other 409s, every 5xx, status 0 (transport death) and every non-verdict 4xx + (auth, rate, timeout) say nothing about delivery and are ``delivery_unknown``; + ANY 404 — every daemon 404 has a body — is the host's own ``not_found`` + (custody untouched: ``daemon_says_absent`` is never consulted here). + """ + status = int(getattr(exc, "status_code", 0) or 0) + code = str(getattr(exc, "code", "") or "") or f"http_{status}" + if status == 409 and code == "idempotency_conflict": + return "rejected", code, refusal_host_code(code) + if status == 409: + return "delivery_unknown", code, None + if status == 404: + return "not_found", code, SUBSTRATE_REFUSAL_CODE + if status in _MESSAGE_PAYLOAD_VERDICT_CODES: + return "rejected", code, AGENT_FAULT_CODE + return "delivery_unknown", code, None + + +def _message_result(ctx: ToolContext, facts: Dict[str, Any], *, outcome: str, + reason: str = "", http_status: int = 0, detail: str = "", + host_code: Optional[str] = None, **engine: Any) -> ToolResult: + """Record the receipt (``delegate_message_outcome``, digest and size, never + the text) and render the typed result. ``host_code`` marks a refusal; + ``delivery_unknown`` and the two positive outcomes are OK observations.""" + _emit(ctx, "delegate_message_outcome", { + **facts, "outcome": outcome, "reason": reason, "http_status": int(http_status), + "attempt_id": str(engine.get("attempt_id") or ""), + "harness_id": str(engine.get("harness_id") or ""), + }) + payload: Dict[str, Any] = { + "status": outcome, "run_id": facts["run_id"], "message_id": facts["message_id"], + "accepted": outcome in ("delivered", "accepted"), "reason": reason or None, + "attempt_id": str(engine.get("attempt_id") or "") or None, + "harness_id": str(engine.get("harness_id") or "") or None, + "live_input": str(engine.get("live_input") or "") or None, + "native_turn_id": str(engine.get("native_turn_id") or "") or None, + "detail": detail, "note": _MESSAGE_NOTES.get(outcome, ""), + } + if host_code: + payload.update({"ok": False, "host_code": host_code}) + return delegate_result(payload) + + +def _delegate_message(ctx: ToolContext, run_id: str, text: Any, + message_id: str = "") -> ToolResult: + """Place one live message into a delegated run's running turn. + + Custody-gated like answer/cancel: only the task that started the run may + speak into it. The outcome is TYPED end to end and mirrors the engine's + ``LiveMessageOutcome`` 1:1 (``delivered`` / ``accepted`` / ``rejected`` / + ``not_active`` / ``unsupported`` / ``delivery_unknown``) plus the host's own + ``not_found``; the engine's ``reason`` rides verbatim. Two host short-circuits + never POST a FRESH message: a settled run (``not_active``) and a route with + no live-input channel (``unsupported``, discovered structurally from the + operation catalog and the route row's ``liveInput`` — no harness-name + branch). The host mints ``message_id`` (= the wire Idempotency-Key) and + returns it; a call carrying a previously returned id SKIPS both + short-circuits and POSTs so the engine replays the stored receipt (A27) — + the recovery for ``delivery_unknown``, and ONLY for it. No retry loop, no + stall detector: the whole call runs under one internal deadline strictly + below its ToolEntry timeout, and no failure reaches the model as a + traceback. + """ + from ouroboros.delegate_progress import poll_bound + from ouroboros.gateways.claudexor import ClaudexorGateway, ClaudexorUnavailable + + rid = str(run_id or "").strip() + if not rid: + return _fail("delegate_message", "missing_run_id", "run_id is required") + body_text = text if isinstance(text, str) else "" + if not body_text.strip(): + return _fail("delegate_message", "message_text_required", + "text is required: the message to place into the run's live " + "session, a non-empty string (nothing is coerced).", run_id=rid) + not_mine, entry = _owned_run(ctx, "delegate_message", rid) + if not_mine or entry is None: + return not_mine or _fail("delegate_message", "run_ownership_unknown", + "custody unresolved", run_id=rid) + mid = str(message_id or "").strip() + replay = bool(mid) + if not replay: + mid = uuid.uuid4().hex + facts = {"run_id": rid, "message_id": mid, "text_chars": len(body_text), + "text_sha256": hashlib.sha256(body_text.encode("utf-8", "replace")).hexdigest()} + if not replay and entry.settled: + return _message_result( + ctx, facts, outcome="not_active", reason="run_settled", + host_code=SUBSTRATE_REFUSAL_CODE, + detail="custody records this run as settled; there is no live turn to steer") + deadline = time.monotonic() + _MESSAGE_DEADLINE_SEC + + def _left() -> float: + return deadline - time.monotonic() + + try: + gateway = ClaudexorGateway() + gateway.handshake(timeout_sec=poll_bound(min(_left(), _ANSWER_HANDSHAKE_MAX_SEC))) + except ClaudexorUnavailable as exc: + return _fail("delegate_message", exc.code, str(exc), run_id=rid, message_id=mid) + try: + live_input = "" + if not replay: + reason, detail, live_input = _live_input_unsupported(gateway, str(entry.route_id or "")) + if reason: + return _message_result(ctx, facts, outcome="unsupported", reason=reason, + host_code=SUBSTRATE_REFUSAL_CODE, detail=detail, + live_input=live_input) + if _left() <= 0: + # Spent before the POST: nothing was sent, and the same-id retry the + # delivery_unknown note prescribes is exactly right (the key is unused). + return _message_result( + ctx, facts, outcome="delivery_unknown", reason="deadline_exhausted", + detail=(f"local time budget ({_MESSAGE_DEADLINE_SEC:.0f}s) spent before " + "the message POST was sent; nothing was sent"), live_input=live_input) + try: + body = gateway.send_run_message(rid, body_text, idempotency_key=mid, + timeout_sec=poll_bound(_left())) + except ClaudexorUnavailable as exc: + outcome, reason, host_code = _message_problem_outcome(exc) + return _message_result( + ctx, facts, outcome=outcome, reason=reason, host_code=host_code, + http_status=int(getattr(exc, "status_code", 0) or 0), detail=str(exc), + live_input=live_input) + outcome = str(body.get("outcome") or "") + reason = str(body.get("reason") or "") + host_code = None + if outcome == "rejected": + host_code = refusal_host_code(reason) + elif outcome in ("not_active", "unsupported"): + host_code = SUBSTRATE_REFUSAL_CODE + return _message_result( + ctx, facts, outcome=outcome, reason=reason, http_status=200, + host_code=host_code, detail=str(body.get("message") or ""), + attempt_id=body.get("attemptId"), harness_id=body.get("harnessId"), + live_input=body.get("liveInput") or live_input, + native_turn_id=body.get("nativeTurnId")) + except Exception as exc: # noqa: BLE001 — F7: never a raw traceback to the model + log.warning("delegate_message failed untyped for %s/%s", rid, mid, exc_info=True) + return _message_result( + ctx, facts, outcome="delivery_unknown", reason="host_exception", + detail=f"{type(exc).__name__}: {exc}") + finally: + gateway.close() diff --git a/ouroboros/delegate_progress.py b/ouroboros/delegate_progress.py index 71e172e5a..647b734b2 100644 --- a/ouroboros/delegate_progress.py +++ b/ouroboros/delegate_progress.py @@ -68,7 +68,9 @@ def _bounded(rows: List[Dict[str, Any]]) -> List[Dict[str, Any]]: for row in rows[-_TIMELINE_TAIL:]: item = {"type": _label(row.get("type")), "title": _label(row.get("title")), "severity": _label(row.get("severity"))} - for key in ("attemptId", "harnessId"): + # ``messageId``/``outcome`` ride the engine's ``message.*`` rows: the + # receipts a recovered nanny reconciles a delegate_message against (A26). + for key in ("attemptId", "harnessId", "messageId", "outcome"): if isinstance(row.get(key), str) and row[key]: item[key] = _label(row[key]) if row.get("textKind") in _TEXT_KINDS and isinstance(row.get("detail"), str): diff --git a/ouroboros/delegate_shared.py b/ouroboros/delegate_shared.py index 5cf5391c8..21551bfe7 100644 --- a/ouroboros/delegate_shared.py +++ b/ouroboros/delegate_shared.py @@ -10,9 +10,9 @@ so every existing reference and monkeypatch target keeps the same objects. This module also owns the external-executor family's RESULT ENVELOPE. Inside the family (``delegate_start``/``delegate_wait``/``delegate_cancel``/ -``delegate_answer``, their producers and their host consumers) a result is a -native ``ToolResult``; only the four registered entries project it back to the -``str`` handler ABI. The envelope is two additive JSON keys — ``ok`` and +``delegate_answer``/``delegate_message``, their producers and their host +consumers) a result is a native ``ToolResult``; only the five registered entries +project it back to the ``str`` handler ABI. The envelope is two additive JSON keys — ``ok`` and ``host_code`` — written beside the domain payload, never instead of it: the domain ``reason`` keeps its own name and its own vocabulary, and nothing here renames it into ``ToolResult.code``. @@ -64,6 +64,10 @@ _AGENT_FAULT_REASONS = frozenset({ "configured_actor_resource_mismatch", "configured_actor_route_mismatch", "empty_prompt", + # The engine's typed rejection of a live message whose message_id was + # replayed with DIFFERENT text: the caller reused an invocation identity. + "idempotency_conflict", + "message_text_required", "missing_interaction_id", "missing_run_id", "payload_binding_mismatch", diff --git a/ouroboros/nanny_pacing.py b/ouroboros/nanny_pacing.py index 42b2383be..e52525bc9 100644 --- a/ouroboros/nanny_pacing.py +++ b/ouroboros/nanny_pacing.py @@ -7,6 +7,7 @@ from typing import Any, Dict, List, Tuple DELEGATE_ACTIVITY_TOOLS = frozenset({ "delegate_start", "delegate_wait", "delegate_cancel", "delegate_answer", + "delegate_message", }) # Only genuine ACTS of delegation reset the burn baseline: starting a physical diff --git a/ouroboros/tool_capabilities.py b/ouroboros/tool_capabilities.py index c04d91b44..2bf5e99f7 100644 --- a/ouroboros/tool_capabilities.py +++ b/ouroboros/tool_capabilities.py @@ -108,9 +108,11 @@ LOCAL_READONLY_SUBAGENT_TOOL_NAMES: frozenset[str] = frozenset({ "tree_note", "tree_read", "override_delegation_constraint", # Nanny verbs. The child gets no shell — it gets the right to ASK the host to run a # session, and the host derives the access profile from THIS task's authority, so a - # read-only child can only ever host a read-only session. delegate_answer speaks - # only to a run this task already owns (custody-gated like cancel). + # read-only child can only ever host a read-only session. delegate_answer and + # delegate_message speak only to a run this task already owns (custody-gated + # like cancel); a message is placed into the run's live turn, never a shell. "delegate_start", "delegate_wait", "delegate_cancel", "delegate_answer", + "delegate_message", "web_search", "browse_page", "browser_action", "analyze_screenshot", "vlm_query", "view_image", # Bounded media projection: writes derived frames only under artifact_store/video_frames. "ocr_pdf", "youtube_transcript", "extract_video_frames", @@ -153,6 +155,7 @@ ACTING_SUBAGENT_TOOL_NAMES: frozenset[str] = frozenset({ # workspace_write session confined to a private snapshot of its own write # root, and explicitly integrates the captured diff (C1). "delegate_start", "delegate_wait", "delegate_cancel", "delegate_answer", + "delegate_message", "integrate_delegated_patch", "web_search", "browse_page", "browser_action", "analyze_screenshot", "vlm_query", "view_image", "ocr_pdf", "youtube_transcript", "extract_video_frames", diff --git a/ouroboros/tools/delegate.py b/ouroboros/tools/delegate.py index e949385a6..6bd865ff0 100644 --- a/ouroboros/tools/delegate.py +++ b/ouroboros/tools/delegate.py @@ -6,10 +6,12 @@ API tokens it starts a Claudexor run, watches it, and brings the result home. Be the nanny IS the host, verification receipts stay host-authored and the harness's output is a claim, not proof. -Four verbs: ``delegate_start``, a time-bounded ``delegate_wait``, -``delegate_cancel``, and ``delegate_answer`` (a run's pending interactive question is -answered by its own nanny — owner decision 7=A, poltergeist phase B). There is still -no ``hurry`` — Claudexor's only control verb is ``cancel``, and cancelling a reviewer +Five verbs: ``delegate_start``, a time-bounded ``delegate_wait``, +``delegate_cancel``, ``delegate_answer`` (a run's pending interactive question is +answered by its own nanny — owner decision 7=A, poltergeist phase B) and +``delegate_message`` (a live message into the run's running turn, gated by the +route's declared ``liveInput`` capability, never by a harness name). There is still +no ``hurry`` — Claudexor's control verb is ``cancel``, and cancelling a reviewer destroys the verdict you wanted. Read-only and mutating children share ONE nanny and ONE transport. The only difference @@ -95,6 +97,7 @@ from ouroboros.delegate_interactions import ( # noqa: F401 _answer_delivery_unknown, _bounded_interactions, _delegate_answer, + _delegate_message, _interactions_are_news, _normalized_answers, _waiting_on_user_payload, @@ -1293,7 +1296,7 @@ def _delegate_cancel(ctx: ToolContext, run_id: str, reason: str = "") -> ToolRes def _published_entry(core: Any) -> Any: - """The family's four REGISTERED entries, wrapped in their one string boundary. + """The family's five REGISTERED entries, wrapped in their one string boundary. Inside the family a result is a native ``ToolResult``; the handler ABI is still ``str``. Publication happens HERE, after every decorator the core ran @@ -1444,7 +1447,8 @@ def get_tools() -> List[ToolEntry]: "continuation=new_physical_run. A " "large terminal result is delivered as a bounded preview plus an " "artifact: read output_delivery and finish reading the artifact before " - "you rely on it." + "you rely on it. A delegate_message receipt is reconciled HERE: the " + "timeline's message.* rows carry messageId and outcome." ), "parameters": {"type": "object", "required": ["run_id"], "properties": { "run_id": {"type": "string", "description": "Run id from delegate_start."}, @@ -1491,8 +1495,8 @@ def get_tools() -> List[ToolEntry]: "— the answer did NOT land; retry the SAME answers after reset_at); " "delivery_unknown " "(transport died mid-answer — re-check with delegate_wait and NEVER " - "post a different answer for the same interaction). Codex-lane runs " - "have no mid-run questions: a run that ENDS needing input " + "post a different answer for the same interaction). A run on a route " + "without a mid-run question channel that ENDS needing input " "(outcome_facts.reason=input_required) is answered with a plain NEW " "delegate_start(subagent_id=..., prompt=...) whose prompt carries the " "assignment plus the answers " @@ -1523,6 +1527,40 @@ def get_tools() -> List[ToolEntry]: "before delivering it."}, }}, }, _published_entry(_delegate_answer), timeout_sec=120), + ToolEntry("delegate_message", { + "name": "delegate_message", + "description": ( + "Place one live message into a delegated run's RUNNING turn (a " + "correction, a new fact, a redirection) without cancelling it. Only the " + "task that started the run may send. Capability-gated, never by harness " + "name: the route's catalog row declares liveInput (mid_turn / " + "next_tool_boundary / none) and the engine must list the operation; " + "otherwise the typed outcome is unsupported and nothing is sent. Typed " + "outcomes mirror the engine: delivered (the harness consumed it in the " + "live turn; obedience unproved); accepted (the acceptance boundary was " + "observed; consumption unproved until a message.delivered timeline row); " + "rejected (an explicit refusal of THIS submission — see reason); " + "not_active (no live target: terminal/settled run, turn gap, attempt " + "mismatch, or a PENDING question — answer that with delegate_answer); " + "unsupported; delivery_unknown (it MAY have landed); not_found. Every " + "result returns message_id, the delivery identity: pass it back ONLY to " + "retry the SAME text after delivery_unknown (the engine replays the " + "stored receipt instead of delivering twice); after any other outcome a " + "new message needs a NEW id (omit message_id). A message steers only the " + "current attempt — a later retry or continuation never re-injects it — and " + "is reconciled on the delegate_wait timeline (message.* rows)." + ), + "parameters": {"type": "object", "required": ["run_id", "text"], "properties": { + "run_id": {"type": "string", "description": "Run id from delegate_start."}, + "text": {"type": "string", "description": + "The message, verbatim, as the harness will read it mid-turn " + "(non-empty; the engine bounds its length)."}, + "message_id": {"type": "string", "description": + "ONLY the message_id a previous delivery_unknown result returned, to " + "replay that exact message under its original key. Omit for a new " + "message; never invent one."}, + }}, + }, _published_entry(_delegate_message), timeout_sec=120), ] diff --git a/tests/system_e2e/interfaces.py b/tests/system_e2e/interfaces.py index 47b3d1aad..34025049b 100644 --- a/tests/system_e2e/interfaces.py +++ b/tests/system_e2e/interfaces.py @@ -9,7 +9,10 @@ digest → 409 ``idempotency_conflict``), run detail with the ``summary`` facts custody settler consumes, the cancel control verb, and the interactive question surface — ``pendingInteractions`` on the detail plus the ``POST /v2/runs/:id/interactions/:iid/answer`` verb ``delegate_answer`` speaks, -with its typed delivered/already_resolved/rejected statuses. Behavior is scripted +with its typed delivered/already_resolved/rejected statuses — and the live-message +surface ``delegate_message`` negotiates: the ``/v2/operations`` catalog row for +``POST /v2/runs/:id/messages``, ``liveInput`` on the harness row, and the route +itself (Idempotency-Key replay, typed outcomes at HTTP 200). Behavior is scripted PER RUN by markers in the POSTed prompt (success / hang / typed refusal / ask) plus the pinned-profile refusal, and the applied facts a WRITING run produces: the edits themselves, made inside the private execution @@ -158,8 +161,12 @@ class FakeClaudexorDaemon: applied_profile: str = "fake-profile-1", ghost_profile: str = "ghost-profile", workspace_edits: Optional[Dict[str, str]] = None, - runs_dir: Optional[pathlib.Path] = None) -> None: + runs_dir: Optional[pathlib.Path] = None, + live_input: str = "mid_turn") -> None: self.harness_id = str(harness_id) + # The harness row's declared live-input capability (``liveInput``); a + # scenario passes "none" to script a route with no mid-run channel. + self.live_input = str(live_input) pin_version, pin_sha = _tree_engine_identity() self.engine_version = str(engine_version or pin_version) self.engine_build_sha = str(engine_build_sha or pin_sha) @@ -295,6 +302,15 @@ class FakeClaudexorDaemon: "sha": self.engine_build_sha}} if method == "GET" and clean == "/v2/agent-capabilities": return 200, {"harnesses": [self._harness_row()]} + if method == "GET" and clean == "/v2/operations": + # The engine's own route catalog (Express-style templates, the shape + # ``run_message_supported`` negotiates against). + return 200, {"protocolMajor": 3, "operations": [ + {"id": "run.message", "method": "POST", "path": "/v2/runs/:id/messages", + "mutability": "mutating", "idempotency": "key_required"}, + {"id": "run.control", "method": "POST", "path": "/v2/runs/:id/control", + "mutability": "mutating", "idempotency": "natural"}, + ]} if method == "GET" and clean == "/v2/harnesses": return 200, {"harnesses": [self._harness_row()]} if method == "GET" and clean == "/v2/quota": @@ -354,6 +370,8 @@ class FakeClaudexorDaemon: run["turn"] += 1 run["pending"] = [_fake_turn_interaction(run["id"], self.harness_id, run["turn"])] return 200, {"accepted": True, "status": "delivered"} + if method == "POST" and len(parts) == 4 and parts[3] == "messages": + return self._run_message(run, record) if method == "POST" and len(parts) == 4 and parts[3] == "control": control = body.get("control") if isinstance(body.get("control"), dict) else {} if str(control.get("kind") or "") == "cancel": @@ -367,7 +385,50 @@ class FakeClaudexorDaemon: def _harness_row(self) -> Dict[str, Any]: return {"id": self.harness_id, "enabled": True, "accessProfilesSupported": ["readonly", "workspace_write", - "external_sandbox_full"]} + "external_sandbox_full"], + "liveInput": self.live_input} + + def _run_message(self, run: Dict[str, Any], record: Dict[str, Any]) -> tuple: + """``POST /v2/runs/:id/messages``: the live-message route, every typed + outcome at HTTP 200 (the deliberate difference from the answer route), + served through the same Idempotency-Key replay as run creation — a + replayed key with the same digest returns the STORED receipt (never a + second delivery), a different digest is 409 ``idempotency_conflict``.""" + body = record["body"] + key = record["idempotency_key"] + if not key: + return 400, {"code": "missing_idempotency_key", + "message": "run message requires Idempotency-Key"} + digest = hashlib.sha256(json.dumps(body, sort_keys=True).encode("utf-8")).hexdigest() + replayed = self._replay.get(key) + if replayed is not None: + if replayed["digest"] != digest: + return 409, {"code": "idempotency_conflict", + "message": "Idempotency-Key replayed with a different request digest"} + return replayed["status"], json.loads(json.dumps(replayed["payload"])) + + def _remember(status: int, payload: Dict[str, Any]) -> tuple: + self._replay[key] = {"digest": digest, "status": status, "payload": payload} + return status, payload + + text = body.get("text") + if not isinstance(text, str) or not text: + return _remember(400, {"code": "invalid_request", "message": "text is required"}) + receipt: Dict[str, Any] = {"runId": run["id"], "messageId": key, + "harnessId": self.harness_id, "liveInput": self.live_input} + if run["state"] in ("succeeded", "cancelled", "failed"): + return _remember(200, {**receipt, "accepted": False, "outcome": "not_active", + "reason": "run_terminal"}) + if run["pending"]: + # INV-048: a steer beside an open question is never written to the harness. + return _remember(200, {**receipt, "accepted": False, "outcome": "not_active", + "reason": "interaction_pending", "attemptId": "a01"}) + if self.live_input == "none": + return _remember(200, {**receipt, "accepted": False, "outcome": "unsupported", + "reason": "no_live_session", "attemptId": "a01"}) + run.setdefault("messages", []).append({"message_id": key, "text": text}) + return _remember(200, {**receipt, "accepted": True, "outcome": "accepted", + "attemptId": "a01"}) def _start_run(self, record: Dict[str, Any]) -> tuple: body = record["body"] diff --git a/tests/test_builtin_refusal_results.py b/tests/test_builtin_refusal_results.py index f27674af5..cb0ae6f8c 100644 --- a/tests/test_builtin_refusal_results.py +++ b/tests/test_builtin_refusal_results.py @@ -144,9 +144,9 @@ def test_existing_warning_and_review_policy_are_not_blanket_reclassified(text, c # --- the external-executor family (owner Q8A) -------------------------------- # -# `delegate_start`/`delegate_wait`/`delegate_cancel`/`delegate_answer` speak a -# native `ToolResult` among themselves and project a `str` at their four -# registered entries. The incident these pin: `_fail` used to render +# `delegate_start`/`delegate_wait`/`delegate_cancel`/`delegate_answer`/ +# `delegate_message` speak a native `ToolResult` among themselves and project a +# `str` at their five registered entries. The incident these pin: `_fail` used to render # `{"status": "refused", ...}` as a plain string, which the registry's legacy # adapter classified as OK — so a refused wait/cancel (daemon unreachable, run # not owned, containment fault, refused cancel) was recorded as a SUCCESSFUL @@ -155,7 +155,7 @@ def test_existing_warning_and_review_policy_are_not_blanket_reclassified(text, c def _delegate_registry(tmp_path, monkeypatch, task_id="t-family"): - """A real registry, with the family's four entries registered as production does.""" + """A real registry, with the family's five entries registered as production does.""" from ouroboros.tools.registry import ToolRegistry import ouroboros.safety as safety @@ -163,8 +163,8 @@ def _delegate_registry(tmp_path, monkeypatch, task_id="t-family"): registry = ToolRegistry(repo_dir=tmp_path, drive_root=tmp_path) registry._ctx.task_id = task_id registry._ctx.task_metadata = {"root_task_id": task_id, "parent_task_id": task_id} - assert {"delegate_start", "delegate_wait", "delegate_cancel", "delegate_answer"} <= set( - registry._entries) + assert {"delegate_start", "delegate_wait", "delegate_cancel", "delegate_answer", + "delegate_message"} <= set(registry._entries) return registry @@ -177,6 +177,8 @@ _PRE_DAEMON_REFUSALS = [ ("delegate_cancel", {"run_id": ""}, "missing_run_id", "TOOL_ARG_ERROR"), ("delegate_answer", {"run_id": "", "interaction_id": "i-1", "answers": [{"question_id": "q"}]}, "missing_run_id", "TOOL_ARG_ERROR"), + ("delegate_message", {"run_id": "run-1", "text": " "}, + "message_text_required", "TOOL_ARG_ERROR"), ] diff --git a/tests/test_delegate_message.py b/tests/test_delegate_message.py new file mode 100644 index 000000000..6a3df8e19 --- /dev/null +++ b/tests/test_delegate_message.py @@ -0,0 +1,709 @@ +"""A live message into a running delegated run: ``delegate_message``. + +The fifth nanny verb places one message into the run's RUNNING turn through the +engine's capability-declared channel (the route row's ``liveInput``, the +``POST /v2/runs/:id/messages`` operation). What these pin: the gateway returns a +typed body at any HTTP status and types every untyped refusal with its code; +the verb refuses malformed calls and foreign runs before any wire call, never +POSTs a FRESH message to a settled run or an incapable route, mirrors the +engine's typed outcomes 1:1 (plus the host's ``not_found``), and keeps custody +of ``message_id`` (the wire Idempotency-Key): returned in every result, replayed +only when handed back, and the replay skips the host short-circuits so the +engine can answer with the stored receipt. +""" + +import json +import queue as stdqueue + +import pytest + + +@pytest.fixture(autouse=True) +def _fresh_custody_memo(): + from ouroboros import delegate_custody as custody + + custody._CUSTODY.clear() + yield + custody._CUSTODY.clear() + + +def _ctx(tmp_path, task_id="t-nanny"): + from ouroboros.contracts.task_constraint import TaskConstraint + from ouroboros.tools.registry import ToolContext + + repo = tmp_path / "repo" + repo.mkdir(exist_ok=True) + ctx = ToolContext(repo_dir=repo, drive_root=tmp_path, + task_constraint=TaskConstraint(mode="local_readonly_subagent")) + ctx.task_id = task_id + ctx.event_queue = stdqueue.Queue() + return ctx + + +def _own_run(run_id="run-1", *, task_id="t-nanny", route_id="some-route", settled=False): + from ouroboros import delegate_custody as custody + + custody._CUSTODY[run_id] = custody.RunCustody( + run_id=run_id, task_id=task_id, route_id=route_id, model="m", + project_id="prj", project_owned=False, settled=settled, + ) + + +def _row(route_id="some-route", live_input="mid_turn"): + row = {"id": route_id, "enabled": True, "accessProfilesSupported": ["readonly"]} + if live_input is not None: + row["liveInput"] = live_input + return row + + +_OPS_WITH_ROUTE = [{"id": "run.message", "method": "POST", "path": "/v2/runs/:id/messages"}] + + +class _Stub: + """The gateway as the verb sees it; every wire call is RECORDED.""" + + engine_version = "3.16.0" + + def __init__(self, *, result=None, error=None, operations=None, harnesses=None, + capability_error=None): + self.result = result + self.error = error + self.operations_rows = _OPS_WITH_ROUTE if operations is None else operations + self.harnesses = [_row()] if harnesses is None else harnesses + self.capability_error = capability_error + self.posts = [] + self.reads = [] + + def handshake(self, **_kw): + return {} + + def operations(self): + self.reads.append("operations") + if self.capability_error is not None: + raise self.capability_error + return list(self.operations_rows) + + def agent_capabilities(self): + self.reads.append("agent_capabilities") + return {"harnesses": list(self.harnesses)} + + def send_run_message(self, rid, text, *, idempotency_key, expected_attempt_id="", + timeout_sec=None): + self.posts.append({"run_id": rid, "text": text, "key": idempotency_key, + "timeout_sec": timeout_sec}) + if self.error is not None: + raise self.error + return dict(self.result) + + def close(self): + pass + + +def _install(monkeypatch, stub): + from ouroboros.gateways import claudexor as gw + + monkeypatch.setattr(gw, "ClaudexorGateway", lambda *a, **k: stub) + return stub + + +def _events(tmp_path): + path = tmp_path / "logs" / "events.jsonl" + if not path.exists(): + return [] + return [json.loads(line) for line in path.read_text(encoding="utf-8").splitlines() if line] + + +def _receipts(tmp_path): + return [e for e in _events(tmp_path) if e["type"] == "delegate_message_outcome"] + + +# -- the gateway: typed bodies at any status, typed problems otherwise ----------- + + +def test_send_run_message_returns_typed_bodies_and_types_every_refusal(monkeypatch): + import httpx + + from ouroboros.gateways import claudexor as cx + + replies = {} + + class _Recorder: + def request(self, method, path, **kwargs): + replies.update(method=method, path=path, json=kwargs.get("json"), + headers=kwargs.get("headers"), timeout=kwargs.get("timeout")) + return httpx.Response(replies["code"], json=replies["body"]) + + gateway = cx.ClaudexorGateway(cx.DaemonEndpoint("127.0.0.1", 1, "secret")) + gateway.close() + gateway._client = _Recorder() + + replies.update(code=200, body={"accepted": True, "outcome": "accepted", "runId": "run 1", + "messageId": "m-1", "attemptId": "a01"}) + body = gateway.send_run_message("run 1", "steer left", idempotency_key="m-1", + timeout_sec=12.0) + assert body["outcome"] == "accepted" + assert replies["method"] == "POST" + assert replies["path"] == "/v2/runs/run%201/messages" + assert replies["json"] == {"text": "steer left"} + # The message identity IS the wire Idempotency-Key, sent verbatim. + assert replies["headers"] == {"Idempotency-Key": "m-1"} + assert replies["timeout"] is not None + + # expectedAttemptId rides only when a caller holds one. + gateway.send_run_message("run 1", "x", idempotency_key="m-2", expected_attempt_id="a01") + assert replies["json"] == {"text": "x", "expectedAttemptId": "a01"} + + # A typed outcome is the ANSWER whatever the status (defensive: the engine + # answers every typed outcome at 200 by contract). + replies.update(code=200, body={"accepted": False, "outcome": "not_active", + "reason": "run_terminal"}) + assert gateway.send_run_message("run 1", "x", idempotency_key="m-3")["reason"] == "run_terminal" + + # The 409 idempotency problems carry their code AND status. + replies.update(code=409, body={"code": "idempotency_conflict", "message": "different digest"}) + with pytest.raises(cx.ClaudexorUnavailable) as exc: + gateway.send_run_message("run 1", "y", idempotency_key="m-3") + assert (exc.value.status_code, exc.value.code) == (409, "idempotency_conflict") + + # The daemon's own 404 (every daemon 404 has a body) stays the typed refusal it is. + replies.update(code=404, body={"error": "no such run"}) + with pytest.raises(cx.ClaudexorUnavailable) as exc: + gateway.send_run_message("run-gone", "x", idempotency_key="m-4") + assert exc.value.status_code == 404 + + # 501: this engine build has no live-message service. + replies.update(code=501, body={"error": "not supported"}) + with pytest.raises(cx.ClaudexorUnavailable) as exc: + gateway.send_run_message("run 1", "x", idempotency_key="m-5") + assert exc.value.status_code == 501 + + # A 2xx without a typed outcome is malformed, never an accepted delivery. + replies.update(code=200, body={"ok": True}) + with pytest.raises(cx.ClaudexorUnavailable) as exc: + gateway.send_run_message("run 1", "x", idempotency_key="m-6") + assert exc.value.code == "malformed_response" + + +def test_send_run_message_transport_death_is_daemon_unreachable(): + import httpx + + from ouroboros.gateways import claudexor as cx + + class _Dead: + def request(self, *a, **k): + raise httpx.ConnectError("refused") + + gateway = cx.ClaudexorGateway(cx.DaemonEndpoint("127.0.0.1", 1, "secret")) + gateway.close() + gateway._client = _Dead() + with pytest.raises(cx.ClaudexorUnavailable) as exc: + gateway.send_run_message("run-1", "x", idempotency_key="m-1") + assert exc.value.code == "daemon_unreachable" and exc.value.status_code == 0 + + +def test_run_message_supported_reads_the_engine_route_catalog(): + from ouroboros.gateways import claudexor as cx + + assert cx.run_message_supported(_OPS_WITH_ROUTE) + assert not cx.run_message_supported([]) + assert not cx.run_message_supported([{"method": "GET", "path": "/v2/runs/:id/messages"}]) + assert not cx.run_message_supported([{"method": "POST", "path": "/v2/runs/:id/control"}]) + assert not cx.run_message_supported(["junk", None]) + + +# -- argument and custody refusals (no wire call) --------------------------------- + + +def test_empty_text_is_an_agent_fault_before_any_wire_call(tmp_path, monkeypatch): + from ouroboros.delegate_interactions import _delegate_message + from ouroboros.delegate_shared import AGENT_FAULT_CODE + + stub = _install(monkeypatch, _Stub()) + ctx = _ctx(tmp_path) + _own_run() + for text in ("", " ", None, 7): + out = _delegate_message(ctx, "run-1", text) + payload = json.loads(out.text) + assert payload["reason"] == "message_text_required" + assert out.code == AGENT_FAULT_CODE and payload["ok"] is False + out = json.loads(_delegate_message(ctx, "", "hello").text) + assert out["reason"] == "missing_run_id" + assert stub.posts == [] and stub.reads == [] + assert _receipts(tmp_path) == [] + + +def test_foreign_and_unknown_runs_are_refused_without_a_wire_call(tmp_path, monkeypatch): + from ouroboros.delegate_interactions import _delegate_message + + stub = _install(monkeypatch, _Stub()) + ctx = _ctx(tmp_path) + _own_run(task_id="t-other") + foreign = json.loads(_delegate_message(ctx, "run-1", "steer").text) + assert foreign["status"] == "refused" and foreign["reason"] == "run_not_owned" + assert foreign["owner_task_id"] == "t-other" + unknown = json.loads(_delegate_message(ctx, "run-nobody", "steer").text) + assert unknown["reason"] == "run_ownership_unknown" + assert stub.posts == [] and stub.reads == [] + + +# -- host short-circuits for a FRESH message ----------------------------------------- + + +def test_a_settled_run_is_not_active_without_a_post(tmp_path, monkeypatch): + from ouroboros.delegate_interactions import _delegate_message + from ouroboros.delegate_shared import SUBSTRATE_REFUSAL_CODE + + stub = _install(monkeypatch, _Stub()) + ctx = _ctx(tmp_path) + _own_run(settled=True) + out = _delegate_message(ctx, "run-1", "steer") + payload = json.loads(out.text) + assert payload["status"] == "not_active" and payload["reason"] == "run_settled" + assert payload["accepted"] is False and out.code == SUBSTRATE_REFUSAL_CODE + assert payload["message_id"] # minted and returned even for a host verdict + assert stub.posts == [] and stub.reads == [] + receipt = _receipts(tmp_path)[0] + assert receipt["outcome"] == "not_active" and receipt["message_id"] == payload["message_id"] + + +@pytest.mark.parametrize("stub_kwargs,reason,live_input", [ + ({"operations": []}, "engine_lacks_operation", None), + ({"harnesses": [_row(live_input="none")]}, "route_live_input_none", "none"), + ({"harnesses": [_row(live_input=None)]}, "route_live_input_none", "none"), + ({"harnesses": [_row(route_id="another-route")]}, "route_not_in_capability_catalog", None), + ({"capability_error": RuntimeError("catalog down")}, "capability_read_failed", "unknown"), +]) +def test_an_incapable_route_is_unsupported_without_a_post(tmp_path, monkeypatch, + stub_kwargs, reason, live_input): + """A18: the operation must be listed AND the route row must declare a + liveInput other than none; a failed read is unsupported too, never a guess + and never a refusal that spends a round.""" + from ouroboros.delegate_interactions import _delegate_message + from ouroboros.delegate_shared import SUBSTRATE_REFUSAL_CODE + + stub = _install(monkeypatch, _Stub(**stub_kwargs)) + ctx = _ctx(tmp_path) + _own_run() + out = _delegate_message(ctx, "run-1", "steer") + payload = json.loads(out.text) + assert payload["status"] == "unsupported" and payload["reason"] == reason + assert payload["live_input"] == live_input + assert out.code == SUBSTRATE_REFUSAL_CODE and payload["ok"] is False + assert stub.posts == [] + assert "delegate_start" in payload["note"] + + +def test_a_capable_route_reads_both_facts_then_posts_once(tmp_path, monkeypatch): + from ouroboros.delegate_interactions import _delegate_message + + stub = _install(monkeypatch, _Stub(result={"accepted": True, "outcome": "accepted", + "attemptId": "a01", "harnessId": "codex", + "liveInput": "mid_turn"})) + ctx = _ctx(tmp_path) + _own_run() + _delegate_message(ctx, "run-1", "steer") + assert stub.reads == ["operations", "agent_capabilities"] + assert len(stub.posts) == 1 + + +# -- typed outcomes ------------------------------------------------------------------- + + +def test_accepted_relays_typed_facts_and_records_a_digest_never_the_text(tmp_path, monkeypatch): + import hashlib + + from ouroboros.delegate_interactions import _delegate_message + + text = "MANGO: stop after step 3 and report" + stub = _install(monkeypatch, _Stub(result={ + "accepted": True, "outcome": "accepted", "runId": "run-1", "messageId": "ignored", + "attemptId": "a01", "harnessId": "codex", "liveInput": "mid_turn", + "nativeTurnId": "turn-7"})) + ctx = _ctx(tmp_path) + _own_run() + out = _delegate_message(ctx, "run-1", text) + payload = json.loads(out.text) + assert out.status == "ok" and "ok" not in payload + assert payload["status"] == "accepted" and payload["accepted"] is True + assert (payload["attempt_id"], payload["harness_id"]) == ("a01", "codex") + assert (payload["live_input"], payload["native_turn_id"]) == ("mid_turn", "turn-7") + # The host-minted identity is what was sent AND what comes back. + assert payload["message_id"] == stub.posts[0]["key"] + assert len(payload["message_id"]) == 32 + assert "delegate_wait" in payload["note"] and "NEW" in payload["note"] + receipt = _receipts(tmp_path)[0] + assert receipt["text_sha256"] == hashlib.sha256(text.encode("utf-8")).hexdigest() + assert receipt["text_chars"] == len(text) + assert receipt["http_status"] == 200 and receipt["attempt_id"] == "a01" + assert text not in json.dumps(receipt) + + +def test_delivered_is_an_ok_observation_with_its_own_note(tmp_path, monkeypatch): + from ouroboros.delegate_interactions import _delegate_message + + _install(monkeypatch, _Stub(result={"accepted": True, "outcome": "delivered", + "nativeTurnId": "t1"})) + ctx = _ctx(tmp_path) + _own_run() + out = _delegate_message(ctx, "run-1", "steer") + payload = json.loads(out.text) + assert out.status == "ok" and payload["status"] == "delivered" + assert "obedience" in payload["note"] + + +def test_engine_rejected_classifies_by_reason(tmp_path, monkeypatch): + """A vendor refusal on an active turn is the substrate saying no; a replayed + message_id with different text is the caller's own defect.""" + from ouroboros.delegate_interactions import _delegate_message + from ouroboros.delegate_shared import AGENT_FAULT_CODE, SUBSTRATE_REFUSAL_CODE + + ctx = _ctx(tmp_path) + _own_run() + _install(monkeypatch, _Stub(result={"accepted": False, "outcome": "rejected", + "reason": "rpc_refused", "message": "turn/steer refused"})) + out = _delegate_message(ctx, "run-1", "steer") + payload = json.loads(out.text) + assert payload["status"] == "rejected" and payload["reason"] == "rpc_refused" + assert out.code == SUBSTRATE_REFUSAL_CODE and payload["detail"] == "turn/steer refused" + assert "NEW message_id" in payload["note"] + + _install(monkeypatch, _Stub(result={"accepted": False, "outcome": "rejected", + "reason": "multi_attempt"})) + assert _delegate_message(ctx, "run-1", "steer").code == SUBSTRATE_REFUSAL_CODE + + from ouroboros.gateways import claudexor as cx + + _install(monkeypatch, _Stub(error=cx.ClaudexorUnavailable( + "idempotency_conflict", "different digest", status_code=409))) + out = _delegate_message(ctx, "run-1", "steer", message_id="m-old") + payload = json.loads(out.text) + assert payload["status"] == "rejected" and payload["reason"] == "idempotency_conflict" + assert out.code == AGENT_FAULT_CODE + + +def test_engine_not_active_and_unsupported_are_substrate_refusals(tmp_path, monkeypatch): + from ouroboros.delegate_interactions import _delegate_message + from ouroboros.delegate_shared import SUBSTRATE_REFUSAL_CODE + + ctx = _ctx(tmp_path) + _own_run() + _install(monkeypatch, _Stub(result={"accepted": False, "outcome": "not_active", + "reason": "interaction_pending", "attemptId": "a01"})) + out = _delegate_message(ctx, "run-1", "steer") + payload = json.loads(out.text) + assert (payload["status"], payload["reason"]) == ("not_active", "interaction_pending") + assert out.code == SUBSTRATE_REFUSAL_CODE and "delegate_answer" in payload["note"] + + _install(monkeypatch, _Stub(result={"accepted": False, "outcome": "unsupported", + "reason": "thread_bound"})) + out = _delegate_message(ctx, "run-1", "steer") + payload = json.loads(out.text) + assert (payload["status"], payload["reason"]) == ("unsupported", "thread_bound") + assert out.code == SUBSTRATE_REFUSAL_CODE + + +def test_payload_verdict_4xx_is_rejected_as_an_agent_fault(tmp_path, monkeypatch): + from ouroboros.delegate_interactions import _delegate_message + from ouroboros.delegate_shared import AGENT_FAULT_CODE + from ouroboros.gateways import claudexor as cx + + ctx = _ctx(tmp_path) + _own_run() + _install(monkeypatch, _Stub(error=cx.ClaudexorUnavailable( + "invalid_request", "text too long", status_code=400))) + out = _delegate_message(ctx, "run-1", "steer") + payload = json.loads(out.text) + assert payload["status"] == "rejected" and payload["reason"] == "invalid_request" + assert out.code == AGENT_FAULT_CODE + assert _receipts(tmp_path)[0]["http_status"] == 400 + + +@pytest.mark.parametrize("code,status", [ + ("daemon_unreachable", 0), # transport death + ("http_503", 503), # any 5xx + ("internal_error", 500), # the receipt-save failure path + ("delivery_in_progress", 409), # the idempotency store: fate unknown + ("delivery_interrupted", 409), + ("http_429", 429), # says nothing about delivery +]) +def test_ambiguous_failures_are_delivery_unknown_naming_the_same_id_retry( + tmp_path, monkeypatch, code, status): + from ouroboros.delegate_interactions import _delegate_message + from ouroboros.gateways import claudexor as cx + + ctx = _ctx(tmp_path) + _own_run() + _install(monkeypatch, _Stub(error=cx.ClaudexorUnavailable(code, "boom", status_code=status))) + out = _delegate_message(ctx, "run-1", "steer") + payload = json.loads(out.text) + # Recorded as an OK observation (the fate is unknown, not refused), with + # the reason relayed and the recovery named: the SAME message_id. + assert out.status == "ok" and "ok" not in payload + assert payload["status"] == "delivery_unknown" and payload["reason"] == code + assert payload["accepted"] is False + assert "SAME message_id" in payload["note"] and "different" in payload["note"] + assert _receipts(tmp_path)[0]["http_status"] == status + + +def test_any_404_after_positive_reads_is_not_found_with_custody_untouched(tmp_path, monkeypatch): + from ouroboros import delegate_custody as custody + from ouroboros.delegate_interactions import _delegate_message + from ouroboros.delegate_shared import SUBSTRATE_REFUSAL_CODE + from ouroboros.gateways import claudexor as cx + + ctx = _ctx(tmp_path) + _own_run() + stub = _install(monkeypatch, _Stub(error=cx.ClaudexorUnavailable( + "http_404", "no such run", status_code=404))) + out = _delegate_message(ctx, "run-1", "steer") + payload = json.loads(out.text) + assert payload["status"] == "not_found" and out.code == SUBSTRATE_REFUSAL_CODE + assert stub.reads == ["operations", "agent_capabilities"] and len(stub.posts) == 1 + # Custody is not closed by this path: the run row stays open and unsettled. + assert custody._CUSTODY["run-1"].settled is False + assert not [e for e in _events(tmp_path) if e["type"].startswith("delegate_run_")] + + +def test_an_untyped_host_failure_is_delivery_unknown_never_a_traceback(tmp_path, monkeypatch): + from ouroboros.delegate_interactions import _delegate_message + + class _Broken(_Stub): + def send_run_message(self, *a, **k): + raise KeyError("wire shape") + + ctx = _ctx(tmp_path) + _own_run() + _install(monkeypatch, _Broken()) + payload = json.loads(_delegate_message(ctx, "run-1", "steer").text) + assert payload["status"] == "delivery_unknown" and payload["reason"] == "host_exception" + assert "KeyError" in payload["detail"] + + +def test_a_dead_daemon_at_handshake_is_a_typed_refusal(tmp_path, monkeypatch): + from ouroboros.delegate_interactions import _delegate_message + from ouroboros.gateways import claudexor as cx + + class _Down(_Stub): + def handshake(self, **_kw): + raise cx.ClaudexorUnavailable("daemon_unreachable", "no socket") + + ctx = _ctx(tmp_path) + _own_run() + stub = _install(monkeypatch, _Down()) + payload = json.loads(_delegate_message(ctx, "run-1", "steer").text) + assert payload["status"] == "refused" and payload["reason"] == "daemon_unreachable" + assert payload["message_id"] and stub.posts == [] + + +def test_a_spent_budget_before_the_post_sends_nothing(tmp_path, monkeypatch): + """The internal deadline sits strictly below the 120 s ToolEntry timeout; + spent before the POST, the answer is typed and nothing was sent — the + same-id retry the note prescribes is exactly right (the key is unused).""" + import ouroboros.delegate_interactions as interactions + from ouroboros.delegate_interactions import _delegate_message + + assert interactions._MESSAGE_DEADLINE_SEC < 120 + monkeypatch.setattr(interactions, "_MESSAGE_DEADLINE_SEC", -1.0) + ctx = _ctx(tmp_path) + _own_run() + stub = _install(monkeypatch, _Stub()) + payload = json.loads(_delegate_message(ctx, "run-1", "steer").text) + assert payload["status"] == "delivery_unknown" and payload["reason"] == "deadline_exhausted" + assert "nothing was sent" in payload["detail"] + assert stub.posts == [] + + +# -- message_id custody (A16/A27) ------------------------------------------------------- + + +def test_a_returned_message_id_replays_under_the_same_key_skipping_the_short_circuits( + tmp_path, monkeypatch): + from ouroboros.delegate_interactions import _delegate_message + + ctx = _ctx(tmp_path) + # The replay must reach the engine even though the host would refuse a FRESH + # message here (settled custody, an engine catalog without the operation): + # the daemon holds the stored receipt and is the only party that can answer. + _own_run(settled=True) + stub = _install(monkeypatch, _Stub(operations=[], result={ + "accepted": True, "outcome": "accepted", "attemptId": "a01"})) + payload = json.loads(_delegate_message(ctx, "run-1", "steer", message_id="m-earlier").text) + assert payload["status"] == "accepted" and payload["message_id"] == "m-earlier" + assert stub.reads == [] + assert stub.posts == [{"run_id": "run-1", "text": "steer", "key": "m-earlier", + "timeout_sec": stub.posts[0]["timeout_sec"]}] + + +def test_a_fresh_call_mints_a_new_id_each_time(tmp_path, monkeypatch): + from ouroboros.delegate_interactions import _delegate_message + + ctx = _ctx(tmp_path) + _own_run() + stub = _install(monkeypatch, _Stub(result={"accepted": True, "outcome": "accepted"})) + first = json.loads(_delegate_message(ctx, "run-1", "same text").text)["message_id"] + second = json.loads(_delegate_message(ctx, "run-1", "same text").text)["message_id"] + # Never a content-stable key: the same text twice is two deliveries. + assert first != second and [p["key"] for p in stub.posts] == [first, second] + + +# -- the registration surfaces -------------------------------------------------------- + + +def test_the_message_verb_is_registered_on_every_contract_surface(): + import ouroboros.delegate_interactions as interactions + from ouroboros.delegate_shared import _AGENT_FAULT_REASONS + from ouroboros.nanny_pacing import BASELINE_RESET_TOOLS, DELEGATE_ACTIVITY_TOOLS + from ouroboros.tool_capabilities import ( + ACTING_SUBAGENT_TOOL_NAMES, + LOCAL_READONLY_SUBAGENT_TOOL_NAMES, + ) + from ouroboros.tools import delegate + + entries = {entry.name: entry for entry in delegate.get_tools()} + assert "delegate_message" in entries + entry = entries["delegate_message"] + assert entry.timeout_sec == 120 and entry.timeout_sec > interactions._MESSAGE_DEADLINE_SEC + schema = entry.schema["parameters"] + assert schema["required"] == ["run_id", "text"] + assert set(schema["properties"]) == {"run_id", "text", "message_id"} + assert "delegate_message" in LOCAL_READONLY_SUBAGENT_TOOL_NAMES + assert "delegate_message" in ACTING_SUBAGENT_TOOL_NAMES + # Supervision activity, never a burn-baseline reset (only start/schedule are). + assert "delegate_message" in DELEGATE_ACTIVITY_TOOLS + assert "delegate_message" not in BASELINE_RESET_TOOLS + assert {"message_text_required", "idempotency_conflict"} <= _AGENT_FAULT_REASONS + assert delegate._delegate_message is interactions._delegate_message + + +def test_descriptions_are_capability_based_not_harness_named(): + from ouroboros.tools import delegate + + text = json.dumps([entry.schema for entry in delegate.get_tools()]) + assert "Codex-lane runs have no mid-run questions" not in text + assert "liveInput" in text + assert "message.* rows" in text # delegate_wait names the timeline as the reconciler + + +def test_the_wait_timeline_keeps_message_receipts(tmp_path): + from ouroboros.delegate_progress import timeline_tail + + rows = timeline_tail({"timeline": [ + {"type": "message.accepted", "title": "Live message accepted (12 bytes)", + "severity": "info", "attemptId": "a01", "messageId": "m-1", "outcome": "accepted"}, + {"type": "message.delivered", "title": "Live message delivered (12 bytes)", + "severity": "info", "messageId": "m-1", "outcome": "delivered"}, + {"type": "harness.started", "title": "started", "severity": "info", "outcome": 7}, + ]}) + assert (rows[0]["messageId"], rows[0]["outcome"], rows[0]["attemptId"]) == ("m-1", "accepted", "a01") + assert (rows[1]["messageId"], rows[1]["outcome"]) == ("m-1", "delivered") + assert "outcome" not in rows[2] and "messageId" not in rows[2] + + +# -- the fake daemon's contract, pinned against the REAL gateway -------------------------- + + +def test_fake_daemon_run_message_contract(tmp_path): + from ouroboros.gateways.claudexor import ( + ClaudexorGateway, + ClaudexorUnavailable, + discover_daemon_at, + run_message_supported, + ) + from tests.system_e2e.interfaces import FAKE_ASK_MARKER, FAKE_HANG_MARKER, FakeClaudexorDaemon + + with FakeClaudexorDaemon(runs_dir=tmp_path / "runs") as daemon: + daemon.install(tmp_path / "cx") + with ClaudexorGateway(discover_daemon_at(tmp_path / "cx")) as gateway: + gateway.handshake() + # Both discovery facts the verb negotiates are served. + assert run_message_supported(gateway.operations()) + row = gateway.agent_capabilities()["harnesses"][0] + assert row["id"] == daemon.harness_id and row["liveInput"] == "mid_turn" + request = { + "prompt": FAKE_HANG_MARKER + " keep running", "instructions": "i", + "authPreference": "subscription", "mode": "ask", + "scope": {"kind": "project", "root": str(tmp_path)}, + "harnesses": [daemon.harness_id], "primaryHarness": daemon.harness_id, + "access": "readonly", "maxSeconds": 60, + } + run_id = str(gateway.start_run(request, idempotency_key="inv-msg-1")["runId"]) + first = gateway.send_run_message(run_id, "steer left", idempotency_key="m-1") + assert first["outcome"] == "accepted" and first["accepted"] is True + assert first["messageId"] == "m-1" and first["runId"] == run_id + assert first["liveInput"] == "mid_turn" and first["attemptId"] == "a01" + # A replay under the same key is the STORED receipt, never a second delivery. + assert gateway.send_run_message(run_id, "steer left", idempotency_key="m-1") == first + assert [m["text"] for m in daemon.runs[run_id]["messages"]] == ["steer left"] + # The same key with different text is the typed 409. + with pytest.raises(ClaudexorUnavailable) as exc: + gateway.send_run_message(run_id, "steer RIGHT", idempotency_key="m-1") + assert (exc.value.status_code, exc.value.code) == (409, "idempotency_conflict") + # The key is REQUIRED on the wire. + with pytest.raises(ClaudexorUnavailable) as exc: + gateway.send_run_message(run_id, "x", idempotency_key="") + assert exc.value.code == "missing_idempotency_key" + # A run parked on a question is never steered beside it (INV-048). + asking = str(gateway.start_run({**request, "prompt": FAKE_ASK_MARKER + " q"}, + idempotency_key="inv-msg-2")["runId"]) + parked = gateway.send_run_message(asking, "steer", idempotency_key="m-2") + assert (parked["outcome"], parked["reason"]) == ("not_active", "interaction_pending") + # An unknown run is the daemon's own 404, with a body. + with pytest.raises(ClaudexorUnavailable) as exc: + gateway.send_run_message("run-nobody", "x", idempotency_key="m-3") + assert exc.value.status_code == 404 + posts = daemon.calls("POST", f"/v2/runs/{run_id}/messages") + assert [p["idempotency_key"] for p in posts] == ["m-1", "m-1", "m-1", ""] + assert posts[0]["body"] == {"text": "steer left"} + + +def test_fake_daemon_route_without_live_input_answers_unsupported(tmp_path): + from ouroboros.gateways.claudexor import ClaudexorGateway, discover_daemon_at + from tests.system_e2e.interfaces import FAKE_HANG_MARKER, FakeClaudexorDaemon + + with FakeClaudexorDaemon(runs_dir=tmp_path / "runs", live_input="none") as daemon: + daemon.install(tmp_path / "cx") + with ClaudexorGateway(discover_daemon_at(tmp_path / "cx")) as gateway: + gateway.handshake() + assert gateway.agent_capabilities()["harnesses"][0]["liveInput"] == "none" + run_id = str(gateway.start_run({ + "prompt": FAKE_HANG_MARKER, "instructions": "i", "authPreference": "subscription", + "mode": "ask", "scope": {"kind": "project", "root": str(tmp_path)}, + "harnesses": [daemon.harness_id], "primaryHarness": daemon.harness_id, + "access": "readonly", "maxSeconds": 60, + }, idempotency_key="inv-none-1")["runId"]) + body = gateway.send_run_message(run_id, "steer", idempotency_key="m-1") + assert (body["outcome"], body["reason"]) == ("unsupported", "no_live_session") + assert "messages" not in daemon.runs[run_id] + + +def test_the_verb_end_to_end_against_the_fake_daemon(tmp_path, monkeypatch): + """The whole verb over the loopback fake: discovery, POST, typed receipt.""" + from ouroboros.delegate_interactions import _delegate_message + from ouroboros.gateways import claudexor as cx + from tests.system_e2e.interfaces import FAKE_HANG_MARKER, FakeClaudexorDaemon + + with FakeClaudexorDaemon(runs_dir=tmp_path / "runs") as daemon: + daemon.install(tmp_path / "cx") + endpoint = cx.discover_daemon_at(tmp_path / "cx") + # The verb builds ``ClaudexorGateway()`` bare; point its discovery at the fake. + monkeypatch.setattr(cx, "discover_daemon", lambda home=None: endpoint) + with cx.ClaudexorGateway(endpoint) as setup: + setup.handshake() + run_id = str(setup.start_run({ + "prompt": FAKE_HANG_MARKER, "instructions": "i", "authPreference": "subscription", + "mode": "ask", "scope": {"kind": "project", "root": str(tmp_path)}, + "harnesses": [daemon.harness_id], "primaryHarness": daemon.harness_id, + "access": "readonly", "maxSeconds": 60, + }, idempotency_key="inv-e2e-1")["runId"]) + ctx = _ctx(tmp_path) + _own_run(run_id, route_id=daemon.harness_id) + payload = json.loads(_delegate_message(ctx, run_id, "steer left").text) + assert payload["status"] == "accepted" and payload["accepted"] is True + assert payload["harness_id"] == daemon.harness_id and payload["live_input"] == "mid_turn" + posts = daemon.calls("POST", f"/v2/runs/{run_id}/messages") + assert len(posts) == 1 and posts[0]["idempotency_key"] == payload["message_id"] + # The replay reaches the daemon and comes back as the stored receipt. + again = json.loads(_delegate_message(ctx, run_id, "steer left", + message_id=payload["message_id"]).text) + assert again["status"] == "accepted" and again["message_id"] == payload["message_id"] + assert len(daemon.runs[run_id]["messages"]) == 1 diff --git a/tests/test_delegated_executor_axis.py b/tests/test_delegated_executor_axis.py index df6118e58..d04a3270b 100644 --- a/tests/test_delegated_executor_axis.py +++ b/tests/test_delegated_executor_axis.py @@ -27,7 +27,8 @@ from tests._delegated_transport_shared import ( # noqa: F401 (autouse fixture ) -NANNY_TOOLS = {"delegate_start", "delegate_wait", "delegate_cancel", "delegate_answer"} +NANNY_TOOLS = {"delegate_start", "delegate_wait", "delegate_cancel", "delegate_answer", + "delegate_message"} def test_subagent_harness_key_stays_out_of_the_model_key_sweep(): diff --git a/tests/test_smoke.py b/tests/test_smoke.py index 0e22b4138..efe47a7d0 100644 --- a/tests/test_smoke.py +++ b/tests/test_smoke.py @@ -115,6 +115,7 @@ EXPECTED_TOOLS = [ "set_next_wakeup", "switch_model", "get_task_result", "wait_task", "wait_tasks", "await_messages", "tree_note", "tree_read", "delegate_start", "delegate_wait", "delegate_cancel", "delegate_answer", + "delegate_message", "read_file", "list_files", "write_file", "edit_text", "apply_patch", "edit_batch", "send_photo", "send_video", "send_file", "send_links", "search_code", "query_code", "escalate", diff --git a/tests/test_tool_capabilities.py b/tests/test_tool_capabilities.py index 23a31a414..398d89612 100644 --- a/tests/test_tool_capabilities.py +++ b/tests/test_tool_capabilities.py @@ -125,6 +125,7 @@ def test_top_level_workspace_focus_has_tool_and_schema_parity(tmp_path, monkeypa assert names == schemas assert { "delegate_start", "delegate_wait", "delegate_cancel", "delegate_answer", + "delegate_message", "switch_model", "send_photo", "send_video", "send_file", "send_links", "commit_reviewed", "promote_chat_to_task", "vcs_restore", } <= names