fix(realtime): resolve API Server session handoff for brokered Hermes turns

The Realtime Agent's brokered Hermes path (hermes_run_task) could fail
two ways when reaching back to the API Server:

- a caller-supplied chat_session_id from another session namespace (the
  gateway/client session store) was passed straight to
  /api/sessions/{id}/chat/stream and rejected with 404 session_not_found
- _create_session() only read a flat id/session_id, but the current API
  Server returns the session nested under {"session": {"id": ...}}, so
  creation raised "Hermes API created a session without an id"

stream_task() now tracks whether it owns the API Server session and, on a
404 session_not_found for a caller-supplied id, mints a fresh API Server
session (emitting a session.bound handoff event) and retries the turn
once — a session it created itself, or a second failure, is not retried,
so there is no loop. Valid existing API sessions are reused untouched.
_create_session() parses both the nested and legacy flat response shapes.

Adds plugin/tests/test_hermes_tool_broker.py (13) covering both parsers
and the namespace-mismatch handoff/retry against a local aiohttp fake
API Server.

Closes #101

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
This commit is contained in:
Bailey Dixon
2026-06-21 21:47:57 -04:00
co-authored by Claude Opus 4.8
parent 0aa1b38a18
commit f6b965a97c
4 changed files with 351 additions and 27 deletions
+1
View File
@@ -32,6 +32,7 @@ The format is based on [Keep a Changelog](https://keepachangelog.com/), and this
- **Hold-to-talk no longer releases on accidental drift (Android).** The mic button holds until the finger genuinely lifts, instead of cancelling when it drifts off the button.
- **Voice overlay is readable (Android).** The voice dropdown panel and its status bubbles are opaque (no bleed-through), and the Focus/Overlay/Exit labels no longer wrap to two lines; invalid engine/route combinations are no longer selectable.
- **Connection status overlay clears faster (Android).** Resolved (error/warning) connection toasts auto-dismiss within ~5s instead of lingering.
- **Realtime Agent: brokered Hermes turns no longer fail with `session_not_found`.** When the Realtime Agent reached back to Hermes for context or tool work, it could hand the API Server a session id that belonged to a different session namespace (the gateway/client store), which the API Server rejected. The broker now mints a valid API Server session and retries the turn once when that happens, and reuses an existing API Server session when the id is already valid. It also reads the API Server's current nested `{"session": {"id": …}}` create-session response (previously only the legacy flat shape), so session creation no longer errored with "created a session without an id." Provider-native turns are unaffected.
## [1.2.0] - 2026-06-20
+12
View File
@@ -1,5 +1,17 @@
# Hermes-Relay — Dev Log
## 2026-06-21 — Realtime Agent API Server session handoff (issue #101)
**Why.** The Realtime Agent's brokered Hermes path (`hermes_run_task`) could fail two ways when reaching back to the API Server. (1) A caller-supplied `chat_session_id` that originated in a different session namespace (the gateway/client session store) was passed straight to `POST /api/sessions/{id}/chat/stream`, which the API Server rejects with `404 session_not_found`. (2) `_create_session()` only read a flat `id`/`session_id`, but the current API Server returns the created session nested under `{"object":"hermes.session","session":{"id":"api_…"}}` — so creation raised "Hermes API created a session without an id."
**Verified against upstream first.** `gateway/platforms/api_server.py` confirms the contract: create-session returns the nested `session` object at status 201 (`_session_response`, line ~1426); `_get_existing_session_or_404` emits `{"error":{"code":"session_not_found"}}` at 404 (line ~1349). Coded to the verified shapes, not the docs.
- **`hermes_tool_broker.py` — nested create-session parse.** Extracted `_session_id_from_create_response()` that accepts top-level `id`/`session_id` *and* nested `session.id`/`session.session_id`, preferring the flat form for back-compat with older/partial builds. `_create_session()` now delegates to it.
- **`hermes_tool_broker.py` — `session_not_found` handoff + single retry.** `stream_task()` tracks whether it owns the API Server session (`api_session_owned`). When a caller-supplied id 404s with `session_not_found` (matched by `_is_session_not_found()`, structured-or-substring), the broker mints a fresh API Server session, emits a second `hermes.session.bound` event with `reason: "session_not_found_handoff"` (so the orchestrator rebinds `session.chat_session_id`), and retries the chat/stream POST once. A session the broker created itself, or a second failure, is not retried — no loop. The 404 is raised before any SSE bytes stream, so the retry never double-emits chat content. Valid existing API sessions are reused untouched.
- **Tests.** New `plugin/tests/test_hermes_tool_broker.py` (13): pure-function coverage for both parsers (nested/flat/precedence/empty, 404-only `session_not_found` detection) plus end-to-end `stream_task` against a local aiohttp `TestServer` fake API Server — no-id-creates-session, existing-session-reused, namespace-mismatch handoff+retry, and single-retry-then-give-up. `aioresponses` isn't installed, so the tests drive the real aiohttp client path against a local server (the repo's existing pattern).
**Verification.** `python -m unittest plugin.tests.test_hermes_tool_broker` → 13/13 green. `plugin.tests.test_realtime_agent_routes` → 34/34 green (no regression). Server-side only; no Android/CLI changes.
## 2026-06-21 — Profile lock + voice fixes (orchestration batch)
**Why.** User-requested batch (TODO User-Added) covering the profile-lock setting and the concrete voice TODOs. Investigated and implemented via a planning→implementation orchestration pass: four read-only investigators, then three disjoint file-ownership implementation lanes. All changes are client-side Kotlin; the server-side realtime-voice half is deferred to TODO. **Unbuilt at time of writing — pending Studio build + `./gradlew lint`.**
@@ -43,8 +43,14 @@ class HermesToolBroker:
try:
async with aiohttp.ClientSession(timeout=timeout) as http:
session_id = request.session_id
# Track whether *we* own the API Server session. A caller-supplied
# session_id may come from a different namespace (gateway/client
# session store) and not exist in the API Server, in which case the
# chat/stream call 404s and we resolve to a fresh API session below.
api_session_owned = False
if not session_id:
session_id = await self._create_session(http, request, headers)
api_session_owned = True
yield {
"type": "hermes.session.bound",
"session_id": session_id,
@@ -65,32 +71,55 @@ class HermesToolBroker:
body["system_message"] = interface_system_message
if request.profile and request.profile != "default":
body["profile"] = request.profile
url = f"{self.webapi_url}/api/sessions/{session_id}/chat/stream"
async with http.post(
url,
json=body,
headers={**headers, "Accept": "text/event-stream"},
) as resp:
if resp.status in (401, 403):
await resp.read()
yield _auth_error_event(
resp.status,
session_id=session_id,
has_bearer=bool(headers.get("Authorization")),
)
return
if resp.status != 200:
text = await resp.text()
yield {
"type": "voice.error",
"message": f"Hermes API error ({resp.status}): {text[:240]}",
"session_id": session_id,
}
return
async for event in _iter_sse_events(resp):
mapped = _map_sse_event(event, session_id)
if mapped is not None:
yield mapped
resolved_via_handoff = False
while True:
url = f"{self.webapi_url}/api/sessions/{session_id}/chat/stream"
async with http.post(
url,
json=body,
headers={**headers, "Accept": "text/event-stream"},
) as resp:
if resp.status in (401, 403):
await resp.read()
yield _auth_error_event(
resp.status,
session_id=session_id,
has_bearer=bool(headers.get("Authorization")),
)
return
if resp.status != 200:
text = await resp.text()
# The supplied session id was created outside the API
# Server (e.g. a gateway/client session). Resolve once
# to a fresh API Server session and retry the turn.
if (
not api_session_owned
and not resolved_via_handoff
and _is_session_not_found(resp.status, text)
):
resolved_via_handoff = True
session_id = await self._create_session(
http, request, headers
)
api_session_owned = True
yield {
"type": "hermes.session.bound",
"session_id": session_id,
"profile": request.profile,
"reason": "session_not_found_handoff",
}
continue
yield {
"type": "voice.error",
"message": f"Hermes API error ({resp.status}): {text[:240]}",
"session_id": session_id,
}
return
async for event in _iter_sse_events(resp):
mapped = _map_sse_event(event, session_id)
if mapped is not None:
yield mapped
break
yield {
"type": "hermes.run.completed",
@@ -183,7 +212,7 @@ class HermesToolBroker:
headers=resp.headers,
)
data = await resp.json()
session_id = str(data.get("id") or data.get("session_id") or "").strip()
session_id = _session_id_from_create_response(data)
if not session_id:
raise aiohttp.ClientError("Hermes API created a session without an id")
return session_id
@@ -199,6 +228,49 @@ def _headers(bearer_token: str | None) -> dict[str, str]:
return headers
def _session_id_from_create_response(data: Any) -> str:
"""Extract a session id from the API Server's create-session response.
The current Hermes API Server returns the session nested under
``{"object": "hermes.session", "session": {"id": "api_..."}}`` while older
or partial builds returned a flat ``{"id": ...}`` / ``{"session_id": ...}``.
Accept both so the broker keeps working across server versions.
"""
if not isinstance(data, dict):
return ""
session_obj = data.get("session") if isinstance(data.get("session"), dict) else {}
return str(
data.get("id")
or data.get("session_id")
or session_obj.get("id")
or session_obj.get("session_id")
or ""
).strip()
def _is_session_not_found(status: int, body_text: str) -> bool:
"""True when a chat/stream response is the API Server's 404 ``session_not_found``.
The API Server emits ``{"error": {"code": "session_not_found", ...}}`` (see
upstream ``_get_existing_session_or_404``). Fall back to substring matching so
a non-JSON body or a slightly different shape is still recognised.
"""
if status != 404:
return False
try:
payload = json.loads(body_text)
except (json.JSONDecodeError, TypeError):
payload = None
if isinstance(payload, dict):
error = payload.get("error")
if isinstance(error, dict):
if str(error.get("code") or "").strip().lower() == "session_not_found":
return True
if "session not found" in str(error.get("message") or "").lower():
return True
return "session_not_found" in body_text or "session not found" in body_text.lower()
def _auth_error_event(
status: int,
*,
+239
View File
@@ -0,0 +1,239 @@
"""Regression tests for the realtime-agent Hermes session/tool broker.
Covers issue #101:
* ``_create_session`` must parse both the legacy flat shape and the current
nested ``{"session": {"id": ...}}`` API Server response.
* ``stream_task`` must recover from a ``404 session_not_found`` when a
caller-supplied ``session_id`` was created outside the API Server, by
minting a fresh API Server session and retrying the turn once.
These exercise the real aiohttp client path against a local fake API Server
(the repo does not ship ``aioresponses``, and the broker creates its own
``aiohttp.ClientSession``).
"""
from __future__ import annotations
import unittest
import aiohttp
from aiohttp import web
from aiohttp.test_utils import TestServer
from plugin.relay.realtime_agent.hermes_tool_broker import (
HermesTaskRequest,
HermesToolBroker,
_is_session_not_found,
_session_id_from_create_response,
)
class FakeHermesApiServer:
"""Minimal stand-in for the Hermes API Server session/chat surface."""
def __init__(self) -> None:
self.known_sessions: set[str] = set()
self.create_bodies: list[dict] = []
self.chat_calls: list[str] = []
self.create_shape = "nested" # nested | flat_id | flat_session_id
self.next_created_id = "api_1700000000_abcd1234"
def _create_payload(self, session_id: str) -> dict:
if self.create_shape == "flat_id":
return {"id": session_id}
if self.create_shape == "flat_session_id":
return {"session_id": session_id}
return {"object": "hermes.session", "session": {"id": session_id}}
async def create_session(self, request: web.Request) -> web.Response:
self.create_bodies.append(await request.json())
session_id = self.next_created_id
self.known_sessions.add(session_id)
return web.json_response(self._create_payload(session_id), status=201)
async def chat_stream(self, request: web.Request) -> web.StreamResponse:
session_id = request.match_info["session_id"]
self.chat_calls.append(session_id)
if session_id not in self.known_sessions:
return web.json_response(
{
"error": {
"message": f"Session not found: {session_id}",
"code": "session_not_found",
}
},
status=404,
)
resp = web.StreamResponse(
status=200, headers={"Content-Type": "text/event-stream"}
)
await resp.prepare(request)
await resp.write(b'event: assistant.delta\ndata: {"delta": "hello"}\n\n')
await resp.write(b"event: run.completed\ndata: {}\n\n")
await resp.write_eof()
return resp
def _make_app(server: FakeHermesApiServer) -> web.Application:
app = web.Application()
app.router.add_post("/api/sessions", server.create_session)
app.router.add_post(
"/api/sessions/{session_id}/chat/stream", server.chat_stream
)
return app
class SessionIdParseTest(unittest.TestCase):
"""Pure-function coverage for the response-shape parsers."""
def test_nested_session_id(self) -> None:
data = {"object": "hermes.session", "session": {"id": "api_123"}}
self.assertEqual(_session_id_from_create_response(data), "api_123")
def test_nested_session_session_id_key(self) -> None:
data = {"session": {"session_id": "api_456"}}
self.assertEqual(_session_id_from_create_response(data), "api_456")
def test_flat_id(self) -> None:
self.assertEqual(_session_id_from_create_response({"id": "api_flat"}), "api_flat")
def test_flat_session_id(self) -> None:
self.assertEqual(
_session_id_from_create_response({"session_id": "api_flat2"}), "api_flat2"
)
def test_top_level_id_wins_over_nested(self) -> None:
data = {"id": "api_top", "session": {"id": "api_nested"}}
self.assertEqual(_session_id_from_create_response(data), "api_top")
def test_missing_id_returns_empty(self) -> None:
self.assertEqual(_session_id_from_create_response({"object": "hermes.session"}), "")
self.assertEqual(_session_id_from_create_response({"session": {}}), "")
self.assertEqual(_session_id_from_create_response(None), "")
def test_session_not_found_detection(self) -> None:
body = '{"error": {"message": "Session not found: x", "code": "session_not_found"}}'
self.assertTrue(_is_session_not_found(404, body))
# Only a 404 should ever be treated as a handoff trigger.
self.assertFalse(_is_session_not_found(500, body))
# Plain-text body still recognised.
self.assertTrue(_is_session_not_found(404, "Session not found: abc"))
# Unrelated 404 is not a session handoff.
self.assertFalse(_is_session_not_found(404, '{"error": {"code": "not_found"}}'))
class HermesToolBrokerStreamTest(unittest.IsolatedAsyncioTestCase):
async def asyncSetUp(self) -> None:
self.fake = FakeHermesApiServer()
self.server = TestServer(_make_app(self.fake))
await self.server.start_server()
base_url = str(self.server.make_url("")).rstrip("/")
self.broker = HermesToolBroker(base_url)
async def asyncTearDown(self) -> None:
await self.server.close()
async def _collect(self, request: HermesTaskRequest) -> list[dict]:
return [event async for event in self.broker.stream_task(request)]
async def test_create_session_parses_nested_shape(self) -> None:
async with aiohttp.ClientSession() as http:
session_id = await self.broker._create_session(
http,
HermesTaskRequest(text="hi", profile=None, session_id=None),
{},
)
self.assertEqual(session_id, self.fake.next_created_id)
async def test_create_session_parses_flat_shape(self) -> None:
self.fake.create_shape = "flat_id"
async with aiohttp.ClientSession() as http:
session_id = await self.broker._create_session(
http,
HermesTaskRequest(text="hi", profile=None, session_id=None),
{},
)
self.assertEqual(session_id, self.fake.next_created_id)
async def test_no_session_id_creates_api_session(self) -> None:
events = await self._collect(
HermesTaskRequest(text="hi", profile=None, session_id=None)
)
bound = [e for e in events if e["type"] == "hermes.session.bound"]
self.assertEqual(len(bound), 1)
self.assertEqual(bound[0]["session_id"], self.fake.next_created_id)
deltas = [e for e in events if e["type"] == "voice.response.delta"]
self.assertEqual("".join(e["delta"] for e in deltas), "hello")
self.assertEqual(self.fake.chat_calls, [self.fake.next_created_id])
self.assertEqual(len(self.fake.create_bodies), 1)
async def test_existing_api_session_is_reused(self) -> None:
self.fake.known_sessions.add("api_existing")
events = await self._collect(
HermesTaskRequest(text="hi", profile=None, session_id="api_existing")
)
# No session creation, no handoff bound event.
self.assertEqual(self.fake.create_bodies, [])
self.assertEqual(
[e for e in events if e["type"] == "hermes.session.bound"], []
)
self.assertEqual(self.fake.chat_calls, ["api_existing"])
deltas = [e for e in events if e["type"] == "voice.response.delta"]
self.assertEqual("".join(e["delta"] for e in deltas), "hello")
async def test_session_not_found_triggers_handoff_and_retry(self) -> None:
# "gateway-xyz" is from another namespace and not in the API Server store.
events = await self._collect(
HermesTaskRequest(text="hi", profile=None, session_id="gateway-xyz")
)
# First chat hit the stale id, then the freshly minted API session.
self.assertEqual(
self.fake.chat_calls, ["gateway-xyz", self.fake.next_created_id]
)
self.assertEqual(len(self.fake.create_bodies), 1)
bound = [e for e in events if e["type"] == "hermes.session.bound"]
self.assertEqual(len(bound), 1)
self.assertEqual(bound[0]["session_id"], self.fake.next_created_id)
self.assertEqual(bound[0].get("reason"), "session_not_found_handoff")
# The turn still completes after the handoff.
deltas = [e for e in events if e["type"] == "voice.response.delta"]
self.assertEqual("".join(e["delta"] for e in deltas), "hello")
self.assertTrue(
any(e["type"] == "hermes.run.completed" for e in events)
)
# No stray voice.error surfaced to the user.
self.assertEqual([e for e in events if e["type"] == "voice.error"], [])
async def test_handoff_retried_only_once(self) -> None:
# Pathological case: the API Server never recognises any session (even the
# one we just minted). The broker must give up after a single resolution
# attempt rather than loop forever creating sessions.
fake = FakeHermesApiServer()
async def always_missing(request: web.Request) -> web.StreamResponse:
fake.chat_calls.append(request.match_info["session_id"])
return web.json_response(
{"error": {"code": "session_not_found"}}, status=404
)
fake.chat_stream = always_missing # type: ignore[assignment]
server = TestServer(_make_app(fake))
await server.start_server()
try:
broker = HermesToolBroker(str(server.make_url("")).rstrip("/"))
events = [
e
async for e in broker.stream_task(
HermesTaskRequest(text="hi", profile=None, session_id="gateway-xyz")
)
]
finally:
await server.close()
# Exactly two chat attempts: original + one handoff retry, then give up.
self.assertEqual(len(fake.chat_calls), 2)
self.assertEqual(len(fake.create_bodies), 1)
self.assertTrue(any(e["type"] == "voice.error" for e in events))
if __name__ == "__main__":
unittest.main()