"""Inbound ``message.send`` handling (app -> agent). Mixin for ``adapter.IrisAdapter``. Echoes the user message to all devices (multi-device sync + ack), resolves ``media_refs`` to ``MessageEvent.media_urls``/``media_types``, applies auto-threading, and hands the ``MessageEvent`` to ``handle_message()`` (the gateway's command pipeline + agent turn). """ import asyncio import logging import threading import time import uuid from typing import Any from gateway.platforms.base import MessageEvent, MessageType from . import protocol from .classify import _derive_thread_name from .mixin_base import IrisAdapterBase logger = logging.getLogger(__name__) class InboundHandlers(IrisAdapterBase): """Inbound message.send (see module docstring).""" async def on_message_send(self, frame: protocol.Frame, device_id: str) -> None: # noqa: PLR0912,PLR0915 """Handle an inbound ``message.send`` frame. Echoes the user message to all devices (multi-device sync + ack), then builds a ``MessageEvent`` and hands it to ``handle_message()`` (the gateway's command pipeline + agent turn). M4: ``media_refs`` reference completed ``media.upload``s; they are resolved to ``MessageEvent.media_urls``/``media_types`` (local paths the agent's vision/audio tools can read) and echoed in the user message's ``media[]`` so every device renders the attachments. """ payload = frame.payload text = payload.get("text") text = text if isinstance(text, str) else "" refs_raw = payload.get("media_refs") media_refs = ( [r for r in refs_raw if isinstance(r, str) and r] if isinstance(refs_raw, list) else [] ) if not text.strip() and not media_refs: await self._reply( device_id, protocol.error( protocol.ERR_UNSUPPORTED, "message.send requires non-empty text", id=frame.id ), ) return chat_id = frame.chat_id or payload.get("chat_id") if not isinstance(chat_id, str) or not chat_id.strip(): chat_id = self.home_channel chat_id = chat_id.strip() # Automation channels are read-only for the user: they only receive # gateway-originated output (cron jobs, webhooks). Reject direct sends # (the app hides the composer for them, this is the server-side # enforcement). target = self._channels.get(chat_id) if target is not None and target.get("automation"): await self._reply( device_id, protocol.error( protocol.ERR_UNSUPPORTED, f"{target.get('name') or chat_id} is an automation channel " "(read-only: cron/webhook output only)", id=frame.id, ), ) return thread_id = frame.thread_id or payload.get("thread_id") if not isinstance(thread_id, str) or not thread_id.strip(): thread_id = None reply_to = payload.get("reply_to") if not isinstance(reply_to, str) or not reply_to.strip(): reply_to = None # Auto-threading (the app's Threads setting, docs/06 ยง6.3): a message # in a channel's flat lane gets its own fresh thread, the way Telegram # topic mode mints a topic per new conversation. The thread is named # instantly from the user's opening message (derived title) and the # LLM upgrades the name in the background. The user echo, the agent # turn, and all streaming frames then carry the new thread_id. # Skipped for slash commands (session-scoped, not conversation # starters) and replies (they continue where the user is). Threading # is only active on the default channel; other channels stay flat. auto_thread = bool(payload.get("auto_thread")) default_entry = self._channels.default() if ( auto_thread and thread_id is None and text.strip() and not text.lstrip().startswith("/") and reply_to is None and default_entry is not None and chat_id == default_entry["chat_id"] ): entry = self._channels.create( name=_derive_thread_name(text), kind="thread", parent_chat_id=chat_id, ) thread_id = entry["chat_id"] # Bare broadcast (like channel.create): the directory is # re-served on hello.ack, so no outbox entry is needed. await self._broadcast_both(protocol.channel_created(entry, auto=True)) self._schedule_thread_title_upgrade(entry["chat_id"], text) # M4: resolve media refs (single-use; unknown ref -> error). media_urls: list[str] = [] media_types: list[str] = [] media_wire: list[dict[str, Any]] = [] for ref in media_refs: entry = self._media.get_inbound(ref) if entry is None: await self._reply( device_id, protocol.error( protocol.ERR_UNSUPPORTED, f"unknown media_ref {ref}", id=frame.id ), ) return media_urls.append(entry.path) media_types.append(entry.mime) media_wire.append( { "media_id": entry.media_id, "kind": entry.kind, "mime": entry.mime, "size": entry.size, "filename": entry.filename, } ) device = self._devices.get(device_id) or {} user_name = device.get("name") or device_id # Echo to all devices: the sender confirms (server-assigned id), # other devices see the message too (single-user, multi-device). # Routed through _broadcast_or_log (not a bare broadcast) so the echo # is appended to the outbox: the app's ChatStore is in-memory only, so # after a process death / activity recreation the only way the user's # own message is restored is via the sync replay. Without this, user # messages vanish on reconnect while bot messages (already parked) # survive. message_id = f"m_{uuid.uuid4().hex[:16]}" echo = protocol.message( chat_id=chat_id, message_id=message_id, role=protocol.ROLE_USER, text=text, thread_id=thread_id, media=media_wire or None, reply_to=reply_to, ts=int(time.time() * 1000), ) await self._broadcast_or_log(chat_id, echo) # Refs are consumed by this message (no replay). for ref in media_refs: self._media.pop_inbound(ref) # M4: a new user turn starts -- stale offer association is dropped. self._last_message_id.pop(chat_id, None) kind = media_wire[0]["kind"] if media_wire else None if kind == "image": message_type = MessageType.PHOTO elif kind == "video": message_type = MessageType.VIDEO elif kind == "audio": message_type = MessageType.AUDIO elif kind == "voice": message_type = MessageType.VOICE elif kind == "document": message_type = MessageType.DOCUMENT else: message_type = MessageType.TEXT source = self.build_source( chat_id=chat_id, chat_name=self._channel_name(chat_id), chat_type="dm", user_id=device_id, user_name=user_name, thread_id=thread_id, ) event = MessageEvent( text=text, message_type=message_type, user_id=device_id, user_name=user_name, source=source, message_id=message_id, reply_to_message_id=reply_to, media_urls=media_urls, media_types=media_types, ) await self.handle_message(event) # M5: acknowledge the user message to the originating device (the # app shows โœ“โœ“) at the moment it is handed to the agent. await self._reply( device_id, protocol.read_receipt(chat_id, message_id), ) def _schedule_thread_title_upgrade(self, thread_id: str, text: str) -> None: """Upgrade an auto-created thread's name with the model's title. Stage 2 of hermes' two-stage session titling (``agent/title_generator .py``): the thread was created with an instant derived name; this background call on the ``title_generation`` auxiliary task replaces it with the model's title and broadcasts ``channel.renamed``. Best-effort โ€” any failure (config, model, network) leaves the derived name in place, and a thread the user already renamed or archived is untouched. """ loop = asyncio.get_running_loop() def _work() -> None: try: from agent.title_generator import generate_title title = generate_title(text) except Exception: logger.debug("Thread title upgrade failed", exc_info=True) return if not title: return entry = self._channels.get(thread_id) if entry is None or entry.get("archived"): return if (entry.get("name") or "") == title: return renamed = self._channels.rename(thread_id, title) if renamed is None: return try: asyncio.run_coroutine_threadsafe( self._broadcast_both(protocol.channel_renamed(renamed)), loop, ) except Exception: logger.debug("Thread title rename broadcast failed", exc_info=True) threading.Thread(target=_work, daemon=True, name="iris-thread-title").start()