Split adapter.py monolith into focused modules; restore Ruff complexity defaults #15
No files matched your search
+103
-2739
File diff suppressed because it is too large.
Load diff
@@ -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
|
||||
# ---------------------------------------------------------------------------
|
||||
@@ -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)
|
||||
@@ -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
|
||||
@@ -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"}
|
||||
@@ -24,7 +24,7 @@ INBOUND_BURST = 40
|
||||
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
|
||||
ignored (forward-compat)."""
|
||||
if frame.type == protocol.TYPE_MESSAGE_SEND:
|
||||
|
||||
@@ -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,
|
||||
)
|
||||
@@ -480,7 +480,7 @@ class HttpServer:
|
||||
|
||||
# ── 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)
|
||||
cursor = _parse_cursor(qs.get("cursor", [None])[0], handler.headers.get("Last-Event-ID"))
|
||||
# Device registration (the HTTP equivalent of the WS hello upsert):
|
||||
@@ -600,7 +600,7 @@ class HttpServer:
|
||||
|
||||
# ── 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.
|
||||
|
||||
Mirrors the WS ``media.upload`` contract (docs/07 §7.2) in one
|
||||
@@ -820,7 +820,7 @@ class _Handler(BaseHTTPRequestHandler):
|
||||
return
|
||||
_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
|
||||
if not hs.enabled:
|
||||
_send_json(self, 503, {"error": "http leg disabled"})
|
||||
|
||||
@@ -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()
|
||||
@@ -184,7 +184,7 @@ class UploadSession:
|
||||
arrive so an over-limit transfer is rejected early.
|
||||
"""
|
||||
|
||||
def __init__(
|
||||
def __init__( # noqa: PLR0913
|
||||
self,
|
||||
media_ref: str,
|
||||
kind: str,
|
||||
@@ -272,7 +272,7 @@ class MediaStore:
|
||||
|
||||
# ── Inbound uploads ───────────────────────────────────────────────────
|
||||
|
||||
def create_upload(
|
||||
def create_upload( # noqa: PLR0913
|
||||
self,
|
||||
device_id: str,
|
||||
media_ref: str,
|
||||
|
||||
@@ -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)
|
||||
@@ -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
|
||||
@@ -171,7 +171,7 @@ class Outbox:
|
||||
|
||||
# ── history (full message history for a chat/thread) ──────────────────
|
||||
|
||||
def history(
|
||||
def history( # noqa: PLR0912
|
||||
self,
|
||||
chat_id: str,
|
||||
thread_id: str | None = None,
|
||||
|
||||
@@ -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)
|
||||
@@ -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
|
||||
@@ -388,7 +388,7 @@ def message_update(
|
||||
)
|
||||
|
||||
|
||||
def message_stop(
|
||||
def message_stop( # noqa: PLR0913
|
||||
chat_id: str,
|
||||
message_id: str,
|
||||
final_text: str,
|
||||
@@ -428,7 +428,7 @@ def message_stop(
|
||||
# ---------------------------------------------------------------------------
|
||||
|
||||
|
||||
def tool_start(
|
||||
def tool_start( # noqa: PLR0913
|
||||
chat_id: str,
|
||||
index: int,
|
||||
name: str,
|
||||
@@ -472,7 +472,7 @@ def tool_progress(
|
||||
)
|
||||
|
||||
|
||||
def tool_end(
|
||||
def tool_end( # noqa: PLR0913
|
||||
chat_id: str,
|
||||
index: int,
|
||||
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,
|
||||
messages: list[dict[str, Any]],
|
||||
has_more: bool,
|
||||
@@ -755,7 +755,7 @@ def message_deleted(
|
||||
# ---------------------------------------------------------------------------
|
||||
|
||||
|
||||
def notification(
|
||||
def notification( # noqa: PLR0913
|
||||
chat_id: str,
|
||||
kind: str,
|
||||
title: str,
|
||||
@@ -807,7 +807,7 @@ def status(state: str) -> Frame:
|
||||
# ---------------------------------------------------------------------------
|
||||
|
||||
|
||||
def media_offer(
|
||||
def media_offer( # noqa: PLR0913
|
||||
media_id: str,
|
||||
kind: str,
|
||||
mime: str,
|
||||
|
||||
@@ -96,7 +96,7 @@ def delete_lane(db_path: Path, chat_id: str, thread_id: str | None = None) -> in
|
||||
conn.close()
|
||||
|
||||
|
||||
def delete_message(
|
||||
def delete_message( # noqa: PLR0913
|
||||
db_path: Path,
|
||||
chat_id: str,
|
||||
thread_id: str | None,
|
||||
|
||||
@@ -63,7 +63,7 @@ class PushBackend:
|
||||
"""True when the backend has credentials to send with."""
|
||||
raise NotImplementedError
|
||||
|
||||
async def send(
|
||||
async def send( # noqa: PLR0913
|
||||
self,
|
||||
*,
|
||||
device_id: str,
|
||||
@@ -126,7 +126,7 @@ class FcmBackend(PushBackend):
|
||||
self._sa_failed = True
|
||||
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
|
||||
OAuth2 access token (JWT-bearer grant, minted with PyJWT)."""
|
||||
if self._server_key:
|
||||
@@ -187,7 +187,7 @@ class FcmBackend(PushBackend):
|
||||
self._token_expiry = now + 3600.0
|
||||
return token
|
||||
|
||||
async def send(
|
||||
async def send( # noqa: PLR0913
|
||||
self,
|
||||
*,
|
||||
device_id: str,
|
||||
@@ -277,7 +277,7 @@ class NtfyBackend(PushBackend):
|
||||
def configured(self) -> bool:
|
||||
return bool(self._topic)
|
||||
|
||||
async def send(
|
||||
async def send( # noqa: PLR0913
|
||||
self,
|
||||
*,
|
||||
device_id: str,
|
||||
@@ -315,7 +315,7 @@ class NtfyBackend(PushBackend):
|
||||
return True
|
||||
|
||||
|
||||
def build_push_backend(
|
||||
def build_push_backend( # noqa: PLR0913
|
||||
name: str | None,
|
||||
*,
|
||||
fcm_service_account: str | None = None,
|
||||
|
||||
@@ -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()
|
||||
@@ -232,7 +232,7 @@ def _version_info(version: int) -> int:
|
||||
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
|
||||
# matrix[r][c] = dark; reserved[r][c] = function module (not data)
|
||||
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
|
||||
|
||||
|
||||
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:
|
||||
return (r + c) % 2 == 0
|
||||
if mask == 1:
|
||||
|
||||
@@ -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
@@ -4,11 +4,17 @@
|
||||
# hermes-agent/.venv/bin/python -m ruff check gateway-plugin
|
||||
#
|
||||
# The rule set is deliberately broad (pycodestyle, pyflakes, isort, pyupgrade,
|
||||
# bugbear, flake8-simplify, pylint, return, comprehensions). Thresholds below
|
||||
# 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
|
||||
# mirror the schema, so the complexity ceilings are set just above the current
|
||||
# maxima rather than an idealized small-function target.
|
||||
# bugbear, flake8-simplify, pylint, return, comprehensions).
|
||||
#
|
||||
# The pylint complexity ceilings (PLR0911/0912/0913/0915) are left at Ruff's
|
||||
# built-in defaults (see [lint.pylint]). We deliberately do NOT raise them to
|
||||
# "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
|
||||
|
||||
@@ -32,12 +38,12 @@ select = [
|
||||
ignore = ["PLC0415"]
|
||||
|
||||
[lint.pylint]
|
||||
# Current maxima in the codebase: 22 branches, 64 statements, 9 returns,
|
||||
# 8 args (protocol.py:252 frame builder is the lone 11-arg outlier, noqa'd).
|
||||
max-branches = 24
|
||||
max-statements = 70
|
||||
max-returns = 9
|
||||
max-args = 8
|
||||
# Ruff's built-in defaults. Existing outliers are noqa'd at the def line
|
||||
# (search for `# noqa: PLR09`), not absorbed into a raised ceiling.
|
||||
max-branches = 12
|
||||
max-statements = 50
|
||||
max-returns = 6
|
||||
max-args = 5
|
||||
|
||||
[lint.per-file-ignores]
|
||||
# The e2e / ws_probe drivers are assertion scripts: scenario numbers and
|
||||
|
||||
@@ -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,
|
||||
query: str,
|
||||
scope: str,
|
||||
@@ -153,7 +153,7 @@ def _fts_query(
|
||||
return [_row_to_hit(r) for r in rows]
|
||||
|
||||
|
||||
def _like_query(
|
||||
def _like_query( # noqa: PLR0913
|
||||
conn: sqlite3.Connection,
|
||||
query: str,
|
||||
scope: str,
|
||||
@@ -197,7 +197,7 @@ def _like_query(
|
||||
return [_row_to_hit(r) for r in rows]
|
||||
|
||||
|
||||
def search(
|
||||
def search( # noqa: PLR0913
|
||||
db_path: Path,
|
||||
query: str,
|
||||
scope: str = "all",
|
||||
|
||||
@@ -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
|
||||
@@ -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")
|
||||
@@ -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()
|
||||
Reference in new issue
Block a user