fix: persist API delegation units once without waking the model
This commit is contained in:
80
evals/api_delegation_http_probe.py
Normal file
80
evals/api_delegation_http_probe.py
Normal file
@@ -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))
|
||||
@@ -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
|
||||
|
||||
@@ -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(
|
||||
|
||||
@@ -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/
|
||||
|
||||
@@ -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:
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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"
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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,
|
||||
|
||||
81
tests/gateway/test_api_delegation_delivery_contract.py
Normal file
81
tests/gateway/test_api_delegation_delivery_contract.py
Normal file
@@ -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()
|
||||
@@ -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)
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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")
|
||||
|
||||
@@ -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
|
||||
|
||||
Reference in New Issue
Block a user