Files
ARIA fb980d12b4
CI / Gateway plugin tests (push) Successful in 5m19s
CI / Kotlin tests (android host + desktop) (push) Successful in 7m0s
Release version management: single VERSION file as source of truth
- VERSION at repo root (0.1.2); bump it to cut a release
- App: generated AppVersion.kt (config-cache-safe Gradle task with
  VERSION as declared input) shown in Settings; sent to the gateway
  via X-Iris-App-Version header on the SSE open
- Gateway: reports its own version in hello.ack server_caps.app_version
  (read from the repo-root VERSION via the plugin symlink); stores the
  app's version in the device registry caps (merge, not overwrite, so
  an old app reconnecting without the header doesn't wipe it)
- Settings: app + gateway version rows, mismatch hint, and a best-effort
  Gitea latest-release check (ReleaseCheck) with an 'update available' hint
- Release workflow: reads VERSION from the repo (no manual input), with
  a guard against an empty file
- Docs: frames.schema.json + 04-wire-protocol.md updated for app_version
2026-08-25 14:42:39 +02:00

861 lines
38 KiB
Python
Raw Permalink Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
"""
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.
Module layout (split out of a single 3.4k-line file, issue #12):
* ``adapter.py`` — ``IrisAdapter`` core: lifecycle, outbound
classification (``send``/``edit_message``), broadcast/outbox/push routing,
plugin registration.
* ``inbound.py`` — ``message.send`` handling (app -> agent).
* ``classify.py`` — outbound frame classification: turn state, tool-line
parsing, reasoning split, previews.
* ``hooks.py`` — plugin hook capture (reasoning, tool results,
runtime metadata).
* ``tool_frames.py`` — tool.start / tool.end lifecycle.
* ``push_frames.py`` — push mirroring + typing/turn tracking.
* ``media_frames.py`` — outbound media (``media.offer``).
* ``picker_frames.py`` — approval / clarify / choice-picker frames.
* ``channel_frames.py`` — ``channel.*`` handlers + directory queries.
* ``query_frames.py`` — catalog / search / sync / history / delete frames.
* ``setup.py`` — interactive setup + config probes + env enablement.
* ``pickers.py`` — picker selection callbacks.
* ``commands.py`` — slash-command catalog.
* ``defaults.py`` — platform defaults.
* ``secrets.py`` — scope-aware secret reads.
Milestones M1–M5 (text round-trip, agent transparency, outbox/search, media,
push/offline) are documented in ``docs/14-milestones.md``; the frame protocol
in ``protocol.py`` (mirrored in ``docs/protocol/frames.schema.json``).
Configuration in config.yaml::
gateway:
platforms:
iris:
enabled: true
extra:
host: 127.0.0.1
http_port: 8791
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_HTTP_HOST, IRIS_HTTP_PORT, IRIS_HOME_CHANNEL,
IRIS_PUSH_BACKEND, IRIS_FCM_SERVICE_ACCOUNT, NTFY_TOPIC, ...
"""
import asyncio
import contextlib
import logging
import os
import queue
import threading
import time
from typing import Any
from gateway.config import Platform
from gateway.platforms.base import BasePlatformAdapter, SendResult
from hermes_constants import get_hermes_home
from . import hooks, protocol
from . import media as media_bridge
from .channel_frames import ChannelFrameHandlers
from .channels import get_directory
from .classify import (
_cron_brief,
_extract_verbose_args, # noqa: F401 (test-facing re-export)
_is_gateway_lifecycle_notice,
_is_tool_progress,
_mint_message_id,
_push_preview,
_short_preview_from_args, # noqa: F401 (test-facing re-export)
_split_reasoning,
_strip_streaming_cursor,
_thread_id_from_metadata,
_TurnState,
)
from .commands import _slash_command_catalog # noqa: F401 (test-facing re-export)
from .defaults import (
DEFAULT_HOME_CHANNEL,
DEFAULT_HOME_CHANNEL_NAME,
DEFAULT_HOST,
DEFAULT_HTTP_PORT,
DEFAULT_MAX_UPLOAD_BYTES,
DEFAULT_OUTBOX_RETENTION_HOURS,
DEFAULT_PUSH_BACKEND,
_truthy,
)
from .hooks import (
_build_runtime_footer,
_home_relative_cwd, # noqa: F401 (test-facing re-export)
_on_post_api_request,
_on_post_tool_call,
_on_stream_delta,
_reset_reasoning,
_reset_tool_results, # noqa: F401 (test-facing re-export)
_resolve_context_length, # noqa: F401 (test-facing re-export)
_take_reasoning,
_take_runtime_meta,
_tool_emoji, # noqa: F401 (test-facing re-export)
_tool_end_fields, # noqa: F401 (test-facing re-export)
_wait_for_reasoning_flushed,
)
from .http_server import HttpServer
from .inbound import InboundHandlers
from .media_frames import MediaHandlers
from .outbox import Outbox
from .pairing import DeviceRegistry
from .picker_frames import PickerHandlers
from .push import NtfyBackend, PushBackend, build_push_backend
from .push_frames import PushHandlers
from .query_frames import QueryFrameHandlers
from .secrets import _get_scoped_secret
from .setup import (
_ensure_verbose_tool_progress,
_env_enablement,
_offer_device_removal, # noqa: F401 (test-facing re-export)
_parse_port,
_parse_target_ref,
_standalone_send,
check_requirements,
interactive_setup,
is_connected,
validate_config,
)
from .tool_frames import ToolProgressHandlers
from .version import plugin_version
logger = logging.getLogger(__name__)
class IrisAdapter(
InboundHandlers,
ToolProgressHandlers,
PushHandlers,
MediaHandlers,
PickerHandlers,
QueryFrameHandlers,
ChannelFrameHandlers,
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).
self.host = os.getenv("IRIS_HTTP_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 (ntfy default, FCM optional) + 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] = {}
# M5: turn-aware push. chat_ids whose agent turn is in flight, per
# the typing indicator (hermes turns typing on at turn start and off
# in the handler's finally at turn end). While a turn is active and
# the device is offline, pushable message/media frames are held back
# in _pending_push instead of pushing per intermediate status
# message; the latest one is pushed when the turn ends, so an
# offline user gets ONE push with the final answer. High-priority
# notifications (approval/clarify/cron) always push immediately.
self._typing_turns: set[str] = set()
self._pending_push: dict[str, tuple[protocol.Frame, int]] = {}
# Live todo lists (lane key -> items): the agent's planning state,
# re-served as a snapshot when a device opens its event stream.
# In-memory only — a gateway restart drops it (the agent re-emits on
# the next todo call).
self._todo_lists: dict[str, list[dict[str, str]]] = {}
# The lane (chat_id, thread_id) of the turn currently in flight — the
# post_tool_call hook is global (no chat id), so the todo emission
# routes to the lane the adapter last saw tool activity for. A
# personal gateway serves one active turn at a time.
self._active_lane: tuple[str, str | None] | None = None
# The adapter's event loop (captured on connect): the synchronous
# plugin hooks schedule broadcasts onto it via run_coroutine_threadsafe.
self._loop: asyncio.AbstractEventLoop | None = None
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()
# The synchronous plugin hooks (post_tool_call) schedule frame
# broadcasts onto this loop; register the live adapter so they can
# reach it (single-adapter personal gateway).
self._loop = asyncio.get_running_loop()
hooks._live_adapter = self
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()
if hooks._live_adapter is self:
hooks._live_adapter = None
self._loop = None
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)
# ── 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)
# ── Live todo list (todo.update) ─────────────────────────────────────
@staticmethod
def _lane_key(chat_id: str, thread_id: str | None) -> str:
"""Lane key for a (chat, thread) pair (mirrors the app's ChatStore)."""
return chat_id if not thread_id else f"{chat_id}::{thread_id}"
async def _emit_todo_update(self, todos: list[dict[str, str]]) -> None:
"""Store + broadcast the agent's current todo list for the active lane.
Scheduled from the synchronous post_tool_call hook (the todo tool
just completed); *todos* is the tool result's full list.
"""
chat_id, thread_id = self._active_lane or (self.home_channel, None)
self._todo_lists[self._lane_key(chat_id, thread_id)] = todos
await self._broadcast_or_log(
chat_id, protocol.todo_update(chat_id, todos, thread_id=thread_id)
)
def todo_snapshot_frames(self) -> list["protocol.Frame"]:
"""``todo.update`` frames for every lane with a non-empty list.
Written by the HTTP server right after the hello frame when a device
opens its event stream, so a (re)connecting app re-learns the agent's
current plan (the frames are ephemeral — never outboxed).
"""
frames: list[protocol.Frame] = []
for lane, todos in self._todo_lists.items():
if not todos:
continue
parts = lane.split("::", 1)
thread_id = parts[1] if len(parts) > 1 else None
frames.append(protocol.todo_update(parts[0], todos, thread_id=thread_id))
return frames
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 -- unless
# this is a normal-priority message/media frame while the
# agent's turn is still in flight: hold it back and push the
# latest one when the turn ends (stop_typing), so an offline
# user gets one push with the final answer instead of one per
# intermediate status message. High-priority notifications
# (approval/clarify/cron) still push immediately.
if chat_id in self._typing_turns and frame.type in (
protocol.TYPE_MESSAGE,
protocol.TYPE_MESSAGE_STOP,
protocol.TYPE_MEDIA_OFFER,
):
self._pending_push[chat_id] = (frame, cursor)
logger.info(
"iris: push deferred for %s (turn active; %s frame cursor=%s)",
chat_id,
frame.type,
cursor,
)
else:
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)
# ── 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)
# Release version from the repo-root VERSION file (version.py);
# the app shows it in Settings and hints on app/gateway mismatch.
"app_version": plugin_version(),
}
# ---------------------------------------------------------------------------
# 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."
),
)