Split adapter.py monolith into focused modules; restore Ruff complexity defaults #15

Merged
Pakobbix merged 1 commits from refactor/adapter-split into master 2026-08-24 19:03:16 +00:00
26 changed files with 3101 additions and 2775 deletions

No files matched your search

+104 -2740
View File
File diff suppressed because it is too large. Load diff
+358
View File
@@ -0,0 +1,358 @@
"""M3: channel-directory frame handlers (``channel.*`` + directory queries).
Mixin for ``adapter.IrisAdapter``. Each request is answered by broadcasting
the matching ``channel.*`` event carrying the request ``id``: the requester's
pending request completes on the id, and every other device reconciles its
local copy from the same frame (single broadcast serves as event + response).
"""
import logging
from typing import Any
from hermes_constants import get_hermes_home
from . import protocol
from . import purge as purge_bridge
from .defaults import DEFAULT_HOME_CHANNEL
from .mixin_base import IrisAdapterBase
logger = logging.getLogger(__name__)
class ChannelFrameHandlers(IrisAdapterBase):
"""Channel directory management (see module docstring)."""
async def on_channel_create(self, frame: protocol.Frame, device_id: str) -> None:
payload = frame.payload
name = payload.get("name")
if not isinstance(name, str) or not name.strip():
await self._reply(
device_id,
protocol.error(
protocol.ERR_UNSUPPORTED, "channel.create requires a name", id=frame.id
),
)
return
kind = payload.get("kind")
kind = kind if kind in ("channel", "thread") else "channel"
parent_chat_id = payload.get("parent_chat_id")
if not isinstance(parent_chat_id, str) or not parent_chat_id.strip():
parent_chat_id = None
if kind == "thread" and not parent_chat_id:
parent_chat_id = frame.chat_id or self.home_channel
try:
entry = self._channels.create(name=name, kind=kind, parent_chat_id=parent_chat_id)
except ValueError as e:
await self._reply(
device_id, protocol.error(protocol.ERR_UNSUPPORTED, str(e), id=frame.id)
)
return
resp = protocol.channel_created(entry)
resp.id = frame.id
await self._http_server.fanout(resp)
# M5: banner + push mirror (parked in the outbox when offline).
await self._broadcast_or_log(
entry["chat_id"],
protocol.notification(
entry["chat_id"],
protocol.NOTIF_CHANNEL_CREATED,
"Channels",
f"New channel: {name}",
),
)
async def on_channel_rename(self, frame: protocol.Frame, device_id: str) -> None:
chat_id = frame.chat_id or frame.payload.get("chat_id")
if not isinstance(chat_id, str) or not chat_id.strip():
await self._reply(
device_id,
protocol.error(
protocol.ERR_NOT_FOUND, "channel.rename requires chat_id", id=frame.id
),
)
return
name = frame.payload.get("name")
if not isinstance(name, str) or not name.strip():
await self._reply(
device_id,
protocol.error(
protocol.ERR_UNSUPPORTED, "channel.rename requires a name", id=frame.id
),
)
return
try:
entry = self._channels.rename(chat_id, name)
except ValueError as e:
await self._reply(
device_id, protocol.error(protocol.ERR_UNSUPPORTED, str(e), id=frame.id)
)
return
if entry is None:
await self._reply(
device_id,
protocol.error(protocol.ERR_NOT_FOUND, f"unknown chat_id {chat_id}", id=frame.id),
)
return
resp = protocol.channel_renamed(entry)
resp.id = frame.id
await self._http_server.fanout(resp)
# M5: banner + push mirror (parked in the outbox when offline).
await self._broadcast_or_log(
chat_id,
protocol.notification(
chat_id,
protocol.NOTIF_CHANNEL_RENAMED,
"Channels",
f"Renamed to {name}",
),
)
async def on_channel_set_default(self, frame: protocol.Frame, device_id: str) -> None:
chat_id = frame.chat_id or frame.payload.get("chat_id")
if not isinstance(chat_id, str) or not chat_id.strip():
await self._reply(
device_id,
protocol.error(
protocol.ERR_NOT_FOUND, "channel.set_default requires chat_id", id=frame.id
),
)
return
entry = self._channels.set_default(chat_id)
if entry is None:
await self._reply(
device_id,
protocol.error(protocol.ERR_NOT_FOUND, f"unknown chat_id {chat_id}", id=frame.id),
)
return
# Reuse the renamed event shape: it carries the full entry (incl. the
# new is_default flag) so every device reconciles the default change.
resp = protocol.channel_renamed(entry)
resp.id = frame.id
await self._http_server.fanout(resp)
async def on_channel_favorite(self, frame: protocol.Frame, device_id: str) -> None:
chat_id = frame.chat_id or frame.payload.get("chat_id")
if not isinstance(chat_id, str) or not chat_id.strip():
await self._reply(
device_id,
protocol.error(
protocol.ERR_NOT_FOUND, "channel.favorite requires chat_id", id=frame.id
),
)
return
on = bool(frame.payload.get("on"))
entry = self._channels.set_favorite(chat_id, on)
if entry is None:
await self._reply(
device_id,
protocol.error(protocol.ERR_NOT_FOUND, f"unknown chat_id {chat_id}", id=frame.id),
)
return
# Reuse the renamed event shape: it carries the full entry (incl. the
# new favorite flag) so every device reconciles the change.
resp = protocol.channel_renamed(entry)
resp.id = frame.id
await self._http_server.fanout(resp)
async def on_channel_icon(self, frame: protocol.Frame, device_id: str) -> None:
chat_id = frame.chat_id or frame.payload.get("chat_id")
if not isinstance(chat_id, str) or not chat_id.strip():
await self._reply(
device_id,
protocol.error(
protocol.ERR_NOT_FOUND, "channel.icon requires chat_id", id=frame.id
),
)
return
payload = frame.payload
icon = payload.get("icon")
icon = icon if isinstance(icon, str) and icon else None
color = payload.get("color")
color = color if isinstance(color, str) and color else None
# Guard against a runaway base64 blob (a channel icon is small).
if icon is not None and len(icon) > 512 * 1024:
await self._reply(
device_id,
protocol.error(protocol.ERR_UNSUPPORTED, "channel icon too large", id=frame.id),
)
return
entry = self._channels.set_icon(chat_id, icon, color)
if entry is None:
await self._reply(
device_id,
protocol.error(protocol.ERR_NOT_FOUND, f"unknown chat_id {chat_id}", id=frame.id),
)
return
resp = protocol.channel_renamed(entry)
resp.id = frame.id
await self._http_server.fanout(resp)
async def on_channel_set_automation(self, frame: protocol.Frame, device_id: str) -> None:
chat_id = frame.chat_id or frame.payload.get("chat_id")
if not isinstance(chat_id, str) or not chat_id.strip():
await self._reply(
device_id,
protocol.error(
protocol.ERR_NOT_FOUND, "channel.set_automation requires chat_id", id=frame.id
),
)
return
on = bool(frame.payload.get("on"))
entry = self._channels.set_automation(chat_id, on)
if entry is None:
await self._reply(
device_id,
protocol.error(
protocol.ERR_NOT_FOUND,
f"cannot set automation on {chat_id} (unknown or default)",
id=frame.id,
),
)
return
# Reuse the renamed event shape: it carries the full entry (incl. the
# new automation flag) so every device reconciles the change.
resp = protocol.channel_renamed(entry)
resp.id = frame.id
await self._http_server.fanout(resp)
async def on_channel_delete(self, frame: protocol.Frame, device_id: str) -> None:
chat_id = frame.chat_id or frame.payload.get("chat_id")
if not isinstance(chat_id, str) or not chat_id.strip():
await self._reply(
device_id,
protocol.error(
protocol.ERR_NOT_FOUND, "channel.delete requires chat_id", id=frame.id
),
)
return
entry = self._channels.delete(chat_id)
if entry is None:
await self._reply(
device_id,
protocol.error(
protocol.ERR_NOT_FOUND,
f"cannot delete {chat_id} (unknown or default)",
id=frame.id,
),
)
return
# Complete deletion: wipe the lane's history from the outbox (so
# ``history`` / ``sync`` can't resurrect it) and from the hermes
# session store (so no search trace survives). A channel delete takes
# its threads with it (thread_id=None); a thread delete is scoped to
# its parent channel + thread_id.
if entry.get("kind") == "thread":
lane_chat_id = entry.get("parent_chat_id") or chat_id
thread_id = chat_id
else:
lane_chat_id = chat_id
thread_id = None
removed_frames = self._outbox.delete_lane(lane_chat_id, thread_id=thread_id)
removed_msgs = purge_bridge.delete_lane(
get_hermes_home() / "state.db", lane_chat_id, thread_id=thread_id
)
logger.info(
"iris: channel.delete %s kind=%s outbox_frames=%s session_msgs=%s",
chat_id,
entry.get("kind"),
removed_frames,
removed_msgs,
)
resp = protocol.channel_deleted(chat_id)
resp.id = frame.id
await self._http_server.fanout(resp)
# M5: banner + push mirror (parked in the outbox when offline).
await self._broadcast_or_log(
chat_id,
protocol.notification(
chat_id,
protocol.NOTIF_CHANNEL_DELETED,
"Channels",
f"{entry.get('name') or chat_id} deleted",
),
)
async def on_channel_list(self, frame: protocol.Frame, device_id: str) -> None:
channels = self._channels.list(include_archived=False)
resp = protocol.channel_list(channels)
resp.id = frame.id
await self._reply(device_id, resp)
# ── Chat info ─────────────────────────────────────────────────────────
def _channel_name(self, chat_id: str) -> str:
"""Channel display name (M3: from the channel directory)."""
if not chat_id:
return "chat"
entry = self._channels.get(chat_id)
if entry is not None:
return entry["name"]
if chat_id in (self.home_channel, DEFAULT_HOME_CHANNEL):
return self.home_channel_name
return chat_id
async def get_chat_info(self, chat_id: str) -> dict[str, Any]:
"""Return ``{name, type, chat_id}`` for a chat (M3: directory-backed)."""
entry = self._channels.get(chat_id)
kind = entry["kind"] if entry else "channel"
return {
"name": self._channel_name(chat_id),
"type": "dm" if kind == "default" else "channel",
"chat_id": chat_id,
}
def channel_list(self) -> list[dict[str, Any]]:
"""Channel directory for ``hello.ack`` (M3: full non-archived list)."""
return self._channels.list(include_archived=False)
# ── M3: core channel-directory hook (cron / send_message name resolution)
async def list_channels(self) -> list[dict[str, Any]]:
"""Expose the directory to the gateway's core channel directory.
``gateway/channel_directory.build_channel_directory`` calls this to
populate ``channel_directory.json``, which ``resolve_channel_name``
reads for friendly-name -> chat_id resolution (cron + send_message).
Threads are addressed via the explicit ``iris:<chat>:<thread>``
syntax (see ``_parse_target_ref``), so only channels are listed here.
"""
out: list[dict[str, Any]] = []
for entry in self._channels.list(include_archived=False):
if entry["kind"] == "thread":
continue
out.append(
{
"id": entry["chat_id"],
"name": entry["name"],
"type": "dm" if entry["kind"] == "default" else "channel",
}
)
return out
# ── M3: thread handoff (gateway create_handoff_thread) ────────────────
async def create_handoff_thread(self, parent_chat_id: str, name: str) -> str | None:
"""Mint a named thread under *parent_chat_id* (gateway handoff path).
Returns the new ``thread_id`` (``t_<n>``) so the handed-off session is
isolated in its own lane, or ``None`` when the parent is unknown.
"""
parent = self._channels.get(parent_chat_id)
if parent is None:
# Unknown parent: still mint a thread under it so the handoff has a
# lane (the directory row is created lazily on first use).
parent_chat_id = parent_chat_id or self.home_channel
try:
entry = self._channels.create(
name=name or "Handoff", kind="thread", parent_chat_id=parent_chat_id
)
except Exception:
logger.warning("iris: create_handoff_thread failed", exc_info=True)
return None
await self._broadcast_both(protocol.channel_created(entry))
return entry["chat_id"]
# ---------------------------------------------------------------------------
# Plugin entry point
# ---------------------------------------------------------------------------
+335
View File
@@ -0,0 +1,335 @@
"""Outbound frame classification (M2): turn state + content heuristics.
The main gateway delivers through the legacy callback path: the stream
consumer calls ``send()`` (first bubble of a segment) and ``edit_message()``
(updates), tool progress flows through ``send()``/``edit_message()`` of an
accumulated line buffer, and interim commentary arrives as a plain
``send()``. The adapter classifies each outbound call into a structured frame
using a per-chat turn state machine + the content markers below:
* ``metadata["expect_edits"] is True`` -> streaming segment start
* ``metadata["notify"] is True`` -> final message (or fallback final)
* tool-progress line format -> tool.start / tool.end
* anything else -> commentary
Verified empirically against the live gateway with ``tests/ws_probe.py``.
"""
import json
import logging
import re
import uuid
from dataclasses import dataclass, field
from typing import Any
logger = logging.getLogger(__name__)
_STREAMING_CURSOR = " ▉"
# Code-style reasoning prefix (gateway/run.py, reasoning_style="code"):
# "💭 **Reasoning:**\n```\n<reasoning>\n```\n\n<response>"
_REASONING_PREFIX = "💭 **Reasoning:**\n```\n"
_REASONING_CLOSE = "\n```\n\n"
# A gateway tool-progress line begins with a (non-ASCII) tool emoji.
_TOOL_LINE_RE = re.compile(r"^(\S+)\s+(.+)$")
_TOOL_NAME_PREVIEW_RE = re.compile(r'^(\S+):\s*"(.*)"\s*$')
_TOOL_NAME_BARE_RE = re.compile(r"^(\S+)\.\.\.\s*$")
_TOOL_NAME_ARGS_RE = re.compile(r"^(\S+)\(([^)]*)\)\s*$")
# Terminal code block: "💻 terminal\n```\n<cmd>\n```"
_TOOL_CODEBLOCK_HEAD_RE = re.compile(r"^(\S+)\s+(\S+)\s*$")
# Reverse map of the gateway's friendly tool verbs (agent/display.py
# _TOOL_VERBS) so a verb-form line ("🔍 Searching the web for …") can be
# recovered to a structured (tool_name, preview). Longest-first matching is
# done at parse time. Verbs shared by several tools map to the most common.
_VERB_TO_TOOL: dict[str, str] = {
"Searching the web": "web_search",
"Searching files": "search_files",
"Searching past sessions": "session_search",
"Running code": "execute_code",
"Running": "terminal",
"Reading skill": "skill_view",
"Reading": "read_file",
"Writing": "write_file",
"Editing": "patch",
"Browsing": "browser_navigate",
"Clicking": "browser_click",
"Typing": "browser_type",
"Generating image": "image_generate",
"Generating video": "video_generate",
"Generating speech": "text_to_speech",
"Looking at the image": "vision_analyze",
"Listing skills": "skills_list",
"Updating skill": "skill_manage",
"Updating memory": "memory",
"Updating tasks": "todo",
"Delegating": "delegate_task",
"Scheduling": "cronjob",
"Asking": "clarify",
}
# Verbs that take a " for " connector before the preview.
_VERB_FOR_CONNECTOR = {"web_search", "search_files"}
def _mint_message_id() -> str:
return f"m_{uuid.uuid4().hex[:16]}"
def _mint_picker_id() -> str:
return f"pc_{uuid.uuid4().hex[:16]}"
def _thread_id_from_metadata(metadata: dict[str, Any] | None) -> str | None:
if not metadata:
return None
tid = metadata.get("thread_id")
if isinstance(tid, str) and tid:
return tid
return None
def _derive_thread_name(text: str) -> str:
"""Instant auto-thread name from the user's opening message (no model).
Reuses hermes' session-title derivation (``agent/title_generator.py``):
a deterministic slice of the user's own words, so the thread is named the
moment it is created. The LLM upgrade (``_schedule_thread_title_upgrade``)
replaces it moments later — the same two-stage titling hermes uses for
sessions (derived < llm < user).
"""
try:
from agent.title_generator import derive_title
title = derive_title(text)
except Exception:
logger.debug("Thread name derivation failed", exc_info=True)
title = None
return (title or "").strip() or "New thread"
def _strip_streaming_cursor(text: str) -> str:
if text and text.endswith(_STREAMING_CURSOR):
return text[: -len(_STREAMING_CURSOR)]
return text
# M5: coalesce back-to-back pushes for the same chat (a cron delivery parks
# a notification frame AND a message frame; only the first should push).
_PUSH_COALESCE_S = 5.0
def _push_preview(text: Any, limit: int = 120) -> str:
"""Short single-line preview for push bodies (lock-screen privacy: no
secrets, no full bodies -- full content arrives via ``sync``)."""
s = " ".join(str(text or "").split())
if len(s) > limit:
s = s[: limit - 1] + "…"
return s
# Cron delivery wrap (cron/scheduler.py ``_deliver_result``,
# cron.wrap_response: true):
# "Cronjob Response: <name>\n(job_id: <id>)\n-------------\n\n<content>\n\n
# To stop or manage this job, send me a new message (e.g. ...)."
_CRON_WRAP_RE = re.compile(r"^Cronjob Response: (.+?)\n\(job_id: [^)]*\)\n-+\n\n")
_CRON_FOOTER = "\n\nTo stop or manage this job"
def _cron_brief(content: str, job_id: str) -> tuple[str, str]:
"""Parse a cron delivery into ``(job_name, inner_text)``.
Falls back to ``(job_id, content)`` when the wrap is disabled
(``cron.wrap_response: false``) or unrecognised.
"""
m = _CRON_WRAP_RE.match(content or "")
if not m:
return str(job_id or "cron"), (content or "").strip()
name = m.group(1).strip()
body = content[m.end() :]
idx = body.rfind(_CRON_FOOTER)
if idx != -1:
body = body[:idx]
return name, body.strip()
def _split_reasoning(text: str) -> tuple[str | None, str]:
"""Split a code-style reasoning prefix off the front of *text*.
Returns ``(reasoning, body)``; ``reasoning`` is ``None`` when no prefix is
present (reasoning off / no reasoning / non-code style). Best-effort parse
of a stable, gateway-owned format: on any mismatch the fallback is
``(None, full text)`` so the answer still renders.
"""
if not text or not text.startswith(_REASONING_PREFIX):
return None, text
close_idx = text.find(_REASONING_CLOSE, len(_REASONING_PREFIX))
if close_idx == -1:
return None, text
reasoning = text[len(_REASONING_PREFIX) : close_idx]
body = text[close_idx + len(_REASONING_CLOSE) :]
return reasoning, body
def _parse_tool_line(line: str) -> tuple[str, str | None] | None: # noqa: PLR0911
"""Parse a single gateway tool-progress line into ``(name, preview)``.
Returns ``None`` when the line is not a tool line. The gateway formats
tool lines as ``<emoji> <name>: "<preview>"``, ``<emoji> <name>...``,
``<emoji> <name>(keys)``, or a friendly verb phrase (``<emoji> <verb> …``).
The verb form is lossy (no tool name), so we surface the verb as the name.
"""
line = line.strip()
if not line:
return None
m = _TOOL_LINE_RE.match(line)
if not m:
return None
emoji, rest = m.group(1), m.group(2)
if emoji.isascii():
return None # a tool line always leads with a non-ASCII emoji
mp = _TOOL_NAME_PREVIEW_RE.match(rest)
if mp:
return mp.group(1), mp.group(2)
mb = _TOOL_NAME_BARE_RE.match(rest)
if mb:
return mb.group(1), None
ma = _TOOL_NAME_ARGS_RE.match(rest)
if ma:
return ma.group(1), None
# Friendly verb phrase: reverse-map to (tool_name, preview).
verb_parsed = _parse_verb_phrase(rest)
if verb_parsed is not None:
return verb_parsed
# Unrecognised: use the phrase as the label.
return rest, None
def _parse_verb_phrase(phrase: str) -> tuple[str, str | None] | None:
"""Reverse-map a friendly verb phrase to ``(tool_name, preview)``.
Matches the longest verb first so "Running code" wins over "Running".
Returns ``None`` when no known verb leads the phrase.
"""
for verb in sorted(_VERB_TO_TOOL, key=len, reverse=True):
tool = _VERB_TO_TOOL[verb]
if phrase == verb:
return tool, None
if tool in _VERB_FOR_CONNECTOR and phrase.startswith(verb + " for "):
return tool, phrase[len(verb) + len(" for ") :].strip() or None
if phrase.startswith(verb + " "):
return tool, phrase[len(verb) + 1 :].strip() or None
return None
def _extract_code_block(content: str) -> str | None:
"""Return the first fenced code block's body in *content*, else ``None``.
Used to recover the terminal command from a tool-progress code block
(``<emoji> terminal`` head line + fenced command).
"""
m = re.search(r"```[^\n]*\n(.*?)\n```", content, re.DOTALL)
if m:
return m.group(1).strip() or None
return None
def _extract_verbose_args(line: str, content: str) -> dict[str, Any] | None:
"""Recover the full args dict from a verbose tool line, else ``None``.
In verbose mode the gateway renders ``<emoji> <name>(keys)`` on one line
and the full args JSON on the line that follows. When *line* is such a
header, return the parsed JSON object from the following line.
"""
parts = line.strip().split(None, 1)
# 2 == "tool name" + "args JSON" on the header line.
if len(parts) < 2 or not _TOOL_NAME_ARGS_RE.match(parts[1]): # noqa: PLR2004
return None
lines = content.splitlines()
for i, ln in enumerate(lines):
if ln.strip() != line.strip():
continue
for follow_line in lines[i + 1 :]:
follow = follow_line.strip()
if not follow:
continue
if follow.startswith("{"):
try:
obj = json.loads(follow)
return obj if isinstance(obj, dict) else None
except Exception:
return None
return None # next non-empty line is not the args JSON
return None
def _short_preview_from_args(args: dict[str, Any], cap: int = 60) -> str | None:
"""Derive a short one-line preview from a verbose args dict.
The verbose line carries no explicit preview, so the Truncated display
would otherwise lose its one-liner. Use the first non-empty string value
(whitespace-collapsed, capped) as a stand-in.
"""
if not isinstance(args, dict):
return None
for value in args.values():
if isinstance(value, str) and value.strip():
s = " ".join(value.split())
return s[: cap - 1] + "…" if len(s) > cap else s
return None
def _is_tool_progress(content: str) -> bool:
"""Heuristic: does *content* look like gateway tool-progress line(s)?
Tool progress is delivered as one or more lines, each led by a tool emoji
(or a terminal code block). Commentary is free-form prose. We classify on
the first non-empty line; subsequent lines of the same bubble are tracked
by message id, not re-classified.
"""
if not content:
return False
lines = [ln for ln in content.splitlines() if ln.strip()]
if not lines:
return False
first = lines[0].strip()
# Terminal code block: "<emoji> terminal" then a fenced command.
if len(lines) > 1 and lines[1].strip().startswith("```"):
return _TOOL_CODEBLOCK_HEAD_RE.match(first) is not None
return _parse_tool_line(first) is not None
def _is_gateway_lifecycle_notice(content: str) -> bool:
"""True for hermes gateway lifecycle notices (restart / shutdown / online).
These are system notices, not tool progress. Their leading ⚠️/♻️ emoji
would otherwise trip the tool-line heuristic and render them as a
never-completing tool card (an endless spinner, since no ``tool.end``
ever arrives for a notice that is not a real tool).
"""
if not content:
return False
c = content.strip()
return any(
marker in c for marker in ("Gateway restarting", "Gateway shutting down", "Gateway online")
)
@dataclass
class _TurnState:
"""Per-chat turn state for outbound frame classification (M2)."""
active: bool = False
# message_id of the currently streaming segment (message.start open).
stream_id: str | None = None
# message_id of the current tool-progress bubble (editable line buffer).
tool_msg_id: str | None = None
# Monotonic per-turn tool counter (start -> end correlation).
tool_index: int = 0
# Index of the most recently started tool (awaiting tool.end).
open_tool_index: int | None = None
# Name of the most recently started tool (matches the post_tool_call
# record when the tool completes, so tool.end can carry its output).
open_tool_name: str | None = None
# Tool lines already emitted as tool.start (dedup across edits).
seen_tool_lines: set = field(default_factory=set)
+62
View File
@@ -0,0 +1,62 @@
"""Slash-command catalog for the app's "/" drawer.
Derived from hermes' central ``COMMAND_REGISTRY`` (``hermes_cli/commands.py``)
— the same source the gateway help text and the Telegram command menu use —
restricted to commands available on gateway surfaces, plus plugin-registered
commands. Never raises: any import/attribute problem (code skew between the
plugin and the hermes checkout) degrades to an empty catalog, so the app's
drawer simply stays closed.
"""
import logging
from typing import Any
logger = logging.getLogger(__name__)
def _slash_command_catalog() -> list[dict[str, Any]]:
try:
from hermes_cli import commands as hermes_commands
except Exception:
logger.warning(
"iris: slash catalog unavailable (hermes_cli.commands import failed)",
exc_info=True,
)
return []
def _entry(
name: str, description: str, args_hint: str, category: str, aliases: list[str]
) -> dict[str, Any]:
return {
"name": f"/{name}",
"description": description,
"args_hint": args_hint or "",
"category": category,
"aliases": [f"/{a}" for a in aliases],
}
entries: list[dict[str, Any]] = []
try:
overrides = hermes_commands._resolve_config_gates()
for cmd in hermes_commands.COMMAND_REGISTRY:
if not hermes_commands._is_gateway_available(cmd, overrides):
continue
entries.append(
_entry(cmd.name, cmd.description, cmd.args_hint, cmd.category, list(cmd.aliases))
)
except Exception:
# Code skew: the private helpers moved. Fall back to the plain
# cli_only filter (config-gated commands are dropped, acceptable).
logger.warning("iris: slash catalog fell back to cli_only filter", exc_info=True)
entries = [
_entry(cmd.name, cmd.description, cmd.args_hint, cmd.category, list(cmd.aliases))
for cmd in hermes_commands.COMMAND_REGISTRY
if not cmd.cli_only
]
try:
for name, description, args_hint in hermes_commands._iter_plugin_command_entries():
entries.append(_entry(name, description, args_hint, "Plugin", []))
except Exception:
# Best-effort: a broken plugin-command registry should not break the
# built-in catalog, so the failure is intentionally swallowed.
logger.debug("iris: plugin command enumeration failed", exc_info=True)
return entries
+14
View File
@@ -0,0 +1,14 @@
"""Platform defaults (config.yaml ``extra`` / env fallbacks)."""
DEFAULT_HOST = "127.0.0.1"
DEFAULT_PORT = 8790
DEFAULT_HTTP_PORT = 8791 # docs/19: HTTP fallback leg
DEFAULT_HOME_CHANNEL = "default"
DEFAULT_HOME_CHANNEL_NAME = "Default"
DEFAULT_PUSH_BACKEND = "ntfy"
DEFAULT_OUTBOX_RETENTION_HOURS = 72
DEFAULT_MAX_UPLOAD_BYTES = 100 * 1024 * 1024 # 100 MB
def _truthy(value: str | None) -> bool:
return (value or "").strip().lower() in {"1", "true", "yes", "on"}
+1 -1
View File
@@ -24,7 +24,7 @@ INBOUND_BURST = 40
MAX_DEVICE_ID_LEN = 128 MAX_DEVICE_ID_LEN = 128
async def dispatch_frame(adapter: Any, frame: protocol.Frame, device_id: str) -> None: async def dispatch_frame(adapter: Any, frame: protocol.Frame, device_id: str) -> None: # noqa: PLR0912
"""Shared inbound frame dispatch (docs/19 §19.4). Unknown types are """Shared inbound frame dispatch (docs/19 §19.4). Unknown types are
ignored (forward-compat).""" ignored (forward-compat)."""
if frame.type == protocol.TYPE_MESSAGE_SEND: if frame.type == protocol.TYPE_MESSAGE_SEND:
+352
View File
@@ -0,0 +1,352 @@
"""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,
)
+3 -3
View File
@@ -480,7 +480,7 @@ class HttpServer:
# ── GET /v1/events (SSE) ────────────────────────────────────────────── # ── GET /v1/events (SSE) ──────────────────────────────────────────────
def _handle_sse(self, handler: BaseHTTPRequestHandler, device_id: str, parsed: Any) -> None: def _handle_sse(self, handler: BaseHTTPRequestHandler, device_id: str, parsed: Any) -> None: # noqa: PLR0912,PLR0915
qs = parse_qs(parsed.query) qs = parse_qs(parsed.query)
cursor = _parse_cursor(qs.get("cursor", [None])[0], handler.headers.get("Last-Event-ID")) cursor = _parse_cursor(qs.get("cursor", [None])[0], handler.headers.get("Last-Event-ID"))
# Device registration (the HTTP equivalent of the WS hello upsert): # Device registration (the HTTP equivalent of the WS hello upsert):
@@ -600,7 +600,7 @@ class HttpServer:
# ── POST /v1/media (upload, docs/19 §19.15) ─────────────────────────── # ── POST /v1/media (upload, docs/19 §19.15) ───────────────────────────
def _handle_media_upload(self, handler: BaseHTTPRequestHandler, device_id: str) -> None: def _handle_media_upload(self, handler: BaseHTTPRequestHandler, device_id: str) -> None: # noqa: PLR0911
"""Whole-file upload: metadata in headers, file bytes as the body. """Whole-file upload: metadata in headers, file bytes as the body.
Mirrors the WS ``media.upload`` contract (docs/07 §7.2) in one Mirrors the WS ``media.upload`` contract (docs/07 §7.2) in one
@@ -820,7 +820,7 @@ class _Handler(BaseHTTPRequestHandler):
return return
_send_json(self, 404, {"error": "not found"}) _send_json(self, 404, {"error": "not found"})
def do_POST(self) -> None: # noqa: N802 def do_POST(self) -> None: # noqa: N802, PLR0911, PLR0912
hs = self.server.http_server hs = self.server.http_server
if not hs.enabled: if not hs.enabled:
_send_json(self, 503, {"error": "http leg disabled"}) _send_json(self, 503, {"error": "http leg disabled"})
+255
View File
@@ -0,0 +1,255 @@
"""Inbound ``message.send`` handling (app -> agent).
Mixin for ``adapter.IrisAdapter``. Echoes the user message to all devices
(multi-device sync + ack), resolves ``media_refs`` to
``MessageEvent.media_urls``/``media_types``, applies auto-threading, and
hands the ``MessageEvent`` to ``handle_message()`` (the gateway's command
pipeline + agent turn).
"""
import asyncio
import logging
import threading
import time
import uuid
from typing import Any
from gateway.platforms.base import MessageEvent, MessageType
from . import protocol
from .classify import _derive_thread_name
from .mixin_base import IrisAdapterBase
logger = logging.getLogger(__name__)
class InboundHandlers(IrisAdapterBase):
"""Inbound message.send (see module docstring)."""
async def on_message_send(self, frame: protocol.Frame, device_id: str) -> None: # noqa: PLR0912,PLR0915
"""Handle an inbound ``message.send`` frame.
Echoes the user message to all devices (multi-device sync + ack),
then builds a ``MessageEvent`` and hands it to ``handle_message()``
(the gateway's command pipeline + agent turn).
M4: ``media_refs`` reference completed ``media.upload``s; they are
resolved to ``MessageEvent.media_urls``/``media_types`` (local paths
the agent's vision/audio tools can read) and echoed in the user
message's ``media[]`` so every device renders the attachments.
"""
payload = frame.payload
text = payload.get("text")
text = text if isinstance(text, str) else ""
refs_raw = payload.get("media_refs")
media_refs = (
[r for r in refs_raw if isinstance(r, str) and r] if isinstance(refs_raw, list) else []
)
if not text.strip() and not media_refs:
await self._reply(
device_id,
protocol.error(
protocol.ERR_UNSUPPORTED, "message.send requires non-empty text", id=frame.id
),
)
return
chat_id = frame.chat_id or payload.get("chat_id")
if not isinstance(chat_id, str) or not chat_id.strip():
chat_id = self.home_channel
chat_id = chat_id.strip()
# Automation channels are read-only for the user: they only receive
# gateway-originated output (cron jobs, webhooks). Reject direct sends
# (the app hides the composer for them, this is the server-side
# enforcement).
target = self._channels.get(chat_id)
if target is not None and target.get("automation"):
await self._reply(
device_id,
protocol.error(
protocol.ERR_UNSUPPORTED,
f"{target.get('name') or chat_id} is an automation channel "
"(read-only: cron/webhook output only)",
id=frame.id,
),
)
return
thread_id = frame.thread_id or payload.get("thread_id")
if not isinstance(thread_id, str) or not thread_id.strip():
thread_id = None
reply_to = payload.get("reply_to")
if not isinstance(reply_to, str) or not reply_to.strip():
reply_to = None
# Auto-threading (the app's Threads setting, docs/06 §6.3): a message
# in a channel's flat lane gets its own fresh thread, the way Telegram
# topic mode mints a topic per new conversation. The thread is named
# instantly from the user's opening message (derived title) and the
# LLM upgrades the name in the background. The user echo, the agent
# turn, and all streaming frames then carry the new thread_id.
# Skipped for slash commands (session-scoped, not conversation
# starters) and replies (they continue where the user is). Threading
# is only active on the default channel; other channels stay flat.
auto_thread = bool(payload.get("auto_thread"))
default_entry = self._channels.default()
if (
auto_thread
and thread_id is None
and text.strip()
and not text.lstrip().startswith("/")
and reply_to is None
and default_entry is not None
and chat_id == default_entry["chat_id"]
):
entry = self._channels.create(
name=_derive_thread_name(text),
kind="thread",
parent_chat_id=chat_id,
)
thread_id = entry["chat_id"]
# Bare broadcast (like channel.create): the directory is
# re-served on hello.ack, so no outbox entry is needed.
await self._broadcast_both(protocol.channel_created(entry, auto=True))
self._schedule_thread_title_upgrade(entry["chat_id"], text)
# M4: resolve media refs (single-use; unknown ref -> error).
media_urls: list[str] = []
media_types: list[str] = []
media_wire: list[dict[str, Any]] = []
for ref in media_refs:
entry = self._media.get_inbound(ref)
if entry is None:
await self._reply(
device_id,
protocol.error(
protocol.ERR_UNSUPPORTED, f"unknown media_ref {ref}", id=frame.id
),
)
return
media_urls.append(entry.path)
media_types.append(entry.mime)
media_wire.append(
{
"media_id": entry.media_id,
"kind": entry.kind,
"mime": entry.mime,
"size": entry.size,
"filename": entry.filename,
}
)
device = self._devices.get(device_id) or {}
user_name = device.get("name") or device_id
# Echo to all devices: the sender confirms (server-assigned id),
# other devices see the message too (single-user, multi-device).
# Routed through _broadcast_or_log (not a bare broadcast) so the echo
# is appended to the outbox: the app's ChatStore is in-memory only, so
# after a process death / activity recreation the only way the user's
# own message is restored is via the sync replay. Without this, user
# messages vanish on reconnect while bot messages (already parked)
# survive.
message_id = f"m_{uuid.uuid4().hex[:16]}"
echo = protocol.message(
chat_id=chat_id,
message_id=message_id,
role=protocol.ROLE_USER,
text=text,
thread_id=thread_id,
media=media_wire or None,
reply_to=reply_to,
ts=int(time.time() * 1000),
)
await self._broadcast_or_log(chat_id, echo)
# Refs are consumed by this message (no replay).
for ref in media_refs:
self._media.pop_inbound(ref)
# M4: a new user turn starts -- stale offer association is dropped.
self._last_message_id.pop(chat_id, None)
kind = media_wire[0]["kind"] if media_wire else None
if kind == "image":
message_type = MessageType.PHOTO
elif kind == "video":
message_type = MessageType.VIDEO
elif kind == "audio":
message_type = MessageType.AUDIO
elif kind == "voice":
message_type = MessageType.VOICE
elif kind == "document":
message_type = MessageType.DOCUMENT
else:
message_type = MessageType.TEXT
source = self.build_source(
chat_id=chat_id,
chat_name=self._channel_name(chat_id),
chat_type="dm",
user_id=device_id,
user_name=user_name,
thread_id=thread_id,
)
event = MessageEvent(
text=text,
message_type=message_type,
user_id=device_id,
user_name=user_name,
source=source,
message_id=message_id,
reply_to_message_id=reply_to,
media_urls=media_urls,
media_types=media_types,
)
await self.handle_message(event)
# M5: acknowledge the user message to the originating device (the
# app shows ✓✓) at the moment it is handed to the agent.
await self._reply(
device_id,
protocol.read_receipt(chat_id, message_id),
)
def _schedule_thread_title_upgrade(self, thread_id: str, text: str) -> None:
"""Upgrade an auto-created thread's name with the model's title.
Stage 2 of hermes' two-stage session titling (``agent/title_generator
.py``): the thread was created with an instant derived name; this
background call on the ``title_generation`` auxiliary task replaces it
with the model's title and broadcasts ``channel.renamed``. Best-effort
— any failure (config, model, network) leaves the derived name in
place, and a thread the user already renamed or archived is untouched.
"""
loop = asyncio.get_running_loop()
def _work() -> None:
try:
from agent.title_generator import generate_title
title = generate_title(text)
except Exception:
logger.debug("Thread title upgrade failed", exc_info=True)
return
if not title:
return
entry = self._channels.get(thread_id)
if entry is None or entry.get("archived"):
return
if (entry.get("name") or "") == title:
return
renamed = self._channels.rename(thread_id, title)
if renamed is None:
return
try:
asyncio.run_coroutine_threadsafe(
self._broadcast_both(protocol.channel_renamed(renamed)),
loop,
)
except Exception:
logger.debug("Thread title rename broadcast failed", exc_info=True)
threading.Thread(target=_work, daemon=True, name="iris-thread-title").start()
+2 -2
View File
@@ -184,7 +184,7 @@ class UploadSession:
arrive so an over-limit transfer is rejected early. arrive so an over-limit transfer is rejected early.
""" """
def __init__( def __init__( # noqa: PLR0913
self, self,
media_ref: str, media_ref: str,
kind: str, kind: str,
@@ -272,7 +272,7 @@ class MediaStore:
# ── Inbound uploads ─────────────────────────────────────────────────── # ── Inbound uploads ───────────────────────────────────────────────────
def create_upload( def create_upload( # noqa: PLR0913
self, self,
device_id: str, device_id: str,
media_ref: str, media_ref: str,
+127
View File
@@ -0,0 +1,127 @@
"""M4: outbound media (agent -> app): ``media.offer`` emission.
Mixin for ``adapter.IrisAdapter``. The gateway's dispatch partition
(gateway/run.py) extracts MEDIA: tags / image URLs from the final response,
filters them through ``filter_media_delivery_paths``, then calls the
``send_*`` overrides with local file paths. We re-validate each path
(defense in depth), register it in the media registry, mint a ``media_id``,
and emit ``media.offer``; the app fetches the bytes via ``media.pull``.
"""
import logging
import os
from typing import Any
from gateway.platforms.base import SendResult, validate_media_delivery_path
from . import media as media_bridge
from . import protocol
from .classify import _thread_id_from_metadata
from .mixin_base import IrisAdapterBase
logger = logging.getLogger(__name__)
class MediaHandlers(IrisAdapterBase):
"""Outbound media (see module docstring)."""
async def _offer_media(
self,
chat_id: str,
path: str,
kind: str,
filename: str | None,
metadata: dict[str, Any] | None,
) -> SendResult:
safe = validate_media_delivery_path(path)
if safe is None:
logger.warning("iris: media path failed delivery validation: %s", path)
return SendResult(success=False, error="iris: media path not deliverable")
try:
size = os.path.getsize(safe)
except OSError as e:
logger.warning("iris: media file unreadable %s: %s", safe, e)
return SendResult(success=False, error="iris: media file unreadable")
entry = self._media.register_outbound(
safe, kind, media_bridge.mime_for_path(safe), filename or os.path.basename(safe), size
)
thread_id = _thread_id_from_metadata(metadata)
frame = protocol.media_offer(
entry.media_id,
entry.kind,
entry.mime,
entry.size,
entry.filename,
chat_id=chat_id,
thread_id=thread_id,
message_id=self._last_message_id.get(chat_id),
)
await self._broadcast_or_log(chat_id, frame)
return SendResult(success=True, message_id=entry.media_id)
async def send_image(
self,
chat_id: str,
image_url: str,
caption: str | None = None,
reply_to: str | None = None,
metadata: dict[str, Any] | None = None,
) -> SendResult:
"""Send an image (M4: local files offered over WS; remote URLs fall
back to the base text rendering)."""
if image_url.startswith("file://"):
from urllib.parse import unquote
return await self._offer_media(chat_id, unquote(image_url[7:]), "image", None, metadata)
return await super().send_image(
chat_id, image_url, caption=caption, reply_to=reply_to, metadata=metadata
)
async def send_image_file(
self,
chat_id: str,
image_path: str,
caption: str | None = None,
reply_to: str | None = None,
metadata: dict[str, Any] | None = None,
**kwargs: Any,
) -> SendResult:
"""Send a local image file (M4)."""
return await self._offer_media(chat_id, image_path, "image", None, metadata)
async def send_video(
self,
chat_id: str,
video_path: str,
caption: str | None = None,
reply_to: str | None = None,
metadata: dict[str, Any] | None = None,
**kwargs: Any,
) -> SendResult:
"""Send a video (M4)."""
return await self._offer_media(chat_id, video_path, "video", None, metadata)
async def send_voice(
self,
chat_id: str,
audio_path: str,
caption: str | None = None,
reply_to: str | None = None,
metadata: dict[str, Any] | None = None,
**kwargs: Any,
) -> SendResult:
"""Send a voice note / audio file (M4)."""
return await self._offer_media(chat_id, audio_path, "voice", None, metadata)
async def send_document( # noqa: PLR0913
self,
chat_id: str,
file_path: str,
caption: str | None = None,
file_name: str | None = None,
reply_to: str | None = None,
metadata: dict[str, Any] | None = None,
**kwargs: Any,
) -> SendResult:
"""Send a document (M4)."""
return await self._offer_media(chat_id, file_path, "document", file_name, metadata)
+54
View File
@@ -0,0 +1,54 @@
"""Shared base for the ``IrisAdapter`` mixin classes.
``adapter.IrisAdapter`` is assembled from several small mixin classes
(``inbound``, ``tool_frames``, ``push_frames``, ...) plus the core state and
lifecycle in ``adapter`` itself. Each mixin references instance attributes and
helper methods that are defined in the core class or in a *sibling* mixin, so a
type checker analysing one mixin in isolation cannot see them.
This base declares those shared names (as ``Any``) so static analysis resolves
``self.<name>`` inside every mixin. The annotations carry no runtime effect;
the real values are set in ``IrisAdapter.__init__`` and the real methods live
in the core class / sibling mixins.
"""
from typing import Any
class IrisAdapterBase:
"""Declaration-only base for the ``IrisAdapter`` mixins (see module doc)."""
# -- shared state (set in ``IrisAdapter.__init__``) -------------------
_channels: Any
_devices: Any
_http_server: Any
_media: Any
_outbox: Any
_push: Any
_active_lane: Any
_last_message_id: Any
_last_push_at: Any
_pending_push: Any
_pending_pickers: Any
_prune_notified_at: Any
_typing_turns: Any
home_channel: Any
home_channel_name: Any
# -- shared helpers (core class or sibling mixins) --------------------
_broadcast_both: Any
_broadcast_or_log: Any
_channel_name: Any
_maybe_push: Any
_offer_media: Any
_parse_tool_line_or_block: Any
_push_summary: Any
_reply: Any
_schedule_thread_title_upgrade: Any
# -- provided by ``BasePlatformAdapter`` / core -----------------------
build_source: Any
handle_message: Any
send: Any
send_image: Any
send_slash_confirm: Any
+1 -1
View File
@@ -171,7 +171,7 @@ class Outbox:
# ── history (full message history for a chat/thread) ────────────────── # ── history (full message history for a chat/thread) ──────────────────
def history( def history( # noqa: PLR0912
self, self,
chat_id: str, chat_id: str,
thread_id: str | None = None, thread_id: str | None = None,
+265
View File
@@ -0,0 +1,265 @@
"""M5: approval / clarify / choice-picker frames (interactive banners).
Mixin for ``adapter.IrisAdapter``. Hermes detects these methods on the
adapter type; each emits a high-priority ``notification`` (pushed even when
a device is live) and, with a live device, an interactive ``picker.choice``
card whose selection runs the stored callback (``pickers.py``).
"""
import time
from typing import Any
from gateway.platforms.base import SendResult
from . import protocol
from .classify import (
_mint_message_id,
_mint_picker_id,
_push_preview,
_thread_id_from_metadata,
)
from .mixin_base import IrisAdapterBase
from .pickers import (
_approval_picker_callback,
_clarify_is_multi,
_clarify_picker_callback,
)
class PickerHandlers(IrisAdapterBase):
"""Interactive pickers + approvals (see module docstring)."""
async def send_slash_confirm( # noqa: PLR0913
self,
chat_id: str,
title: str,
message: str,
session_key: str,
confirm_id: str,
metadata: dict[str, Any] | None = None,
) -> SendResult:
"""Banner + push for a slash-command approval prompt.
The gateway's text fallback still renders the actionable prompt (the
app has no inline buttons yet); the notification is the push-visible
signal (high priority: pushed even when a device is live).
"""
thread_id = _thread_id_from_metadata(metadata)
await self._broadcast_or_log(
chat_id,
protocol.notification(
chat_id,
protocol.NOTIF_APPROVAL,
title or "Approval needed",
_push_preview(message),
thread_id=thread_id,
),
)
return await super().send_slash_confirm(
chat_id, title, message, session_key, confirm_id, metadata=metadata
)
async def send_exec_approval( # noqa: PLR0913
self,
chat_id: str,
command: str,
session_key: str,
description: str = "dangerous command",
metadata: dict[str, Any] | None = None,
allow_permanent: bool = True,
allow_session: bool = True,
smart_denied: bool = False,
) -> SendResult:
"""Interactive exec-approval picker (buttons) for a dangerous command.
Hermes calls this (detected on the adapter type) when the agent wants
to run a command that needs approval; the agent thread blocks until the
user decides. With a live device we render the same choice set as the
native adapters (Allow Once / Session / Always / Deny, gated by the
same flags) as a ``picker.choice`` card, reusing the clarify/slash
picker mechanism. A tap resolves via ``resolve_gateway_approval``
(the same primitive the text ``/approve`` / ``/deny`` handlers use),
unblocking the agent, and a short confirmation is delivered as a
normal message. A high-priority ``approval`` notification is also
emitted so a backgrounded device is woken (pushed even when live).
With no live device the picker could never be answered, so report
failure and let hermes fall back to the text ``/approve`` prompt.
"""
if not self._http_server.has_devices():
return SendResult(success=False, error="no live devices for approval picker")
thread_id = _thread_id_from_metadata(metadata)
# High-priority banner + push (wakes a backgrounded device).
await self._broadcast_or_log(
chat_id,
protocol.notification(
chat_id,
protocol.NOTIF_APPROVAL,
"Approval needed",
_push_preview(description or command),
thread_id=thread_id,
),
)
# Choice set mirrors the native adapters (telegram/relay).
frame_choices = [{"value": "once", "label": "\u2705 Allow Once", "is_current": False}]
if not smart_denied and allow_session:
frame_choices.append(
{"value": "session", "label": "\u2705 Allow Session", "is_current": False}
)
if allow_permanent:
frame_choices.append(
{"value": "always", "label": "\u2705 Always Allow", "is_current": False}
)
frame_choices.append({"value": "deny", "label": "\u274c Deny", "is_current": False})
cmd_preview = command if len(command) <= 1500 else command[:1500] + "\u2026"
title = (
"\u26a0\ufe0f **Command approval required**\n\n"
f"```\n{cmd_preview}\n```\n\n"
f"Reason: {description}"
)
if smart_denied:
title += "\n\n**Smart DENY:** owner override applies to this one operation only."
picker_id = _mint_picker_id()
self._pending_pickers[picker_id] = {
"chat_id": chat_id,
"thread_id": thread_id,
"on_choice_selected": _approval_picker_callback(session_key),
}
await self._broadcast_or_log(
chat_id,
protocol.picker_choice(picker_id, title, frame_choices, chat_id, thread_id=thread_id),
)
return SendResult(success=True, message_id=picker_id)
async def send_choice_picker( # noqa: PLR0913
self,
chat_id: str,
title: str,
choices: list,
session_key: str,
on_choice_selected,
metadata: dict[str, Any] | None = None,
) -> SendResult:
"""Send an interactive choice picker (one tap → one value).
The generic companion to Telegram's inline-keyboard pickers, used by
``/reasoning``, ``/fast``, and any future finite-choice slash command
(hermes detects this method on the adapter type). Emits a
``picker.choice`` frame; the app answers with ``picker.select``,
which runs ``on_choice_selected(chat_id, value)`` and delivers the
returned text as a normal message. Outboxed, so a reconnecting
device re-renders a still-pending picker.
With no live device the picker could never be answered, so report
failure and let hermes fall back to the text status card.
"""
if not self._http_server.has_devices():
return SendResult(success=False, error="no live devices for picker")
thread_id = _thread_id_from_metadata(metadata)
picker_id = _mint_picker_id()
self._pending_pickers[picker_id] = {
"chat_id": chat_id,
"thread_id": thread_id,
"on_choice_selected": on_choice_selected,
}
await self._broadcast_or_log(
chat_id,
protocol.picker_choice(picker_id, title, choices, chat_id, thread_id=thread_id),
)
return SendResult(success=True, message_id=picker_id)
async def send_clarify( # noqa: PLR0913
self,
chat_id: str,
question: str,
choices: list | None,
clarify_id: str,
session_key: str,
metadata: dict[str, Any] | None = None,
) -> SendResult:
"""Banner + push for a clarify prompt.
Single-select clarifies with a live device render as an interactive
``picker.choice`` card (one tap per option + an "Other" free-text
button), reusing the slash-command picker mechanism. A real pick
resolves via ``resolve_gateway_clarify`` (the agent then continues and
replies); "Other" flips the entry to text-capture. Multi-select,
open-ended, and no-live-device clarifies fall back to a numbered text
list whose reply the gateway's text-intercept captures via
``mark_awaiting_text``.
"""
thread_id = _thread_id_from_metadata(metadata)
await self._broadcast_or_log(
chat_id,
protocol.notification(
chat_id,
protocol.NOTIF_CLARIFY,
"Question",
_push_preview(question),
thread_id=thread_id,
),
)
# Single-select + live device → interactive picker card.
if choices and not _clarify_is_multi(clarify_id) and self._http_server.has_devices():
picker_id = _mint_picker_id()
self._pending_pickers[picker_id] = {
"chat_id": chat_id,
"thread_id": thread_id,
"on_choice_selected": _clarify_picker_callback(
clarify_id, [str(c) for c in choices]
),
}
frame_choices = [
{"value": f"c{i}", "label": str(c)[:75], "is_current": False}
for i, c in enumerate(choices)
]
frame_choices.append(
{"value": "other", "label": "✏️ Other (type your answer)", "is_current": False}
)
await self._broadcast_or_log(
chat_id,
protocol.picker_choice(
picker_id, f"❓ {question}", frame_choices, chat_id, thread_id=thread_id
),
)
return SendResult(success=True, message_id=picker_id)
# Text fallback (multi-select / open-ended / no live device).
if choices:
lines = [f"❓ {question}", ""]
for i, choice in enumerate(choices, start=1):
lines.append(f" {i}. {choice}")
lines.append("")
if _clarify_is_multi(clarify_id):
lines.append(
"Multiple selections allowed — reply with the numbers "
'separated by commas or spaces (e.g. "1, 3"), the option '
"text, or your own answer."
)
else:
lines.append("Reply with the number, the option text, or your own answer.")
text = "\n".join(lines)
# Text fallback: enable text-capture so the gateway intercept
# picks up the user's typed reply (e.g. "2" or choice text).
from tools.clarify_gateway import mark_awaiting_text
mark_awaiting_text(clarify_id)
else:
text = f"❓ {question}"
message_id = _mint_message_id()
await self._broadcast_or_log(
chat_id,
protocol.message(
chat_id=chat_id,
message_id=message_id,
role=protocol.ROLE_ASSISTANT,
text=text,
thread_id=thread_id,
ts=int(time.time() * 1000),
),
)
return SendResult(success=True, message_id=message_id)
+83
View File
@@ -0,0 +1,83 @@
"""Choice-picker callbacks: clarify + exec approval.
Build the ``on_choice_selected`` closures the adapter stores in
``_pending_pickers``; a ``picker.select`` from the app runs the closure and
delivers its reply text as a normal message.
"""
def _clarify_is_multi(clarify_id: str) -> bool:
"""True when the pending clarify [clarify_id] allows multiple selections.
The flag lives on the gateway's pending entry; a missing/expired entry (or
any lookup error) is treated as single-select.
"""
try:
from tools import clarify_gateway as _cg
with _cg._lock:
_entry = _cg._entries.get(clarify_id)
return bool(_entry and getattr(_entry, "multi_select", False))
except Exception:
return False
def _clarify_picker_callback(clarify_id: str, choices: list[str]):
"""Build the ``on_choice_selected`` callback for a clarify picker.
Option values are positional (``c0``..``cN``) plus an ``other`` sentinel;
the closure maps them back to the real choice strings. A real pick resolves
the clarify (the agent then continues and replies); "Other" flips the entry
to text-capture so the next typed message is the answer. An unmappable value
also flips to text so a clarify never dead-ends.
"""
async def on_choice_selected(chat_id: str, value: str) -> str | None:
from tools.clarify_gateway import mark_awaiting_text, resolve_gateway_clarify
if value == "other":
mark_awaiting_text(clarify_id)
return "✏️ Type your answer:"
try:
idx = int(value[1:]) if value.startswith("c") else -1
except ValueError:
idx = -1
if 0 <= idx < len(choices):
resolve_gateway_clarify(clarify_id, choices[idx])
return None
mark_awaiting_text(clarify_id)
return "✏️ Type your answer:"
return on_choice_selected
# The four exec-approval outcomes hermes understands (tools.approval).
_APPROVAL_CHOICES = ("once", "session", "always", "deny")
def _approval_picker_callback(session_key: str):
"""Build the ``on_choice_selected`` callback for an exec-approval picker.
The option values are the raw hermes approval outcomes (``once`` /
``session`` / ``always`` / ``deny``); a tap resolves the waiting agent
thread via ``resolve_gateway_approval`` (the same primitive the text
``/approve`` / ``/deny`` handlers use) and returns a short confirmation
label, which the picker handler delivers as a normal message. An unknown
value is treated as a deny so a stray tap never approves a command.
"""
async def on_choice_selected(chat_id: str, value: str) -> str | None:
from tools.approval import resolve_gateway_approval
choice = value if value in _APPROVAL_CHOICES else "deny"
count = resolve_gateway_approval(session_key, choice)
label = {
"once": "✅ Approved once",
"session": "✅ Approved for this session",
"always": "✅ Approved permanently",
"deny": "❌ Denied",
}[choice]
if not count:
label = "⌛ Approval expired — no command was waiting."
return label
return on_choice_selected
+6 -6
View File
@@ -388,7 +388,7 @@ def message_update(
) )
def message_stop( def message_stop( # noqa: PLR0913
chat_id: str, chat_id: str,
message_id: str, message_id: str,
final_text: str, final_text: str,
@@ -428,7 +428,7 @@ def message_stop(
# --------------------------------------------------------------------------- # ---------------------------------------------------------------------------
def tool_start( def tool_start( # noqa: PLR0913
chat_id: str, chat_id: str,
index: int, index: int,
name: str, name: str,
@@ -472,7 +472,7 @@ def tool_progress(
) )
def tool_end( def tool_end( # noqa: PLR0913
chat_id: str, chat_id: str,
index: int, index: int,
name: str, name: str,
@@ -694,7 +694,7 @@ def sync_done(cursor: int, *, id: int | None = None) -> Frame:
# --------------------------------------------------------------------------- # ---------------------------------------------------------------------------
def history( def history( # noqa: PLR0913
chat_id: str, chat_id: str,
messages: list[dict[str, Any]], messages: list[dict[str, Any]],
has_more: bool, has_more: bool,
@@ -755,7 +755,7 @@ def message_deleted(
# --------------------------------------------------------------------------- # ---------------------------------------------------------------------------
def notification( def notification( # noqa: PLR0913
chat_id: str, chat_id: str,
kind: str, kind: str,
title: str, title: str,
@@ -807,7 +807,7 @@ def status(state: str) -> Frame:
# --------------------------------------------------------------------------- # ---------------------------------------------------------------------------
def media_offer( def media_offer( # noqa: PLR0913
media_id: str, media_id: str,
kind: str, kind: str,
mime: str, mime: str,
+1 -1
View File
@@ -96,7 +96,7 @@ def delete_lane(db_path: Path, chat_id: str, thread_id: str | None = None) -> in
conn.close() conn.close()
def delete_message( def delete_message( # noqa: PLR0913
db_path: Path, db_path: Path,
chat_id: str, chat_id: str,
thread_id: str | None, thread_id: str | None,
+5 -5
View File
@@ -63,7 +63,7 @@ class PushBackend:
"""True when the backend has credentials to send with.""" """True when the backend has credentials to send with."""
raise NotImplementedError raise NotImplementedError
async def send( async def send( # noqa: PLR0913
self, self,
*, *,
device_id: str, device_id: str,
@@ -126,7 +126,7 @@ class FcmBackend(PushBackend):
self._sa_failed = True self._sa_failed = True
return None return None
async def _authorization(self, client: httpx.AsyncClient) -> str | None: async def _authorization(self, client: httpx.AsyncClient) -> str | None: # noqa: PLR0911
"""Bearer token: the legacy server key, or a cached service-account """Bearer token: the legacy server key, or a cached service-account
OAuth2 access token (JWT-bearer grant, minted with PyJWT).""" OAuth2 access token (JWT-bearer grant, minted with PyJWT)."""
if self._server_key: if self._server_key:
@@ -187,7 +187,7 @@ class FcmBackend(PushBackend):
self._token_expiry = now + 3600.0 self._token_expiry = now + 3600.0
return token return token
async def send( async def send( # noqa: PLR0913
self, self,
*, *,
device_id: str, device_id: str,
@@ -277,7 +277,7 @@ class NtfyBackend(PushBackend):
def configured(self) -> bool: def configured(self) -> bool:
return bool(self._topic) return bool(self._topic)
async def send( async def send( # noqa: PLR0913
self, self,
*, *,
device_id: str, device_id: str,
@@ -315,7 +315,7 @@ class NtfyBackend(PushBackend):
return True return True
def build_push_backend( def build_push_backend( # noqa: PLR0913
name: str | None, name: str | None,
*, *,
fcm_service_account: str | None = None, fcm_service_account: str | None = None,
+199
View File
@@ -0,0 +1,199 @@
"""M5: push mirroring, typing indicators, turn-aware push.
Mixin for ``adapter.IrisAdapter``. Frames with no live subscriber wake the
device via the push backend (``push.py``); high-priority kinds push even
when a device is live. While the agent's turn is in flight (typing on),
normal-priority message/media frames are held back and the latest one is
pushed when the turn ends.
"""
import logging
import time
from typing import Any
from . import protocol
from .classify import _PUSH_COALESCE_S, _push_preview
from .mixin_base import IrisAdapterBase
logger = logging.getLogger(__name__)
# How often (seconds) the outbox-prune "storage reclaimed" notice may repeat.
_PRUNE_NOTIFY_INTERVAL_S = 3600.0
class PushHandlers(IrisAdapterBase):
"""Push + typing/turn tracking (see module docstring)."""
def _push_summary(self, frame: "protocol.Frame") -> tuple[str, str, str, str] | None:
"""``(title, body, kind, priority)`` for a pushable frame, else None.
Only terminal/interesting frames wake a device: intermediate
streaming and tool frames are replayed by ``sync`` without a push
(no notification spam per turn).
"""
t = frame.type
p = frame.payload
if t == protocol.TYPE_MESSAGE:
return (
self._channel_name(frame.chat_id or ""),
_push_preview(p.get("text")),
"message",
"normal",
)
if t == protocol.TYPE_MESSAGE_STOP:
return (
self._channel_name(frame.chat_id or ""),
_push_preview(p.get("final_text")),
"message",
"normal",
)
if t == protocol.TYPE_NOTIFICATION:
kind = str(p.get("kind") or protocol.NOTIF_GENERIC)
priority = "high" if kind in protocol.HIGH_PRIORITY_NOTIF_KINDS else "normal"
return (
str(p.get("title") or "Iris"),
str(p.get("body") or ""),
kind,
priority,
)
if t == protocol.TYPE_MEDIA_OFFER:
return (
self._channel_name(frame.chat_id or ""),
f"New {p.get('kind') or 'media'}: {p.get('filename') or ''}".strip(),
"media",
"normal",
)
return None
async def _maybe_push(self, chat_id: str, frame: "protocol.Frame", cursor: int) -> None:
"""Fire the configured push backend for a parked (or high-priority)
frame. Best-effort: failures are logged, never raised."""
summary = self._push_summary(frame)
if summary is None:
return
# M5: coalesce back-to-back pushes for the same chat (cron delivery
# = notification frame + message frame). The suppressed frame is
# still synced when the app reconnects.
now = time.time()
if now - self._last_push_at.get(chat_id, 0.0) < _PUSH_COALESCE_S:
logger.info(
"iris: push coalesced for %s (%s frame within %.0fs of last push)",
chat_id,
frame.type,
_PUSH_COALESCE_S,
)
return
title, body, kind, priority = summary
backend = self._push
if backend is None or not backend.token_field:
return
devices = self._devices.list()
if not backend.configured() and not any(d.get(backend.token_field) for d in devices):
return
data: dict[str, Any] = {"chat_id": chat_id, "kind": kind, "cursor": str(cursor)}
if frame.thread_id:
data["thread_id"] = frame.thread_id
message_id = frame.payload.get("message_id")
if isinstance(message_id, str) and message_id:
data["message_id"] = message_id
for device in devices:
device_id = device.get("device_id")
if not device_id:
continue
token = device.get(backend.token_field)
if not token:
continue
try:
ok = await backend.send(
device_id=device_id,
chat_id=chat_id,
title=title,
body=body,
data=data,
token=token,
priority=priority,
)
except Exception:
logger.warning("iris: push via %s failed", backend.name, exc_info=True)
continue
if ok:
# M5: remember that this cursor reached the device via push,
# so the app can dedupe it on the next sync replay.
try:
self._devices.update_push_cursor(device_id, cursor)
except Exception:
logger.warning(
"iris: push cursor update failed for %s", device_id, exc_info=True
)
self._last_push_at[chat_id] = time.time()
logger.info(
"iris: push via %s -> %s (%s, chat=%s)",
backend.name,
device_id,
frame.type,
chat_id,
)
async def _maybe_notify_outbox_prune(self, chat_id: str) -> None:
"""When the outbox row cap pruned old frames, tell the app (throttled
to once per hour so a full box doesn't banner per frame)."""
pruned = self._outbox.take_overflow_pruned()
if pruned <= 0:
return
now = time.time()
if now - self._prune_notified_at < _PRUNE_NOTIFY_INTERVAL_S:
return
self._prune_notified_at = now
await self._broadcast_or_log(
chat_id,
protocol.notification(
chat_id,
protocol.NOTIF_GENERIC,
"Outbox",
f"{pruned} older message(s) pruned",
),
)
async def send_typing(self, chat_id: str, metadata: dict[str, Any] | None = None) -> None:
"""Send a typing indicator (``typing`` frame, on=true)."""
thread_id = None
if metadata:
tid = metadata.get("thread_id")
if isinstance(tid, str) and tid:
thread_id = tid
# M5: typing on marks the chat's agent turn as in flight (hermes
# turns typing on at turn start, before the first output).
self._typing_turns.add(chat_id)
frame = protocol.typing(chat_id, True, thread_id=thread_id)
await self._http_server.fanout(frame, cursor=None)
async def stop_typing(self, chat_id: str) -> None:
"""Clear the typing indicator (``typing`` frame, on=false)."""
frame = protocol.typing(chat_id, False)
await self._http_server.fanout(frame, cursor=None)
# M5: turn ended (hermes fires stop_typing in the handler's finally,
# after the final send). Flush the held-back push -- the final
# answer -- but only while the device is still offline; a live
# device already got the frames via its event stream / sync.
# Idempotent: hermes may call stop_typing more than once per turn.
self._typing_turns.discard(chat_id)
pending = self._pending_push.pop(chat_id, None)
if pending is None:
return
if self._http_server.has_devices():
logger.info("iris: deferred push dropped for %s (device back online)", chat_id)
return
held_frame, held_cursor = pending
await self._maybe_push(chat_id, held_frame, held_cursor)
def on_device_online(self) -> None:
"""A device opened its event stream (SSE/long-poll): it will sync
the outbox, so drop any held-back pushes -- flushing them later
would duplicate what the app already shows. Called from the HTTP
server's handler thread; dict.clear() is atomic under the GIL."""
if self._pending_push:
logger.info(
"iris: device online; dropping %d deferred push(es)",
len(self._pending_push),
)
self._pending_push.clear()
+2 -2
View File
@@ -232,7 +232,7 @@ def _version_info(version: int) -> int:
return _bch(version, 12, 0x1F25) return _bch(version, 12, 0x1F25)
def _build_matrix(version: int, level: str, codewords: list[int], mask: int) -> list[list[bool]]: def _build_matrix(version: int, level: str, codewords: list[int], mask: int) -> list[list[bool]]: # noqa: PLR0912,PLR0915
size = 17 + 4 * version size = 17 + 4 * version
# matrix[r][c] = dark; reserved[r][c] = function module (not data) # matrix[r][c] = dark; reserved[r][c] = function module (not data)
matrix = [[False] * size for _ in range(size)] matrix = [[False] * size for _ in range(size)]
@@ -355,7 +355,7 @@ def _build_matrix(version: int, level: str, codewords: list[int], mask: int) ->
return matrix return matrix
def _mask_bit(mask: int, r: int, c: int) -> bool: def _mask_bit(mask: int, r: int, c: int) -> bool: # noqa: PLR0911
if mask == 0: if mask == 0:
return (r + c) % 2 == 0 return (r + c) % 2 == 0
if mask == 1: if mask == 1:
+266
View File
@@ -0,0 +1,266 @@
"""Inbound query frames: catalog, search, sync, history, delete, push tokens,
picker selection.
Mixin for ``adapter.IrisAdapter``. These are the read/catch-up half of the
inbound surface (``message.send`` lives in ``inbound.py``): they answer
point-to-point (``_reply``) or broadcast the matching event frame.
"""
import logging
from hermes_constants import get_hermes_home
from . import protocol
from . import purge as purge_bridge
from . import search as search_bridge
from .commands import _slash_command_catalog
logger = logging.getLogger(__name__)
class QueryFrameHandlers:
"""Inbound query frames (see module docstring)."""
# ── Slash-command catalog (app's "/" drawer) ──────────────────────────
async def on_commands_catalog(self, frame: protocol.Frame, device_id: str) -> None:
"""Handle an inbound ``commands.catalog`` request: reply with the
gateway's slash-command catalog (hermes ``COMMAND_REGISTRY``,
gateway-available subset + plugin commands). The app fuzzy-matches
the typed prefix client-side; the catalog is static per gateway run,
so no caching is needed here."""
resp = protocol.commands_catalog(_slash_command_catalog(), id=frame.id)
await self._reply(device_id, resp)
# ── M3: search (app -> agent) ─────────────────────────────────────────
async def on_search(self, frame: protocol.Frame, device_id: str) -> None:
payload = frame.payload
query = payload.get("query")
if not isinstance(query, str) or not query.strip():
await self._reply(
device_id,
protocol.error(protocol.ERR_UNSUPPORTED, "search requires a query", id=frame.id),
)
return
scope = payload.get("scope")
scope = scope if scope in ("all", "chat") else "all"
chat_id = payload.get("chat_id") or frame.chat_id
if not isinstance(chat_id, str) or not chat_id.strip():
chat_id = None
thread_id = payload.get("thread_id") or frame.thread_id
if not isinstance(thread_id, str) or not thread_id.strip():
thread_id = None
limit = payload.get("limit")
try:
limit = int(limit) if limit is not None else 20
except (TypeError, ValueError):
limit = 20
db_path = get_hermes_home() / "state.db"
hits = search_bridge.search(
db_path, query, scope=scope, chat_id=chat_id, thread_id=thread_id, limit=limit
)
resp = protocol.search_results(query, scope, hits, id=frame.id)
await self._reply(device_id, resp)
# ── M3: sync (reconnect catch-up) ─────────────────────────────────────
async def on_sync(self, frame: protocol.Frame, device_id: str) -> None:
payload = frame.payload
cursor = payload.get("cursor")
try:
cursor = int(cursor) if cursor is not None else 0
except (TypeError, ValueError):
cursor = 0
for e in self._outbox.replay(cursor):
raw = e["frame"]
replayed = protocol.Frame(
type=raw.get("type", ""),
payload=raw.get("payload", {}) if isinstance(raw.get("payload"), dict) else {},
id=raw.get("id") if isinstance(raw.get("id"), int) else None,
chat_id=(
raw.get("chat_id") if isinstance(raw.get("chat_id"), str) else e.get("chat_id")
),
thread_id=raw.get("thread_id") if isinstance(raw.get("thread_id"), str) else None,
# M5: tag replayed frames with their outbox cursor so the app
# can skip re-notifying frames that already woke the device
# via push (cursor <= last_pushed_cursor, docs/08 §8.7).
cursor=e.get("cursor"),
v=raw.get("v") if isinstance(raw.get("v"), int) else protocol.PROTOCOL_VERSION,
)
await self._reply(device_id, replayed)
done = protocol.sync_done(self._outbox.latest_cursor(), id=frame.id)
await self._reply(device_id, done)
# ── Full message history (initial channel open / scroll-up) ───────────
async def on_history(self, frame: protocol.Frame, device_id: str) -> None:
"""Handle an inbound ``history`` request.
``sync`` only replays the outbox delta since the device's cursor, so
after a process death the app's in-memory ChatStore is empty and the
delta does not cover older messages. ``history`` loads the full
message list for a chat/thread (reconstructed from the outbox log) so
the app can populate the view on first open / restart.
"""
payload = frame.payload
chat_id = frame.chat_id or payload.get("chat_id")
logger.info("iris: history request from %s chat_id=%r", device_id, chat_id)
if not isinstance(chat_id, str) or not chat_id.strip():
await self._reply(
device_id,
protocol.error(protocol.ERR_UNSUPPORTED, "history requires a chat_id", id=frame.id),
)
return
chat_id = chat_id.strip()
thread_id = frame.thread_id or payload.get("thread_id")
if not isinstance(thread_id, str) or not thread_id.strip():
thread_id = None
before = payload.get("before_message_id")
if not isinstance(before, str) or not before.strip():
before = None
limit_raw = payload.get("limit")
try:
limit = int(limit_raw) if limit_raw is not None else 50
except (TypeError, ValueError):
limit = 50
page = self._outbox.history(
chat_id,
thread_id=thread_id,
before_message_id=before,
limit=limit,
)
resp = protocol.history(
chat_id,
page["messages"],
page["has_more"],
thread_id=thread_id,
oldest_message_id=page["oldest_message_id"],
id=frame.id,
)
await self._reply(device_id, resp)
# ── Message deletion (app -> agent) ───────────────────────────────────
async def on_message_delete(self, frame: protocol.Frame, device_id: str) -> None:
"""Handle an inbound ``message.delete`` request.
Completely deletes the requested message(s): they are removed from the
outbox (so ``history`` and ``sync`` no longer return them) **and** from
the hermes session store (so no search trace survives and they are not
recoverable). ``message.deleted`` is broadcast to every device
(outboxed too, so an offline device learns of the deletion on its next
``sync``). Deleting is idempotent: a message that is already gone
(pruned by retention) simply yields 0 removed rows, and the
``message.deleted`` broadcast is still emitted so live caches drop it.
"""
payload = frame.payload
chat_id = frame.chat_id or payload.get("chat_id")
if not isinstance(chat_id, str) or not chat_id.strip():
await self._reply(
device_id,
protocol.error(
protocol.ERR_NOT_FOUND, "message.delete requires chat_id", id=frame.id
),
)
return
chat_id = chat_id.strip()
thread_id = frame.thread_id or payload.get("thread_id")
if not isinstance(thread_id, str) or not thread_id.strip():
thread_id = None
message_ids = payload.get("message_ids")
if not isinstance(message_ids, list):
message_ids = [payload.get("message_id")] if payload.get("message_id") else []
message_ids = [m for m in message_ids if isinstance(m, str) and m.strip()]
if not message_ids:
await self._reply(
device_id,
protocol.error(
protocol.ERR_UNSUPPORTED, "message.delete requires message_ids", id=frame.id
),
)
return
removed = 0
purged = 0
db_path = get_hermes_home() / "state.db"
for mid in message_ids:
# Read the final frame data first (role / text / ts) so the
# session-store row can be matched, then drop the outbox frames.
info = self._outbox.message_info(chat_id, mid, thread_id=thread_id)
removed += self._outbox.delete_message(chat_id, mid, thread_id=thread_id)
if info:
purged += purge_bridge.delete_message(
db_path,
chat_id,
thread_id,
info.get("role") or "",
info.get("text") or "",
info.get("ts"),
)
logger.info(
"iris: message.delete from %s chat_id=%r thread_id=%r ids=%s removed=%s purged=%s",
device_id,
chat_id,
thread_id,
message_ids,
removed,
purged,
)
resp = protocol.message_deleted(chat_id, message_ids, thread_id=thread_id)
resp.id = frame.id
await self._broadcast_or_log(chat_id, resp)
# ── M5: push token registration ───────────────────────────────────────
async def on_fcm_register(self, frame: protocol.Frame, device_id: str) -> None:
"""Update the device's push tokens (FCM rotation / ntfy topic).
Persists to the device registry so the next push targets the current
token without a stale read.
"""
fcm_token = frame.payload.get("fcm_token")
ntfy_topic = frame.payload.get("ntfy_topic")
fcm_token = fcm_token if isinstance(fcm_token, str) and fcm_token else None
ntfy_topic = ntfy_topic if isinstance(ntfy_topic, str) and ntfy_topic else None
if fcm_token is None and ntfy_topic is None:
return
try:
self._devices.update_push_tokens(device_id, fcm_token=fcm_token, ntfy_topic=ntfy_topic)
except Exception:
logger.warning("iris: fcm.register update failed", exc_info=True)
return
logger.info("iris: push tokens updated for %s", device_id)
# ── Interactive pickers (slash-command choice menus) ─────────────────
async def on_picker_select(self, frame: protocol.Frame, device_id: str) -> None:
"""Resolve a pending choice picker (``picker.select`` from the app).
Runs the command's selection callback and delivers its reply text as
a normal final message in the picker's chat. Unknown/expired picker
ids (gateway restart, double tap) are a no-op — the app already
marked the card resolved locally.
"""
picker_id = frame.payload.get("picker_id")
value = frame.payload.get("value")
if not isinstance(picker_id, str) or not isinstance(value, str):
return
state = self._pending_pickers.pop(picker_id, None)
if state is None:
logger.info("iris: picker.select for unknown/expired picker %s", picker_id)
return
callback = state.get("on_choice_selected")
if callback is None:
return
try:
result_text = await callback(state["chat_id"], value)
except Exception:
logger.error("iris: picker selection failed for %s", picker_id, exc_info=True)
return
if not result_text:
return
await self.send(
state["chat_id"],
str(result_text),
metadata={"notify": True, "thread_id": state.get("thread_id")},
)
+17 -11
View File
@@ -4,11 +4,17 @@
# hermes-agent/.venv/bin/python -m ruff check gateway-plugin # hermes-agent/.venv/bin/python -m ruff check gateway-plugin
# #
# The rule set is deliberately broad (pycodestyle, pyflakes, isort, pyupgrade, # The rule set is deliberately broad (pycodestyle, pyflakes, isort, pyupgrade,
# bugbear, flake8-simplify, pylint, return, comprehensions). Thresholds below # bugbear, flake8-simplify, pylint, return, comprehensions).
# reflect the plugin's real shape: it is a single large dispatch surface #
# (adapter.py) plus a wire-protocol layer (protocol.py) whose frame builders # The pylint complexity ceilings (PLR0911/0912/0913/0915) are left at Ruff's
# mirror the schema, so the complexity ceilings are set just above the current # built-in defaults (see [lint.pylint]). We deliberately do NOT raise them to
# maxima rather than an idealized small-function target. # "just above the current maxima": that ratchets the bar down every time code
# grows (LLM maintenance adds functions, it does not refactor them), so new
# complex code would silently pass. Instead, the handful of genuinely complex
# functions that already exist (frame builders that mirror 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.
line-length = 100 line-length = 100
@@ -32,12 +38,12 @@ select = [
ignore = ["PLC0415"] ignore = ["PLC0415"]
[lint.pylint] [lint.pylint]
# Current maxima in the codebase: 22 branches, 64 statements, 9 returns, # Ruff's built-in defaults. Existing outliers are noqa'd at the def line
# 8 args (protocol.py:252 frame builder is the lone 11-arg outlier, noqa'd). # (search for `# noqa: PLR09`), not absorbed into a raised ceiling.
max-branches = 24 max-branches = 12
max-statements = 70 max-statements = 50
max-returns = 9 max-returns = 6
max-args = 8 max-args = 5
[lint.per-file-ignores] [lint.per-file-ignores]
# The e2e / ws_probe drivers are assertion scripts: scenario numbers and # The e2e / ws_probe drivers are assertion scripts: scenario numbers and
+3 -3
View File
@@ -116,7 +116,7 @@ def _row_to_hit(row: sqlite3.Row) -> dict[str, Any]:
} }
def _fts_query( def _fts_query( # noqa: PLR0913
conn: sqlite3.Connection, conn: sqlite3.Connection,
query: str, query: str,
scope: str, scope: str,
@@ -153,7 +153,7 @@ def _fts_query(
return [_row_to_hit(r) for r in rows] return [_row_to_hit(r) for r in rows]
def _like_query( def _like_query( # noqa: PLR0913
conn: sqlite3.Connection, conn: sqlite3.Connection,
query: str, query: str,
scope: str, scope: str,
@@ -197,7 +197,7 @@ def _like_query(
return [_row_to_hit(r) for r in rows] return [_row_to_hit(r) for r in rows]
def search( def search( # noqa: PLR0913
db_path: Path, db_path: Path,
query: str, query: str,
scope: str = "all", scope: str = "all",
+29
View File
@@ -0,0 +1,29 @@
"""Scope-aware secret reads for the iris plugin.
Shared by ``adapter.py`` (token / TLS / FCM credentials) and ``setup.py``
(config probes + env enablement).
"""
import os
from agent.secret_scope import UnscopedSecretError as _UnscopedSecretError
from agent.secret_scope import get_secret as _scoped_get_secret
def _get_scoped_secret(name, default=None):
"""Scope-aware credential read with the default-profile startup fallback.
Secondary profiles construct their adapters under a profile secret scope
-- the scope is authoritative and a scoped miss returns ``default`` (no
cross-profile borrow from ``os.environ``, which may hold another
profile's value). The DEFAULT profile's adapter constructs and sends
*unscoped* under multiplexing, where a bare ``get_secret`` would raise
``UnscopedSecretError`` and crash this path; there ``os.environ`` is that
profile's own value, so fall back to it. Same pattern as the IRC
``IRC_SERVER_PASSWORD`` read (``plugins/platforms/irc/adapter.py``).
"""
try:
val = _scoped_get_secret(name, default)
except _UnscopedSecretError:
val = os.getenv(name)
return val if val is not None else default
+400
View File
@@ -0,0 +1,400 @@
"""Interactive setup, passive config probes, env-driven auto-configuration.
``interactive_setup`` is the ``hermes gateway setup`` flow (token, host,
port, push backend, pairing QR, device removal). ``check_requirements`` /
``validate_config`` / ``is_connected`` are the passive probes the platform
registry calls from status displays. ``_env_enablement`` seeds
``PlatformConfig.extra`` from env vars before adapter construction.
"""
import logging
import os
import re
import time
from typing import Any
from hermes_constants import get_hermes_home
from . import qr
from .channels import get_directory
from .defaults import (
DEFAULT_HOME_CHANNEL_NAME,
DEFAULT_HOST,
DEFAULT_HTTP_PORT,
DEFAULT_PORT,
DEFAULT_PUSH_BACKEND,
)
from .pairing import (
DeviceRegistry,
advertise_host,
generate_token,
pairing_url,
qr_payload,
)
from .secrets import _get_scoped_secret
logger = logging.getLogger(__name__)
# ---------------------------------------------------------------------------
# Passive / config probes (called from status displays -- no side effects)
# ---------------------------------------------------------------------------
def check_requirements() -> bool:
"""PASSIVE dependency probe: token set.
Must be side-effect free (called from ``hermes setup`` / ``status`` /
dashboard readiness). Never installs. The HTTP transport is stdlib-only,
so there is no extra dependency to probe.
"""
return bool(_get_scoped_secret("IRIS_TOKEN"))
def validate_config(config) -> bool:
"""Given a PlatformConfig, is the platform properly configured?"""
extra = getattr(config, "extra", {}) or {}
token = _get_scoped_secret("IRIS_TOKEN") or extra.get("token", "")
return bool(token)
def is_connected(config) -> bool:
"""Is the platform configured (env or config.yaml)?"""
return validate_config(config)
# ---------------------------------------------------------------------------
# Env-driven auto-configuration (seeds PlatformConfig.extra pre-adapter)
# ---------------------------------------------------------------------------
def _env_enablement() -> dict | None:
"""Seed ``PlatformConfig.extra`` from env vars during gateway config load.
Called by the platform registry's env-enablement hook BEFORE adapter
construction, so ``gateway status`` and ``get_connected_platforms()``
reflect env-only configuration without instantiating the adapter.
Returns ``None`` when the platform isn't minimally configured (no token);
the caller then skips auto-enabling.
The special ``home_channel`` key in the returned dict is handled by the
core hook -- it becomes a proper ``HomeChannel`` dataclass on the
``PlatformConfig`` rather than being merged into ``extra``.
"""
token = _get_scoped_secret("IRIS_TOKEN", "")
if not token:
return None
# Seed ONLY explicitly-set env vars: the core commits this seed on top of
# config.yaml (``extra.update(seed)``), so default values here would
# clobber user YAML. Unset keys fall through to config.yaml / adapter
# defaults.
seed: dict[str, Any] = {}
host = os.getenv("IRIS_WS_HOST", "").strip()
if host:
seed["host"] = host
http_port_raw = os.getenv("IRIS_HTTP_PORT", "").strip()
if http_port_raw:
seed["http_port"] = _parse_port(http_port_raw)
push = os.getenv("IRIS_PUSH_BACKEND", "").strip().lower()
if push:
seed["push_backend"] = push
home = os.getenv("IRIS_HOME_CHANNEL", "").strip()
if home:
seed["home_channel"] = {
"chat_id": home,
"name": os.getenv("IRIS_HOME_CHANNEL_NAME", "").strip() or DEFAULT_HOME_CHANNEL_NAME,
}
return seed
def _parse_port(raw: str) -> int:
try:
return int((raw or "").strip())
except (ValueError, TypeError):
return DEFAULT_PORT
# ---------------------------------------------------------------------------
# Target parsing: "<chat_id>[:<thread>]" (platform prefix stripped by core)
# ---------------------------------------------------------------------------
def _parse_target_ref(target_ref: str) -> tuple | None: # noqa: PLR0911
"""Parse a raw target string into ``(chat_id, thread_id)`` or ``None``.
The core strips the platform prefix before calling us, so the native
syntax is simply ``<chat_id>[:<thread>]`` (e.g. ``chan_7`` or
``chan_7:t_31``); the home channel is ``default``. Chat ids are direct
(no embedded platform prefix), so a cron delivery reads
``iris:chan_7`` end to end. A bare friendly name (e.g. ``Cron Reports``)
is resolved against the channel directory so cron / ``send_message`` can
target a channel by name immediately, without waiting for the core
directory's refresh timer. Returns ``None`` for anything unrecognised so
the target proceeds to the core channel-directory resolution.
"""
if not target_ref:
return None
t = target_ref.strip()
if not t:
return None
thread_id: str | None = None
if ":" in t:
head, tail = t.rsplit(":", 1)
if head and tail.startswith("t_"):
thread_id = tail
t = head
else:
# Not a <chat>:<thread> pair -- treat the whole string as a name.
t = target_ref.strip()
if not t:
return None
# Native chat id (default / chan_<n>) or any id known to the directory
# (covers custom IRIS_HOME_CHANNEL values).
try:
known = get_directory().get(t) is not None
except Exception:
known = False
if t == "default" or re.fullmatch(r"chan_\d+", t) or known:
return (t, thread_id)
# Bare friendly name -> resolve via the channel directory. A thread resolves
# to its session lane (parent_chat_id + thread_id); a channel/default to
# its chat_id.
try:
entry = get_directory().resolve_entry(t)
except Exception:
entry = None
if entry is not None:
if entry["kind"] == "thread":
return (entry["parent_chat_id"], entry["chat_id"])
return (entry["chat_id"], None)
return None
# ---------------------------------------------------------------------------
# Standalone (out-of-process) send -- best-effort, stretch for v1
# ---------------------------------------------------------------------------
async def _standalone_send( # noqa: PLR0913
pconfig,
chat_id: str,
message: str,
*,
thread_id: str | None = None,
media_files: list[str] | None = None,
force_document: bool = False,
) -> dict[str, Any]:
"""Out-of-process delivery for cron jobs that run separately from the
gateway.
The outbox is served by the *running* gateway, so standalone delivery
while the gateway process is fully down is best-effort only (see
``docs/00-overview.md`` "Out of scope"). For M1 this is a stub that
reports the gateway is required; the real implementation lands with the
outbox (M3/M5).
"""
return {
"error": (
"iris standalone send: the running gateway is required to serve "
"the outbox (standalone delivery is best-effort only)"
)
}
# ---------------------------------------------------------------------------
# Verbose tool progress (full args on the progress line)
# ---------------------------------------------------------------------------
def _ensure_verbose_tool_progress() -> None:
"""Ensure the iris platform renders tool progress in ``verbose`` mode.
Verbose mode makes the gateway's tool-progress line carry the FULL
argument JSON (not just a ~40-char preview), which the adapter parses
into the ``tool.start`` frame's ``args`` field; the app then decides how
much to show (Settings → Tool detail). The tool *output* is captured
separately via the ``post_tool_call`` hook (verbose mode does not stream
it).
Best-effort and idempotent: writes
``display.platforms.iris.tool_progress: verbose`` to config.yaml only
when it isn't already set. The gateway's config cache is mtime-keyed, so
the write takes effect on the next turn without a restart. Never raises.
"""
try:
from hermes_cli.config import load_config_readonly
cfg = load_config_readonly() or {}
display = cfg.get("display") or {}
platforms = display.get("platforms") or {}
iris_cfg = platforms.get("iris") or {}
if iris_cfg.get("tool_progress") == "verbose":
return # already set
from utils import atomic_roundtrip_yaml_update
atomic_roundtrip_yaml_update(
get_hermes_home() / "config.yaml",
"display.platforms.iris.tool_progress",
"verbose",
)
logger.info("iris: set display.platforms.iris.tool_progress=verbose")
except Exception:
logger.debug("iris: could not ensure verbose tool_progress", exc_info=True)
# ---------------------------------------------------------------------------
# Interactive setup (hermes gateway setup flow)
# ---------------------------------------------------------------------------
def _offer_device_removal() -> None:
"""Setup-flow device management (docs/09 §9.3): if devices are already
paired, offer to revoke one. Revocation is server-side — no access to
the device is needed: its per-device token is deleted and its id is
denylisted, so even the shared token no longer authenticates it.
Flow: ask (default No) → numbered select menu (last option = exit the
removal loop, NOT the setup) → confirmation → back to the menu, so
several devices can be removed in a row.
"""
try:
from hermes_cli.cli_output import (
print_info,
print_success,
prompt,
prompt_yes_no,
)
except Exception:
return
try:
reg = DeviceRegistry(get_hermes_home() / "iris" / "devices.db")
except Exception:
return
try:
devices = reg.list()
if not devices:
return
if not prompt_yes_no("Remove a paired device?", default=False):
return
while True:
print_info("Paired devices:")
for i, d in enumerate(devices, 1):
last_seen = time.strftime("%Y-%m-%d %H:%M", time.localtime(d["last_seen"]))
print_info(f" {i}. {d['name']} ({d['device_id']}) last seen {last_seen}")
exit_idx = len(devices) + 1
print_info(f" {exit_idx}. Exit")
# Default = exit: pressing Enter leaves the removal loop (and
# continues the setup) without removing anything.
choice = prompt("Select a device to remove", default=str(exit_idx))
idx = int(choice) if choice.isdigit() else exit_idx
if idx < 1 or idx >= exit_idx:
return
target = devices[idx - 1]
if not prompt_yes_no(
f"Remove {target['name']} ({target['device_id']})? It will no longer "
"be able to connect (shared token included).",
default=False,
):
continue # back to the select menu
reg.revoke(target["device_id"])
devices = [d for d in devices if d["device_id"] != target["device_id"]]
print_success(f"Removed {target['device_id']} \u2014 it can no longer connect.")
if not devices:
print_info("No paired devices left.")
return
finally:
reg.close()
def interactive_setup() -> None:
"""Prompt for the pairing token / host / port / push backend.
M1: token generation, host/port/push prompts, and the pairing QR payload
(``iris://pair?...``) + app URL printed for the Connect screen.
"""
try:
from hermes_cli.cli_output import (
print_info,
print_success,
print_warning,
prompt,
)
from hermes_cli.config import get_env_value, save_env_value
except Exception:
print("iris: setup helpers unavailable; set IRIS_TOKEN in ~/.hermes/.env")
return
print_info("📱 Android / Desktop (Iris x Hermes)")
token = get_env_value("IRIS_TOKEN") or ""
if not token:
generated = generate_token()
save_env_value("IRIS_TOKEN", generated)
print_success(f"Generated pairing token: {generated}")
print_warning("Keep this secret -- the app presents it on connect.")
else:
print_info("Existing IRIS_TOKEN found (not shown).")
# Device management (docs/09 §9.3): on an existing setup, offer to cut
# off a lost/compromised device before continuing with the config.
_offer_device_removal()
host = prompt("Bind host", default=get_env_value("IRIS_WS_HOST") or DEFAULT_HOST)
save_env_value("IRIS_WS_HOST", host or DEFAULT_HOST)
# _parse_port falls back to DEFAULT_PORT (8790) for empty input, so the
# HTTP default must be applied explicitly (docs/19: 8791).
http_port_raw = (get_env_value("IRIS_HTTP_PORT") or "").strip()
port = prompt(
"HTTP port",
default=str(int(http_port_raw) if http_port_raw.isdigit() else DEFAULT_HTTP_PORT),
)
save_env_value("IRIS_HTTP_PORT", str(_parse_port(port)))
backend = prompt(
"Push backend (ntfy/fcm)",
default=get_env_value("IRIS_PUSH_BACKEND") or DEFAULT_PUSH_BACKEND,
)
backend = (backend or DEFAULT_PUSH_BACKEND).strip().lower()
save_env_value("IRIS_PUSH_BACKEND", backend)
if backend == "fcm":
print_warning(
"FCM push metadata (notification title, device token) is routed "
"through Google's servers. For truly private communication use "
"ntfy (self-hosted) instead."
)
# Pairing payload for the app's Connect screen (manual entry + QR scan).
# Advertise a routable host: a bind wildcard (0.0.0.0/127.0.0.1) is
# replaced by the default-route LAN IP so the QR points somewhere a phone
# can actually reach (the user can still override the Server URL in-app).
advertised = advertise_host(host or DEFAULT_HOST)
url = pairing_url(advertised, _parse_port(port))
pairing = qr_payload(advertised, _parse_port(port), token)
print_info("Pair your device (enter this on the app's Connect screen):")
print_info(f"Pairing URL: {pairing}")
print_info(f"Server URL: {url}")
if advertised != (host or DEFAULT_HOST):
print_info(
f"QR points to {advertised} (your default LAN address). If your "
"phone is on a different network, change the Server URL in the app."
)
# Scannable QR (docs/20): the same payload as a terminal QR. The URL text
# lines stay — the QR is a convenience, not a replacement (non-UTF-8
# terminals still work, and the text is copy-pasteable). render_qr returns
# '' (not an exception) when the payload is too long to encode.
qr_block = qr.render_qr(pairing)
if qr_block:
print_info("Scan with the Iris app (Connect → Scan QR) or any camera app:")
print(qr_block)
else:
print_warning("QR too large to render; use the pairing URL above.")
# Always render tool progress verbosely so the app receives the full tool
# call args (it decides how much to show via Settings → Tool detail).
_ensure_verbose_tool_progress()
print_success("Iris configuration saved to ~/.hermes/.env")
print_info("Restart the gateway for changes to take effect: hermes gateway restart")
+157
View File
@@ -0,0 +1,157 @@
"""Tool-progress frame emission (M2): the tool.start / tool.end lifecycle.
Mixin for ``adapter.IrisAdapter``. The gateway accumulates tool lines in one
editable bubble; on an edit the full buffer is re-sent, so new lines are
diffed against ``seen_tool_lines`` and each new tool closes the previously
open one (attaching the output/duration captured by the ``post_tool_call``
hook).
"""
from typing import Any
from gateway.platforms.base import SendResult
from . import protocol
from .classify import (
_extract_code_block,
_extract_verbose_args,
_mint_message_id,
_parse_tool_line,
_short_preview_from_args,
_TurnState,
)
from .hooks import _reset_tool_results, _tool_emoji, _tool_end_fields
from .mixin_base import IrisAdapterBase
class ToolProgressHandlers(IrisAdapterBase):
"""Tool-progress lifecycle (see module docstring)."""
async def _emit_tool_lines(
self,
chat_id: str,
content: str,
state: _TurnState,
thread_id: str | None,
*,
is_edit: bool,
) -> SendResult:
"""Emit ``tool.start`` for each NEW tool line in *content*.
The gateway accumulates tool lines in one editable bubble; on an edit
the full buffer is re-sent, so we diff against ``seen_tool_lines`` to
emit only the new ones. A new tool closes the previously-open tool.
"""
message_id = state.tool_msg_id or _mint_message_id()
state.tool_msg_id = message_id
state.active = True
# Tool activity marks this lane as the turn in flight: the global
# post_tool_call hook (no chat id) routes todo emissions here.
self._active_lane = (chat_id, thread_id)
lines = [ln for ln in content.splitlines() if ln.strip()]
for line in lines:
key = line.strip()
if key in state.seen_tool_lines:
continue
state.seen_tool_lines.add(key)
parsed = self._parse_tool_line_or_block(line, content)
if parsed is None:
continue
name, preview, args = parsed
# A new tool begins: close the previously-open one, attaching the
# output/duration/ok captured by the post_tool_call hook.
if state.open_tool_index is not None:
extra = _tool_end_fields(state.open_tool_name or "")
await self._broadcast_or_log(
chat_id,
protocol.tool_end(
chat_id,
state.open_tool_index,
state.open_tool_name or "",
ok=extra.get("ok", True),
duration=extra.get("duration"),
output_preview=extra.get("output_preview"),
thread_id=thread_id,
),
)
state.tool_index += 1
state.open_tool_index = state.tool_index
state.open_tool_name = name
await self._broadcast_or_log(
chat_id,
protocol.tool_start(
chat_id,
state.tool_index,
name,
preview=preview,
args=args,
emoji=_tool_emoji(name),
thread_id=thread_id,
),
)
return SendResult(success=True, message_id=message_id)
@staticmethod
def _parse_tool_line_or_block(
line: str, content: str
) -> tuple[str, str | None, dict[str, Any] | None] | None:
"""Parse a tool line into ``(name, preview, args)``.
Expands a terminal code block to its command, and a verbose header
(``<emoji> <name>(keys)``) to its full args JSON (the JSON sits on the
following line). ``args`` is ``None`` unless the line is a verbose
header with a parseable JSON body.
"""
parsed = _parse_tool_line(line)
if parsed is None:
return None
name, preview = parsed
# Terminal code block: the command lives in the fenced lines that
# follow the "<emoji> terminal" head line.
if name == "terminal" and preview is None and "```" in content:
cmd = _extract_code_block(content)
if cmd:
return name, cmd, None
# Verbose mode: recover the full args from the following JSON line.
args = _extract_verbose_args(line, content)
if args is not None and preview is None:
preview = _short_preview_from_args(args)
return name, preview, args
async def _close_open_tool(
self, chat_id: str, state: _TurnState, thread_id: str | None
) -> None:
"""Emit ``tool.end`` for the currently-open tool, if any.
A tool is considered complete when the next tool starts OR a new
content segment begins (the model only produces content after the
tool it was waiting on has returned).
"""
if state.open_tool_index is not None:
extra = _tool_end_fields(state.open_tool_name or "")
await self._broadcast_or_log(
chat_id,
protocol.tool_end(
chat_id,
state.open_tool_index,
state.open_tool_name or "",
ok=extra.get("ok", True),
duration=extra.get("duration"),
output_preview=extra.get("output_preview"),
thread_id=thread_id,
),
)
state.open_tool_index = None
state.open_tool_name = None
def _reset_tool_state(self, state: _TurnState) -> None:
"""Clear per-turn tool bookkeeping (called at turn finalization)."""
state.tool_msg_id = None
state.seen_tool_lines = set()
state.tool_index = 0
state.open_tool_index = None
state.open_tool_name = None
# Drop any captured tool results not consumed by a tool.end this turn
# (e.g. tool_progress off) so they can't leak into the next turn.
_reset_tool_results()