fix(gateway): require wake-capable session provenance for background delegate_task (#98619)
Rebuilt against the post-refactor owners: the chat-completions route now
lives in gateway/platforms/api_server_openai_routes.py, the wake gate in
tools/delegate_tool_dispatch.py, and session_context.py was reshaped —
the original patch aimed at code main no longer has.
Header-less OpenAI-compatible clients get a fingerprint-derived session
id bound as the api_server chat_id. delegate_task's background gate
treated ANY bound session id as wake-capable and dispatched detached
subagents, but for derived ids the wake self-post lands in a session
whose history the client never reloads — the result is undeliverable by
construction.
Bind a wake_capable provenance flag at session-bind time, DEFAULT-DENY
at the central boundary (set_session_vars and _bind_api_server_session
both treat an omitted declaration as denied; a binder that says nothing
grants no wake authority): "1" only from audited producers whose client
can address the id again (explicit X-Hermes-Session-Id — 403-gated on
API_SERVER_KEY — native /api/sessions/{id} routes, /v1/runs); "" for
fingerprint-derived ids. The delegate gate requires the flag (fail-closed,
captured pre-child-construction alongside origin_wake_sid) and keeps the
forced-sync fallback with its honest note.
Reported-and-investigated-by: shojikumaru (Sho + Alpha) via #98619
This commit is contained in:
@@ -3024,7 +3024,11 @@ class APIServerAdapter(OpenAICompatRoutesMixin, BasePlatformAdapter):
|
||||
gateway_session_key=gateway_session_key, route=route, session_model=session_model,
|
||||
requested_runtime=runtime_request.get("requested") or {},
|
||||
route_source=runtime_request.get("route_source") or "global",
|
||||
confirmed_runtime_lock=lock_active, **agent_overrides)
|
||||
confirmed_runtime_lock=lock_active,
|
||||
# #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)
|
||||
return {
|
||||
"gateway_session_key": gateway_session_key, "session_id": session_id, "body": body,
|
||||
"user_message": user_message, "runtime_request": runtime_request,
|
||||
@@ -3536,12 +3540,20 @@ class APIServerAdapter(OpenAICompatRoutesMixin, BasePlatformAdapter):
|
||||
@staticmethod
|
||||
def _bind_api_server_session(
|
||||
*, chat_id: str = "", session_key: str = "", session_id: str = "",
|
||||
browser_control_principal: str = "", browser_control_transport_family: str = "") -> list:
|
||||
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).
|
||||
|
||||
``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.
|
||||
"""
|
||||
from gateway.session_context import set_session_vars
|
||||
@@ -3549,7 +3561,7 @@ class APIServerAdapter(OpenAICompatRoutesMixin, BasePlatformAdapter):
|
||||
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="")
|
||||
async_delivery=False, cron_session="", wake_capable=wake_capable)
|
||||
|
||||
def _turn_runtime_metadata(
|
||||
self, agent: Any, *, route: Optional[Dict[str, Any]], requested_runtime: Optional[Dict[str, Any]],
|
||||
@@ -3623,11 +3635,15 @@ class APIServerAdapter(OpenAICompatRoutesMixin, BasePlatformAdapter):
|
||||
requested_provider: Optional[str] = None, model_options: Optional[Dict[str, Any]] = None,
|
||||
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) -> tuple:
|
||||
confirmed_runtime_lock: bool = False, bind_declared_conversation: bool = False,
|
||||
wake_capable: 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."""
|
||||
provider/model must match or the turn fails; ``runtime`` metadata is attached.
|
||||
``wake_capable`` 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()
|
||||
# ContextVars do not follow run_in_executor threads: capture here, re-enter in _run().
|
||||
request_profile = _api_request_profile.get()
|
||||
@@ -3641,7 +3657,8 @@ class APIServerAdapter(OpenAICompatRoutesMixin, BasePlatformAdapter):
|
||||
chat_id=session_id or "", session_key=gateway_session_key or session_id or "",
|
||||
session_id=session_id or "",
|
||||
browser_control_principal=request_browser_control_principal,
|
||||
browser_control_transport_family=request_browser_control_transport_family)
|
||||
browser_control_transport_family=request_browser_control_transport_family,
|
||||
wake_capable=wake_capable)
|
||||
agent = None
|
||||
try:
|
||||
agent = self._create_agent(
|
||||
|
||||
@@ -494,7 +494,13 @@ class OpenAICompatRoutesMixin:
|
||||
run_kwargs = dict(
|
||||
user_message=user_message, conversation_history=history,
|
||||
ephemeral_system_prompt=system_prompt, session_id=session_id,
|
||||
gateway_session_key=gateway_session_key, **agent_overrides, route=route)
|
||||
gateway_session_key=gateway_session_key, **agent_overrides, route=route,
|
||||
# #98619: only an explicitly provided X-Hermes-Session-Id is wake-capable (the
|
||||
# header is 403-gated on API_SERVER_KEY, so the wake self-post can authenticate
|
||||
# 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 ""))
|
||||
if stream:
|
||||
_stream_q = ThreadSafeAsyncQueue()
|
||||
# tool_call_ids with an emitted "running": a "completed" without one (internal/
|
||||
|
||||
@@ -486,7 +486,12 @@ def _run_agent_sync(self, run: _RunLaunch, agent, approval_notify, *, _api_serve
|
||||
session_tokens = self._bind_api_server_session(
|
||||
chat_id=session_id or "", session_key=run.approval_session_key, session_id=session_id or "",
|
||||
browser_control_principal=run.browser_control_principal,
|
||||
browser_control_transport_family=run.browser_control_transport_family)
|
||||
browser_control_transport_family=run.browser_control_transport_family,
|
||||
# #98619 audited opt-in: every /v1/runs session id is one the client can address
|
||||
# again (body session_id, a response chain id, a declared session key, or the
|
||||
# run_id fallback returned by the create response) — and a wake-injected turn
|
||||
# lands in session history a later run on the same id will load.
|
||||
wake_capable="1")
|
||||
if session_tokens:
|
||||
resets.append((session_tokens, clear_session_vars))
|
||||
if run.agent_kwargs["room_dispatch"] is not None:
|
||||
|
||||
@@ -52,6 +52,18 @@ _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)
|
||||
|
||||
# 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)
|
||||
_CRON_AUTO_DELIVER_CHAT_ID = ContextVar("HERMES_CRON_AUTO_DELIVER_CHAT_ID", default=_UNSET)
|
||||
@@ -115,10 +127,17 @@ 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,
|
||||
) -> 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)."""
|
||||
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:
|
||||
``"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
|
||||
cannot grant wake authority."""
|
||||
global _session_context_engaged
|
||||
_session_context_engaged = True
|
||||
values = (
|
||||
@@ -128,6 +147,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))
|
||||
_runtime_cwd("set_session_cwd", cwd)
|
||||
return tokens
|
||||
|
||||
@@ -135,10 +155,13 @@ def set_session_vars(
|
||||
def clear_session_vars(tokens: list) -> None:
|
||||
"""Mark session context variables as explicitly cleared (``""``, not ``_UNSET``), so
|
||||
``get_session_env`` returns empty instead of stale ``os.environ`` values. Async-delivery
|
||||
goes back to ``_UNSET``: a cleared context is default-supported, not opted-out."""
|
||||
goes back to ``_UNSET``: a cleared context is default-supported, not opted-out. Wake
|
||||
capability goes back to ``_UNSET`` too — but for the opposite reason: a cleared context has
|
||||
declared nothing, and an undeclared capability FAILS CLOSED (#98619)."""
|
||||
for var in _SESSION_VARS:
|
||||
var.set("")
|
||||
_SESSION_ASYNC_DELIVERY.set(_UNSET)
|
||||
_SESSION_WAKE_CAPABLE.set(_UNSET)
|
||||
_runtime_cwd("clear_session_cwd")
|
||||
|
||||
|
||||
@@ -146,10 +169,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`` (outside ``_VAR_MAP``) is reset explicitly too."""
|
||||
identity. ``_SESSION_ASYNC_DELIVERY`` and ``_SESSION_WAKE_CAPABLE`` (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)
|
||||
_runtime_cwd("clear_session_cwd")
|
||||
|
||||
|
||||
@@ -198,3 +223,12 @@ def async_delivery_supported() -> bool:
|
||||
return False
|
||||
value = _SESSION_ASYNC_DELIVERY.get()
|
||||
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"
|
||||
|
||||
@@ -2385,6 +2385,40 @@ 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"] == ""
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# X-Hermes-Session-Key header (long-term memory scoping)
|
||||
|
||||
@@ -43,6 +43,7 @@ class _Batch:
|
||||
origin_ui_session_id: str
|
||||
origin_owner_transport: Any
|
||||
origin_owner_session_record: Any
|
||||
origin_wake_capable: 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
|
||||
@@ -66,17 +67,22 @@ def _announce_batch(parent_agent, n_tasks: int, live_deleg_id: Optional[str]) ->
|
||||
_hdr = f" 🔀 [{format_batch_tag(live_deleg_id, parent_agent)}] delegating {n_tasks} tasks"
|
||||
_print_completion_line(parent_agent, getattr(parent_agent, "_delegate_spinner", None), _hdr, console_line=_hdr)
|
||||
|
||||
def _capture_origin() -> tuple[str, str, Any, Any]:
|
||||
"""``(wake_sid, ui_session_id, owner_transport, owner_session_record)`` of the
|
||||
def _capture_origin() -> tuple[str, str, Any, Any, bool]:
|
||||
"""``(wake_sid, ui_session_id, owner_transport, owner_session_record, wake_capable)`` of the
|
||||
ORIGINATING session, captured BEFORE building any child: AIAgent construction
|
||||
clobbers the HERMES_SESSION_ID ContextVar/os.environ with the subagent's id."""
|
||||
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
|
||||
reason — and fails closed: a binding that never declared it (or a read error) leaves the
|
||||
session treated as non-wake-capable (#98619)."""
|
||||
from tools.async_delegation import _current_origin_session_id
|
||||
_origin_wake_sid = _current_origin_session_id()
|
||||
_origin_ui_session_id = ""
|
||||
_origin_wake_capable = False
|
||||
with _quiet(None):
|
||||
from gateway.session_context import get_session_env
|
||||
from gateway.session_context import get_session_env, wake_capable_session
|
||||
_origin_ui_session_id = get_session_env("HERMES_UI_SESSION_ID", "")
|
||||
return (_origin_wake_sid, _origin_ui_session_id, *_capture_gateway_steer_authority(_origin_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)
|
||||
|
||||
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
|
||||
@@ -204,14 +210,16 @@ 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) -> Optional[str]:
|
||||
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.
|
||||
|
||||
Finite sessions (stateless HTTP requests, one-shot Kanban workers) cannot route a detached result back after their
|
||||
turn/process ends — but if a raw session id is bound (the API server always binds one), gateway.wake can still
|
||||
reach it by self-POSTing /v1/chat/completions, so only fall back to sync when there is truly no session id to
|
||||
wake. Uses the origin captured BEFORE child construction — HERMES_SESSION_ID here would be the subagent's internal
|
||||
id.
|
||||
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.
|
||||
"""
|
||||
try:
|
||||
# Finite sessions cannot route a detached subagent result back to the agent after their turn/process
|
||||
@@ -223,13 +231,19 @@ def _resolve_async_wake_sid(origin_wake_sid: str) -> Optional[str]:
|
||||
return ""
|
||||
except Exception:
|
||||
return ""
|
||||
if origin_wake_sid:
|
||||
if origin_wake_sid and origin_wake_capable:
|
||||
logger.info(
|
||||
"delegate_task: async delivery unsupported on this session, but a 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: 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,
|
||||
)
|
||||
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,
|
||||
)
|
||||
return None
|
||||
|
||||
def _resolve_async_session_key(parent_agent: Any, origin_ui_session_id: str) -> tuple[str, str]:
|
||||
@@ -366,7 +380,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)
|
||||
wake_sid = _resolve_async_wake_sid(batch.origin_wake_sid, batch.origin_wake_capable)
|
||||
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")
|
||||
|
||||
Reference in New Issue
Block a user