adapter.py was a 3,493-line monolith. Split it into focused modules with clear separation of responsibilities, bringing it down to ~857 lines: - Module-level helpers: hooks, classify, pickers, commands, setup, defaults, secrets - Frame-handler mixins: inbound, tool_frames, push_frames, media_frames, picker_frames, channel_frames, query_frames - mixin_base: IrisAdapterBase (declaration-only base for shared attrs) - adapter.py now holds only IrisAdapter (the composition of the 7 mixins + BasePlatformAdapter), register(), and test-facing re-exports The mixins come before BasePlatformAdapter in the MRO so their methods override the base; super() calls (e.g. send_image) still resolve to BasePlatformAdapter. No circular imports; dispatch.py and http_server.py (instance-method callers) are unaffected. Ruff complexity ceilings (PLR0911/0912/0913/0915) restored to Ruff's built-in defaults (12/50/6/5) instead of "just above the current maxima", which ratchets the bar down as code grows. The existing genuinely-complex functions (frame builders mirroring the wire schema, the QR matrix builder, the dispatch table) carry an explicit `# noqa: PLR09xx` marking them as reviewed, frozen exceptions; new code is held to the default ceilings. All 125 tests green (94 test_android + 31 test_android_http); no new ruff errors introduced.
353 lines
14 KiB
Python
353 lines
14 KiB
Python
"""Plugin hook capture: reasoning, tool results, runtime metadata.
|
|
|
|
The gateway's plugin hooks (``on_stream_delta``, ``post_tool_call``,
|
|
``post_api_request``) fire on gateway worker threads; these module-level
|
|
buffers accumulate the per-turn data the adapter attaches to outbound frames
|
|
(reasoning on ``message.stop``, tool output on ``tool.end``, the ``runtime``
|
|
footer on final messages). A personal iris gateway serves one active turn at
|
|
a time, so global buffers suffice; each is reset at the turn boundary.
|
|
"""
|
|
|
|
import asyncio
|
|
import contextlib
|
|
import json
|
|
import logging
|
|
import os
|
|
import threading
|
|
import time
|
|
from collections import deque
|
|
from typing import TYPE_CHECKING, Any
|
|
|
|
from . import protocol
|
|
|
|
if TYPE_CHECKING: # pragma: no cover - typing only
|
|
from .adapter import IrisAdapter
|
|
|
|
logger = logging.getLogger(__name__)
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# M2 — reasoning capture (streaming)
|
|
#
|
|
# The gateway streams only ``content`` to the platform and suppresses the
|
|
# final send (which would carry the prepended reasoning), so the model's
|
|
# separate ``reasoning_content`` is otherwise lost in the streaming case.
|
|
# hermes exposes a plugin ``on_stream_delta`` hook that fires reasoning
|
|
# deltas with ``kind="reasoning"`` (gated by ``plugins.stream_reasoning_deltas``).
|
|
# We accumulate those deltas here and attach the result to the turn's
|
|
# ``message.stop`` frame. Single-chat for now (the default home channel), so a
|
|
# module-level buffer suffices; it is reset at each turn start.
|
|
# ---------------------------------------------------------------------------
|
|
|
|
_reasoning_parts: list[str] = []
|
|
_reasoning_lock = threading.Lock()
|
|
# Barrier: set by the hook worker once it has processed the first content
|
|
# delta (kind="text"). The worker drains a FIFO queue and reasoning deltas are
|
|
# enqueued before content deltas, so at that point every reasoning delta has
|
|
# already been appended -- a reliable "reasoning flushed" signal that avoids
|
|
# racing message.stop against the async hook thread.
|
|
_reasoning_flushed = threading.Event()
|
|
|
|
|
|
def _on_stream_delta(**kwargs: Any) -> None:
|
|
"""Plugin hook: capture reasoning deltas (kind="reasoning")."""
|
|
kind = kwargs.get("kind")
|
|
if kind == "reasoning":
|
|
delta = kwargs.get("delta") or ""
|
|
if delta:
|
|
with _reasoning_lock:
|
|
_reasoning_parts.append(delta)
|
|
elif kind == "text":
|
|
_reasoning_flushed.set()
|
|
|
|
|
|
async def _wait_for_reasoning_flushed(timeout: float = 0.3) -> None:
|
|
"""Wait (without blocking the event loop) until the hook worker has
|
|
processed all reasoning deltas, or *timeout* seconds elapse."""
|
|
loop = asyncio.get_running_loop()
|
|
deadline = loop.time() + timeout
|
|
while loop.time() < deadline:
|
|
if _reasoning_flushed.is_set():
|
|
return
|
|
await asyncio.sleep(0.01)
|
|
|
|
|
|
def _take_reasoning() -> str:
|
|
"""Drain and return the accumulated reasoning (empty string if none)."""
|
|
with _reasoning_lock:
|
|
parts = _reasoning_parts[:]
|
|
_reasoning_parts.clear()
|
|
_reasoning_flushed.clear()
|
|
return "".join(parts).strip()
|
|
|
|
|
|
def _reset_reasoning() -> None:
|
|
with _reasoning_lock:
|
|
_reasoning_parts.clear()
|
|
_reasoning_flushed.clear()
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# M2 — tool-result capture (post_tool_call hook)
|
|
#
|
|
# The gateway renders tool *progress* lines to the platform but never streams
|
|
# the tool *output* (it is the agent's concern, persisted to history, not
|
|
# presentation). To let the app show the full call + result on demand
|
|
# (Settings → Tool detail), we capture each completed tool call via the
|
|
# ``post_tool_call`` hook and attach it to the ``tool.end`` frame.
|
|
#
|
|
# Global FIFO (like the reasoning buffer): a personal iris gateway serves
|
|
# one active turn at a time, and records are matched to the open tool by name
|
|
# in completion order. Bounded so a runaway turn can't grow it without limit.
|
|
# ---------------------------------------------------------------------------
|
|
|
|
_tool_results: "deque[dict[str, Any]]" = deque()
|
|
_tool_results_lock = threading.Lock()
|
|
_MAX_TOOL_RESULTS = 200
|
|
_MAX_OUTPUT_PREVIEW = 8000
|
|
|
|
|
|
def _on_post_tool_call(**kwargs: Any) -> None:
|
|
"""Plugin hook: capture a completed tool call's result + timing."""
|
|
result = kwargs.get("result")
|
|
record = {
|
|
"tool_name": kwargs.get("tool_name") or "",
|
|
"result": (str(result) if result is not None else "")[:_MAX_OUTPUT_PREVIEW],
|
|
"duration_ms": kwargs.get("duration_ms") or 0,
|
|
"status": kwargs.get("status") or "ok",
|
|
}
|
|
with _tool_results_lock:
|
|
_tool_results.append(record)
|
|
while len(_tool_results) > _MAX_TOOL_RESULTS:
|
|
_tool_results.popleft()
|
|
# Live todo list: the todo tool's result is the authoritative full list,
|
|
# so emit it the moment the call completes — the tool.end frame only
|
|
# arrives when the NEXT tool starts or the turn ends, which would lag the
|
|
# app's strip behind the agent's actual progress. Best-effort: a parse
|
|
# failure (truncated preview) or a missing live adapter just skips it.
|
|
if kwargs.get("tool_name") == "todo" and record["status"] == "ok":
|
|
todos = _parse_todo_result(record["result"])
|
|
adapter = _live_adapter
|
|
if todos is not None and adapter is not None and adapter._loop is not None:
|
|
with contextlib.suppress(Exception):
|
|
asyncio.run_coroutine_threadsafe(adapter._emit_todo_update(todos), adapter._loop)
|
|
|
|
|
|
def _take_tool_result(tool_name: str) -> dict[str, Any] | None:
|
|
"""Pop the first completed record matching *tool_name* (FIFO), else None."""
|
|
if not tool_name:
|
|
return None
|
|
with _tool_results_lock:
|
|
for i, rec in enumerate(_tool_results):
|
|
if rec["tool_name"] == tool_name:
|
|
del _tool_results[i]
|
|
return rec
|
|
return None
|
|
|
|
|
|
def _reset_tool_results() -> None:
|
|
"""Clear captured records (turn boundary — drop anything unconsumed)."""
|
|
with _tool_results_lock:
|
|
_tool_results.clear()
|
|
|
|
|
|
def _tool_end_fields(tool_name: str) -> dict[str, Any]:
|
|
"""Build the ``tool.end`` enrichment (ok/duration/output_preview) from the
|
|
captured hook record for *tool_name*; empty dict when none is available
|
|
(e.g. tool_progress off, or the call came from another session)."""
|
|
rec = _take_tool_result(tool_name)
|
|
if rec is None:
|
|
return {}
|
|
fields: dict[str, Any] = {
|
|
"ok": rec["status"] == "ok",
|
|
"output_preview": rec["result"] or None,
|
|
}
|
|
if rec["duration_ms"]:
|
|
fields["duration"] = round(rec["duration_ms"] / 1000.0, 3)
|
|
return fields
|
|
|
|
|
|
# The live adapter instance (module-level so the synchronous plugin hooks
|
|
# below can reach it). A personal iris gateway runs exactly one adapter;
|
|
# set on connect, cleared on disconnect.
|
|
_live_adapter: "IrisAdapter | None" = None
|
|
|
|
# Valid todo item statuses (hermes tools/todo_tool.py VALID_STATUSES).
|
|
_TODO_STATUSES = ("pending", "in_progress", "completed", "cancelled")
|
|
|
|
|
|
def _parse_todo_result(result: Any) -> list[dict[str, str]] | None:
|
|
"""Parse the ``todo`` tool's result into a clean item list.
|
|
|
|
The tool returns ``{"todos": [...], "summary": {...}}`` — the FULL current
|
|
list, which is authoritative even for ``merge`` writes (whose args carry
|
|
only the changed items) and read-only calls. Returns ``None`` when the
|
|
result is not a parseable todo list (error string, truncated preview, …)
|
|
so the caller skips the emission instead of broadcasting garbage.
|
|
"""
|
|
try:
|
|
data = json.loads(result)
|
|
except (TypeError, ValueError):
|
|
return None
|
|
items = data.get("todos") if isinstance(data, dict) else None
|
|
if not isinstance(items, list):
|
|
return None
|
|
todos: list[dict[str, str]] = []
|
|
for it in items:
|
|
if not isinstance(it, dict):
|
|
continue
|
|
content = str(it.get("content") or "").strip()
|
|
status = str(it.get("status") or "")
|
|
if not content or status not in _TODO_STATUSES:
|
|
continue
|
|
todos.append({"id": str(it.get("id") or ""), "content": content, "status": status})
|
|
return todos
|
|
|
|
|
|
def _tool_emoji(tool_name: str) -> str | None:
|
|
"""Cosmetic per-tool emoji for the ``tool.start`` frame.
|
|
|
|
Resolved via hermes' own display layer (``agent.display.get_tool_emoji``):
|
|
active-skin ``tool_emojis`` overrides first, then the tool registry's
|
|
per-tool ``emoji`` field — so icons track the user's hermes theme and
|
|
new/plugin tools get their registered glyph for free. Returns ``None``
|
|
when the tool is unknown (or the import fails) so the frame omits the
|
|
field and the app falls back to its own default glyph.
|
|
"""
|
|
try:
|
|
from agent.display import get_tool_emoji
|
|
|
|
return get_tool_emoji(tool_name, default="") or None
|
|
except Exception:
|
|
return None
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Runtime-metadata footer (post_api_request hook)
|
|
#
|
|
# The app renders a Telegram-style footer under final assistant messages
|
|
# (model, context %, cwd, latency, cost). Display is controlled by the APP
|
|
# (Settings → Runtime footer), not hermes config — so the gateway ALWAYS
|
|
# sends the data. hermes core only appends its own *text* footer when
|
|
# ``display.runtime_footer.enabled`` is set, and the adapter has no access to
|
|
# the gateway's ``agent_result``, so we capture the same facts ourselves via
|
|
# the ``post_api_request`` plugin hook (fires after every provider call with
|
|
# model + usage):
|
|
#
|
|
# * model — the turn's latest model (failover-aware)
|
|
# * prompt_tokens — the latest call's prompt size (context occupancy)
|
|
# * turn start — the first API call of the turn (latency baseline)
|
|
#
|
|
# Global buffer (same pattern as the reasoning/tool buffers): a personal
|
|
# iris gateway serves one active turn at a time. The hook fires for every
|
|
# platform, so we only record when the turn's platform is iris.
|
|
# ---------------------------------------------------------------------------
|
|
|
|
_runtime_meta: dict[str, Any] = {}
|
|
_runtime_meta_lock = threading.Lock()
|
|
# Per-model context-window cache. Resolution may probe endpoints on first use
|
|
# (slow); the cache is process-lifetime so each model resolves at most once.
|
|
_context_length_cache: dict[str, int] = {}
|
|
# Upper bound (seconds) on context-window resolution during a final send, so a
|
|
# slow first-use probe never delays the reply. The worker thread keeps running
|
|
# and populates the cache, so the next turn is fast.
|
|
_CTX_RESOLVE_TIMEOUT_S = 3.0
|
|
|
|
|
|
def _on_post_api_request(**kwargs: Any) -> None:
|
|
"""Plugin hook: capture per-turn runtime metadata (model, prompt tokens)."""
|
|
platform = kwargs.get("platform")
|
|
if platform and platform != "iris":
|
|
return
|
|
model = kwargs.get("model") or ""
|
|
usage = kwargs.get("usage") or {}
|
|
prompt_tokens = usage.get("prompt_tokens") or 0
|
|
with _runtime_meta_lock:
|
|
if model:
|
|
_runtime_meta["model"] = model
|
|
if prompt_tokens:
|
|
_runtime_meta["prompt_tokens"] = prompt_tokens
|
|
if "turn_start" not in _runtime_meta:
|
|
_runtime_meta["turn_start"] = time.monotonic()
|
|
|
|
|
|
def _take_runtime_meta() -> dict[str, Any]:
|
|
"""Drain the captured turn metadata (turn boundary)."""
|
|
with _runtime_meta_lock:
|
|
meta = dict(_runtime_meta)
|
|
_runtime_meta.clear()
|
|
return meta
|
|
|
|
|
|
def _resolve_context_length(model: str) -> int | None:
|
|
"""Best-effort context window for *model* (cached; None on failure).
|
|
|
|
Runs in a worker thread (may probe endpoints on first use). The cache is
|
|
populated even if the caller's asyncio task times out, so subsequent
|
|
turns resolve instantly.
|
|
"""
|
|
if not model:
|
|
return None
|
|
cached = _context_length_cache.get(model)
|
|
if cached:
|
|
return cached
|
|
try:
|
|
from agent.model_metadata import get_model_context_length
|
|
|
|
ctx = get_model_context_length(model)
|
|
if ctx and ctx > 0:
|
|
_context_length_cache[model] = int(ctx)
|
|
return int(ctx)
|
|
except Exception:
|
|
logger.debug("iris: context-length resolution failed for %s", model, exc_info=True)
|
|
return None
|
|
|
|
|
|
def _home_relative_cwd(cwd: str) -> str:
|
|
"""Collapse ``$HOME`` to ``~`` (matches hermes' runtime footer)."""
|
|
if not cwd:
|
|
return ""
|
|
try:
|
|
home = os.path.expanduser("~")
|
|
p = os.path.abspath(cwd)
|
|
if home and (p == home or p.startswith(home + os.sep)):
|
|
return "~" + p[len(home) :]
|
|
return p
|
|
except Exception:
|
|
return cwd
|
|
|
|
|
|
async def _build_runtime_footer(meta: dict[str, Any]) -> dict[str, Any]:
|
|
"""Build the ``runtime`` footer object from captured turn metadata.
|
|
|
|
Called on every final send (the app decides what to show). Fields without
|
|
data are omitted. ``meta`` is the drained turn buffer (model,
|
|
prompt_tokens, turn_start).
|
|
"""
|
|
# Lazy import: the tests monkeypatch ``adapter._resolve_context_length``,
|
|
# so resolve it through the adapter module at call time (avoids a
|
|
# top-level circular import adapter -> hooks -> adapter).
|
|
from .adapter import _resolve_context_length
|
|
|
|
model = (meta.get("model") or "").rsplit("/", 1)[-1]
|
|
prompt_tokens = meta.get("prompt_tokens") or 0
|
|
context_pct = None
|
|
if prompt_tokens and model:
|
|
try:
|
|
ctx_len = await asyncio.wait_for(
|
|
asyncio.to_thread(_resolve_context_length, model),
|
|
timeout=_CTX_RESOLVE_TIMEOUT_S,
|
|
)
|
|
except (asyncio.TimeoutError, Exception):
|
|
ctx_len = None
|
|
if ctx_len:
|
|
context_pct = round(prompt_tokens / ctx_len * 100)
|
|
turn_start = meta.get("turn_start")
|
|
latency = (time.monotonic() - turn_start) if turn_start else None
|
|
cwd = _home_relative_cwd(os.environ.get("TERMINAL_CWD", ""))
|
|
return protocol.runtime_footer(
|
|
model=model or None,
|
|
context_pct=context_pct,
|
|
cwd=cwd or None,
|
|
latency=latency,
|
|
)
|