adapter.py was a 3,493-line monolith. Split it into focused modules with clear separation of responsibilities, bringing it down to ~857 lines: - Module-level helpers: hooks, classify, pickers, commands, setup, defaults, secrets - Frame-handler mixins: inbound, tool_frames, push_frames, media_frames, picker_frames, channel_frames, query_frames - mixin_base: IrisAdapterBase (declaration-only base for shared attrs) - adapter.py now holds only IrisAdapter (the composition of the 7 mixins + BasePlatformAdapter), register(), and test-facing re-exports The mixins come before BasePlatformAdapter in the MRO so their methods override the base; super() calls (e.g. send_image) still resolve to BasePlatformAdapter. No circular imports; dispatch.py and http_server.py (instance-method callers) are unaffected. Ruff complexity ceilings (PLR0911/0912/0913/0915) restored to Ruff's built-in defaults (12/50/6/5) instead of "just above the current maxima", which ratchets the bar down as code grows. The existing genuinely-complex functions (frame builders mirroring the wire schema, the QR matrix builder, the dispatch table) carry an explicit `# noqa: PLR09xx` marking them as reviewed, frozen exceptions; new code is held to the default ceilings. All 125 tests green (94 test_android + 31 test_android_http); no new ruff errors introduced.
858 lines
38 KiB
Python
858 lines
38 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.
|
||
|
||
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
|
||
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 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
|
||
|
||
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). 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 (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)
|
||
}
|
||
|
||
|
||
# ---------------------------------------------------------------------------
|
||
# 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."
|
||
),
|
||
)
|