Files
hermes-agent/tui_gateway/hosted_room_member_activity.py
teknium1 ebe8cda8ea feat(tui_gateway): real JSON-RPC server→client requests replace the *.request / *.respond notification pair
The backend never sent a JSON-RPC request; when it needed an answer from the
renderer it hand-correlated a `*.request` notification with a later `*.respond`
method through four module-level dicts, a timeout thread and 13 derived
`*.expire` names, plus a separate reconnect snapshot per prompt kind. That is a
second request/response layer built on a protocol that already has one.

`tui_gateway/server_requests.py` sends `{id: "srq-…", method, params}` and
blocks on the response frame with that id (string ids never collide with the
clients' integer ids). One `request.cancel {id, method, reason}` notification
withdraws a request on timeout / interrupt / session close. `open_requests` on
`session.resume` / `session.activate` / `session.events.since` re-delivers
unanswered requests after a reconnect; the shared TypeScript channel does that
itself before the caller sees the result. Batch clarify keeps its per-question
locks as a normal `clarify.lock` RPC (the last lock resolves the request).
Approvals stay queue-backed (`tools.approval` owns the timeout, `/approve all`,
coalescing): the request resolves the queue entry and the entry's own
resolution withdraws the request through `register_gateway_settle`.

Deleted: `_block`, `_respond`, `_pending`, `_answers`,
`_pending_prompt_payloads`, `_batch_clarify`, `_EXPIRING_REQUESTS`, the
`*.respond` methods, every `*.request` / `*.expire` event, `pending_clarify`.
Compute-host (turn isolation) mirrors the child's open request and relays the
response frame / lock to it. Desktop, TUI and shared clients register
`onRequest` handlers where they used to switch on `*.request` events; answers
are response frames over the socket the request arrived on, so #91684's
owner-routing class cannot recur for prompts.
2026-09-14 06:02:05 -07:00

70 lines
3.3 KiB
Python

"""Project a hosted room member's live runtime events to the ``on_room_member_activity`` plugin hook.
A member turn runs on a hidden ``room_plumbing`` session with no client transport, so the tool /
approval / streaming frames the turn loop already emits bottom out at stdio and are lost. Between
``turn.started`` and ``turn.settled`` in the durable room log a client sees nothing. This module
re-routes those frames, stamped with the room coordinates the session carries in
``_hosted_room_task``, to plugins — off the token path, through the same bounded per-consumer
queues the ``on_stream_*`` observers use. Nothing is written to the room log: deltas at room-log
byte budgets would exhaust a room in minutes, and checkpoint replay must stay a pure function of
the durable events.
"""
from __future__ import annotations
from collections.abc import Mapping
from typing import Any
HOOK_NAME = "on_room_member_activity"
# Session event frame ``type`` -> room activity ``kind``. Frames not listed (session.info,
# status.update, message.start/complete, ...) are session chrome, not member activity.
KIND_BY_FRAME_TYPE: Mapping[str, str] = {
"tool.start": "tool.started",
"tool.complete": "tool.completed",
"tool.output_risk": "tool.output_risk",
"message.delta": "message.delta",
"message.interim": "message.interim",
"reasoning.delta": "reasoning.delta",
"error": "turn.error",
}
_COORDINATE_FIELDS = ("room_id", "thread_id", "member_id", "turn_id", "task_id", "execution_generation")
def emit_room_member_activity(hosted_task: Mapping[str, Any], *, kind: str, payload: Mapping[str, Any] | None,
seq: int | None = None) -> bool:
"""Queue one activity event for every registered consumer; False when nobody listens."""
from agent.plugin_stream_hooks import enqueue_plugin_stream_hook
coordinates = {field: hosted_task.get(field) for field in _COORDINATE_FIELDS}
return enqueue_plugin_stream_hook(HOOK_NAME, **coordinates, kind=kind, seq=seq, payload=dict(payload or {}))
# Server→client request frames (``{"id": "srq-…", "method": …}``) that are member activity.
KIND_BY_REQUEST_METHOD: Mapping[str, str] = {"approval": "request.opened"}
def project_room_member_activity(frame: Mapping[str, Any], sessions: Mapping[str, Mapping[str, Any]]) -> bool:
"""Fire the hook for an outgoing session frame (event notification or server→client request) when its
session is running a room turn."""
params = frame.get("params")
if not isinstance(params, Mapping):
return False
if frame.get("method") == "event":
kind = KIND_BY_FRAME_TYPE.get(str(params.get("type") or ""))
payload = params.get("payload")
elif "id" in frame:
kind = KIND_BY_REQUEST_METHOD.get(str(frame.get("method") or ""))
payload = {k: v for k, v in params.items() if k != "session_id"}
else:
return False
if kind is None:
return False
session = sessions.get(str(params.get("session_id") or ""))
hosted_task = session.get("_hosted_room_task") if isinstance(session, Mapping) else None
if not isinstance(hosted_task, Mapping):
return False
return emit_room_member_activity(
hosted_task, kind=kind, payload=payload if isinstance(payload, Mapping) else None, seq=params.get("seq"))