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.
267 lines
12 KiB
Python
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")},
|
|
)
|