"""Inbound query frames: catalog, search, sync, history, delete, push tokens, picker selection. Mixin for ``adapter.IrisAdapter``. These are the read/catch-up half of the inbound surface (``message.send`` lives in ``inbound.py``): they answer point-to-point (``_reply``) or broadcast the matching event frame. """ import logging from hermes_constants import get_hermes_home from . import protocol from . import purge as purge_bridge from . import search as search_bridge from .commands import _slash_command_catalog logger = logging.getLogger(__name__) class QueryFrameHandlers: """Inbound query frames (see module docstring).""" # ── Slash-command catalog (app's "/" drawer) ────────────────────────── async def on_commands_catalog(self, frame: protocol.Frame, device_id: str) -> None: """Handle an inbound ``commands.catalog`` request: reply with the gateway's slash-command catalog (hermes ``COMMAND_REGISTRY``, gateway-available subset + plugin commands). The app fuzzy-matches the typed prefix client-side; the catalog is static per gateway run, so no caching is needed here.""" resp = protocol.commands_catalog(_slash_command_catalog(), id=frame.id) await self._reply(device_id, resp) # ── M3: search (app -> agent) ───────────────────────────────────────── async def on_search(self, frame: protocol.Frame, device_id: str) -> None: payload = frame.payload query = payload.get("query") if not isinstance(query, str) or not query.strip(): await self._reply( device_id, protocol.error(protocol.ERR_UNSUPPORTED, "search requires a query", id=frame.id), ) return scope = payload.get("scope") scope = scope if scope in ("all", "chat") else "all" chat_id = payload.get("chat_id") or frame.chat_id if not isinstance(chat_id, str) or not chat_id.strip(): chat_id = None thread_id = payload.get("thread_id") or frame.thread_id if not isinstance(thread_id, str) or not thread_id.strip(): thread_id = None limit = payload.get("limit") try: limit = int(limit) if limit is not None else 20 except (TypeError, ValueError): limit = 20 db_path = get_hermes_home() / "state.db" hits = search_bridge.search( db_path, query, scope=scope, chat_id=chat_id, thread_id=thread_id, limit=limit ) resp = protocol.search_results(query, scope, hits, id=frame.id) await self._reply(device_id, resp) # ── M3: sync (reconnect catch-up) ───────────────────────────────────── async def on_sync(self, frame: protocol.Frame, device_id: str) -> None: payload = frame.payload cursor = payload.get("cursor") try: cursor = int(cursor) if cursor is not None else 0 except (TypeError, ValueError): cursor = 0 for e in self._outbox.replay(cursor): raw = e["frame"] replayed = protocol.Frame( type=raw.get("type", ""), payload=raw.get("payload", {}) if isinstance(raw.get("payload"), dict) else {}, id=raw.get("id") if isinstance(raw.get("id"), int) else None, chat_id=( raw.get("chat_id") if isinstance(raw.get("chat_id"), str) else e.get("chat_id") ), thread_id=raw.get("thread_id") if isinstance(raw.get("thread_id"), str) else None, # M5: tag replayed frames with their outbox cursor so the app # can skip re-notifying frames that already woke the device # via push (cursor <= last_pushed_cursor, docs/08 §8.7). cursor=e.get("cursor"), v=raw.get("v") if isinstance(raw.get("v"), int) else protocol.PROTOCOL_VERSION, ) await self._reply(device_id, replayed) done = protocol.sync_done(self._outbox.latest_cursor(), id=frame.id) await self._reply(device_id, done) # ── Full message history (initial channel open / scroll-up) ─────────── async def on_history(self, frame: protocol.Frame, device_id: str) -> None: """Handle an inbound ``history`` request. ``sync`` only replays the outbox delta since the device's cursor, so after a process death the app's in-memory ChatStore is empty and the delta does not cover older messages. ``history`` loads the full message list for a chat/thread (reconstructed from the outbox log) so the app can populate the view on first open / restart. """ payload = frame.payload chat_id = frame.chat_id or payload.get("chat_id") logger.info("iris: history request from %s chat_id=%r", device_id, chat_id) if not isinstance(chat_id, str) or not chat_id.strip(): await self._reply( device_id, protocol.error(protocol.ERR_UNSUPPORTED, "history requires a chat_id", id=frame.id), ) return chat_id = chat_id.strip() thread_id = frame.thread_id or payload.get("thread_id") if not isinstance(thread_id, str) or not thread_id.strip(): thread_id = None before = payload.get("before_message_id") if not isinstance(before, str) or not before.strip(): before = None limit_raw = payload.get("limit") try: limit = int(limit_raw) if limit_raw is not None else 50 except (TypeError, ValueError): limit = 50 page = self._outbox.history( chat_id, thread_id=thread_id, before_message_id=before, limit=limit, ) resp = protocol.history( chat_id, page["messages"], page["has_more"], thread_id=thread_id, oldest_message_id=page["oldest_message_id"], id=frame.id, ) await self._reply(device_id, resp) # ── Message deletion (app -> agent) ─────────────────────────────────── async def on_message_delete(self, frame: protocol.Frame, device_id: str) -> None: """Handle an inbound ``message.delete`` request. Completely deletes the requested message(s): they are removed from the outbox (so ``history`` and ``sync`` no longer return them) **and** from the hermes session store (so no search trace survives and they are not recoverable). ``message.deleted`` is broadcast to every device (outboxed too, so an offline device learns of the deletion on its next ``sync``). Deleting is idempotent: a message that is already gone (pruned by retention) simply yields 0 removed rows, and the ``message.deleted`` broadcast is still emitted so live caches drop it. """ payload = frame.payload chat_id = frame.chat_id or payload.get("chat_id") if not isinstance(chat_id, str) or not chat_id.strip(): await self._reply( device_id, protocol.error( protocol.ERR_NOT_FOUND, "message.delete requires chat_id", id=frame.id ), ) return chat_id = chat_id.strip() thread_id = frame.thread_id or payload.get("thread_id") if not isinstance(thread_id, str) or not thread_id.strip(): thread_id = None message_ids = payload.get("message_ids") if not isinstance(message_ids, list): message_ids = [payload.get("message_id")] if payload.get("message_id") else [] message_ids = [m for m in message_ids if isinstance(m, str) and m.strip()] if not message_ids: await self._reply( device_id, protocol.error( protocol.ERR_UNSUPPORTED, "message.delete requires message_ids", id=frame.id ), ) return removed = 0 purged = 0 db_path = get_hermes_home() / "state.db" for mid in message_ids: # Read the final frame data first (role / text / ts) so the # session-store row can be matched, then drop the outbox frames. info = self._outbox.message_info(chat_id, mid, thread_id=thread_id) removed += self._outbox.delete_message(chat_id, mid, thread_id=thread_id) if info: purged += purge_bridge.delete_message( db_path, chat_id, thread_id, info.get("role") or "", info.get("text") or "", info.get("ts"), ) logger.info( "iris: message.delete from %s chat_id=%r thread_id=%r ids=%s removed=%s purged=%s", device_id, chat_id, thread_id, message_ids, removed, purged, ) resp = protocol.message_deleted(chat_id, message_ids, thread_id=thread_id) resp.id = frame.id await self._broadcast_or_log(chat_id, resp) # ── M5: push token registration ─────────────────────────────────────── async def on_fcm_register(self, frame: protocol.Frame, device_id: str) -> None: """Update the device's push tokens (FCM rotation / ntfy topic). Persists to the device registry so the next push targets the current token without a stale read. """ fcm_token = frame.payload.get("fcm_token") ntfy_topic = frame.payload.get("ntfy_topic") fcm_token = fcm_token if isinstance(fcm_token, str) and fcm_token else None ntfy_topic = ntfy_topic if isinstance(ntfy_topic, str) and ntfy_topic else None if fcm_token is None and ntfy_topic is None: return try: self._devices.update_push_tokens(device_id, fcm_token=fcm_token, ntfy_topic=ntfy_topic) except Exception: logger.warning("iris: fcm.register update failed", exc_info=True) return logger.info("iris: push tokens updated for %s", device_id) # ── Interactive pickers (slash-command choice menus) ───────────────── async def on_picker_select(self, frame: protocol.Frame, device_id: str) -> None: """Resolve a pending choice picker (``picker.select`` from the app). Runs the command's selection callback and delivers its reply text as a normal final message in the picker's chat. Unknown/expired picker ids (gateway restart, double tap) are a no-op — the app already marked the card resolved locally. """ picker_id = frame.payload.get("picker_id") value = frame.payload.get("value") if not isinstance(picker_id, str) or not isinstance(value, str): return state = self._pending_pickers.pop(picker_id, None) if state is None: logger.info("iris: picker.select for unknown/expired picker %s", picker_id) return callback = state.get("on_choice_selected") if callback is None: return try: result_text = await callback(state["chat_id"], value) except Exception: logger.error("iris: picker selection failed for %s", picker_id, exc_info=True) return if not result_text: return await self.send( state["chat_id"], str(result_text), metadata={"notify": True, "thread_id": state.get("thread_id")}, )