From d801a18db5ca41a1efd7bcc684d0344efa58e728 Mon Sep 17 00:00:00 2001 From: ARIA Date: Sun, 23 Aug 2026 14:32:40 +0200 Subject: [PATCH] fix FCM push notifications --- docs/08-push.md | 41 ++++++++++++++-- gateway-plugin/adapter.py | 62 ++++++++++++++++++++++- gateway-plugin/http_server.py | 7 +++ gateway-plugin/tests/test_android.py | 73 ++++++++++++++++++++++++++++ 4 files changed, 177 insertions(+), 6 deletions(-) diff --git a/docs/08-push.md b/docs/08-push.md index b1960b9..718dde7 100644 --- a/docs/08-push.md +++ b/docs/08-push.md @@ -5,13 +5,44 @@ relay. **Decision: FCM primary, ntfy fallback** (`IRIS_PUSH_BACKEND`). ## 8.1 When push fires -- A frame targets a `chat_id` whose device is **disconnected** (WS closed) → - drop to **outbox** + fire **push**. +- A frame targets a `chat_id` whose device is **disconnected** (no live + SSE/long-poll subscriber) → drop to **outbox** + fire **push** — with one + refinement, *turn-aware push* (below). - Also fire push for high-priority foreground events the user should see even if the app is backgrounded (approvals, clarifies, cron completions) — the app decides whether to also show an in-app banner. - If the device is **connected**, no push (the live frame is enough). +### 8.1.1 Turn-aware push (one push per turn, final answer as body) + +An agent turn can span minutes and emit several completed status messages +("researching X…", "found Y…", "writing findings…", final answer). Pushing each +parked message would spam an offline user with the steps in between. So while +the agent's turn is **in flight** (hermes holds the typing indicator on from +turn start until the handler's `finally` at turn end), normal-priority +`message` / `message.stop` / `media.offer` frames that park with no live +device are **held back** per chat instead of pushing; the latest one is pushed +when the turn ends (`stop_typing`), so the offline user gets **one push with +the final answer**. Details: + +- Turn state is tracked per `chat_id` from the typing indicator + (`send_typing` → in flight, `stop_typing` → ended; hermes fires + `stop_typing` in the handler's `finally`, after the final send, so the + flush always sees the final frame). +- The held-back frame is still parked in the outbox — sync catch-up is + unaffected; only the push is deferred. +- **High-priority notifications** (approval/clarify/cron) push immediately, + even mid-turn — they need user action. +- If the device **reconnects mid-turn** (SSE/long-poll open), the held-back + push is dropped: the app syncs the parked frames and must not get a +duplicate push at turn end. +- If the turn ends while the device is live, nothing is pushed (the frames + were delivered live / synced). +- Turns without a typing indicator (e.g. typing disabled in config) and + non-turn deliveries (cron, standalone sends) push immediately as before. +- Best-effort: a gateway crash mid-turn loses the held-back push (the frames + remain in the outbox and sync on reconnect). + ## 8.2 `PushBackend` interface (`push.py`) ```python @@ -26,9 +57,10 @@ class PushBackend(Protocol): Selected at adapter init by `IRIS_PUSH_BACKEND` (`fcm` default, `ntfy`). ### 8.2.1 `FcmBackend` (primary) + - **FCM HTTP v1 API** via `httpx` (core dep). Auth = Firebase **service account** (`IRIS_FCM_SERVICE_ACCOUNT` JSON path) → mint a short-lived - OAuth2 access token (cached, refreshed before expiry). + OAuth2 access token (cached, refreshed before expiry). - Fallback: legacy **server key** (`IRIS_FCM_SERVER_KEY`) if no service account (simpler, but legacy). - Target = the device's **FCM token** (registered via `hello` / @@ -43,6 +75,7 @@ Selected at adapter init by `IRIS_PUSH_BACKEND` (`fcm` default, `ntfy`). devices). ### 8.2.2 `NtfyBackend` (fallback, self-host friendly) + - Reuses hermes's existing ntfy publish path (hermes ships an ntfy adapter). - Publish to `NTFY_TOPIC` on `NTFY_SERVER_URL` (default `https://ntfy.sh`) via `httpx` POST, with an `X-Title` / `X-Message` / `X-Tag` / `X-Priority` and a @@ -132,4 +165,4 @@ device push watermark: Residual edge: FCM is at-least-once, so a lost device ack can still produce a duplicate *system-displayed* notification (two `FCM-Notification:*` ids). The designed evolution is the data-only push option (§8.2.1), which moves display -into the app and lets it use a stable per-message notification id. \ No newline at end of file +into the app and lets it use a stable per-message notification id. diff --git a/gateway-plugin/adapter.py b/gateway-plugin/adapter.py index 65cda54..4071d5c 100644 --- a/gateway-plugin/adapter.py +++ b/gateway-plugin/adapter.py @@ -1360,6 +1360,16 @@ class IrisAdapter(BasePlatformAdapter): # M5: per-chat push throttle (epoch seconds of the last successful # push). The coalesced frame still reaches the app via sync. self._last_push_at: dict[str, float] = {} + # M5: turn-aware push. chat_ids whose agent turn is in flight, per + # the typing indicator (hermes turns typing on at turn start and off + # in the handler's finally at turn end). While a turn is active and + # the device is offline, pushable message/media frames are held back + # in _pending_push instead of pushing per intermediate status + # message; the latest one is pushed when the turn ends, so an + # offline user gets ONE push with the final answer. High-priority + # notifications (approval/clarify/cron) always push immediately. + self._typing_turns: set[str] = set() + self._pending_push: dict[str, tuple[protocol.Frame, int]] = {} # Live todo lists (lane key -> items): the agent's planning state, # re-served as a snapshot when a device opens its event stream. # In-memory only — a gateway restart drops it (the agent re-emits on @@ -1948,8 +1958,27 @@ class IrisAdapter(BasePlatformAdapter): frame.type, cursor, ) - # M5: wake the offline device(s) via the push backend. - await self._maybe_push(chat_id, frame, cursor) + # M5: wake the offline device(s) via the push backend -- unless + # this is a normal-priority message/media frame while the + # agent's turn is still in flight: hold it back and push the + # latest one when the turn ends (stop_typing), so an offline + # user gets one push with the final answer instead of one per + # intermediate status message. High-priority notifications + # (approval/clarify/cron) still push immediately. + if chat_id in self._typing_turns and frame.type in ( + protocol.TYPE_MESSAGE, + protocol.TYPE_MESSAGE_STOP, + protocol.TYPE_MEDIA_OFFER, + ): + self._pending_push[chat_id] = (frame, cursor) + logger.info( + "iris: push deferred for %s (turn active; %s frame cursor=%s)", + chat_id, + frame.type, + cursor, + ) + else: + await self._maybe_push(chat_id, frame, cursor) await self._maybe_notify_outbox_prune(chat_id) elif ( frame.type == protocol.TYPE_NOTIFICATION @@ -2099,6 +2128,9 @@ class IrisAdapter(BasePlatformAdapter): 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) @@ -2106,6 +2138,32 @@ class IrisAdapter(BasePlatformAdapter): """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() # ── M4: outbound media (agent -> app) ───────────────────────────────── # diff --git a/gateway-plugin/http_server.py b/gateway-plugin/http_server.py index e93181e..a2e495d 100644 --- a/gateway-plugin/http_server.py +++ b/gateway-plugin/http_server.py @@ -275,6 +275,13 @@ class HttpServer: def _add_sub(self, sub: _Subscriber) -> None: with self._subs_lock: self._subs.setdefault(sub.device_id, []).append(sub) + # M5: a live subscriber will sync the outbox -- tell the adapter to + # drop any held-back (deferred) pushes so the turn-end flush doesn't + # duplicate what the app already shows. getattr-guard: test doubles + # may use a bare adapter stub. + on_online = getattr(self._adapter, "on_device_online", None) + if on_online is not None: + on_online() def _remove_sub(self, sub: _Subscriber) -> None: with self._subs_lock: diff --git a/gateway-plugin/tests/test_android.py b/gateway-plugin/tests/test_android.py index f257a9c..3006b50 100644 --- a/gateway-plugin/tests/test_android.py +++ b/gateway-plugin/tests/test_android.py @@ -1729,6 +1729,79 @@ async def test_push_skipped_when_backend_unconfigured(adapter): assert fake.calls == [] +@pytest.mark.asyncio +async def test_push_deferred_while_turn_active_flushed_on_stop_typing(adapter): + """Turn-aware push: while the agent's turn is in flight (typing on) and + the device is offline, intermediate status messages park without a push; + the turn-end (stop_typing) flushes ONE push with the latest (final) + message. Regression: a long research turn pushed every status report.""" + fake = _FakePush() + adapter._push = fake + adapter._devices.upsert(DEVICE_ID, "Test", {}, fcm_token="tok-1") + + await adapter.send_typing("default") # turn starts + await adapter.send("default", "researching XXXX", metadata={"notify": True}) + assert fake.calls == [] # intermediate: held back + await adapter.send( + "default", "found XXX, cross-referencing", metadata={"notify": True} + ) + assert fake.calls == [] # still held back; latest replaces the pending + await adapter.send("default", "done, findings ready", metadata={"notify": True}) + assert fake.calls == [] + + await adapter.stop_typing("default") # turn ends + assert len(fake.calls) == 1 + call = fake.calls[0] + assert call["chat_id"] == "default" + assert call["data"]["kind"] == "message" + assert "done, findings ready" in call["body"] + # All frames are parked in the outbox for sync catch-up. + assert adapter._outbox.latest_cursor() == 3 + + +@pytest.mark.asyncio +async def test_deferred_push_dropped_when_device_comes_online(adapter, ws_client): + """If the device reconnects mid-turn it syncs the parked frames, so the + turn-end flush must not push a duplicate.""" + ws, _ = ws_client + fake = _FakePush() + adapter._push = fake + adapter._devices.upsert(DEVICE_ID, "Test", {}, fcm_token="tok-1") + + await adapter.send_typing("default") + await adapter.send("default", "intermediate", metadata={"notify": True}) + assert fake.calls == [] + # Device reconnects mid-turn (SSE open -> on_device_online clears the + # pending push); the parked frame arrives via the event stream. + frames = await recv_until(ws, lambda f: f.get("type") == "message") + assert frames[-1]["payload"]["text"] == "intermediate" + await adapter.stop_typing("default") + assert fake.calls == [] + + + + + +@pytest.mark.asyncio +async def test_high_priority_notification_pushes_immediately_during_turn(plugin, adapter): + """Approval/clarify/cron notifications need user action: they push + immediately even while a turn is in flight and the device is offline.""" + fake = _FakePush() + adapter._push = fake + adapter._devices.upsert(DEVICE_ID, "Test", {}, fcm_token="tok-1") + + await adapter.send_typing("default") + await adapter._broadcast_or_log( + "default", + plugin.protocol.notification( + "default", plugin.protocol.NOTIF_APPROVAL, "Approve", "run rm -rf?" + ), + ) + assert len(fake.calls) == 1 + assert fake.calls[0]["priority"] == "high" + assert fake.calls[0]["data"]["kind"] == "approval" + + # ── M5: fcm.register ───────────────────────────────────────────────────────