Files
iris_x_hermes/gateway-plugin/query_frames.py
T
ARIA b8e756c3dd
CI / Gateway plugin tests (pull_request) Successful in 4m59s
CI / Kotlin tests (android host + desktop) (pull_request) Successful in 7m5s
Split adapter.py monolith into focused modules; restore Ruff complexity defaults (issue #12)
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.
2026-08-24 20:55:00 +02:00

267 lines
12 KiB
Python

"""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")},
)