"""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()