4 Commits
Author SHA1 Message Date
ARIA 70282dfb65 fix(release): use Forgejo-style /assets and /tags API routes
CI / Gateway plugin tests (push) Successful in 4m55s
CI / Kotlin tests (android host + desktop) (push) Successful in 6m48s
The server (gitea.zephyre.one) exposes a Forgejo-compatible API:
- release attachments live at POST /releases/{id}/assets, not /attachments
- tag deletion is DELETE /tags/{tag}, not DELETE /git/refs/tags/{tag}
(verified against the live API: /attachments 404s, /assets and /tags exist)
2026-08-23 19:53:17 +02:00
ARIA 44e8c7322e fix(release): delete leftover tag on re-run; support \\n in changelog input
CI / Gateway plugin tests (push) Successful in 5m0s
CI / Kotlin tests (android host + desktop) (push) Successful in 6m43s
- Gitea's DELETE /releases/:id does not remove the tag, so re-running the
  workflow for the same version failed with 409 (curl exit 22). The
  re-run safety block now also deletes the tag via git/refs/tags.
- Replace curl -sf with an api() wrapper that prints Gitea's error body
  on HTTP >= 400 instead of failing silently.
- The workflow_dispatch changelog input is single-line (Gitea has no
  multiline input type); convert literal \\n to real newlines and
  document it in the input description.
2026-08-23 19:03:59 +02:00
ARIA 29d0c1a73f Fix two gateway test failures: outbox lane scoping + SSE teardown race
CI / Gateway plugin tests (push) Successful in 5m46s
CI / Kotlin tests (android host + desktop) (push) Successful in 6m57s
- outbox: delete_message/message_info now match the exact lane first
  (a flat-lane delete/lookup with thread_id=None sees only frames with
  no thread_id) and fall back to the message_id across all lanes only
  when the exact lane matches nothing. Previously lane=None meant
  'any lane' in the first pass, so a flat-lane delete also removed
  same-id frames from threads (test expected 3 removed, got 4).

- http_server: the SSE live loop skipped queued frames when stop() set
  sub.closed before the handler thread reached the loop (descheduled
  under load between the initial hello/status writes and the loop).
  The loop now drains frames queued before the close, so the
  status{restarting} teardown broadcast always reaches the client
  before EOF (test_disconnect_broadcasts_status_restarting was flaky
  ~70% under CPU load).
2026-08-23 14:45:09 +02:00
ARIA d801a18db5 fix FCM push notifications
CI / Gateway plugin tests (push) Failing after 6m27s
CI / Kotlin tests (android host + desktop) (push) Successful in 7m2s
2026-08-23 14:32:40 +02:00
6 changed files with 238 additions and 30 deletions

No files matched your search

+28 -8
View File
@@ -8,7 +8,7 @@ on:
required: true required: true
type: string type: string
changelog: changelog:
description: "Release notes (markdown, shown on the release page)" description: "Release notes (markdown, shown on the release page). Single-line field — use literal \\n for line breaks."
required: false required: false
type: string type: string
@@ -178,19 +178,38 @@ jobs:
REPO="${GITEA_REPOSITORY:-$GITHUB_REPOSITORY}" REPO="${GITEA_REPOSITORY:-$GITHUB_REPOSITORY}"
TOKEN="${RELEASE_TOKEN:-$GITHUB_TOKEN}" TOKEN="${RELEASE_TOKEN:-$GITHUB_TOKEN}"
VERSION=$(jq -r '.inputs.version' "$GITHUB_EVENT_PATH") VERSION=$(jq -r '.inputs.version' "$GITHUB_EVENT_PATH")
CHANGELOG=$(jq -r '.inputs.changelog // ""' "$GITHUB_EVENT_PATH") # The dispatch input is a single-line field; turn literal \n into real newlines.
CHANGELOG=$(jq -r '.inputs.changelog // ""' "$GITHUB_EVENT_PATH" | sed 's/\\n/\n/g')
TAG="v$VERSION" TAG="v$VERSION"
API="$SERVER/api/v1/repos/$REPO" API="$SERVER/api/v1/repos/$REPO"
AUTH="Authorization: token $TOKEN" AUTH="Authorization: token $TOKEN"
# Re-run safety: drop a previous release (and its tag) for this version. # curl wrapper: on HTTP >= 400, print the response body (Gitea's error
OLD_ID=$(curl -sf -H "$AUTH" "$API/releases/tags/$TAG" | jq -r '.id // empty') # message) before failing — plain `curl -f` hides it (exit 22).
api() {
local code body
body=$(mktemp)
code=$(curl -s -o "$body" -w '%{http_code}' "$@") || { cat "$body"; rm -f "$body"; return 1; }
if [ "${code:0:1}" != "2" ]; then
echo "API error $code: $(cat "$body")" >&2
rm -f "$body"
return 1
fi
cat "$body"
rm -f "$body"
}
# Re-run safety: drop a previous release AND its tag for this version.
# (Gitea's DELETE /releases/:id does NOT remove the tag; a leftover tag
# makes the POST below fail with 409.)
OLD_ID=$(api -H "$AUTH" "$API/releases/tags/$TAG" | jq -r '.id // empty') || true
if [ -n "$OLD_ID" ]; then if [ -n "$OLD_ID" ]; then
curl -sf -X DELETE -H "$AUTH" "$API/releases/$OLD_ID" > /dev/null api -X DELETE -H "$AUTH" "$API/releases/$OLD_ID" > /dev/null
fi fi
api -X DELETE -H "$AUTH" "$API/tags/$TAG" > /dev/null || true
# Gitea creates the tag at the default branch HEAD automatically. # Gitea creates the tag at the default branch HEAD automatically.
RELEASE_ID=$(curl -sf -X POST -H "$AUTH" -H "Content-Type: application/json" \ RELEASE_ID=$(api -X POST -H "$AUTH" -H "Content-Type: application/json" \
"$API/releases" \ "$API/releases" \
-d "$(jq -n --arg tag "$TAG" --arg title "Iris $VERSION" --arg body "$CHANGELOG" \ -d "$(jq -n --arg tag "$TAG" --arg title "Iris $VERSION" --arg body "$CHANGELOG" \
'{tag_name:$tag, title:$title, body:$body}')" \ '{tag_name:$tag, title:$title, body:$body}')" \
@@ -200,7 +219,8 @@ jobs:
for f in "$GITHUB_WORKSPACE"/iris-android-v* "$GITHUB_WORKSPACE"/iris-desktop-*; do for f in "$GITHUB_WORKSPACE"/iris-android-v* "$GITHUB_WORKSPACE"/iris-desktop-*; do
[ -f "$f" ] || continue [ -f "$f" ] || continue
echo "Uploading $(basename "$f")" echo "Uploading $(basename "$f")"
curl -sf -X POST -H "$AUTH" -F "attachment=@$f" \ # Forgejo-style API: release assets live under /assets, not /attachments.
"$API/releases/$RELEASE_ID/attachments" > /dev/null api -X POST -H "$AUTH" -F "attachment=@$f" \
"$API/releases/$RELEASE_ID/assets" > /dev/null
done done
echo "Done: $SERVER/$REPO/releases/tag/$TAG" echo "Done: $SERVER/$REPO/releases/tag/$TAG"
+37 -4
View File
@@ -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,9 +57,10 @@ 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).
- Fallback: legacy **server key** (`IRIS_FCM_SERVER_KEY`) if no service - Fallback: legacy **server key** (`IRIS_FCM_SERVER_KEY`) if no service
account (simpler, but legacy). account (simpler, but legacy).
- Target = the device's **FCM token** (registered via `hello` / - 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). 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
@@ -132,4 +165,4 @@ device push watermark:
Residual edge: FCM is at-least-once, so a lost device ack can still produce a 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 duplicate *system-displayed* notification (two `FCM-Notification:*` ids). The
designed evolution is the data-only push option (§8.2.1), which moves display 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. into the app and lets it use a stable per-message notification id.
+60 -2
View File
@@ -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) ─────────────────────────────────
# #
+26 -6
View File
@@ -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:
@@ -516,12 +523,25 @@ class HttpServer:
for snap in self._adapter.todo_snapshot_frames(): for snap in self._adapter.todo_snapshot_frames():
self._write_sse(handler, "frame", None, snap.to_json()) self._write_sse(handler, "frame", None, snap.to_json())
# 3. Live frames (cursor=None frames have no id). # 3. Live frames (cursor=None frames have no id).
while not sub.closed.is_set(): while True:
try: if sub.closed.is_set():
item = sub.q.get(timeout=SSE_HEARTBEAT_S) # stop() can land between the initial writes above and
except queue.Empty: # this loop (the handler thread is descheduled under
self._write_raw(handler, ": hb\n\n") # load): drain the frames queued before the close — e.g.
continue # the status{restarting} teardown broadcast — so the
# client sees them before EOF instead of losing them to
# the closed check.
try:
item = sub.q.get_nowait()
except queue.Empty:
reason = "stopped"
break
else:
try:
item = sub.q.get(timeout=SSE_HEARTBEAT_S)
except queue.Empty:
self._write_raw(handler, ": hb\n\n")
continue
if item is _STOP: if item is _STOP:
reason = "stopped" reason = "stopped"
break break
+14 -10
View File
@@ -303,10 +303,11 @@ class Outbox:
frame of a streamed reply -- or ``None`` when the message is not in the frame of a streamed reply -- or ``None`` when the message is not in the
outbox (e.g. already pruned by retention). outbox (e.g. already pruned by retention).
The lane is matched exactly first; when that finds nothing the lookup The lane is matched exactly first (a flat-lane lookup, ``thread_id
falls back to the ``message_id`` alone (it is a unique uuid4), so a = None``, sees only frames with no ``thread_id``); when that finds
stale/missing ``thread_id`` on the request still resolves the row. nothing the lookup falls back to the ``message_id`` alone across all
(``lane=None`` in the scan means "any lane".) lanes (it is a unique uuid4), so a stale/missing ``thread_id`` on the
request still resolves the row.
""" """
if not message_id: if not message_id:
return None return None
@@ -315,7 +316,7 @@ class Outbox:
"SELECT frame FROM outbox WHERE chat_id = ?", (chat_id,) "SELECT frame FROM outbox WHERE chat_id = ?", (chat_id,)
).fetchall() ).fetchall()
def scan(lane: str | None) -> dict[str, Any] | None: def scan(lane: str | None, exact: bool) -> dict[str, Any] | None:
msg_frame: dict[str, Any] | None = None msg_frame: dict[str, Any] | None = None
stop_frame: dict[str, Any] | None = None stop_frame: dict[str, Any] | None = None
for r in rows: for r in rows:
@@ -325,7 +326,7 @@ class Outbox:
continue continue
if not isinstance(frame, dict): if not isinstance(frame, dict):
continue continue
if lane is not None and _frame_thread_id(frame) != lane: if exact and _frame_thread_id(frame) != lane:
continue continue
payload = frame.get("payload") payload = frame.get("payload")
if not isinstance(payload, dict) or payload.get("message_id") != message_id: if not isinstance(payload, dict) or payload.get("message_id") != message_id:
@@ -345,7 +346,7 @@ class Outbox:
} }
return msg_frame or stop_frame return msg_frame or stop_frame
return scan(thread_id) or scan(None) return scan(thread_id, exact=True) or scan(None, exact=False)
def delete_message( def delete_message(
self, self,
@@ -376,7 +377,7 @@ class Outbox:
"SELECT cursor, frame FROM outbox WHERE chat_id = ?", (chat_id,) "SELECT cursor, frame FROM outbox WHERE chat_id = ?", (chat_id,)
).fetchall() ).fetchall()
def cursors_for(lane: str | None) -> list[int]: def cursors_for(lane: str | None, exact: bool) -> list[int]:
out: list[int] = [] out: list[int] = []
for r in rows: for r in rows:
try: try:
@@ -385,14 +386,17 @@ class Outbox:
continue continue
if not isinstance(frame, dict): if not isinstance(frame, dict):
continue continue
if lane is not None and _frame_thread_id(frame) != lane: if exact and _frame_thread_id(frame) != lane:
continue continue
payload = frame.get("payload") payload = frame.get("payload")
if isinstance(payload, dict) and payload.get("message_id") == message_id: if isinstance(payload, dict) and payload.get("message_id") == message_id:
out.append(int(r["cursor"])) out.append(int(r["cursor"]))
return out return out
cursors = cursors_for(thread_id) or cursors_for(None) # Exact lane first (a flat-lane delete must not reach into
# threads); fall back to the message_id across all lanes only
# when the exact lane matches nothing (stale/missing thread_id).
cursors = cursors_for(thread_id, exact=True) or cursors_for(thread_id, exact=False)
if not cursors: if not cursors:
return 0 return 0
# One bound-parameter delete per cursor (a message spans only a few # One bound-parameter delete per cursor (a message spans only a few
+73
View File
@@ -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 ───────────────────────────────────────────────────────