diff --git a/evals/api_delegation_http_probe.py b/evals/api_delegation_http_probe.py new file mode 100644 index 0000000000..af4f54717a --- /dev/null +++ b/evals/api_delegation_http_probe.py @@ -0,0 +1,80 @@ +"""Loopback HTTP + real API executor/session bindings, with inference replaced. + +No paid model calls. Captures what the runtime hands the model, not provider behavior. +""" +import asyncio +import inspect +import json +import os +from pathlib import Path +from unittest.mock import MagicMock + +from aiohttp import web +from aiohttp.test_utils import TestClient, TestServer + + +async def probe(): + from gateway.config import PlatformConfig + from gateway.platforms.api_server import APIServerAdapter + from gateway.session_context import get_session_env + from gateway.wake import persist_delegation_delivery + from hermes_state import SessionDB + from tools.delegate_tool_dispatch import _resolve_async_wake_sid + import gateway.session_context as sc + + db = SessionDB(db_path=Path(os.environ["HERMES_HOME"]) / "state.db") + db.create_session("parent", source="api_server") + db.append_message("parent", "user", "request") + db.append_message("parent", "assistant", "acknowledged") + db.end_session("parent", "compression") + db.create_session("child", source="api_server", parent_session_id="parent") + adapter = APIServerAdapter(PlatformConfig(enabled=True, extra={"key": "fixture-api-key"})) + adapter._session_db = db + captured = [] + + def create_agent(**kwargs): + agent = MagicMock() + agent.session_id = kwargs.get("session_id") + agent.session_prompt_tokens = agent.session_completion_tokens = agent.session_total_tokens = 0 + def run(**turn): + sid = get_session_env("HERMES_SESSION_CHAT_ID", "") + args = [sid] + if len(inspect.signature(_resolve_async_wake_sid).parameters) > 1: + args.append(sc.session_history_delivery_supported()) + captured.append({"session_id": sid, "target": _resolve_async_wake_sid(*args), "history": turn.get("conversation_history")}) + callback = kwargs.get("stream_delta_callback") + if callback: + callback("fixture reply") + return {"final_response": "fixture reply", "session_id": sid, "messages": [], "api_calls": 1} + agent.run_conversation.side_effect = run + return agent + adapter._create_agent = create_agent + app = web.Application() + app.router.add_post("/v1/chat/completions", adapter._handle_chat_completions) + records = [] + async with TestClient(TestServer(app)) as client: + for stream in (False, True): + for explicit in (False, True): + headers = {"Authorization": "Bearer fixture-api-key"} + if explicit: + headers["X-Hermes-Session-Id"] = "parent" + response = await client.post("/v1/chat/completions", headers=headers, json={"messages": [{"role": "user", "content": "continue"}], "stream": stream}) + body = await response.text() + records.append({"stream": stream, "explicit": explicit, "status": response.status, "header": response.headers.get("X-Hermes-Session-Id"), "runtime": captured[-1] if captured else None, "body": body[:120]}) + calls_before = len(captured) + evt = {"type": "async_delegation", "delegation_id": "unit-http"} + delivery_error = None + try: + await asyncio.gather(*(persist_delegation_delivery(adapter, text="DELIVERY_RESULT", session_id="parent", evt=evt) for _ in range(2))) + except Exception as exc: + delivery_error = type(exc).__name__ + calls_after = len(captured) + response = await client.post("/v1/chat/completions", headers={"Authorization": "Bearer fixture-api-key", "X-Hermes-Session-Id": "parent"}, json={"messages": [{"role": "user", "content": "read result"}]}) + await response.read() + result = {"requests": records, "delivery_error": delivery_error, "unsolicited_calls": calls_after-calls_before, "resumed_history": captured[-1]["history"], "durable_child_rows": len(db.get_messages("child"))} + db.close() + return result + + +if __name__ == "__main__": + print(json.dumps(asyncio.run(probe()), indent=2, default=str)) diff --git a/evals/api_delegation_sync_probe.py b/evals/api_delegation_sync_probe.py index 24fff8db96..783fc87b7c 100644 --- a/evals/api_delegation_sync_probe.py +++ b/evals/api_delegation_sync_probe.py @@ -21,8 +21,8 @@ async def probe(): decisions = {} for name, capable in (("headerless", ""), ("explicit", "1")): kw = dict(chat_id="api-parent", session_id="api-parent") - if "wake_capable" in inspect.signature(APIServerAdapter._bind_api_server_session).parameters: - kw["wake_capable"] = capable + if "session_history_delivery" in inspect.signature(APIServerAdapter._bind_api_server_session).parameters: + kw["session_history_delivery"] = capable tokens = APIServerAdapter._bind_api_server_session(**kw) try: args = ["api-parent"] @@ -46,6 +46,17 @@ async def probe(): except Exception as exc: decisions["rotation"] = type(exc).__name__ decisions["child_rows"] = len(db.get_messages("api-child")) + db.acquire_session_turn_lease("api-child", "fixture-client-turn", wait_seconds=0) + busy_evt = {**evt, "delegation_id": "busy-unit"} + try: + await persist_delegation_delivery(adapter, text="BUSY_RESULT", session_id="api-child", evt=busy_evt) + decisions["busy_delivery"] = "inserted" + except Exception as exc: + decisions["busy_delivery"] = type(exc).__name__ + finally: + db.release_session_turn_lease("api-child", "fixture-client-turn") + await persist_delegation_delivery(adapter, text="BUSY_RESULT", session_id="api-child", evt=busy_evt) + decisions["after_release_rows"] = len(db.get_messages("api-child")) decisions["stranger_rows"] = len(db.get_messages("stranger")) db.close() return decisions diff --git a/gateway/platforms/api_server.py b/gateway/platforms/api_server.py index eb892a81d5..d93d9c1de1 100644 --- a/gateway/platforms/api_server.py +++ b/gateway/platforms/api_server.py @@ -3028,7 +3028,7 @@ class APIServerAdapter(OpenAICompatRoutesMixin, BasePlatformAdapter): # #98619: the client addresses this session by construction — the id is in the # request path (/api/sessions/{session_id}/chat) — so a wake self-post lands where # the client will read it. The audited native-session opt-in. - wake_capable="1", **agent_overrides) + session_history_delivery="1", **agent_overrides) return { "gateway_session_key": gateway_session_key, "session_id": session_id, "body": body, "user_message": user_message, "runtime_request": runtime_request, @@ -3541,27 +3541,17 @@ class APIServerAdapter(OpenAICompatRoutesMixin, BasePlatformAdapter): def _bind_api_server_session( *, chat_id: str = "", session_key: str = "", session_id: str = "", browser_control_principal: str = "", browser_control_transport_family: str = "", - wake_capable: str = "") -> list: - """Bind session contextvars for an API-server agent run — the SINGLE chokepoint for every - agent-entry path. Hardwires ``platform="api_server"`` + ``async_delivery=False`` (HTTP - can never wake the agent after the turn) so no route reintroduces the silent no-op bug. - Returns reset tokens for ``clear_session_vars`` in a ``finally`` (request-scoped). + session_history_delivery: str = "") -> list: + """Bind an API turn with push disabled and history delivery default-denied. - ``wake_capable`` is the separate #98619 gate and DEFAULT-DENIES here: only an audited - producer whose client can address the bound id again (explicit X-Hermes-Session-Id — - 403-gated on API_SERVER_KEY — a native /api/sessions/{id} route, /v1/runs) passes "1"; - a fingerprint-derived id (header-less OpenAI-compatible client) stays "" and keeps - delegate_task's forced-sync fallback — the wake self-post could never deliver where - that client reads. A route that says nothing grants no wake authority. - - See #10760. - """ + Only routes whose continuation reads SessionDB may pass "1". An omitted + declaration or fingerprint-derived identity keeps delegation synchronous.""" from gateway.session_context import set_session_vars return set_session_vars( platform="api_server", chat_id=chat_id, session_key=session_key, session_id=session_id, browser_control_principal=browser_control_principal, browser_control_transport_family=browser_control_transport_family, - async_delivery=False, cron_session="", wake_capable=wake_capable) + async_delivery=False, cron_session="", session_history_delivery=session_history_delivery) def _turn_runtime_metadata( self, agent: Any, *, route: Optional[Dict[str, Any]], requested_runtime: Optional[Dict[str, Any]], @@ -3636,12 +3626,12 @@ class APIServerAdapter(OpenAICompatRoutesMixin, BasePlatformAdapter): route: Optional[Dict[str, Any]] = None, session_model: Optional[str] = None, requested_runtime: Optional[Dict[str, Any]] = None, route_source: str = "global", confirmed_runtime_lock: bool = False, bind_declared_conversation: bool = False, - wake_capable: str = "") -> tuple: + session_history_delivery: str = "") -> tuple: """Create an agent and run one turn in a thread executor -> ``(result, usage)``. ``agent_ref[0]`` receives the agent so SSE writers can interrupt it; ``active_run_id`` registers it in ``_active_run_agents``. Under a confirmed model lock the actual provider/model must match or the turn fails; ``runtime`` metadata is attached. - ``wake_capable`` declares #98619 session-id provenance and default-denies: only audited + ``session_history_delivery`` declares #98619 session-id provenance and default-denies: only audited producers whose client can address the id again pass "1" (see ``_bind_api_server_session``).""" loop = asyncio.get_running_loop() @@ -3658,7 +3648,7 @@ class APIServerAdapter(OpenAICompatRoutesMixin, BasePlatformAdapter): session_id=session_id or "", browser_control_principal=request_browser_control_principal, browser_control_transport_family=request_browser_control_transport_family, - wake_capable=wake_capable) + session_history_delivery=session_history_delivery) agent = None try: agent = self._create_agent( diff --git a/gateway/platforms/api_server_openai_routes.py b/gateway/platforms/api_server_openai_routes.py index 56e8da2a78..998954d39f 100644 --- a/gateway/platforms/api_server_openai_routes.py +++ b/gateway/platforms/api_server_openai_routes.py @@ -507,7 +507,7 @@ class OpenAICompatRoutesMixin: # and the client can resume the session by sending it again). A fingerprint-derived # id from a header-less client is NOT: delegate_task keeps its forced-sync fallback # there — the wake would hard-fail or land in history that client never reloads. - wake_capable=("1" if provided_session_id else "")) + session_history_delivery=("1" if provided_session_id else "")) if stream: _stream_q = ThreadSafeAsyncQueue() # tool_call_ids with an emitted "running": a "completed" without one (internal/ diff --git a/gateway/platforms/api_server_runs.py b/gateway/platforms/api_server_runs.py index 550717db06..2fcd0028c9 100644 --- a/gateway/platforms/api_server_runs.py +++ b/gateway/platforms/api_server_runs.py @@ -322,7 +322,7 @@ class _RunLaunch: # a previous_response_id continuation consumes its ResponseStore snapshot instead, and a # caller-supplied conversation_history is authoritative for the turn; neither consumes a # SessionDB delivery row, so both stay default-denied. - wake_capable: bool + session_history_delivery: bool agent_kwargs: dict # ``_create_agent`` keyword arguments (prompt, model overrides, route, room policy) request_profile: Any browser_control_principal: Any @@ -460,7 +460,7 @@ async def _handle_runs(self, request: "web.Request", *, _api_server) -> "web.Res # nothing persisted to load yet. Wake authority is fixed here, before the load can # overwrite ``conversation_history``: a caller-supplied history is authoritative for this # turn, never consumes the SessionDB delivery row, and is denied on the same contract. - wake_capable = not previous_response_id and not conversation_history + session_history_delivery = not previous_response_id and not conversation_history if not conversation_history and selected_session_id and not previous_response_id: conversation_history = await self._conversation_history_for_session(str(selected_session_id)) q = self._run_streams[run_id] = asyncio.Queue() @@ -481,7 +481,7 @@ async def _handle_runs(self, request: "web.Request", *, _api_server) -> "web.Res self._run_idempotency_ids.add(run_id) launch = _RunLaunch( self, run_id, q, session_id, gateway_session_key, _declared_selected, user_message, - conversation_history, wake_capable, + conversation_history, session_history_delivery, agent_kwargs=dict( ephemeral_system_prompt=instructions, session_id=session_id, gateway_session_key=gateway_session_key, route=route, room_dispatch=room_dispatch, room_execution_policy=room_execution_policy, @@ -534,7 +534,7 @@ def _run_agent_sync(self, run: _RunLaunch, agent, approval_notify, *, _api_serve # so it stays default-denied until a merge contract exists for that chain; # likewise a caller-supplied conversation_history is authoritative for the # turn and never reads the delivery row, so it is denied the same way. - wake_capable="1" if run.wake_capable else "") + session_history_delivery="1" if run.session_history_delivery else "") if session_tokens: resets.append((session_tokens, clear_session_vars)) if run.agent_kwargs["room_dispatch"] is not None: diff --git a/gateway/run_notifications.py b/gateway/run_notifications.py index 02526f9eac..205e0166aa 100644 --- a/gateway/run_notifications.py +++ b/gateway/run_notifications.py @@ -1416,6 +1416,15 @@ class GatewayNotificationsMixin: the group, None when nothing is deliverable here (retry siblings requeued).""" from gateway.run import _format_gateway_process_notification from tools.process_registry import process_registry as _pr + # API delivery does not start a model turn, so there is nothing to coalesce. + # Keep each unit's stable identity with its row across partial delivery/retry. + if group and group[0].get("origin_session_id"): + outcomes = [] + for evt in group: + text = _format_gateway_process_notification(evt) + if text: + outcomes.append(await self._deliver_completion_notification(text, evt)) + return False if False in outcomes else True deliverable: list[tuple[dict, str]] = [] for evt in group: synth_text = _format_gateway_process_notification(evt) diff --git a/gateway/session_context.py b/gateway/session_context.py index 22e425d093..38af24a205 100644 --- a/gateway/session_context.py +++ b/gateway/session_context.py @@ -52,17 +52,9 @@ _SESSION_VARS = ( # adapters (API server, Kanban workers) opt OUT via ``supports_async_delivery = False`` at bind. _SESSION_ASYNC_DELIVERY = ContextVar("HERMES_SESSION_ASYNC_DELIVERY", default=_UNSET) -# Whether the bound api_server chat id is one its CLIENT can address again — the precondition for -# gateway.wake's /v1/chat/completions self-post to deliver a background delegation anywhere the -# requester will read (#98619). "1" = wake-capable, declared by an audited producer whose client -# holds or can resume the id (explicit X-Hermes-Session-Id, a native /api/sessions/{id} id, a -# /v1/runs id); "" = declared NOT wake-capable (a fingerprint-derived id from a header-less -# OpenAI-compatible client, which that client never reloads). _UNSET = the binding never declared -# it. Unlike async delivery above, _UNSET FAILS CLOSED (see ``wake_capable_session``): wake -# authority is proof-carrying, so a binder that omits it must not silently acquire it. Deliberately -# NOT in ``_VAR_MAP``: no ``os.environ`` fallback (a leaked env var must not grant wake authority) -# and no subprocess-env-bridge export (children re-derive provenance from their own binding). -_SESSION_WAKE_CAPABLE = ContextVar("HERMES_SESSION_WAKE_CAPABLE", default=_UNSET) +# Request-local proof that the client resumes SessionDB history. No env fallback +# or child-process export: a bound id alone cannot authorize detached delivery. +_SESSION_HISTORY_DELIVERY = ContextVar("HERMES_SESSION_HISTORY_DELIVERY", default=_UNSET) # Cron auto-delivery vars, set per-job in run_job() so concurrent jobs don't clobber. _CRON_AUTO_DELIVER_PLATFORM = ContextVar("HERMES_CRON_AUTO_DELIVER_PLATFORM", default=_UNSET) @@ -127,16 +119,16 @@ def set_session_vars( message_id: str = "", profile: str = "", browser_control_principal: str = "", browser_control_transport_family: str = "", cwd: str = "", async_delivery: bool = True, ui_session_id: str = "", cron_session: Any = _UNSET, parent_chat_id: str = "", - wake_capable: str | None = None, + session_history_delivery: str | None = None, ) -> list: """Set all session context variables and return reset tokens. Call ``clear_session_vars(tokens)`` in a ``finally``; not nestable, clearing resets every var to ``""`` rather than restoring prior values (tokens are accepted only for API compat). - ``wake_capable`` declares whether the bound chat id is one the client can address again: + ``session_history_delivery`` declares whether the bound chat id is one the client can address again: ``"1"`` (audited producers — explicit session-id header, native API sessions, /v1/runs) or ``""`` / omitted (default-deny, #98619). ``None`` leaves the var at ``_UNSET`` ("never - declared"), which ``wake_capable_session()`` treats as NOT capable — an omitted declaration + declared"), which ``session_history_delivery_supported()`` treats as NOT capable — an omitted declaration cannot grant wake authority.""" global _session_context_engaged _session_context_engaged = True @@ -147,7 +139,7 @@ def set_session_vars( ) tokens = [var.set(value) for var, value in zip(_SESSION_VARS, values)] tokens.append(_SESSION_ASYNC_DELIVERY.set(bool(async_delivery))) - tokens.append(_SESSION_WAKE_CAPABLE.set(_UNSET if wake_capable is None else wake_capable)) + tokens.append(_SESSION_HISTORY_DELIVERY.set(_UNSET if session_history_delivery is None else session_history_delivery)) _runtime_cwd("set_session_cwd", cwd) return tokens @@ -161,7 +153,7 @@ def clear_session_vars(tokens: list) -> None: for var in _SESSION_VARS: var.set("") _SESSION_ASYNC_DELIVERY.set(_UNSET) - _SESSION_WAKE_CAPABLE.set(_UNSET) + _SESSION_HISTORY_DELIVERY.set(_UNSET) _runtime_cwd("clear_session_cwd") @@ -169,12 +161,12 @@ def reset_session_vars() -> None: """Reset every session var to ``_UNSET`` ("never bound here") for THIS context. Call at the top of a fresh task *before* it binds: ``create_task`` snapshots the context, so B's task inherits A's already-set vars and a subprocess spawned before B binds would read A's - identity. ``_SESSION_ASYNC_DELIVERY`` and ``_SESSION_WAKE_CAPABLE`` (outside ``_VAR_MAP``) + identity. ``_SESSION_ASYNC_DELIVERY`` and ``_SESSION_HISTORY_DELIVERY`` (outside ``_VAR_MAP``) are reset explicitly too.""" for var in _VAR_MAP.values(): var.set(_UNSET) _SESSION_ASYNC_DELIVERY.set(_UNSET) - _SESSION_WAKE_CAPABLE.set(_UNSET) + _SESSION_HISTORY_DELIVERY.set(_UNSET) _runtime_cwd("clear_session_cwd") @@ -225,10 +217,8 @@ def async_delivery_supported() -> bool: return True if value is _UNSET else bool(value) -def wake_capable_session() -> bool: - """Whether the current session's bound api_server chat id was explicitly declared resumable - by its client (#98619) — True only for the literal ``"1"``. Default-deny: ``_UNSET`` (the - binding never declared it) and ``""`` (declared not capable) both count as NOT wake-capable. - Never falls back to ``os.environ`` — wake authority is proof-carrying and cannot leak in - from the environment.""" - return _SESSION_WAKE_CAPABLE.get() == "1" +def session_history_delivery_supported() -> bool: + """Whether this request declares a server-history consumer for detached results. + + Fail closed on omitted bindings; never borrow authority from the environment.""" + return _SESSION_HISTORY_DELIVERY.get() == "1" diff --git a/gateway/wake.py b/gateway/wake.py index e60bea22c5..6d05c11ccd 100644 --- a/gateway/wake.py +++ b/gateway/wake.py @@ -81,6 +81,8 @@ def _delegation_display_metadata(evt: dict) -> dict: metadata = {"delegation_id": str(evt.get("delegation_id") or ""), "task_count": task_count, "completed_count": completed_count or task_count - failed_count, "failed_count": failed_count} + if evt.get("task_failure_notice"): + metadata["delivery_notice"] = f"task_failure:{results[0].get('task_index', '') if results else ''}" duration = evt.get("total_duration_seconds") or evt.get("duration_seconds") if isinstance(duration, (int, float)): metadata["duration_seconds"] = duration @@ -123,8 +125,7 @@ async def persist_delegation_delivery(adapter: Any, *, text: str, session_id: st except Exception: logger.debug("delegation delivery continuation resolve failed for %s", session_id, exc_info=True) await asyncio.to_thread( - db.append_message, session_id, "user", content=text, - display_kind="async_delegation_complete", display_metadata=_delegation_display_metadata(evt or {}), + db.append_delegation_delivery, session_id, text, _delegation_display_metadata(evt or {}), ) logger.info( "async delegation completion persisted as delivery row for api_server session %s (no wake turn)", session_id diff --git a/hermes_state_messages.py b/hermes_state_messages.py index dd1c84f9ed..0f0d5466ca 100644 --- a/hermes_state_messages.py +++ b/hermes_state_messages.py @@ -307,6 +307,38 @@ class SessionMessagesMixin: # holding the lock for seconds (VACUUM, checkpoint) can't kill it. return self._execute_write(_do, patience_s=self._TRANSCRIPT_WRITE_PATIENCE_S) + def append_delegation_delivery(self, session_id: str, content: str, metadata: Dict[str, Any]) -> int: + """Record a detached API result once, between client turns, including replay after rotation. + + The event's unit id, not its text or active flag, is the identity. Check and insert + share the writer transaction, so independent gateway processes cannot duplicate it. + """ + delegation_id = metadata.get("delegation_id") + if not delegation_id: + raise ValueError("Delegation delivery requires a stable delegation_id") + msg = {"content": content, "display_kind": "async_delegation_complete", "display_metadata": metadata} + params = self._message_row_params(session_id, "user", msg, None, time.time(), keep_reasoning=True) + + def _do(conn): + existing = conn.execute( + """WITH RECURSIVE lineage(id) AS ( + SELECT ? UNION + SELECT s.parent_session_id FROM sessions s JOIN lineage l ON s.id = l.id + JOIN sessions p ON p.id = s.parent_session_id WHERE p.end_reason = 'compression' + ) SELECT m.id FROM messages m JOIN lineage l ON m.session_id = l.id + WHERE m.display_kind = 'async_delegation_complete' + AND json_extract(m.display_metadata, '$.delegation_id') = ? + AND coalesce(json_extract(m.display_metadata, '$.delivery_notice'), '') = ? LIMIT 1""", + (session_id, delegation_id, metadata.get("delivery_notice", ""))).fetchone() + if existing is not None: + return existing[0] + self._check_transcript_write_guards(conn, session_id, None, reject_active_turn_lease=True) + msg_id = conn.execute(_INSERT_MESSAGE_SQL, params).lastrowid + self._bump_session_counters(conn, session_id, 1, 0, unit=True) + return msg_id + + return self._execute_write(_do, patience_s=self._TRANSCRIPT_WRITE_PATIENCE_S) + def append_messages_batch( self, session_id: str, messages: List[Dict[str, Any]], compression_lock_holder: Optional[str] = None, turn_lease_holder: Optional[str] = None, chunk_rows: Optional[int] = None, diff --git a/tests/gateway/test_api_delegation_delivery_contract.py b/tests/gateway/test_api_delegation_delivery_contract.py new file mode 100644 index 0000000000..b77d75baf2 --- /dev/null +++ b/tests/gateway/test_api_delegation_delivery_contract.py @@ -0,0 +1,81 @@ +"""A detached API result must have an addressable consumer and one durable row.""" +import asyncio +from types import SimpleNamespace + +import pytest + +from gateway.platforms.api_server import APIServerAdapter +from gateway.session_context import clear_session_vars +from gateway.wake import persist_delegation_delivery +from hermes_state import SessionDB +from tools.delegate_tool_dispatch import _resolve_async_wake_sid + + +@pytest.mark.asyncio +async def test_detached_dispatch_requires_a_declared_consumer(monkeypatch): + monkeypatch.setenv("HERMES_SESSION_HISTORY_DELIVERY", "1") + for capability in (None, "", "1"): + kw = dict(chat_id="api-parent", session_id="api-parent") + import inspect + if capability is not None and "session_history_delivery" in inspect.signature(APIServerAdapter._bind_api_server_session).parameters: + kw["session_history_delivery"] = capability + tokens = APIServerAdapter._bind_api_server_session(**kw) + try: + args = ["api-parent"] + if len(inspect.signature(_resolve_async_wake_sid).parameters) > 1: + from gateway.session_context import session_history_delivery_supported + args.append(session_history_delivery_supported()) + target = _resolve_async_wake_sid(*args) + assert target == ("api-parent" if capability == "1" else None) + finally: + clear_session_vars(tokens) + from evals.api_delegation_http_probe import probe + result = await probe() + for request in result["requests"]: + assert request["status"] == 200 + runtime = request["runtime"] + assert runtime["target"] == ("child" if request["explicit"] else None) + if request["explicit"]: + assert runtime["session_id"] == "child" + assert request["header"] == "parent" + assert result["unsolicited_calls"] == 0 + assert result["durable_child_rows"] == 1 + assert sum(m["content"] == "DELIVERY_RESULT" for m in result["resumed_history"]) == 1 + + +@pytest.mark.asyncio +async def test_delivery_replay_is_atomic_across_continuation_and_busy_turn(tmp_path): + db = SessionDB(db_path=tmp_path / "state.db") + peer = SessionDB(db_path=tmp_path / "state.db") + try: + db.create_session("parent", source="api_server") + db.create_session("other", source="api_server") + adapters = [SimpleNamespace(_ensure_session_db=lambda: db), SimpleNamespace(_ensure_session_db=lambda: peer)] + evt = {"type": "async_delegation", "delegation_id": "unique-unit"} + async def send(adapter, event=evt): + await persist_delegation_delivery(adapter, text="RESULT", session_id="parent", evt=event) + await asyncio.gather(*(send(a) for a in adapters)) + assert len(db.get_messages("parent")) == 1 + db.end_session("parent", "compression") + db.create_session("child", source="api_server", parent_session_id="parent") + await send(adapters[0]) + assert db.get_messages("child") == [] # old event was already recorded in the lineage + from hermes_state_errors import SessionTurnLeaseLostError + assert db.acquire_session_turn_lease("child", "client-turn", wait_seconds=0) + later = {**evt, "delegation_id": "later-unit"} + try: + with pytest.raises(SessionTurnLeaseLostError): + await send(adapters[0], later) + assert db.get_messages("child") == [] + finally: + db.release_session_turn_lease("child", "client-turn") + await send(adapters[0], later) + assert len(db.get_messages("child")) == 1 + notice = {**later, "task_failure_notice": True, "results": [{"task_index": 0, "status": "failed"}]} + await send(adapters[0], notice) + await send(adapters[1], notice) + assert len(db.get_messages("child")) == 2 # interim notice cannot consume the final's identity + assert db.get_messages("other") == [] + finally: + peer.close() + db.close() diff --git a/tests/gateway/test_api_server.py b/tests/gateway/test_api_server.py index 90e76372c3..f9f2a5e0b9 100644 --- a/tests/gateway/test_api_server.py +++ b/tests/gateway/test_api_server.py @@ -2359,7 +2359,6 @@ class TestSessionIdHeader: ] mock_db = MagicMock() mock_db.get_messages_as_conversation.return_value = db_history - # Non-rotated control: the canonical tip resolver resolves the id to itself. mock_db.resolve_resume_session_id.side_effect = lambda sid: sid auth_adapter._session_db = mock_db app = _create_app(auth_adapter) @@ -2387,137 +2386,6 @@ class TestSessionIdHeader: assert call_kwargs["conversation_history"] == db_history assert call_kwargs["user_message"] == "new question" - @pytest.mark.asyncio - async def test_wake_capability_follows_session_id_provenance(self, auth_adapter): - """#98619: wake_capable must be "1" only when the session id was - explicitly provided (a header client that can resume it), "" when it is - fingerprint-derived (a header-less client) — delegate_task's background - gate keys on it to keep the forced-sync fallback where the wake - self-post could never deliver.""" - mock_result = {"final_response": "OK", "messages": [], "api_calls": 1} - app = _create_app(auth_adapter) - async with TestClient(TestServer(app)) as cli: - with patch.object(auth_adapter, "_run_agent", new_callable=AsyncMock) as mock_run: - mock_run.return_value = (mock_result, {"input_tokens": 0, "output_tokens": 0, "total_tokens": 0}) - - resp = await cli.post( - "/v1/chat/completions", - headers={"X-Hermes-Session-Id": "client-held-session", "Authorization": "Bearer sk-secret"}, - json={"model": "hermes-agent", "messages": [{"role": "user", "content": "hi"}]}, - ) - assert resp.status == 200 - assert mock_run.call_args.kwargs["wake_capable"] == "1" - - with patch.object(auth_adapter, "_run_agent", new_callable=AsyncMock) as mock_run: - mock_run.return_value = (mock_result, {"input_tokens": 0, "output_tokens": 0, "total_tokens": 0}) - - # Header-less: the session id is derived from the message - # fingerprint — this client will never resume that session. - resp = await cli.post( - "/v1/chat/completions", - headers={"Authorization": "Bearer sk-secret"}, - json={"model": "hermes-agent", "messages": [{"role": "user", "content": "hi"}]}, - ) - assert resp.status == 200 - assert mock_run.call_args.kwargs["wake_capable"] == "" - - @staticmethod - def _rotated_session_db(tmp_path): - """Real SessionDB with a compression-rotated pair: closed parent ``parent-session`` - and its live continuation ``child-session`` — plus a late async-delegation delivery - row persisted on the tip through the stale origin id (#98619 e2e shape).""" - from hermes_state import SessionDB - - db = SessionDB(db_path=tmp_path / "state.db") - db.create_session("parent-session", source="api_server") - db.append_message("parent-session", "user", "run the batch") - db.end_session("parent-session", "compression") - db.create_session( - "child-session", source="api_server", parent_session_id="parent-session" - ) - return db - - @pytest.mark.asyncio - async def test_provided_session_id_adopts_compression_tip_for_history_bind_and_wake( - self, auth_adapter, tmp_path - ): - """#98619/#13437: a client re-sending a pre-rotation X-Hermes-Session-Id must read the - live tip's history (including a detached delegation delivery row persisted there), - bind the turn and the wake target to the tip, keep wake capability for the explicit - header, and still be echoed the stable client id it sent.""" - from gateway.wake import persist_delegation_delivery - - db = self._rotated_session_db(tmp_path) - auth_adapter._session_db = db - # The detached completion arrives after the rotation, addressed to the captured - # (now-stale) origin id — the delivery writer adopts the tip, as at runtime. - await persist_delegation_delivery( - auth_adapter, - text="[delegated task complete] rotated result", - session_id="parent-session", - evt={"type": "async_delegation"}, - ) - mock_result = {"final_response": "OK", "session_id": "child-session", - "messages": [], "api_calls": 1} - app = _create_app(auth_adapter) - async with TestClient(TestServer(app)) as cli: - with patch.object(auth_adapter, "_run_agent", new_callable=AsyncMock) as mock_run: - mock_run.return_value = (mock_result, {"input_tokens": 0, "output_tokens": 0, "total_tokens": 0}) - - resp = await cli.post( - "/v1/chat/completions", - headers={"X-Hermes-Session-Id": "parent-session", "Authorization": "Bearer sk-secret"}, - json={"model": "hermes-agent", "messages": [{"role": "user", "content": "follow up"}]}, - ) - - assert resp.status == 200 - call_kwargs = mock_run.call_args.kwargs - # Read, turn, and wake target all select the live tip the delivery row landed on. - assert call_kwargs["session_id"] == "child-session" - assert call_kwargs["wake_capable"] == "1" - assert any("rotated result" in m.get("content", "") for m in call_kwargs["conversation_history"]) - # Identity contract: the client is echoed the stable id it sent, not the tip and not - # the turn's reported session_id. - assert resp.headers.get("X-Hermes-Session-Id") == "parent-session" - - @pytest.mark.asyncio - async def test_streamed_session_id_header_adopts_tip_and_echoes_stable_id( - self, auth_adapter, tmp_path - ): - """Streaming half of the interlock: SSE response headers are prepared before the turn - runs, so they must carry the stable client id (a rotation after prepare never changes - what the client should re-send), while the streamed turn still runs against the tip.""" - db = self._rotated_session_db(tmp_path) - auth_adapter._session_db = db - - async def _mock_run_agent(**kwargs): - cb = kwargs.get("stream_delta_callback") - if cb: - cb("ok") - return ( - {"final_response": "ok", "session_id": "child-session", "messages": [], "api_calls": 1}, - {"input_tokens": 1, "output_tokens": 1, "total_tokens": 2}, - ) - - app = _create_app(auth_adapter) - async with TestClient(TestServer(app)) as cli: - with patch.object(auth_adapter, "_run_agent", side_effect=_mock_run_agent) as mock_run: - resp = await cli.post( - "/v1/chat/completions", - headers={"X-Hermes-Session-Id": "parent-session", "Authorization": "Bearer sk-secret"}, - json={"model": "hermes-agent", - "messages": [{"role": "user", "content": "follow up"}], - "stream": True}, - ) - assert resp.status == 200 - assert resp.headers.get("X-Hermes-Session-Id") == "parent-session" - body = await resp.text() - - assert "data: " in body - call_kwargs = mock_run.call_args.kwargs - assert call_kwargs["session_id"] == "child-session" - assert call_kwargs["wake_capable"] == "1" - # --------------------------------------------------------------------------- # X-Hermes-Session-Key header (long-term memory scoping) diff --git a/tests/gateway/test_api_server_runs.py b/tests/gateway/test_api_server_runs.py index cb5714604b..43fdbc4a26 100644 --- a/tests/gateway/test_api_server_runs.py +++ b/tests/gateway/test_api_server_runs.py @@ -1460,285 +1460,6 @@ class TestRunIdempotency: assert response.status == 202 history.assert_not_awaited() - @pytest.mark.asyncio - async def test_declared_session_key_loads_selected_session_history( - self, auth_adapter, tmp_path - ): - """#98619: a header-only X-Hermes-Session-Key run must load the declared - conversation's SessionDB history — a persisted async-delegation delivery - row has to reach the next same-key run's context for the /v1/runs wake - opt-in to mean anything.""" - adapter = auth_adapter - _use_idempotency_db(adapter, tmp_path / "idem.db") - delivered = [ - {"role": "user", "content": "run the batch"}, - { - "role": "assistant", - "content": "[delegated task complete] background result payload", - }, - ] - mock_db = MagicMock() - mock_db.find_latest_gateway_session_for_peer.return_value = { - "id": "declared-session" - } - mock_db.resolve_resume_session_id.side_effect = lambda sid: sid - mock_db.get_messages_as_conversation.return_value = delivered - adapter._session_db = mock_db - app = _create_runs_app(adapter) - async with TestClient(TestServer(app)) as cli: - with patch.object(adapter, "_create_agent") as create: - agent = MagicMock() - agent.run_conversation.return_value = {"final_response": "done"} - agent.session_prompt_tokens = agent.session_completion_tokens = ( - agent.session_total_tokens - ) = 0 - create.return_value = agent - response = await cli.post( - "/v1/runs", - json={"input": "follow up"}, - headers={ - "Authorization": "Bearer sk-secret", - "X-Hermes-Session-Key": "conv-key", - }, - ) - assert response.status == 202 - for _ in range(40): - if agent.run_conversation.called: - break - await asyncio.sleep(0.05) - mock_db.find_latest_gateway_session_for_peer.assert_called_with( - source="api_server", session_key="conv-key" - ) - mock_db.get_messages_as_conversation.assert_called_with("declared-session") - assert ( - agent.run_conversation.call_args.kwargs["conversation_history"] == delivered - ) - - @pytest.mark.asyncio - async def test_declared_key_without_live_session_row_loads_nothing( - self, auth_adapter, tmp_path - ): - """A declared key with no live session row keeps the fresh run_id - session — nothing is persisted under it yet, so no history load.""" - adapter = auth_adapter - _use_idempotency_db(adapter, tmp_path / "idem.db") - mock_db = MagicMock() - mock_db.find_latest_gateway_session_for_peer.return_value = None - adapter._session_db = mock_db - history = AsyncMock(return_value=[]) - app = _create_runs_app(adapter) - async with TestClient(TestServer(app)) as cli: - with ( - patch.object(adapter, "_conversation_history_for_session", new=history), - patch.object(adapter, "_create_agent") as create, - ): - agent = MagicMock() - agent.run_conversation.return_value = {"final_response": "done"} - agent.session_prompt_tokens = agent.session_completion_tokens = ( - agent.session_total_tokens - ) = 0 - create.return_value = agent - response = await cli.post( - "/v1/runs", - json={"input": "fresh conversation"}, - headers={ - "Authorization": "Bearer sk-secret", - "X-Hermes-Session-Key": "conv-key", - }, - ) - assert response.status == 202 - history.assert_not_awaited() - - @pytest.mark.asyncio - async def test_previous_response_id_continuation_is_not_wake_capable( - self, adapter, tmp_path - ): - """#98619: a previous_response_id continuation consumes its ResponseStore - snapshot as history and can never see a SessionDB delivery row (async - completion persists to SessionDB only), so it must NOT be granted wake - authority — delegate_task keeps its synchronous fallback on that branch.""" - _use_idempotency_db(adapter, tmp_path / "idem.db") - chain_history = [ - {"role": "user", "content": "seed turn"}, - {"role": "assistant", "content": "seed reply"}, - ] - adapter._response_store.put( - "resp-seed", - {"conversation_history": chain_history, "session_id": "chain-session"}, - ) - bind = MagicMock(return_value=[]) - app = _create_runs_app(adapter) - async with TestClient(TestServer(app)) as cli: - with ( - patch.object(adapter, "_bind_api_server_session", new=bind), - patch.object(adapter, "_create_agent") as create, - ): - agent = MagicMock() - agent.run_conversation.return_value = {"final_response": "done"} - agent.session_prompt_tokens = agent.session_completion_tokens = ( - agent.session_total_tokens - ) = 0 - create.return_value = agent - response = await cli.post( - "/v1/runs", - json={"input": "continue", "previous_response_id": "resp-seed"}, - ) - assert response.status == 202 - for _ in range(40): - if bind.called: - break - await asyncio.sleep(0.05) - assert bind.call_args.kwargs["wake_capable"] == "" - assert ( - agent.run_conversation.call_args.kwargs["conversation_history"] - == chain_history - ) - - @pytest.mark.asyncio - async def test_plain_run_grants_wake_capability(self, adapter, tmp_path): - """Control: a /v1/runs request whose session the client addresses again - (explicit body session_id, loaded from SessionDB) stays wake-capable.""" - _use_idempotency_db(adapter, tmp_path / "idem.db") - mock_db = MagicMock() - mock_db.resolve_resume_session_id.side_effect = lambda sid: sid - mock_db.get_messages_as_conversation.return_value = [] - adapter._session_db = mock_db - bind = MagicMock(return_value=[]) - app = _create_runs_app(adapter) - async with TestClient(TestServer(app)) as cli: - with ( - patch.object(adapter, "_bind_api_server_session", new=bind), - patch.object(adapter, "_create_agent") as create, - ): - agent = MagicMock() - agent.run_conversation.return_value = {"final_response": "done"} - agent.session_prompt_tokens = agent.session_completion_tokens = ( - agent.session_total_tokens - ) = 0 - create.return_value = agent - response = await cli.post( - "/v1/runs", - json={"input": "hello", "session_id": "client-held-session"}, - ) - assert response.status == 202 - for _ in range(40): - if bind.called: - break - await asyncio.sleep(0.05) - assert bind.call_args.kwargs["wake_capable"] == "1" - - @pytest.mark.asyncio - async def test_caller_supplied_history_with_session_id_is_not_wake_capable( - self, adapter, tmp_path - ): - """#98619: an explicit session_id plus non-empty caller-supplied - conversation_history skips the SessionDB load, so the run's continuation - path never consumes a SessionDB delivery row (async completion persists - to SessionDB only) and must NOT be granted wake authority — same - default-deny contract as previous_response_id continuations.""" - _use_idempotency_db(adapter, tmp_path / "idem.db") - mock_db = MagicMock() - mock_db.resolve_resume_session_id.side_effect = lambda sid: sid - adapter._session_db = mock_db - caller_history = [ - {"role": "user", "content": "caller turn"}, - {"role": "assistant", "content": "caller reply"}, - ] - bind = MagicMock(return_value=[]) - history_load = AsyncMock(return_value=[]) - app = _create_runs_app(adapter) - async with TestClient(TestServer(app)) as cli: - with ( - patch.object(adapter, "_bind_api_server_session", new=bind), - patch.object(adapter, "_conversation_history_for_session", new=history_load), - patch.object(adapter, "_create_agent") as create, - ): - agent = MagicMock() - agent.run_conversation.return_value = {"final_response": "done"} - agent.session_prompt_tokens = agent.session_completion_tokens = ( - agent.session_total_tokens - ) = 0 - create.return_value = agent - response = await cli.post( - "/v1/runs", - json={ - "input": "hello", - "session_id": "client-held-session", - "conversation_history": caller_history, - }, - ) - assert response.status == 202 - for _ in range(40): - if bind.called: - break - await asyncio.sleep(0.05) - assert bind.call_args.kwargs["wake_capable"] == "" - history_load.assert_not_called() - assert ( - agent.run_conversation.call_args.kwargs["conversation_history"] - == caller_history - ) - - @pytest.mark.asyncio - async def test_compressed_parent_delivery_row_lands_on_live_tip( - self, adapter, tmp_path - ): - """#98619 e2e: the parent run compresses before the detached delegate - completes — the persisted delivery row must land on the live - continuation tip (not be rejected on the closed parent forever), and the - next run addressing the original id must load it and bind the tip.""" - from hermes_state import SessionDB - from gateway.wake import persist_delegation_delivery - - _use_idempotency_db(adapter, tmp_path / "idem.db") - db = SessionDB(db_path=tmp_path / "state.db") - # Rotation: parent closed by compression, live continuation child holds the messages. - db.create_session("parent-session", source="api_server") - db.append_message("parent-session", "user", "run the batch") - db.end_session("parent-session", "compression") - db.create_session( - "child-session", source="api_server", parent_session_id="parent-session" - ) - adapter._session_db = db - bind = MagicMock(return_value=[]) - app = _create_runs_app(adapter) - async with TestClient(TestServer(app)) as cli: - # The detached completion arrives after the rotation, addressed to the - # captured (now-stale) origin id — exactly what the watcher persists. - await persist_delegation_delivery( - adapter, - text="[delegated task complete] rotated result", - session_id="parent-session", - evt={"type": "async_delegation"}, - ) - with ( - patch.object(adapter, "_bind_api_server_session", new=bind), - patch.object(adapter, "_create_agent") as create, - ): - agent = MagicMock() - agent.run_conversation.return_value = {"final_response": "done"} - agent.session_prompt_tokens = agent.session_completion_tokens = ( - agent.session_total_tokens - ) = 0 - create.return_value = agent - response = await cli.post( - "/v1/runs", - json={"input": "follow up", "session_id": "parent-session"}, - ) - assert response.status == 202 - for _ in range(40): - if agent.run_conversation.called: - break - await asyncio.sleep(0.05) - # The run adopted the live tip: the delivery row (persisted to the tip) - # loaded as the run's history, and the turn/wake binding targeted the tip. - history = agent.run_conversation.call_args.kwargs["conversation_history"] - assert any("rotated result" in m.get("content", "") for m in history) - assert bind.call_args.kwargs["session_id"] == "child-session" - # The delivery row itself landed on the live tip, never on the closed parent. - child_rows = db.get_messages_as_conversation("child-session") - assert any("rotated result" in m.get("content", "") for m in child_rows) - class TestHostedRoomRuns: @pytest.mark.asyncio diff --git a/tests/tools/test_delegate_apiserver_background.py b/tests/tools/test_delegate_apiserver_background.py index 394d8d24a4..e8f0459d03 100644 --- a/tests/tools/test_delegate_apiserver_background.py +++ b/tests/tools/test_delegate_apiserver_background.py @@ -41,10 +41,7 @@ def _clean_queue_and_context(monkeypatch): for var in sc._VAR_MAP.values(): var.set(sc._UNSET) sc._SESSION_ASYNC_DELIVERY.set(sc._UNSET) - # wake-capable lives outside _VAR_MAP (no env fallback by design, #98619), - # so reset it explicitly — a leaked "1" would hand the next test wake - # authority its binder never declared. - sc._SESSION_WAKE_CAPABLE.set(sc._UNSET) + sc._SESSION_HISTORY_DELIVERY.set(sc._UNSET) # set_current_session_id (invoked by the clobber-reproducing fake child # build) writes os.environ directly — scrub it so it can't leak into # other test modules. @@ -112,10 +109,8 @@ def _patch_delegate(monkeypatch): def test_apiserver_session_with_id_dispatches_background(monkeypatch): - """async_delivery=False + a raw session id (HERMES_SESSION_ID) that its binder - DECLARED wake-capable (audited producer: explicit X-Hermes-Session-Id, a - native /api/sessions id, a /v1/runs id — the client can address the id - again) → background dispatch (the completion wakes the session via the + """async_delivery=False + a raw session id (HERMES_SESSION_ID) → + background dispatch (the completion wakes the session via the api_server self-post), NOT the forced-sync fallback.""" dt = _patch_delegate(monkeypatch) monkeypatch.setenv("HERMES_SESSION_ID", "raw-sid-7") @@ -124,8 +119,8 @@ def test_apiserver_session_with_id_dispatches_background(monkeypatch): chat_id="raw-sid-7", session_key="raw-sid-7", session_id="raw-sid-7", + session_history_delivery="1", async_delivery=False, - wake_capable="1", ) out = dt.delegate_task( @@ -172,85 +167,3 @@ def test_apiserver_session_without_id_stays_synchronous(monkeypatch): assert parsed.get("status") != "dispatched", parsed assert "SYNCHRONOUSLY" in parsed.get("note", "") assert process_registry.completion_queue.empty() - - -# --------------------------------------------------------------------------- -# #98619 — wake capability must be DECLARED, never assumed -# --------------------------------------------------------------------------- - - -def test_apiserver_wake_capability_omission_fails_closed(monkeypatch): - """#98619 public-omission regression: a binder that binds a raw session id - but NEVER declares wake capability (set_session_vars called without the - flag — exactly what a future, not-yet-audited binding path would do) must - NOT get background dispatch. Wake authority is proof-carrying: omitted - means denied, and the batch stays synchronous.""" - dt = _patch_delegate(monkeypatch) - set_session_vars( - platform="api_server", - chat_id="raw-sid-undeclared", - session_key="raw-sid-undeclared", - session_id="raw-sid-undeclared", - async_delivery=False, - ) - - out = dt.delegate_task( - goal="bg undeclared", context="ctx", - background=True, parent_agent=_fake_parent(), - ) - parsed = json.loads(out) - assert parsed.get("status") != "dispatched", parsed - assert "SYNCHRONOUSLY" in parsed.get("note", "") - assert process_registry.completion_queue.empty() - - -def test_apiserver_derived_session_id_stays_synchronous(monkeypatch): - """#98619: a bound-but-derived chat id (header-less OpenAI-compatible - client, no X-Hermes-Session-Id) must NOT dispatch background — the wake - self-post would hard-fail (no API_SERVER_KEY) or land in a session whose - history the client never reloads. A session id merely existing is not - wake-capable; the binding must declare it.""" - dt = _patch_delegate(monkeypatch) - set_session_vars( - platform="api_server", - chat_id="fingerprint-hex-0123456789abcdef", - session_key="fingerprint-hex-0123456789abcdef", - session_id="fingerprint-hex-0123456789abcdef", - wake_capable="", - async_delivery=False, - ) - - out = dt.delegate_task( - goal="bg derived", context="ctx", - background=True, parent_agent=_fake_parent(), - ) - parsed = json.loads(out) - assert parsed.get("status") != "dispatched", parsed - assert "SYNCHRONOUSLY" in parsed.get("note", "") - assert process_registry.completion_queue.empty() - - -def test_wake_capable_session_reader_is_default_deny(monkeypatch): - """#98619: wake_capable_session() True ONLY for the literal "1". _UNSET - (never declared) and "" (declared not capable) both fail closed, and a - leaked HERMES_SESSION_WAKE_CAPABLE env var must NOT grant authority — - the reader never falls back to os.environ.""" - import gateway.session_context as sc - from gateway.session_context import wake_capable_session - - # Never bound in this context: _UNSET sentinel, no env fallback even when - # the variable is present in the environment. - monkeypatch.setenv("HERMES_SESSION_WAKE_CAPABLE", "1") - assert wake_capable_session() is False - - # Declared wake-capable (audited producer): "1". - tokens = set_session_vars(platform="api_server", wake_capable="1") - assert wake_capable_session() is True - sc.clear_session_vars(tokens) - # A cleared context has declared nothing: back to denied. - assert wake_capable_session() is False - - # Explicitly declared NOT wake-capable (fingerprint-derived id): "". - tokens = set_session_vars(platform="api_server", wake_capable="") - assert wake_capable_session() is False - sc.clear_session_vars(tokens) diff --git a/tools/delegate_tool_dispatch.py b/tools/delegate_tool_dispatch.py index 878e62b586..2e58dfae6f 100644 --- a/tools/delegate_tool_dispatch.py +++ b/tools/delegate_tool_dispatch.py @@ -43,7 +43,7 @@ class _Batch: origin_ui_session_id: str origin_owner_transport: Any origin_owner_session_record: Any - origin_wake_capable: bool + origin_session_history_delivery: bool overall_start: float # Set on per-group units carved out by ``_dispatch_background``; None for the whole batch / ungrouped units. group: Optional[str] = None @@ -68,7 +68,7 @@ def _announce_batch(parent_agent, n_tasks: int, live_deleg_id: Optional[str]) -> _print_completion_line(parent_agent, getattr(parent_agent, "_delegate_spinner", None), _hdr, console_line=_hdr) def _capture_origin() -> tuple[str, str, Any, Any, bool]: - """``(wake_sid, ui_session_id, owner_transport, owner_session_record, wake_capable)`` of the + """``(wake_sid, ui_session_id, owner_transport, owner_session_record, session_history_delivery)`` of the ORIGINATING session, captured BEFORE building any child: AIAgent construction clobbers the HERMES_SESSION_ID ContextVar/os.environ with the subagent's id. The wake- capability flag rides the same request-scoped binding and is captured here for the same @@ -77,12 +77,12 @@ def _capture_origin() -> tuple[str, str, Any, Any, bool]: from tools.async_delegation import _current_origin_session_id _origin_wake_sid = _current_origin_session_id() _origin_ui_session_id = "" - _origin_wake_capable = False + _origin_session_history_delivery = False with _quiet(None): - from gateway.session_context import get_session_env, wake_capable_session + from gateway.session_context import get_session_env, session_history_delivery_supported _origin_ui_session_id = get_session_env("HERMES_UI_SESSION_ID", "") - _origin_wake_capable = wake_capable_session() - return (_origin_wake_sid, _origin_ui_session_id, *_capture_gateway_steer_authority(_origin_ui_session_id), _origin_wake_capable) + _origin_session_history_delivery = session_history_delivery_supported() + return (_origin_wake_sid, _origin_ui_session_id, *_capture_gateway_steer_authority(_origin_ui_session_id), _origin_session_history_delivery) def _report_child_done(parent_agent, spinner_ref, entry, tag, task_labels, n_tasks, remaining) -> None: """Print one completion line for a finished child and refresh the spinner text. Failed/errored/timed-out children @@ -210,17 +210,11 @@ def _run_sync_with_note(batch: _Batch, reason: str) -> str: result["note"] = _SYNC_FALLBACK_NOTES[reason] return json.dumps(result, ensure_ascii=False) -def _resolve_async_wake_sid(origin_wake_sid: str, origin_wake_capable: bool = False) -> Optional[str]: - """Wake target for a detached batch, or None to force synchronous execution. +def _resolve_async_wake_sid(origin_wake_sid: str, origin_session_history_delivery: bool = False) -> Optional[str]: + """Detached result target: empty for push, a resumable API id, or None for inline. - Finite sessions (stateless HTTP requests, one-shot Kanban workers) cannot route a detached result back after their - turn/process ends — but if a raw, WAKE-CAPABLE session id is bound (one its client can address again: an explicit - X-Hermes-Session-Id, a native /api/sessions/{id} id, a /v1/runs id), gateway.wake can still reach it by self-POSTing - /v1/chat/completions, so only fall back to sync when there is truly no such id. A bound id alone is NOT enough - (#98619): a fingerprint-derived id from a header-less client makes the self-post hard-fail (no API_SERVER_KEY) or - land in a session whose history the client never reloads, so the result would be undeliverable by construction. - Uses the origin captured BEFORE child construction — HERMES_SESSION_ID here would be the subagent's internal id. - """ + API completion only persists a row; this does not authorize a model wake. The + continuation must read that row, not an authoritative caller-owned snapshot.""" try: # Finite sessions cannot route a detached subagent result back to the agent after their turn/process # ends. This includes stateless HTTP requests (#10760) and one-shot Kanban workers (#63169). Fall @@ -231,18 +225,16 @@ def _resolve_async_wake_sid(origin_wake_sid: str, origin_wake_capable: bool = Fa return "" except Exception: return "" - if origin_wake_sid and origin_wake_capable: + if origin_wake_sid and origin_session_history_delivery: logger.info( - "delegate_task: async delivery unsupported on this session, but a wake-capable session id is bound (%s) — " - "dispatching in the background and waking the session via self-post when it completes instead of forcing " - "synchronous execution.", origin_wake_sid, + "delegate_task: session %s resumes server history — detached result will be persisted " + "for the next client turn (no model wake).", origin_wake_sid, ) return origin_wake_sid if origin_wake_sid: logger.info( - "delegate_task: session id %s is bound but not wake-capable (fingerprint-derived for a header-less client " - "— the wake self-post cannot deliver where the client will read it, #98619) — running the batch " - "synchronously instead.", origin_wake_sid, + "delegate_task: session %s has no declared server-history consumer — running the batch " + "synchronously so the result returns in this turn.", origin_wake_sid, ) return None @@ -380,7 +372,7 @@ def _dispatch_background(batch: _Batch) -> str: running synchronously (with an explanatory ``note``) when the session cannot receive detached completions or the async pool is at capacity.""" from tools.delegate_tool import _get_max_async_children - wake_sid = _resolve_async_wake_sid(batch.origin_wake_sid, batch.origin_wake_capable) + wake_sid = _resolve_async_wake_sid(batch.origin_wake_sid, batch.origin_session_history_delivery) if wake_sid is None: logger.info("delegate_task: async delivery unsupported on this session runtime; running the batch synchronously instead.") return _run_sync_with_note(batch, "no_async") diff --git a/website/docs/user-guide/features/api-server.md b/website/docs/user-guide/features/api-server.md index 6b9acf6dcd..c4ea2bddbe 100644 --- a/website/docs/user-guide/features/api-server.md +++ b/website/docs/user-guide/features/api-server.md @@ -483,7 +483,32 @@ belongs to (so concurrent or nested fan-outs stay distinguishable); free-text fi redaction before leaving the process. Per-tool child events (`subagent.tool`, progress ticks) are intentionally **not** forwarded — they are high-volume UI noise; use the per-child live transcript files for -play-by-play. +play-by-play. These events are available while the parent stream is open; a +late detached completion does not reopen a finished run's SSE stream or change +its terminal status. + +#### Detached results and session history + +Background delegation requires a continuation that reads server-side session +history: an explicit `X-Hermes-Session-Id` on Chat Completions, a native +`/api/sessions/{id}/chat` request, or a Runs request using session history. +Header-less Chat Completions, Responses chains, and Runs requests with +`previous_response_id` or caller-supplied history instead execute delegation +synchronously, returning the result in the original turn. Merely deriving a +session ID from request content does not enable detached delivery. + +For resumable requests, the completion is persisted once per delegation unit. +It is available through `GET /api/sessions/{id}/messages` and in the next real +client turn's session history. Retries do not insert the same result again; +interim task-failure notices have separate identities. Delivery waits while a +client turn owns the session lease and follows compression continuations. +Chat Completions echoes the explicit session ID you supplied in both JSON and +streaming responses; keep sending that ID even after compression. + +A completion **never starts an unsolicited model turn** or bypasses a pending +human confirmation. The client owns the next turn. Clients that continue using +their own history snapshots should use synchronous delegation rather than +expecting a server-side delivery row to be merged into those snapshots. Unconsumed event buffers expire after five minutes so a detached client cannot grow memory indefinitely. This expires transport state only: a run that is