fix FCM push notifications
This commit is contained in:
1 parent
560d9c19b3
commit
d801a18db5
4 files changed
+175
-4
No files matched your search
+35
-2
@@ -5,13 +5,44 @@ relay. **Decision: FCM primary, ntfy fallback** (`IRIS_PUSH_BACKEND`).
|
|||||||
|
|
||||||
## 8.1 When push fires
|
## 8.1 When push fires
|
||||||
|
|
||||||
- A frame targets a `chat_id` whose device is **disconnected** (WS closed) →
|
- A frame targets a `chat_id` whose device is **disconnected** (no live
|
||||||
drop to **outbox** + fire **push**.
|
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
|
- Also fire push for high-priority foreground events the user should see even if
|
||||||
the app is backgrounded (approvals, clarifies, cron completions) — the app
|
the app is backgrounded (approvals, clarifies, cron completions) — the app
|
||||||
decides whether to also show an in-app banner.
|
decides whether to also show an in-app banner.
|
||||||
- If the device is **connected**, no push (the live frame is enough).
|
- 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`)
|
## 8.2 `PushBackend` interface (`push.py`)
|
||||||
|
|
||||||
```python
|
```python
|
||||||
@@ -26,6 +57,7 @@ class PushBackend(Protocol):
|
|||||||
Selected at adapter init by `IRIS_PUSH_BACKEND` (`fcm` default, `ntfy`).
|
Selected at adapter init by `IRIS_PUSH_BACKEND` (`fcm` default, `ntfy`).
|
||||||
|
|
||||||
### 8.2.1 `FcmBackend` (primary)
|
### 8.2.1 `FcmBackend` (primary)
|
||||||
|
|
||||||
- **FCM HTTP v1 API** via `httpx` (core dep). Auth = Firebase **service
|
- **FCM HTTP v1 API** via `httpx` (core dep). Auth = Firebase **service
|
||||||
account** (`IRIS_FCM_SERVICE_ACCOUNT` JSON path) → mint a short-lived
|
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).
|
||||||
@@ -43,6 +75,7 @@ Selected at adapter init by `IRIS_PUSH_BACKEND` (`fcm` default, `ntfy`).
|
|||||||
devices).
|
devices).
|
||||||
|
|
||||||
### 8.2.2 `NtfyBackend` (fallback, self-host friendly)
|
### 8.2.2 `NtfyBackend` (fallback, self-host friendly)
|
||||||
|
|
||||||
- Reuses hermes's existing ntfy publish path (hermes ships an ntfy adapter).
|
- 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
|
- 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
|
`httpx` POST, with an `X-Title` / `X-Message` / `X-Tag` / `X-Priority` and a
|
||||||
|
|||||||
@@ -1360,6 +1360,16 @@ class IrisAdapter(BasePlatformAdapter):
|
|||||||
# M5: per-chat push throttle (epoch seconds of the last successful
|
# M5: per-chat push throttle (epoch seconds of the last successful
|
||||||
# push). The coalesced frame still reaches the app via sync.
|
# push). The coalesced frame still reaches the app via sync.
|
||||||
self._last_push_at: dict[str, float] = {}
|
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,
|
# Live todo lists (lane key -> items): the agent's planning state,
|
||||||
# re-served as a snapshot when a device opens its event stream.
|
# 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
|
# In-memory only — a gateway restart drops it (the agent re-emits on
|
||||||
@@ -1948,8 +1958,27 @@ class IrisAdapter(BasePlatformAdapter):
|
|||||||
frame.type,
|
frame.type,
|
||||||
cursor,
|
cursor,
|
||||||
)
|
)
|
||||||
# M5: wake the offline device(s) via the push backend.
|
# M5: wake the offline device(s) via the push backend -- unless
|
||||||
await self._maybe_push(chat_id, frame, cursor)
|
# 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)
|
await self._maybe_notify_outbox_prune(chat_id)
|
||||||
elif (
|
elif (
|
||||||
frame.type == protocol.TYPE_NOTIFICATION
|
frame.type == protocol.TYPE_NOTIFICATION
|
||||||
@@ -2099,6 +2128,9 @@ class IrisAdapter(BasePlatformAdapter):
|
|||||||
tid = metadata.get("thread_id")
|
tid = metadata.get("thread_id")
|
||||||
if isinstance(tid, str) and tid:
|
if isinstance(tid, str) and tid:
|
||||||
thread_id = 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)
|
frame = protocol.typing(chat_id, True, thread_id=thread_id)
|
||||||
await self._http_server.fanout(frame, cursor=None)
|
await self._http_server.fanout(frame, cursor=None)
|
||||||
|
|
||||||
@@ -2106,6 +2138,32 @@ class IrisAdapter(BasePlatformAdapter):
|
|||||||
"""Clear the typing indicator (``typing`` frame, on=false)."""
|
"""Clear the typing indicator (``typing`` frame, on=false)."""
|
||||||
frame = protocol.typing(chat_id, False)
|
frame = protocol.typing(chat_id, False)
|
||||||
await self._http_server.fanout(frame, cursor=None)
|
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) ─────────────────────────────────
|
# ── M4: outbound media (agent -> app) ─────────────────────────────────
|
||||||
#
|
#
|
||||||
|
|||||||
@@ -275,6 +275,13 @@ class HttpServer:
|
|||||||
def _add_sub(self, sub: _Subscriber) -> None:
|
def _add_sub(self, sub: _Subscriber) -> None:
|
||||||
with self._subs_lock:
|
with self._subs_lock:
|
||||||
self._subs.setdefault(sub.device_id, []).append(sub)
|
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:
|
def _remove_sub(self, sub: _Subscriber) -> None:
|
||||||
with self._subs_lock:
|
with self._subs_lock:
|
||||||
|
|||||||
@@ -1729,6 +1729,79 @@ async def test_push_skipped_when_backend_unconfigured(adapter):
|
|||||||
assert fake.calls == []
|
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 ───────────────────────────────────────────────────────
|
# ── M5: fcm.register ───────────────────────────────────────────────────────
|
||||||
|
|
||||||
|
|
||||||
|
|||||||
Reference in new issue
Block a user