Slash commands with a finite set of options (/reasoning, /fast, ...) now
render a tappable card with buttons (2 per row, ✓ on the current value)
instead of a plain text status card. The mechanism is generic: any command
that calls the adapter's send_choice_picker() gets a picker automatically.
Wire protocol (docs/04, frames.schema.json):
- picker.choice (server→app): {picker_id, title, choices[]}
- picker.select (app→server): {picker_id, value}
- pickers capability flag now True in server_caps
gateway-plugin:
- protocol.py: picker.choice/picker.select frame types + picker_choice()
- dispatch.py: route picker.select → adapter.on_picker_select
- adapter.py: send_choice_picker() (fails cleanly with no live device so
hermes falls back to text), on_picker_select(), in-memory pending pickers
(gateway restart expires them; stale select is a no-op), pickers=True
app (KMP):
- Protocol.kt: PickerChoice/PickerChoicePayload + pickerSelectFrame()
- ChatStore.kt: PickerItem + onPickerChoice (idempotent) + resolvePicker
(optimistic, one-shot)
- ChatDb.kt: persist PickerItem in the messages table (polymorphic decode)
- IrisController.kt: picker.choice routing + selectPicker() action
- ChatScreen.kt: PickerCard composable (locks after selection)
Tests:
- python: 3 picker tests (roundtrip, no-device fallback, stale-select noop)
- kotlin: ChatStorePickerTest (add/idempotent/resolve/one-shot/noop/serialize)
- fixture fix: clear leaked IRIS_HTTP_PORT/IRIS_WS_HOST env so the adapter
binds the ephemeral port (a prior test's interactive_setup() polluted the
process env, colliding with a live gateway on 8791)
Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
3083 lines
127 KiB
Python
3083 lines
127 KiB
Python
"""
|
|
Iris Platform Adapter for Hermes Agent (Iris x Hermes).
|
|
|
|
A plugin-based gateway adapter that runs an HTTP server *inside* the
|
|
``hermes gateway`` process. The native Android / Desktop app connects to it
|
|
with a pairing token and talks to the agent over a single HTTP transport
|
|
(chat, streaming, tools, media, pairing, push-token): JSON frames via
|
|
``POST /v1/frame``, events via SSE ``GET /v1/events`` (or long-poll), and
|
|
media via ``POST /v1/media`` / ``GET /v1/media/{id}`` (docs/19).
|
|
|
|
Zero new Python dependencies: ``httpx`` is a hermes core dep. Zero
|
|
hermes-core changes.
|
|
|
|
Milestone M1: the gateway core loop (text round-trip). The server binds and
|
|
authenticates devices (constant-time token check), the adapter emits
|
|
``message`` frames from ``send()`` and turns inbound ``message.send`` frames
|
|
into ``MessageEvent``s for ``handle_message()``.
|
|
|
|
Milestone M2: agent transparency. ``send()``/``edit_message()`` are mapped to
|
|
``message.start``/``message.update``/``message.stop`` (streaming), tool
|
|
progress is classified into structured ``tool.start``/``tool.end`` frames,
|
|
interim commentary becomes ``commentary`` frames, and the code-style
|
|
reasoning prefix is split into a ``reasoning`` field. Outbox and search land
|
|
in M3; media, push, and desktop land in later milestones (see
|
|
``docs/14-milestones.md``).
|
|
|
|
Milestone M4: media. Inbound uploads (``POST /v1/media``) are streamed to a
|
|
temp file, verified (size + sha256), re-sniffed, and cached via hermes
|
|
``cache_*_from_bytes``; the resulting refs attach to the next
|
|
``message.send`` as ``MessageEvent.media_urls``. Outbound ``send_*`` calls
|
|
register the (delivery-validated) file in the media registry and emit
|
|
``media.offer``; ``GET /v1/media/{id}`` streams the file back, re-checking
|
|
``validate_media_delivery_path`` at pull time.
|
|
|
|
Milestone M5: push + offline. Frames with no live subscriber are parked in
|
|
the outbox (M3) AND wake the device via the push backend (``push.py``: FCM
|
|
HTTP v1 primary, ntfy fallback, selected by ``IRIS_PUSH_BACKEND``).
|
|
``notification`` frames render in-app banners and mirror to push (channel
|
|
events, cron deliveries, approvals, clarifies); high-priority kinds push even
|
|
when a device is live. ``fcm.register`` rotates push tokens (registry + live
|
|
connection). The outbox enforces a row cap with a throttled prune notice.
|
|
|
|
Configuration in config.yaml::
|
|
|
|
gateway:
|
|
platforms:
|
|
iris:
|
|
enabled: true
|
|
extra:
|
|
host: 127.0.0.1
|
|
port: 8790
|
|
home_channel: default
|
|
push_backend: fcm
|
|
outbox_retention_hours: 72
|
|
max_upload_bytes: 104857600
|
|
|
|
Or via environment variables (overrides config.yaml; secrets live in .env):
|
|
IRIS_TOKEN, IRIS_WS_HOST, IRIS_WS_PORT, IRIS_HOME_CHANNEL,
|
|
IRIS_PUSH_BACKEND, IRIS_FCM_SERVICE_ACCOUNT, NTFY_TOPIC, ...
|
|
"""
|
|
|
|
import asyncio
|
|
import contextlib
|
|
import json
|
|
import logging
|
|
import os
|
|
import queue
|
|
import re
|
|
import threading
|
|
import time
|
|
import uuid
|
|
from collections import deque
|
|
from dataclasses import dataclass, field
|
|
from typing import Any
|
|
|
|
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
|
|
|
|
|
|
logger = logging.getLogger(__name__)
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Lazy import: BasePlatformAdapter and friends live in the main repo.
|
|
# We import at module level (as the bundled plugins do) but guard the heavy
|
|
# gateway imports so the plugin can be discovered before the gateway is fully
|
|
# initialised.
|
|
# ---------------------------------------------------------------------------
|
|
|
|
from gateway.config import Platform # noqa: E402
|
|
from gateway.platforms.base import ( # noqa: E402
|
|
BasePlatformAdapter,
|
|
MessageEvent,
|
|
MessageType,
|
|
SendResult,
|
|
validate_media_delivery_path,
|
|
)
|
|
from hermes_constants import get_hermes_home # noqa: E402
|
|
|
|
from . import media as media_bridge # noqa: E402
|
|
from . import ( # noqa: E402
|
|
protocol,
|
|
qr,
|
|
)
|
|
from . import purge as purge_bridge # noqa: E402
|
|
from . import search as search_bridge # noqa: E402
|
|
from .channels import get_directory # noqa: E402
|
|
from .http_server import HttpServer # noqa: E402
|
|
from .outbox import Outbox # noqa: E402
|
|
from .pairing import ( # noqa: E402
|
|
DeviceRegistry,
|
|
advertise_host,
|
|
generate_token,
|
|
pairing_url,
|
|
qr_payload,
|
|
)
|
|
from .push import NtfyBackend, PushBackend, build_push_backend # noqa: E402
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Slash-command catalog (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.
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
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
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# 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()
|
|
|
|
|
|
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
|
|
|
|
|
|
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).
|
|
"""
|
|
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,
|
|
)
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Defaults
|
|
# ---------------------------------------------------------------------------
|
|
|
|
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 = "fcm"
|
|
DEFAULT_OUTBOX_RETENTION_HOURS = 72
|
|
DEFAULT_MAX_UPLOAD_BYTES = 100 * 1024 * 1024 # 100 MB
|
|
|
|
# How often (seconds) the outbox-prune "storage reclaimed" notice may repeat.
|
|
_PRUNE_NOTIFY_INTERVAL_S = 3600.0
|
|
|
|
|
|
def _truthy(value: str | None) -> bool:
|
|
return (value or "").strip().lower() in {"1", "true", "yes", "on"}
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# M2 — turn state + outbound classification
|
|
#
|
|
# 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()``. We classify each outbound call into a structured frame using a
|
|
# per-chat turn state machine + content markers:
|
|
#
|
|
# * ``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``.
|
|
# ---------------------------------------------------------------------------
|
|
|
|
# Streaming cursor the gateway appends to in-progress edits (" ▉"). Stripped
|
|
# before we forward text to the app (the app renders its own live indicator).
|
|
_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:
|
|
"""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)
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# 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:
|
|
"""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(
|
|
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 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).")
|
|
|
|
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 (fcm/ntfy)",
|
|
default=get_env_value("IRIS_PUSH_BACKEND") or DEFAULT_PUSH_BACKEND,
|
|
)
|
|
save_env_value("IRIS_PUSH_BACKEND", (backend or DEFAULT_PUSH_BACKEND).strip().lower())
|
|
|
|
# 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")
|
|
|
|
|
|
# ---------------------------------------------------------------------------
|
|
# Iris Adapter
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
class IrisAdapter(BasePlatformAdapter):
|
|
"""HTTP-backed adapter for the native Iris app (Android / Desktop).
|
|
|
|
The HTTP server (``http_server.HttpServer``) authenticates devices with
|
|
the pairing token, the device registry tracks live subscribers, ``send()``
|
|
emits ``message`` frames, and inbound ``message.send`` frames become
|
|
``MessageEvent``s for ``handle_message()``.
|
|
"""
|
|
|
|
# The HTTP transport has no per-message size limit. The stream consumer
|
|
# resolves its per-chat chunking budget via ``max_message_length_for_chat``
|
|
# -> this attribute (defaulting to 4096 when unset), which would split
|
|
# long replies — and complete HTML artifacts — across multiple
|
|
# fence-reopened messages. A large cap disables that chunking so a reply
|
|
# arrives as a single message.
|
|
MAX_MESSAGE_LENGTH = 1_000_000
|
|
|
|
def __init__(self, config, **kwargs):
|
|
platform = Platform("iris")
|
|
super().__init__(config=config, platform=platform)
|
|
|
|
# Ensure verbose tool progress (full args on the progress line) so the
|
|
# app can show the full tool call on demand. Best-effort; idempotent.
|
|
_ensure_verbose_tool_progress()
|
|
|
|
extra = getattr(config, "extra", {}) or {}
|
|
|
|
# Connection settings (env vars override config.yaml). The bind host
|
|
# is shared with the (legacy) WS-era env var name for compatibility.
|
|
self.host = os.getenv("IRIS_WS_HOST", "").strip() or extra.get("host", DEFAULT_HOST)
|
|
# docs/19: HTTP transport (the only device-facing transport; optional TLS).
|
|
self.http_port = _parse_port(
|
|
os.getenv("IRIS_HTTP_PORT", "") or str(extra.get("http_port", DEFAULT_HTTP_PORT))
|
|
)
|
|
self.token = _get_scoped_secret("IRIS_TOKEN") or extra.get("token", "")
|
|
self.push_backend = os.getenv("IRIS_PUSH_BACKEND", "").strip().lower() or extra.get(
|
|
"push_backend", DEFAULT_PUSH_BACKEND
|
|
)
|
|
self.outbox_retention_hours = int(
|
|
extra.get("outbox_retention_hours", DEFAULT_OUTBOX_RETENTION_HOURS)
|
|
)
|
|
self.max_upload_bytes = int(extra.get("max_upload_bytes", DEFAULT_MAX_UPLOAD_BYTES))
|
|
self._gateway_status = protocol.STATUS_ONLINE
|
|
|
|
# Home channel: the core hook turns the env-seeded ``home_channel``
|
|
# dict into a HomeChannel dataclass on the config; config.yaml may
|
|
# also put it in extra (dict or bare string).
|
|
home = getattr(config, "home_channel", None)
|
|
if home is not None and getattr(home, "chat_id", None):
|
|
self.home_channel = str(home.chat_id)
|
|
self.home_channel_name = str(getattr(home, "name", "") or DEFAULT_HOME_CHANNEL_NAME)
|
|
else:
|
|
hc = extra.get("home_channel")
|
|
if isinstance(hc, dict) and hc.get("chat_id"):
|
|
self.home_channel = str(hc["chat_id"])
|
|
self.home_channel_name = str(hc.get("name") or DEFAULT_HOME_CHANNEL_NAME)
|
|
elif isinstance(hc, str) and hc.strip():
|
|
self.home_channel = hc.strip()
|
|
self.home_channel_name = DEFAULT_HOME_CHANNEL_NAME
|
|
else:
|
|
self.home_channel = DEFAULT_HOME_CHANNEL
|
|
self.home_channel_name = DEFAULT_HOME_CHANNEL_NAME
|
|
|
|
# TLS (optional)
|
|
self.http_cert = _get_scoped_secret("IRIS_HTTP_CERT") or extra.get("http_cert", "")
|
|
self.http_key = _get_scoped_secret("IRIS_HTTP_KEY") or extra.get("http_key", "")
|
|
|
|
# Auth
|
|
allowed = os.getenv("IRIS_ALLOWED_USERS", "").strip()
|
|
self.allowed_users: list[str] = (
|
|
[u.strip() for u in allowed.split(",") if u.strip()] if allowed else []
|
|
)
|
|
self.allow_all = _truthy(os.getenv("IRIS_ALLOW_ALL_USERS"))
|
|
|
|
# Runtime state
|
|
self._devices = DeviceRegistry(get_hermes_home() / "iris" / "devices.db")
|
|
# docs/19: HTTP transport (the only device-facing transport).
|
|
self._http_server = HttpServer(self, self._devices)
|
|
# docs/19 §19.7: reply sinks for in-flight HTTP requests — while a
|
|
# POST /v1/frame is being dispatched, the handler's point-to-point
|
|
# replies are captured here and returned as the HTTP response.
|
|
# Entry shape: (sink queue, abandoned event).
|
|
self._http_reply_sinks: dict[str, tuple[queue.Queue, threading.Event]] = {}
|
|
self._connected = False
|
|
# M2: per-chat turn state for outbound frame classification.
|
|
self._turns: dict[str, _TurnState] = {}
|
|
# M3: channel directory (shared singleton) + offline outbox.
|
|
self._channels = get_directory()
|
|
self._outbox = Outbox(
|
|
get_hermes_home() / "iris" / "outbox.db",
|
|
retention_hours=self.outbox_retention_hours,
|
|
)
|
|
# M4: media registry (inbound upload refs + outbound offers) and the
|
|
# last finalized assistant message id per chat (offer association).
|
|
self._media = media_bridge.MediaStore(get_hermes_home())
|
|
self._last_message_id: dict[str, str] = {}
|
|
# Interactive choice pickers (slash-command menus, e.g. /reasoning,
|
|
# /fast): picker_id -> pending state. In-memory only — a gateway
|
|
# restart expires them (a stale picker.select is a no-op).
|
|
self._pending_pickers: dict[str, dict] = {}
|
|
# M5: push backend (FCM primary, ntfy fallback) + the throttle for
|
|
# the outbox-prune banner.
|
|
self._push: PushBackend = build_push_backend(
|
|
self.push_backend,
|
|
fcm_service_account=_get_scoped_secret("IRIS_FCM_SERVICE_ACCOUNT"),
|
|
fcm_server_key=_get_scoped_secret("IRIS_FCM_SERVER_KEY"),
|
|
ntfy_topic=_get_scoped_secret("NTFY_TOPIC"),
|
|
ntfy_server_url=os.getenv("NTFY_SERVER_URL", "").strip() or None,
|
|
ntfy_auth_token=_get_scoped_secret("NTFY_AUTH_TOKEN"),
|
|
)
|
|
self._prune_notified_at = 0.0
|
|
# M5: per-chat push throttle (epoch seconds of the last successful
|
|
# push). The coalesced frame still reaches the app via sync.
|
|
self._last_push_at: dict[str, float] = {}
|
|
|
|
def _turn_state(self, chat_id: str) -> _TurnState:
|
|
st = self._turns.get(chat_id)
|
|
if st is None:
|
|
st = _TurnState()
|
|
self._turns[chat_id] = st
|
|
return st
|
|
|
|
@property
|
|
def name(self) -> str:
|
|
return "Iris"
|
|
|
|
# ── Connection lifecycle ──────────────────────────────────────────────
|
|
|
|
async def connect(self, *, is_reconnect: bool = False) -> bool:
|
|
"""Bring the platform up: bind the HTTP server on host:http_port."""
|
|
if not self.token:
|
|
logger.error("iris: IRIS_TOKEN must be set")
|
|
self._set_fatal_error(
|
|
"config_missing",
|
|
"IRIS_TOKEN must be set",
|
|
retryable=False,
|
|
)
|
|
return False
|
|
|
|
# The HTTP server is the only device-facing transport, so a bind
|
|
# failure is fatal (the app has no other way to reach the gateway).
|
|
# start() never raises; it disables the leg and logs on failure.
|
|
await self._http_server.start()
|
|
if not self._http_server.enabled:
|
|
logger.error("iris: HTTP server failed to bind %s:%s", self.host, self.http_port)
|
|
self._set_fatal_error(
|
|
"bind_failed",
|
|
f"HTTP port {self.http_port} unavailable",
|
|
retryable=False,
|
|
)
|
|
return False
|
|
|
|
# M5: announce gateway health to connected clients (none yet at
|
|
# startup; the frame + plumbing exist for future transitions).
|
|
# Reset in case this adapter instance previously went down (the
|
|
# gateway may reconnect the same adapter after a fatal error).
|
|
self._gateway_status = protocol.STATUS_ONLINE
|
|
await self._http_server.fanout(protocol.status(self._gateway_status), cursor=None)
|
|
|
|
# M3: ensure the default (home) channel exists in the directory so the
|
|
# app's channel list and cron home delivery have a stable anchor.
|
|
try:
|
|
self._channels.ensure_default(self.home_channel, self.home_channel_name)
|
|
except Exception:
|
|
logger.warning("iris: ensure_default failed", exc_info=True)
|
|
|
|
# M5: push backend status (degrade gracefully when unconfigured).
|
|
if not self._push.configured():
|
|
logger.warning(
|
|
"iris: push backend %r not configured (no credentials) -- "
|
|
"offline devices will not be woken; outbox + sync still apply",
|
|
self.push_backend,
|
|
)
|
|
else:
|
|
logger.info("iris: push backend: %s", self._push.name)
|
|
|
|
self._connected = True
|
|
self._mark_connected()
|
|
logger.info("iris: connected; HTTP server on %s:%s", self.host, self.http_port)
|
|
return True
|
|
|
|
async def disconnect(self) -> None:
|
|
"""Tear down the platform: stop the server, close device streams."""
|
|
# Tell live clients the gateway is going away (restart/shutdown) so
|
|
# the app can distinguish a clean gateway teardown from a plain
|
|
# network drop: the "Gateway restarting" chat notice is shown only
|
|
# when this frame was received (docs/04 §status).
|
|
self._gateway_status = protocol.STATUS_RESTARTING
|
|
with contextlib.suppress(Exception):
|
|
await self._http_server.fanout(protocol.status(self._gateway_status), cursor=None)
|
|
try:
|
|
await self._http_server.stop()
|
|
except Exception:
|
|
logger.warning("iris: HTTP server stop failed", exc_info=True)
|
|
# Best-effort shutdown: a close failure on an already-closed store is
|
|
# not actionable at disconnect time.
|
|
with contextlib.suppress(Exception):
|
|
self._devices.close()
|
|
with contextlib.suppress(Exception):
|
|
self._outbox.close()
|
|
self._connected = False
|
|
self._mark_disconnected()
|
|
logger.info("iris: disconnected")
|
|
|
|
# ── Outbound (agent -> app) ───────────────────────────────────────────
|
|
|
|
async def send(
|
|
self,
|
|
chat_id: str,
|
|
content: str,
|
|
reply_to: str | None = None,
|
|
metadata: dict[str, Any] | None = None,
|
|
) -> SendResult:
|
|
"""Send a message to a chat.
|
|
|
|
M2: classify the outbound call into a structured frame using the
|
|
per-chat turn state machine (see module docstring):
|
|
|
|
* ``metadata["expect_edits"]`` -> ``message.start`` (streaming segment)
|
|
* ``metadata["notify"]`` -> ``message`` / ``message.stop`` (final)
|
|
* tool-progress line format -> ``tool.start`` (first tool bubble)
|
|
* anything else -> ``commentary``
|
|
|
|
With no live devices the frame is dropped here (the outbox + push
|
|
replay lands in M3/M5).
|
|
"""
|
|
content = content or ""
|
|
meta = metadata or {}
|
|
thread_id = _thread_id_from_metadata(meta)
|
|
state = self._turn_state(chat_id)
|
|
|
|
# 1. Streaming segment start (stream consumer first send).
|
|
if bool(meta.get("expect_edits")):
|
|
# A new content segment means the tool the model was waiting on
|
|
# has returned -> close it before the segment opens.
|
|
await self._close_open_tool(chat_id, state, thread_id)
|
|
message_id = _mint_message_id()
|
|
state.active = True
|
|
state.stream_id = message_id
|
|
await self._broadcast_or_log(
|
|
chat_id,
|
|
protocol.message_start(
|
|
chat_id, message_id, protocol.ROLE_ASSISTANT, thread_id=thread_id
|
|
),
|
|
)
|
|
return SendResult(success=True, message_id=message_id)
|
|
|
|
# 2. Cron delivery (the cron scheduler passes ``job_id`` in metadata):
|
|
# a final assistant message + a high-priority notification banner
|
|
# (pushed even when a device is live, docs/08 §8.1).
|
|
if meta.get("job_id"):
|
|
reasoning, _ = _split_reasoning(content)
|
|
if reasoning:
|
|
_reset_reasoning()
|
|
else:
|
|
await _wait_for_reasoning_flushed()
|
|
reasoning = _take_reasoning() or None
|
|
name, inner = _cron_brief(content, str(meta.get("job_id")))
|
|
message_id = _mint_message_id()
|
|
await self._broadcast_or_log(
|
|
chat_id,
|
|
protocol.notification(
|
|
chat_id,
|
|
protocol.NOTIF_CRON,
|
|
f"Cron: {name}",
|
|
_push_preview(inner),
|
|
thread_id=thread_id,
|
|
),
|
|
)
|
|
await self._broadcast_or_log(
|
|
chat_id,
|
|
protocol.message(
|
|
chat_id=chat_id,
|
|
message_id=message_id,
|
|
role=protocol.ROLE_ASSISTANT,
|
|
text=inner,
|
|
thread_id=thread_id,
|
|
reasoning=reasoning,
|
|
ts=int(time.time() * 1000),
|
|
),
|
|
)
|
|
self._last_message_id[chat_id] = message_id
|
|
state.active = False
|
|
return SendResult(success=True, message_id=message_id)
|
|
|
|
# 3. Final message (non-streaming final, or streaming fallback final).
|
|
if bool(meta.get("notify")):
|
|
reasoning, body = _split_reasoning(content)
|
|
# Non-streaming: reasoning is prepended to content (split above).
|
|
# Streaming fallback: content has no reasoning, so use the
|
|
# reasoning captured via the on_stream_delta hook (wait for the
|
|
# async hook worker to flush it first).
|
|
if not reasoning:
|
|
await _wait_for_reasoning_flushed()
|
|
reasoning = _take_reasoning() or None
|
|
else:
|
|
_reset_reasoning()
|
|
# Runtime-metadata footer (app-controlled display): always attach
|
|
# the structured ``runtime`` object so the app can render its
|
|
# footer (Settings → Runtime footer). Independent of hermes'
|
|
# ``display.runtime_footer`` config.
|
|
runtime = await _build_runtime_footer(_take_runtime_meta())
|
|
if state.stream_id:
|
|
# Fallback final: close the open streaming segment in place.
|
|
message_id = state.stream_id
|
|
state.stream_id = None
|
|
await self._broadcast_or_log(
|
|
chat_id,
|
|
protocol.message_stop(
|
|
chat_id,
|
|
message_id,
|
|
body,
|
|
reasoning=reasoning,
|
|
thread_id=thread_id,
|
|
runtime=runtime or None,
|
|
ts=int(time.time() * 1000),
|
|
),
|
|
)
|
|
else:
|
|
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=body,
|
|
thread_id=thread_id,
|
|
reasoning=reasoning,
|
|
reply_to=reply_to,
|
|
runtime=runtime or None,
|
|
ts=int(time.time() * 1000),
|
|
),
|
|
)
|
|
# M4: media offers emitted after this final associate with it.
|
|
self._last_message_id[chat_id] = message_id
|
|
await self._close_open_tool(chat_id, state, thread_id)
|
|
self._reset_tool_state(state)
|
|
state.active = False
|
|
return SendResult(success=True, message_id=message_id)
|
|
|
|
# 4. Tool progress (first tool bubble of an editable line buffer).
|
|
# Gateway lifecycle notices (restart/shutdown/online) are system
|
|
# notices, not tool progress — skip the tool-line heuristic so they
|
|
# fall through to commentary instead of a never-completing card.
|
|
if not _is_gateway_lifecycle_notice(content) and _is_tool_progress(content):
|
|
return await self._emit_tool_lines(chat_id, content, state, thread_id, is_edit=False)
|
|
|
|
# 5. Commentary (interim assistant beat).
|
|
message_id = _mint_message_id()
|
|
state.active = True
|
|
await self._broadcast_or_log(
|
|
chat_id,
|
|
protocol.commentary(chat_id, message_id, content, thread_id=thread_id),
|
|
)
|
|
return SendResult(success=True, message_id=message_id)
|
|
|
|
async def edit_message(
|
|
self,
|
|
chat_id: str,
|
|
message_id: str,
|
|
content: str,
|
|
*,
|
|
finalize: bool = False,
|
|
metadata: dict[str, Any] | None = None,
|
|
) -> SendResult:
|
|
"""Edit a previously sent message (M2: drives streaming + tool updates).
|
|
|
|
* ``message_id == state.stream_id`` -> ``message.update``
|
|
(``finalize=True`` -> ``message.stop``).
|
|
* ``message_id == state.tool_msg_id`` -> tool-progress update
|
|
(new lines -> ``tool.start``).
|
|
* unknown id -> best-effort ``message.update``.
|
|
"""
|
|
content = content or ""
|
|
thread_id = _thread_id_from_metadata(metadata)
|
|
state = self._turn_state(chat_id)
|
|
|
|
if message_id and message_id == state.stream_id:
|
|
if finalize:
|
|
reasoning, body = _split_reasoning(_strip_streaming_cursor(content))
|
|
# Streaming: the gateway drops the model's separate
|
|
# reasoning_content (final send suppressed), so attach the
|
|
# reasoning we captured via the on_stream_delta hook (wait
|
|
# for the async hook worker to flush it first).
|
|
if not reasoning:
|
|
await _wait_for_reasoning_flushed()
|
|
reasoning = _take_reasoning() or None
|
|
else:
|
|
_reset_reasoning()
|
|
# Runtime-metadata footer (app-controlled display).
|
|
runtime = await _build_runtime_footer(_take_runtime_meta())
|
|
state.stream_id = None
|
|
await self._broadcast_or_log(
|
|
chat_id,
|
|
protocol.message_stop(
|
|
chat_id,
|
|
message_id,
|
|
body,
|
|
reasoning=reasoning,
|
|
thread_id=thread_id,
|
|
runtime=runtime or None,
|
|
ts=int(time.time() * 1000),
|
|
),
|
|
)
|
|
# M4: media offers emitted after this final associate with it.
|
|
self._last_message_id[chat_id] = message_id
|
|
await self._close_open_tool(chat_id, state, thread_id)
|
|
self._reset_tool_state(state)
|
|
state.active = False
|
|
else:
|
|
await self._broadcast_or_log(
|
|
chat_id,
|
|
protocol.message_update(
|
|
chat_id,
|
|
message_id,
|
|
_strip_streaming_cursor(content),
|
|
thread_id=thread_id,
|
|
),
|
|
)
|
|
return SendResult(success=True, message_id=message_id)
|
|
|
|
if message_id and message_id == state.tool_msg_id:
|
|
return await self._emit_tool_lines(chat_id, content, state, thread_id, is_edit=True)
|
|
|
|
# Unknown id: treat as a streaming update (best effort).
|
|
if finalize:
|
|
runtime = await _build_runtime_footer(_take_runtime_meta())
|
|
await self._broadcast_or_log(
|
|
chat_id,
|
|
protocol.message_stop(
|
|
chat_id,
|
|
message_id,
|
|
_strip_streaming_cursor(content),
|
|
thread_id=thread_id,
|
|
runtime=runtime or None,
|
|
ts=int(time.time() * 1000),
|
|
),
|
|
)
|
|
else:
|
|
await self._broadcast_or_log(
|
|
chat_id,
|
|
protocol.message_update(
|
|
chat_id,
|
|
message_id,
|
|
_strip_streaming_cursor(content),
|
|
thread_id=thread_id,
|
|
),
|
|
)
|
|
return SendResult(success=True, message_id=message_id)
|
|
|
|
# ── M2: tool-progress helpers ─────────────────────────────────────────
|
|
|
|
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
|
|
|
|
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()
|
|
|
|
# ── docs/19: HTTP-leg reply routing ───────────────────────────────────
|
|
|
|
def _http_register_sink(
|
|
self, device_id: str, entry: tuple[queue.Queue, threading.Event]
|
|
) -> None:
|
|
self._http_reply_sinks[device_id] = entry
|
|
|
|
def _http_pop_sink(self, device_id: str) -> tuple[queue.Queue, threading.Event] | None:
|
|
return self._http_reply_sinks.pop(device_id, None)
|
|
|
|
def _http_pop_sink_if(
|
|
self, device_id: str, sink: queue.Queue
|
|
) -> tuple[queue.Queue, threading.Event] | None:
|
|
"""Pop the sink entry only if it is still ours (a newer request from
|
|
the same device may have replaced it)."""
|
|
entry = self._http_reply_sinks.get(device_id)
|
|
if entry is None or entry[0] is not sink:
|
|
return None
|
|
return self._http_reply_sinks.pop(device_id, None)
|
|
|
|
async def _broadcast_both(self, frame: "protocol.Frame") -> None:
|
|
"""Bare (non-outbox) broadcast to live subscribers (docs/19): the
|
|
frame reaches every live SSE/long-poll subscriber."""
|
|
await self._http_server.fanout(frame, cursor=None)
|
|
|
|
async def _reply(self, device_id: str, frame: "protocol.Frame") -> None:
|
|
"""Point-to-point reply with broadcast fallback (docs/19 §19.7).
|
|
|
|
For an in-flight HTTP request (a reply sink is registered) the frame
|
|
goes into the HTTP response. Otherwise it is broadcast so the
|
|
device's SSE stream delivers it (single-user model).
|
|
"""
|
|
entry = self._http_reply_sinks.get(device_id)
|
|
if entry is not None:
|
|
entry[0].put(frame)
|
|
return
|
|
await self._broadcast_both(frame)
|
|
|
|
async def _broadcast_or_log(self, chat_id: str, frame: "protocol.Frame") -> None:
|
|
# M3/M5: always append to the outbox so a reconnecting app can catch
|
|
# up on *all* recent frames, not just the ones that were parked. This
|
|
# covers the case where the app's in-memory ChatStore is reset (e.g.
|
|
# activity recreation / composition recompose) while the process stays
|
|
# alive: the recreated controller re-syncs from its stale cursor and
|
|
# replays the frames it missed. The push below still only fires when
|
|
# there is no live subscriber (or for high-priority events).
|
|
try:
|
|
cursor = self._outbox.append(chat_id, frame.to_json())
|
|
except Exception:
|
|
logger.warning("iris: outbox append failed", exc_info=True)
|
|
return
|
|
# docs/19 §19.8: a device reading SSE/long-poll IS a live subscriber
|
|
# — count it in the delivery total or every message would push AND
|
|
# stream to a device that is already receiving it.
|
|
delivered = await self._http_server.fanout(frame, cursor)
|
|
if delivered == 0:
|
|
logger.info(
|
|
"iris: no live devices for %s; %s frame parked in outbox (cursor=%s)",
|
|
chat_id,
|
|
frame.type,
|
|
cursor,
|
|
)
|
|
# M5: wake the offline device(s) via the push backend.
|
|
await self._maybe_push(chat_id, frame, cursor)
|
|
await self._maybe_notify_outbox_prune(chat_id)
|
|
elif (
|
|
frame.type == protocol.TYPE_NOTIFICATION
|
|
and frame.payload.get("kind") in protocol.HIGH_PRIORITY_NOTIF_KINDS
|
|
):
|
|
# M5: high-priority events (approval/clarify/cron) push even when
|
|
# a device is live -- the app may be backgrounded and decides
|
|
# whether to also show an in-app banner (docs/08 §8.1).
|
|
await self._maybe_push(chat_id, frame, cursor)
|
|
|
|
# ── M5: push ───────────────────────────────────────────────────────────
|
|
|
|
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
|
|
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)
|
|
|
|
# ── M4: outbound media (agent -> app) ─────────────────────────────────
|
|
#
|
|
# 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 these ``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``.
|
|
|
|
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(
|
|
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)
|
|
|
|
# ── Inbound (app -> agent) ────────────────────────────────────────────
|
|
|
|
async def on_message_send(self, frame: protocol.Frame, device_id: str) -> None:
|
|
"""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()
|
|
|
|
# ── M3: channel directory management (app -> agent) ───────────────────
|
|
#
|
|
# 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).
|
|
|
|
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)
|
|
|
|
# ── 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")},
|
|
)
|
|
|
|
# ── M5: approval / clarify banners ────────────────────────────────────
|
|
|
|
async def send_slash_confirm(
|
|
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_choice_picker(
|
|
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(
|
|
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.
|
|
|
|
Renders the prompt as a proper ``message`` frame (the base text
|
|
fallback would route through ``send()`` and be misclassified as
|
|
commentary/tool progress) and keeps the gateway's text intercept
|
|
working 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,
|
|
),
|
|
)
|
|
if choices:
|
|
# Multi-select clarifies register their flag on the pending entry;
|
|
# look it up by id (mirrors the base text fallback).
|
|
_is_multi = False
|
|
try:
|
|
from tools import clarify_gateway as _cg
|
|
|
|
with _cg._lock:
|
|
_entry = _cg._entries.get(clarify_id)
|
|
_is_multi = bool(_entry and getattr(_entry, "multi_select", False))
|
|
except Exception:
|
|
_is_multi = False
|
|
lines = [f"❓ {question}", ""]
|
|
for i, choice in enumerate(choices, start=1):
|
|
lines.append(f" {i}. {choice}")
|
|
lines.append("")
|
|
if _is_multi:
|
|
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)
|
|
|
|
# ── 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,
|
|
}
|
|
|
|
# ── hello.ack helpers ─────────────────────────────────────────────────
|
|
|
|
def gateway_status(self) -> str:
|
|
"""Current gateway health state (sent to each pairing connection)."""
|
|
return self._gateway_status
|
|
|
|
def server_caps(self) -> dict[str, Any]:
|
|
"""Capability flags advertised in ``hello.ack`` (M4 surface)."""
|
|
return {
|
|
"streaming": True, # M2: message.start/update/stop
|
|
"reasoning": True, # M2: reasoning field on message / message.stop
|
|
"tools": True, # M2: tool.start/progress/end
|
|
"media": True, # M4: media.upload/offer/pull
|
|
"search": True, # M3: search frame
|
|
"commands_catalog": True, # commands.catalog frame (the "/" drawer)
|
|
"push": self.push_backend,
|
|
# M5: ntfy server URL (app listener discovery; "" when not ntfy).
|
|
"push_ntfy_server": (
|
|
self._push.server_url if isinstance(self._push, NtfyBackend) else ""
|
|
),
|
|
"pickers": True, # picker.choice / picker.select (slash-command menus)
|
|
}
|
|
|
|
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
|
|
# ---------------------------------------------------------------------------
|
|
|
|
|
|
def register(ctx):
|
|
"""Plugin entry point: called by the Hermes plugin system."""
|
|
# M2: capture the model's separate reasoning_content during streaming so
|
|
# it can be attached to the turn's message.stop frame (the gateway
|
|
# otherwise drops it when streaming suppresses the final send).
|
|
try:
|
|
ctx.register_hook("on_stream_delta", _on_stream_delta)
|
|
except Exception:
|
|
logger.debug("iris: on_stream_delta hook registration failed", exc_info=True)
|
|
# M2: capture each completed tool call's result + timing so the tool.end
|
|
# frame can carry the output (the gateway never streams tool output to
|
|
# platforms). The app shows it on demand (Settings → Tool detail).
|
|
try:
|
|
ctx.register_hook("post_tool_call", _on_post_tool_call)
|
|
except Exception:
|
|
logger.debug("iris: post_tool_call hook registration failed", exc_info=True)
|
|
# Runtime-metadata footer: capture the turn's model + prompt tokens (per
|
|
# provider call) so the final message can carry a structured ``runtime``
|
|
# object. The app decides whether/what to show (Settings → Runtime
|
|
# footer); the gateway always sends the data (independent of hermes
|
|
# ``display.runtime_footer`` config).
|
|
try:
|
|
ctx.register_hook("post_api_request", _on_post_api_request)
|
|
except Exception:
|
|
logger.debug("iris: post_api_request hook registration failed", exc_info=True)
|
|
ctx.register_platform(
|
|
name="iris",
|
|
label="Iris",
|
|
adapter_factory=IrisAdapter,
|
|
check_fn=check_requirements,
|
|
validate_config=validate_config,
|
|
is_connected=is_connected,
|
|
required_env=["IRIS_TOKEN"],
|
|
install_hint="No extra packages needed (httpx is a core dep)",
|
|
setup_fn=interactive_setup,
|
|
# Env-driven auto-configuration: seeds PlatformConfig.extra with
|
|
# host/port/push_backend + home_channel so env-only setups show up in
|
|
# gateway status without instantiating the adapter.
|
|
env_enablement_fn=_env_enablement,
|
|
# Cron home-channel delivery support (deliver=iris:<chat_id>[:<thread>]).
|
|
cron_deliver_env_var="IRIS_HOME_CHANNEL",
|
|
# Out-of-process cron delivery (best-effort; outbox is gateway-served).
|
|
standalone_sender_fn=_standalone_send,
|
|
# Native target syntax: "iris:<chat_id>[:<thread>]" (chat ids are direct, e.g. iris:chan_7).
|
|
parse_target_ref_fn=_parse_target_ref,
|
|
# Auth env vars for _is_user_authorized() integration.
|
|
allowed_users_env="IRIS_ALLOWED_USERS",
|
|
allow_all_env="IRIS_ALLOW_ALL_USERS",
|
|
# WS has no message-size limit.
|
|
max_message_length=0,
|
|
# Display.
|
|
emoji="📱",
|
|
pii_safe=False,
|
|
allow_update_command=True,
|
|
# LLM guidance.
|
|
platform_hint=(
|
|
"You are chatting with the user through their native Iris app "
|
|
"(Android/Desktop). It renders Markdown, inline code, images, "
|
|
"audio and video, and shows your reasoning and tool activity. "
|
|
"Conversations are organized into channels and optional threads. "
|
|
"Keep formatting rich but readable. "
|
|
"You can send media files natively: to deliver a file to the user, "
|
|
"include MEDIA:/absolute/path/to/file in your response. Images "
|
|
"(.png, .jpg, .webp) appear as photos, audio (.ogg, .mp3, .m4a) "
|
|
"plays inline, videos (.mp4, .webm, .mov) play inline, and other "
|
|
"files arrive as downloadable documents."
|
|
),
|
|
)
|