Files
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

200 lines
8.0 KiB
Python

"""M5: push mirroring, typing indicators, turn-aware push.
Mixin for ``adapter.IrisAdapter``. Frames with no live subscriber wake the
device via the push backend (``push.py``); high-priority kinds push even
when a device is live. While the agent's turn is in flight (typing on),
normal-priority message/media frames are held back and the latest one is
pushed when the turn ends.
"""
import logging
import time
from typing import Any
from . import protocol
from .classify import _PUSH_COALESCE_S, _push_preview
from .mixin_base import IrisAdapterBase
logger = logging.getLogger(__name__)
# How often (seconds) the outbox-prune "storage reclaimed" notice may repeat.
_PRUNE_NOTIFY_INTERVAL_S = 3600.0
class PushHandlers(IrisAdapterBase):
"""Push + typing/turn tracking (see module docstring)."""
def _push_summary(self, frame: "protocol.Frame") -> tuple[str, str, str, str] | None:
"""``(title, body, kind, priority)`` for a pushable frame, else None.
Only terminal/interesting frames wake a device: intermediate
streaming and tool frames are replayed by ``sync`` without a push
(no notification spam per turn).
"""
t = frame.type
p = frame.payload
if t == protocol.TYPE_MESSAGE:
return (
self._channel_name(frame.chat_id or ""),
_push_preview(p.get("text")),
"message",
"normal",
)
if t == protocol.TYPE_MESSAGE_STOP:
return (
self._channel_name(frame.chat_id or ""),
_push_preview(p.get("final_text")),
"message",
"normal",
)
if t == protocol.TYPE_NOTIFICATION:
kind = str(p.get("kind") or protocol.NOTIF_GENERIC)
priority = "high" if kind in protocol.HIGH_PRIORITY_NOTIF_KINDS else "normal"
return (
str(p.get("title") or "Iris"),
str(p.get("body") or ""),
kind,
priority,
)
if t == protocol.TYPE_MEDIA_OFFER:
return (
self._channel_name(frame.chat_id or ""),
f"New {p.get('kind') or 'media'}: {p.get('filename') or ''}".strip(),
"media",
"normal",
)
return None
async def _maybe_push(self, chat_id: str, frame: "protocol.Frame", cursor: int) -> None:
"""Fire the configured push backend for a parked (or high-priority)
frame. Best-effort: failures are logged, never raised."""
summary = self._push_summary(frame)
if summary is None:
return
# M5: coalesce back-to-back pushes for the same chat (cron delivery
# = notification frame + message frame). The suppressed frame is
# still synced when the app reconnects.
now = time.time()
if now - self._last_push_at.get(chat_id, 0.0) < _PUSH_COALESCE_S:
logger.info(
"iris: push coalesced for %s (%s frame within %.0fs of last push)",
chat_id,
frame.type,
_PUSH_COALESCE_S,
)
return
title, body, kind, priority = summary
backend = self._push
if backend is None or not backend.token_field:
return
devices = self._devices.list()
if not backend.configured() and not any(d.get(backend.token_field) for d in devices):
return
data: dict[str, Any] = {"chat_id": chat_id, "kind": kind, "cursor": str(cursor)}
if frame.thread_id:
data["thread_id"] = frame.thread_id
message_id = frame.payload.get("message_id")
if isinstance(message_id, str) and message_id:
data["message_id"] = message_id
for device in devices:
device_id = device.get("device_id")
if not device_id:
continue
token = device.get(backend.token_field)
if not token:
continue
try:
ok = await backend.send(
device_id=device_id,
chat_id=chat_id,
title=title,
body=body,
data=data,
token=token,
priority=priority,
)
except Exception:
logger.warning("iris: push via %s failed", backend.name, exc_info=True)
continue
if ok:
# M5: remember that this cursor reached the device via push,
# so the app can dedupe it on the next sync replay.
try:
self._devices.update_push_cursor(device_id, cursor)
except Exception:
logger.warning(
"iris: push cursor update failed for %s", device_id, exc_info=True
)
self._last_push_at[chat_id] = time.time()
logger.info(
"iris: push via %s -> %s (%s, chat=%s)",
backend.name,
device_id,
frame.type,
chat_id,
)
async def _maybe_notify_outbox_prune(self, chat_id: str) -> None:
"""When the outbox row cap pruned old frames, tell the app (throttled
to once per hour so a full box doesn't banner per frame)."""
pruned = self._outbox.take_overflow_pruned()
if pruned <= 0:
return
now = time.time()
if now - self._prune_notified_at < _PRUNE_NOTIFY_INTERVAL_S:
return
self._prune_notified_at = now
await self._broadcast_or_log(
chat_id,
protocol.notification(
chat_id,
protocol.NOTIF_GENERIC,
"Outbox",
f"{pruned} older message(s) pruned",
),
)
async def send_typing(self, chat_id: str, metadata: dict[str, Any] | None = None) -> None:
"""Send a typing indicator (``typing`` frame, on=true)."""
thread_id = None
if metadata:
tid = metadata.get("thread_id")
if isinstance(tid, str) and tid:
thread_id = tid
# M5: typing on marks the chat's agent turn as in flight (hermes
# turns typing on at turn start, before the first output).
self._typing_turns.add(chat_id)
frame = protocol.typing(chat_id, True, thread_id=thread_id)
await self._http_server.fanout(frame, cursor=None)
async def stop_typing(self, chat_id: str) -> None:
"""Clear the typing indicator (``typing`` frame, on=false)."""
frame = protocol.typing(chat_id, False)
await self._http_server.fanout(frame, cursor=None)
# M5: turn ended (hermes fires stop_typing in the handler's finally,
# after the final send). Flush the held-back push -- the final
# answer -- but only while the device is still offline; a live
# device already got the frames via its event stream / sync.
# Idempotent: hermes may call stop_typing more than once per turn.
self._typing_turns.discard(chat_id)
pending = self._pending_push.pop(chat_id, None)
if pending is None:
return
if self._http_server.has_devices():
logger.info("iris: deferred push dropped for %s (device back online)", chat_id)
return
held_frame, held_cursor = pending
await self._maybe_push(chat_id, held_frame, held_cursor)
def on_device_online(self) -> None:
"""A device opened its event stream (SSE/long-poll): it will sync
the outbox, so drop any held-back pushes -- flushing them later
would duplicate what the app already shows. Called from the HTTP
server's handler thread; dict.clear() is atomic under the GIL."""
if self._pending_push:
logger.info(
"iris: device online; dropping %d deferred push(es)",
len(self._pending_push),
)
self._pending_push.clear()