From 9f3f9842c88b4b0efb4d62311bc61447bab751d5 Mon Sep 17 00:00:00 2001 From: ARIA Date: Fri, 21 Aug 2026 17:21:54 +0200 Subject: [PATCH] =?UTF-8?q?Push=20notification=20dedupe:=20one=20message?= =?UTF-8?q?=20=3D=20one=20notification=20=E2=80=94=20an=20offline=20messag?= =?UTF-8?q?e=20was=20notified=20twice=20(FCM=20push,=20then=20again=20when?= =?UTF-8?q?=20the=20app=20synced=20the=20outbox=20and=20mirrored=20the=20r?= =?UTF-8?q?eplayed=20frames).=20Fix:=20the=20gateway=20records=20the=20hig?= =?UTF-8?q?hest=20outbox=20cursor=20delivered=20per=20device=20via=20push?= =?UTF-8?q?=20(devices.last=5Fpushed=5Fcursor,=20advanced=20only=20on=20su?= =?UTF-8?q?ccessful=20send)=20and=20returns=20it=20in=20hello.ack;=20sync-?= =?UTF-8?q?replayed=20frames=20carry=20their=20outbox=20cursor=20in=20the?= =?UTF-8?q?=20envelope;=20the=20app=20skips=20system=20notifications=20for?= =?UTF-8?q?=20replayed=20frames=20at/below=20the=20watermark=20(live=20fra?= =?UTF-8?q?mes=20never=20suppressed=20=E2=80=94=20that=20is=20the=20case?= =?UTF-8?q?=20where=20no=20push=20fired).=20Also:=205s=20per-chat=20push?= =?UTF-8?q?=20coalescing=20so=20a=20cron=20delivery=20(notification=20fram?= =?UTF-8?q?e=20+=20message=20frame)=20pushes=20once,=20and=20the=20FCM=20h?= =?UTF-8?q?andler=20no=20longer=20posts=20a=20redundant=20notification=20(?= =?UTF-8?q?skips=20when=20WS=20is=20connected=20or=20FCM=20already=20displ?= =?UTF-8?q?ayed=20the=20notification=20payload;=20data-only=20messages=20a?= =?UTF-8?q?re=20the=20exception).=20Docs:=20frames.schema.json,=2004-wire-?= =?UTF-8?q?protocol.md,=2008-push.md=20=C2=A78.8?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- .../platform/IrisFirebaseMessagingService.kt | 15 +- .../kotlin/iris/net/GatewayClient.kt | 11 +- .../kotlin/iris/protocol/Protocol.kt | 9 + .../kotlin/iris/state/IrisController.kt | 27 ++- docs/04-wire-protocol.md | 9 + docs/08-push.md | 32 ++- docs/protocol/frames.schema.json | 2 + gateway-plugin/adapter.py | 31 +++ gateway-plugin/pairing.py | 65 +++++- gateway-plugin/protocol.py | 195 +++++++++++------- gateway-plugin/ws_server.py | 3 + 11 files changed, 301 insertions(+), 98 deletions(-) diff --git a/app/shared/src/androidMain/kotlin/iris/platform/IrisFirebaseMessagingService.kt b/app/shared/src/androidMain/kotlin/iris/platform/IrisFirebaseMessagingService.kt index 466f9f6..a6128bc 100644 --- a/app/shared/src/androidMain/kotlin/iris/platform/IrisFirebaseMessagingService.kt +++ b/app/shared/src/androidMain/kotlin/iris/platform/IrisFirebaseMessagingService.kt @@ -4,6 +4,7 @@ import android.content.pm.PackageManager import androidx.core.content.ContextCompat import com.google.firebase.messaging.FirebaseMessagingService import com.google.firebase.messaging.RemoteMessage +import iris.net.GatewayClient /** * M5: FCM handler (docs/08 §8.1). @@ -30,13 +31,23 @@ class IrisFirebaseMessagingService : FirebaseMessagingService() { } override fun onMessageReceived(message: RemoteMessage) { + // Foreground + live WS: the in-app banner already showed this. + if (AppBridge.foreground) return + // Live WS: the frame arrives over the socket and the controller + // mirrors it to a system notification itself — posting here would + // duplicate it (docs/08 §8.7). + if (AppBridge.controller?.client?.state?.value is GatewayClient.State.Connected) return + // Backgrounded/killed: FCM already displayed the `notification` + // payload on our behalf (the data payload only carries sync + // metadata). Posting again would show a second notification with a + // different id. Data-only messages (no notification payload) are the + // exception: the app must display them itself. + if (message.notification != null) return val data = message.data val chatId = data["chat_id"] ?: "android:default" val threadId = data["thread_id"] val title = data["title"] ?: "Iris" val body = data["body"] ?: data["title"].orEmpty() - // Foreground + live WS: the in-app banner already showed this. - if (AppBridge.foreground) return if (ContextCompat.checkSelfPermission(applicationContext, android.Manifest.permission.POST_NOTIFICATIONS) != PackageManager.PERMISSION_GRANTED ) { diff --git a/app/shared/src/commonMain/kotlin/iris/net/GatewayClient.kt b/app/shared/src/commonMain/kotlin/iris/net/GatewayClient.kt index e9c073c..2ed3491 100644 --- a/app/shared/src/commonMain/kotlin/iris/net/GatewayClient.kt +++ b/app/shared/src/commonMain/kotlin/iris/net/GatewayClient.kt @@ -70,7 +70,14 @@ class GatewayClient( sealed interface State { data object Disconnected : State data object Connecting : State - data class Connected(val caps: ServerCaps, val channels: List) : State + data class Connected( + val caps: ServerCaps, + val channels: List, + /** M5: highest outbox cursor already pushed to this device + * (from hello.ack; 0 = never). Sync-replayed frames at/below + * it must not re-post system notifications (docs/08 §8.7). */ + val lastPushedCursor: Long = 0, + ) : State data object Reconnecting : State data class AuthFailed(val message: String) : State } @@ -275,7 +282,7 @@ class GatewayClient( helloAck.invokeOnCompletion { e -> if (e == null) { val ack = helloAck.getCompleted() - _state.value = State.Connected(ack.serverCaps, ack.channels) + _state.value = State.Connected(ack.serverCaps, ack.channels, ack.lastPushedCursor) // M5: reconnect catch-up — replay frames parked while offline. val local = store.syncCursor if (local < ack.syncCursor) { diff --git a/app/shared/src/commonMain/kotlin/iris/protocol/Protocol.kt b/app/shared/src/commonMain/kotlin/iris/protocol/Protocol.kt index fa20fc9..c5c5b0c 100644 --- a/app/shared/src/commonMain/kotlin/iris/protocol/Protocol.kt +++ b/app/shared/src/commonMain/kotlin/iris/protocol/Protocol.kt @@ -115,6 +115,11 @@ data class Frame( val type: String, @SerialName("chat_id") val chatId: String? = null, @SerialName("thread_id") val threadId: String? = null, + /** Outbox cursor this frame was parked under. Set only on frames + * replayed by `sync` (live frames carry none) — the app skips + * re-notifying replayed frames with `cursor <= lastPushedCursor` + * (they already woke the device via push, docs/08 §8.7). */ + val cursor: Long? = null, val payload: JsonElement = JsonObject(emptyMap()), ) { /** Payload as a JSON object (the wire format); parse per-type with @@ -185,6 +190,10 @@ data class HelloAckPayload( @SerialName("server_caps") val serverCaps: ServerCaps = ServerCaps(), @SerialName("sync_cursor") val syncCursor: Long = 0, val channels: List = emptyList(), + /** M5: highest outbox cursor already delivered to THIS device via the + * push backend (0 = never). Sync-replayed frames at/below it must not + * re-post system notifications (dedupe, docs/08 §8.7). */ + @SerialName("last_pushed_cursor") val lastPushedCursor: Long = 0, ) // ── message (server -> app) ───────────────────────────────────────────── diff --git a/app/shared/src/commonMain/kotlin/iris/state/IrisController.kt b/app/shared/src/commonMain/kotlin/iris/state/IrisController.kt index 6bf5baf..e8869ec 100644 --- a/app/shared/src/commonMain/kotlin/iris/state/IrisController.kt +++ b/app/shared/src/commonMain/kotlin/iris/state/IrisController.kt @@ -136,6 +136,19 @@ class IrisController( private val _gatewayStatus = MutableStateFlow(null) val gatewayStatus: StateFlow = _gatewayStatus.asStateFlow() + /** M5: highest outbox cursor already delivered to this device via the + * push backend (from hello.ack; 0 = never). Sync-replayed frames with + * `cursor <= lastPushedCursor` already woke the device via push, so the + * app must not post a second system notification for them (docs/08 + * §8.7). In-memory: hello.ack refreshes it on every (re)connect. */ + @Volatile + private var lastPushedCursor: Long = 0 + + /** True while [frame] is a sync replay that already reached the device + * via push (live frames carry no cursor and are never suppressed). */ + private fun isPushedReplay(frame: iris.protocol.Frame): Boolean = + frame.cursor?.let { it <= lastPushedCursor } ?: false + // ── M3: threads toggle (per-app for now; per-channel lands later) ───── // Persisted (Settings → "Threads"). private val _threadsEnabled = MutableStateFlow(store.threadsEnabled) @@ -393,13 +406,15 @@ class IrisController( when (frame.type) { TYPE_MESSAGE_STOP -> { frame.payloadAs()?.let { - notifyMessageIfBackgrounded(frame.chatId, frame.threadId, it.finalText) + if (!isPushedReplay(frame)) { + notifyMessageIfBackgrounded(frame.chatId, frame.threadId, it.finalText) + } } } TYPE_MESSAGE -> { frame.payloadAs()?.let { - if (it.role == ROLE_ASSISTANT) { + if (it.role == ROLE_ASSISTANT && !isPushedReplay(frame)) { notifyMessageIfBackgrounded(frame.chatId, frame.threadId, it.text) } } @@ -493,8 +508,10 @@ class IrisController( // M5: WS is live but the app is backgrounded — the // in-app banner is invisible, so mirror to a system // notification (the push backend only fires when - // there is no live subscriber). - if (!isAppForeground()) { + // there is no live subscriber). Suppressed for + // sync replays that already woke the device via + // push (docs/08 §8.7). + if (!isAppForeground() && !isPushedReplay(frame)) { postSystemNotification(p.chatId, null, p.title, p.body, p.threadId) } } @@ -540,6 +557,8 @@ class IrisController( prevState = s if (s is GatewayClient.State.Connected) { channels.setAll(s.channels) + // M5: refresh the push-dedupe watermark (docs/08 §8.7). + lastPushedCursor = s.lastPushedCursor val home = s.channels.firstOrNull { it.isDefault }?.chatId if (home != null) { _homeChannel.value = home diff --git a/docs/04-wire-protocol.md b/docs/04-wire-protocol.md index 7a9bbc3..188cf4a 100644 --- a/docs/04-wire-protocol.md +++ b/docs/04-wire-protocol.md @@ -21,6 +21,10 @@ Every frame: - `v` — protocol version (currently `1`). Server rejects unknown major versions. - `id` — request id (client-chosen). Responses/acks echo it. Events have no `id`. - `chat_id` / `thread_id` — top-level for convenience; may also be in `payload`. +- `cursor` — outbox cursor the frame was parked under. Present **only** on + frames replayed by `sync` (live frames carry none). The app compares it + against `last_pushed_cursor` from `hello.ack` to skip re-notifying frames + that already woke the device via push (`08-push.md` §8.7). - Unknown `type`s are ignored (forward-compat); unknown `payload` fields ignored. **Binary media frames** are not JSON. A media transfer is: one JSON header frame @@ -36,9 +40,14 @@ Pairing succeeded. "server_caps":{"streaming":true,"reasoning":true,"tools":true,"media":true, "search":true,"push":"fcm","pickers":true}, "sync_cursor":1042, + "last_pushed_cursor":1040, "channels":[{"chat_id":"android:default","name":"Default","kind":"default","is_default":true}] }} ``` +`last_pushed_cursor` is the highest outbox cursor already delivered to THIS +device via the push backend (0 = never). The app skips system notifications +for sync-replayed frames with `cursor <= last_pushed_cursor` — they already +woke the device via push (dedupe, `08-push.md` §8.7). ### `message` A final / standalone message. diff --git a/docs/08-push.md b/docs/08-push.md index 4f356a1..0b3bb7a 100644 --- a/docs/08-push.md +++ b/docs/08-push.md @@ -102,4 +102,34 @@ persist until acted on. short preview (privacy on lock screen). Full content is fetched via `sync` over the authenticated WS. - ntfy: use a **private topic + auth token** for any real trust boundary (hermes - ntfy adapter guidance). \ No newline at end of file + ntfy adapter guidance). + +## 8.8 Notification dedupe (push vs. sync) + +A message sent while the device is offline is notified **twice** by naive +design: once by the push (FCM displays the `notification` payload), and again +when the app reconnects, syncs the outbox, and mirrors the replayed frames to +system notifications (the background-mirror path, §8.5). The fix is a per- +device push watermark: + +- **Gateway** records the highest outbox cursor delivered to each device via + the push backend (`devices.last_pushed_cursor`, advanced only on a + *successful* send) and returns it in `hello.ack` as `last_pushed_cursor`. +- **Gateway** coalesces back-to-back pushes per chat (5 s window): a cron + delivery parks a notification frame AND a message frame, and only the first + pushes — the second reaches the app via sync (tap the first notification). +- **Gateway** tags every `sync`-replayed frame with its outbox cursor in the + frame envelope (`cursor`; live frames carry none). +- **App** skips system notifications for replayed frames with + `cursor <= lastPushedCursor` (they already woke the device). Live frames are + never suppressed — that is exactly the case where no push fired and the app + must notify itself. +- **App** FCM handler (`onMessageReceived`) posts nothing when the WS is + connected (the background-mirror path handles it) and nothing when the + message carried a `notification` payload (FCM already displayed it); + data-only messages are the exception (the app must display them itself). + +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 diff --git a/docs/protocol/frames.schema.json b/docs/protocol/frames.schema.json index f0de422..850c7a2 100644 --- a/docs/protocol/frames.schema.json +++ b/docs/protocol/frames.schema.json @@ -12,6 +12,7 @@ "type": { "type": "string", "description": "Frame type (see frame_types)." }, "chat_id": { "type": "string", "description": "Optional chat scope (e.g. android:default, android:chan_7)." }, "thread_id": { "type": "string", "description": "Optional thread scope within a chat_id." }, + "cursor": { "type": "integer", "description": "Outbox cursor the frame was parked under. Present ONLY on frames replayed by sync (live frames carry none). The app skips re-notifying replayed frames with cursor <= last_pushed_cursor (docs/08 §8.7)." }, "payload": { "type": "object", "description": "Type-specific payload." } } }, @@ -22,6 +23,7 @@ "payload": { "server_caps": { "type": "object", "properties": { "streaming": {"type":"boolean"}, "reasoning": {"type":"boolean"}, "tools": {"type":"boolean"}, "media": {"type":"boolean"}, "search": {"type":"boolean"}, "push": {"type":"string","enum":["fcm","ntfy","none"]}, "push_ntfy_server": {"type":"string","description":"ntfy server URL for the app's listener; empty string when the backend is not ntfy."}, "pickers": {"type":"boolean"} } }, "sync_cursor": { "type": "integer" }, + "last_pushed_cursor": { "type": "integer", "description": "Highest outbox cursor already delivered to THIS device via the push backend (0 = never). The app skips system notifications for sync-replayed frames at/below it (dedupe, docs/08 §8.7)." }, "channels": { "type": "array", "items": { "$ref": "#/definitions/channel" } } } }, diff --git a/gateway-plugin/adapter.py b/gateway-plugin/adapter.py index 03eb676..307383b 100644 --- a/gateway-plugin/adapter.py +++ b/gateway-plugin/adapter.py @@ -439,6 +439,11 @@ def _strip_streaming_cursor(text: str) -> str: return text +# M5: coalesce back-to-back pushes for the same chat (a cron delivery parks +# a notification frame AND a message frame; only the first should push). +_PUSH_COALESCE_S = 5.0 + + def _push_preview(text: Any, limit: int = 120) -> str: """Short single-line preview for push bodies (lock-screen privacy: no secrets, no full bodies -- full content arrives via ``sync``).""" @@ -1019,6 +1024,9 @@ class AndroidAdapter(BasePlatformAdapter): ntfy_auth_token=_get_scoped_secret("NTFY_AUTH_TOKEN"), ) self._prune_notified_at = 0.0 + # 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] = {} def _turn_state(self, chat_id: str) -> _TurnState: st = self._turns.get(chat_id) @@ -1537,6 +1545,16 @@ class AndroidAdapter(BasePlatformAdapter): 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( + "android: 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: @@ -1578,6 +1596,15 @@ class AndroidAdapter(BasePlatformAdapter): logger.warning("android: 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( + "android: push cursor update failed for %s", device_id, exc_info=True + ) + self._last_push_at[chat_id] = time.time() logger.info( "android: push via %s -> %s (%s, chat=%s)", backend.name, device_id, frame.type, chat_id, @@ -2362,6 +2389,10 @@ class AndroidAdapter(BasePlatformAdapter): id=raw.get("id") if isinstance(raw.get("id"), int) else None, chat_id=raw.get("chat_id") if isinstance(raw.get("chat_id"), str) else e.get("chat_id"), thread_id=raw.get("thread_id") if isinstance(raw.get("thread_id"), str) else None, + # M5: tag replayed frames with their outbox cursor so the app + # can skip re-notifying frames that already woke the device + # via push (cursor <= last_pushed_cursor, docs/08 §8.7). + cursor=e.get("cursor"), v=raw.get("v") if isinstance(raw.get("v"), int) else protocol.PROTOCOL_VERSION, ) await self._ws_server.send_to(device_id, replayed) diff --git a/gateway-plugin/pairing.py b/gateway-plugin/pairing.py index 85327e8..52c55ac 100644 --- a/gateway-plugin/pairing.py +++ b/gateway-plugin/pairing.py @@ -17,7 +17,7 @@ import sqlite3 import threading import time from pathlib import Path -from typing import Any, Dict, List, Optional +from typing import Any from urllib.parse import quote logger = logging.getLogger(__name__) @@ -31,7 +31,7 @@ def generate_token() -> str: return secrets.token_hex(TOKEN_BYTES) -def verify_token(provided: Optional[str], expected: Optional[str]) -> bool: +def verify_token(provided: str | None, expected: str | None) -> bool: """Constant-time token comparison (never time-leaks the token).""" if not provided or not expected: return False @@ -65,6 +65,7 @@ def pairing_url(host: str, port: int, secure: bool = False) -> str: # Device registry (SQLite) # --------------------------------------------------------------------------- + class DeviceRegistry: """Persistent device registry under ``get_hermes_home()/"android"``. @@ -88,20 +89,32 @@ class DeviceRegistry: caps TEXT NOT NULL DEFAULT '{}', fcm_token TEXT, ntfy_topic TEXT, + last_pushed_cursor INTEGER NOT NULL DEFAULT 0, last_seen REAL NOT NULL DEFAULT 0, created REAL NOT NULL DEFAULT 0 ) """ ) + # M5: migrate pre-push-cursor databases (the column carries the + # highest outbox cursor already delivered to the device via the + # push backend; hello.ack returns it for notification dedupe). + cols = { + r["name"] + for r in self._conn.execute("PRAGMA table_info(devices)").fetchall() + } + if "last_pushed_cursor" not in cols: + self._conn.execute( + "ALTER TABLE devices ADD COLUMN last_pushed_cursor INTEGER NOT NULL DEFAULT 0" + ) self._conn.commit() def upsert( self, device_id: str, name: str, - caps: Optional[Dict[str, Any]] = None, - fcm_token: Optional[str] = None, - ntfy_topic: Optional[str] = None, + caps: dict[str, Any] | None = None, + fcm_token: str | None = None, + ntfy_topic: str | None = None, ) -> None: now = time.time() caps_json = json.dumps(caps or {}, separators=(",", ":")) @@ -125,8 +138,8 @@ class DeviceRegistry: def update_push_tokens( self, device_id: str, - fcm_token: Optional[str] = None, - ntfy_topic: Optional[str] = None, + fcm_token: str | None = None, + ntfy_topic: str | None = None, ) -> None: with self._lock: self._conn.execute( @@ -149,14 +162,43 @@ class DeviceRegistry: ) self._conn.commit() - def get(self, device_id: str) -> Optional[Dict[str, Any]]: + def update_push_cursor(self, device_id: str, cursor: int) -> None: + """Advance the device's last-pushed cursor (monotonic; never + regresses). Called after a successful push send.""" + try: + cursor = max(0, int(cursor or 0)) + except (TypeError, ValueError): + return + with self._lock: + self._conn.execute( + """ + UPDATE devices SET last_pushed_cursor = MAX(last_pushed_cursor, ?) + WHERE device_id = ? + """, + (cursor, device_id), + ) + self._conn.commit() + + def last_pushed_cursor(self, device_id: str) -> int: + """Highest outbox cursor pushed to this device (0 = never/unknown).""" + with self._lock: + row = self._conn.execute( + "SELECT last_pushed_cursor FROM devices WHERE device_id = ?", + (device_id,), + ).fetchone() + try: + return int(row["last_pushed_cursor"]) if row else 0 + except (TypeError, ValueError, KeyError, IndexError): + return 0 + + def get(self, device_id: str) -> dict[str, Any] | None: with self._lock: row = self._conn.execute( "SELECT * FROM devices WHERE device_id = ?", (device_id,) ).fetchone() return _row_to_device(row) if row else None - def list(self) -> List[Dict[str, Any]]: + def list(self) -> list[dict[str, Any]]: with self._lock: rows = self._conn.execute( "SELECT * FROM devices ORDER BY last_seen DESC" @@ -171,7 +213,7 @@ class DeviceRegistry: pass -def _row_to_device(row: sqlite3.Row) -> Dict[str, Any]: +def _row_to_device(row: sqlite3.Row) -> dict[str, Any]: try: caps = json.loads(row["caps"] or "{}") if not isinstance(caps, dict): @@ -184,6 +226,7 @@ def _row_to_device(row: sqlite3.Row) -> Dict[str, Any]: "caps": caps, "fcm_token": row["fcm_token"], "ntfy_topic": row["ntfy_topic"], + "last_pushed_cursor": row["last_pushed_cursor"] or 0, "last_seen": row["last_seen"], "created": row["created"], - } \ No newline at end of file + } diff --git a/gateway-plugin/protocol.py b/gateway-plugin/protocol.py index 5a71b71..855b1c5 100644 --- a/gateway-plugin/protocol.py +++ b/gateway-plugin/protocol.py @@ -17,7 +17,7 @@ Milestone M5: notification, fcm.register, read.receipt, status. import json from dataclasses import dataclass, field -from typing import Any, Dict, List, Optional +from typing import Any, Optional PROTOCOL_VERSION = 1 @@ -146,30 +146,38 @@ STATUS_DEGRADED = "degraded" # Envelope # --------------------------------------------------------------------------- + @dataclass class Frame: """One wire frame. - ``v`` is always serialised; ``id``/``chat_id``/``thread_id`` are - omitted when ``None`` (events carry no ``id``; chat-scoped frames carry - ``chat_id``/``thread_id`` at the top level for convenience). + ``v`` is always serialised; ``id``/``chat_id``/``thread_id``/``cursor`` + are omitted when ``None`` (events carry no ``id``; chat-scoped frames + carry ``chat_id``/``thread_id`` at the top level for convenience). + + ``cursor`` is set only on frames replayed by ``sync``: the outbox cursor + the frame was parked under. The app uses it to skip re-notifying frames + that already woke the device via push (docs/08 §8.7). """ type: str - payload: Dict[str, Any] = field(default_factory=dict) - id: Optional[int] = None - chat_id: Optional[str] = None - thread_id: Optional[str] = None + payload: dict[str, Any] = field(default_factory=dict) + id: int | None = None + chat_id: str | None = None + thread_id: str | None = None + cursor: int | None = None v: int = PROTOCOL_VERSION - def to_dict(self) -> Dict[str, Any]: - d: Dict[str, Any] = {"v": self.v, "type": self.type} + def to_dict(self) -> dict[str, Any]: + d: dict[str, Any] = {"v": self.v, "type": self.type} if self.id is not None: d["id"] = self.id if self.chat_id is not None: d["chat_id"] = self.chat_id if self.thread_id is not None: d["thread_id"] = self.thread_id + if self.cursor is not None: + d["cursor"] = self.cursor d["payload"] = self.payload return d @@ -202,17 +210,29 @@ class Frame: thread_id = data.get("thread_id") if not isinstance(thread_id, str): thread_id = None - return cls(type=ftype, payload=payload, id=fid, chat_id=chat_id, thread_id=thread_id) + cursor = data.get("cursor") + if not isinstance(cursor, int) or isinstance(cursor, bool): + cursor = None + return cls( + type=ftype, + payload=payload, + id=fid, + chat_id=chat_id, + thread_id=thread_id, + cursor=cursor, + ) # --------------------------------------------------------------------------- # Frame constructors (server -> app) # --------------------------------------------------------------------------- + def hello_ack( - server_caps: Dict[str, Any], + server_caps: dict[str, Any], sync_cursor: int = 0, - channels: Optional[list] = None, + channels: list | None = None, + last_pushed_cursor: int = 0, ) -> Frame: return Frame( type=TYPE_HELLO_ACK, @@ -220,6 +240,11 @@ def hello_ack( "server_caps": server_caps, "sync_cursor": sync_cursor, "channels": channels or [], + # M5: highest outbox cursor already delivered to THIS device via + # the push backend (0 = never). The app skips system + # notifications for sync-replayed frames at/below it (dedupe, + # docs/08 §8.7). + "last_pushed_cursor": last_pushed_cursor, }, ) @@ -230,15 +255,15 @@ def message( role: str, text: str, *, - thread_id: Optional[str] = None, - reasoning: Optional[str] = None, - media: Optional[list] = None, - reply_to: Optional[str] = None, - model: Optional[str] = None, - tokens: Optional[int] = None, - ts: Optional[int] = None, + thread_id: str | None = None, + reasoning: str | None = None, + media: list | None = None, + reply_to: str | None = None, + model: str | None = None, + tokens: int | None = None, + ts: int | None = None, ) -> Frame: - payload: Dict[str, Any] = { + payload: dict[str, Any] = { "message_id": message_id, "role": role, "text": text, @@ -255,10 +280,12 @@ def message( payload["tokens"] = tokens if ts is not None: payload["ts"] = ts - return Frame(type=TYPE_MESSAGE, chat_id=chat_id, thread_id=thread_id, payload=payload) + return Frame( + type=TYPE_MESSAGE, chat_id=chat_id, thread_id=thread_id, payload=payload + ) -def typing(chat_id: str, on: bool = True, *, thread_id: Optional[str] = None) -> Frame: +def typing(chat_id: str, on: bool = True, *, thread_id: str | None = None) -> Frame: return Frame( type=TYPE_TYPING, chat_id=chat_id, @@ -271,12 +298,13 @@ def typing(chat_id: str, on: bool = True, *, thread_id: Optional[str] = None) -> # Streaming frames (M2) # --------------------------------------------------------------------------- + def message_start( chat_id: str, message_id: str, role: str = ROLE_ASSISTANT, *, - thread_id: Optional[str] = None, + thread_id: str | None = None, ) -> Frame: """Open a streaming bubble.""" return Frame( @@ -292,7 +320,7 @@ def message_update( message_id: str, text: str, *, - thread_id: Optional[str] = None, + thread_id: str | None = None, ) -> Frame: """Replace the live bubble text (full snapshot).""" return Frame( @@ -308,14 +336,14 @@ def message_stop( message_id: str, final_text: str, *, - thread_id: Optional[str] = None, - reasoning: Optional[str] = None, - model: Optional[str] = None, - tokens: Optional[int] = None, - ts: Optional[int] = None, + thread_id: str | None = None, + reasoning: str | None = None, + model: str | None = None, + tokens: int | None = None, + ts: int | None = None, ) -> Frame: """Finalize a streaming bubble.""" - payload: Dict[str, Any] = { + payload: dict[str, Any] = { "message_id": message_id, "final_text": final_text, } @@ -339,16 +367,17 @@ def message_stop( # Tool activity frames (M2) # --------------------------------------------------------------------------- + def tool_start( chat_id: str, index: int, name: str, *, - thread_id: Optional[str] = None, - preview: Optional[str] = None, - args: Optional[Dict[str, Any]] = None, + thread_id: str | None = None, + preview: str | None = None, + args: dict[str, Any] | None = None, ) -> Frame: - payload: Dict[str, Any] = {"index": index, "name": name} + payload: dict[str, Any] = {"index": index, "name": name} if preview: payload["preview"] = preview if args: @@ -366,10 +395,10 @@ def tool_progress( index: int, name: str, *, - thread_id: Optional[str] = None, - note: Optional[str] = None, + thread_id: str | None = None, + note: str | None = None, ) -> Frame: - payload: Dict[str, Any] = {"index": index, "name": name} + payload: dict[str, Any] = {"index": index, "name": name} if note: payload["note"] = note return Frame( @@ -385,12 +414,12 @@ def tool_end( index: int, name: str, *, - thread_id: Optional[str] = None, + thread_id: str | None = None, ok: bool = True, - duration: Optional[float] = None, - output_preview: Optional[str] = None, + duration: float | None = None, + output_preview: str | None = None, ) -> Frame: - payload: Dict[str, Any] = {"index": index, "name": name, "ok": ok} + payload: dict[str, Any] = {"index": index, "name": name, "ok": ok} if duration is not None: payload["duration"] = duration if output_preview: @@ -407,12 +436,13 @@ def tool_end( # Commentary frame (M2) # --------------------------------------------------------------------------- + def commentary( chat_id: str, message_id: str, text: str, *, - thread_id: Optional[str] = None, + thread_id: str | None = None, ) -> Frame: """An intermediate assistant beat (between tool iterations).""" return Frame( @@ -427,9 +457,10 @@ def commentary( # Channel directory frames (M3) # --------------------------------------------------------------------------- -def _channel_payload(entry: Dict[str, Any]) -> Dict[str, Any]: + +def _channel_payload(entry: dict[str, Any]) -> dict[str, Any]: """Project a directory entry onto the wire shape.""" - payload: Dict[str, Any] = { + payload: dict[str, Any] = { "chat_id": entry.get("chat_id"), "name": entry.get("name"), "kind": entry.get("kind", "channel"), @@ -451,7 +482,7 @@ def _channel_payload(entry: Dict[str, Any]) -> Dict[str, Any]: return payload -def channel_created(entry: Dict[str, Any], auto: bool = False) -> Frame: +def channel_created(entry: dict[str, Any], auto: bool = False) -> Frame: """Broadcast: a channel/thread was created. ``auto=True`` marks a thread the gateway minted itself for an incoming @@ -465,7 +496,7 @@ def channel_created(entry: Dict[str, Any], auto: bool = False) -> Frame: return Frame(type=TYPE_CHANNEL_CREATED, payload=payload) -def channel_renamed(entry: Dict[str, Any]) -> Frame: +def channel_renamed(entry: dict[str, Any]) -> Frame: """Broadcast: a channel/thread was renamed.""" return Frame(type=TYPE_CHANNEL_RENAMED, payload=_channel_payload(entry)) @@ -475,7 +506,7 @@ def channel_deleted(chat_id: str) -> Frame: return Frame(type=TYPE_CHANNEL_DELETED, payload={"chat_id": chat_id}) -def channel_list(channels: List[Dict[str, Any]]) -> Frame: +def channel_list(channels: list[dict[str, Any]]) -> Frame: """Full directory (response to a ``channel.list`` request).""" return Frame( type=TYPE_CHANNEL_LIST, @@ -487,12 +518,13 @@ def channel_list(channels: List[Dict[str, Any]]) -> Frame: # Search frames (M3) # --------------------------------------------------------------------------- + def search_results( query: str, scope: str, - hits: List[Dict[str, Any]], + hits: list[dict[str, Any]], *, - id: Optional[int] = None, + id: int | None = None, ) -> Frame: """Response to a ``search`` request. @@ -509,7 +541,8 @@ def search_results( # Slash-command catalog frames # --------------------------------------------------------------------------- -def commands_catalog(commands: List[Dict[str, Any]], *, id: Optional[int] = None) -> Frame: + +def commands_catalog(commands: list[dict[str, Any]], *, id: int | None = None) -> Frame: """Response to a ``commands.catalog`` request: the gateway's slash-command catalog for the app's ``/`` drawer. @@ -527,7 +560,8 @@ def commands_catalog(commands: List[Dict[str, Any]], *, id: Optional[int] = None # Sync frames (M3 outbox) # --------------------------------------------------------------------------- -def sync_done(cursor: int, *, id: Optional[int] = None) -> Frame: + +def sync_done(cursor: int, *, id: int | None = None) -> Frame: """Terminal frame of a ``sync`` replay: the new cursor to persist.""" return Frame(type=TYPE_SYNC_DONE, id=id, payload={"cursor": cursor}) @@ -536,20 +570,21 @@ def sync_done(cursor: int, *, id: Optional[int] = None) -> Frame: # History frame (full message history for a chat/thread) # --------------------------------------------------------------------------- + def history( chat_id: str, - messages: List[Dict[str, Any]], + messages: list[dict[str, Any]], has_more: bool, *, - thread_id: Optional[str] = None, - oldest_message_id: Optional[str] = None, - id: Optional[int] = None, + thread_id: str | None = None, + oldest_message_id: str | None = None, + id: int | None = None, ) -> Frame: """Response to a ``history`` request: a page of final messages for a chat/thread, ordered oldest → newest. ``has_more`` signals older pages exist; ``oldest_message_id`` is the ``before_message_id`` for the next (older) page.""" - payload: Dict[str, Any] = { + payload: dict[str, Any] = { "messages": messages, "has_more": has_more, } @@ -568,12 +603,13 @@ def history( # Message deletion frames # --------------------------------------------------------------------------- + def message_deleted( chat_id: str, - message_ids: List[str], + message_ids: list[str], *, - thread_id: Optional[str] = None, - id: Optional[int] = None, + thread_id: str | None = None, + id: int | None = None, ) -> Frame: """Broadcast: the given message(s) were deleted from a chat/thread. @@ -595,18 +631,19 @@ def message_deleted( # Push / notification frames (M5) # --------------------------------------------------------------------------- + def notification( chat_id: str, kind: str, title: str, body: str, *, - thread_id: Optional[str] = None, - ts: Optional[int] = None, + thread_id: str | None = None, + ts: int | None = None, ) -> Frame: """Event: a transient in-app banner (and a push mirror when the device is offline). ``kind`` is one of the ``NOTIF_*`` constants.""" - payload: Dict[str, Any] = {"kind": kind, "title": title, "body": body} + payload: dict[str, Any] = {"kind": kind, "title": title, "body": body} if ts is not None: payload["ts"] = ts return Frame( @@ -614,11 +651,9 @@ def notification( ) -def fcm_register( - fcm_token: Optional[str] = None, ntfy_topic: Optional[str] = None -) -> Frame: +def fcm_register(fcm_token: str | None = None, ntfy_topic: str | None = None) -> Frame: """Request: update the device's push tokens (FCM rotation / ntfy topic).""" - payload: Dict[str, Any] = {} + payload: dict[str, Any] = {} if fcm_token: payload["fcm_token"] = fcm_token if ntfy_topic: @@ -640,6 +675,7 @@ def read_receipt(chat_id: str, message_id: str) -> Frame: # Gateway health frame (M5) # --------------------------------------------------------------------------- + def status(state: str) -> Frame: """Gateway health state (``state`` is one of the ``STATUS_*`` constants).""" return Frame(type=TYPE_STATUS, payload={"state": state}) @@ -649,6 +685,7 @@ def status(state: str) -> Frame: # Media frames (M4) # --------------------------------------------------------------------------- + def media_offer( media_id: str, kind: str, @@ -656,16 +693,16 @@ def media_offer( size: int, filename: str, *, - chat_id: Optional[str] = None, - thread_id: Optional[str] = None, - message_id: Optional[str] = None, + chat_id: str | None = None, + thread_id: str | None = None, + message_id: str | None = None, ) -> Frame: """Event: the agent produced media the app can fetch via ``media.pull``. ``message_id`` (optional) associates the offer with the assistant message it belongs to (the app falls back to the lane's last assistant message). """ - payload: Dict[str, Any] = { + payload: dict[str, Any] = { "media_id": media_id, "kind": kind, "mime": mime, @@ -674,15 +711,17 @@ def media_offer( } if message_id: payload["message_id"] = message_id - return Frame(type=TYPE_MEDIA_OFFER, chat_id=chat_id, thread_id=thread_id, payload=payload) + return Frame( + type=TYPE_MEDIA_OFFER, chat_id=chat_id, thread_id=thread_id, payload=payload + ) -def media_pull_end(ok: bool, *, id: Optional[int] = None) -> Frame: +def media_pull_end(ok: bool, *, id: int | None = None) -> Frame: """Terminal frame of a ``media.pull`` binary stream.""" return Frame(type=TYPE_MEDIA_PULL_END, id=id, payload={"ok": ok}) -def media_upload_ack(ok: bool, media_ref: str, *, id: Optional[int] = None) -> Frame: +def media_upload_ack(ok: bool, media_ref: str, *, id: int | None = None) -> Frame: """Response to ``media.upload.end``: the ref is cached and may be used in a ``message.send`` ``media_refs``. Failures use ``error`` frames instead.""" return Frame( @@ -692,12 +731,12 @@ def media_upload_ack(ok: bool, media_ref: str, *, id: Optional[int] = None) -> F ) -def error(code: str, message: str, *, id: Optional[int] = None) -> Frame: +def error(code: str, message: str, *, id: int | None = None) -> Frame: return Frame(type=TYPE_ERROR, id=id, payload={"code": code, "message": message}) -def pong(ts: Optional[int] = None) -> Frame: - payload: Dict[str, Any] = {} +def pong(ts: int | None = None) -> Frame: + payload: dict[str, Any] = {} if ts is not None: payload["ts"] = ts - return Frame(type=TYPE_PONG, payload=payload) \ No newline at end of file + return Frame(type=TYPE_PONG, payload=payload) diff --git a/gateway-plugin/ws_server.py b/gateway-plugin/ws_server.py index 24c8ae6..8f65070 100644 --- a/gateway-plugin/ws_server.py +++ b/gateway-plugin/ws_server.py @@ -290,6 +290,9 @@ class WsServer: server_caps=self._adapter.server_caps(), sync_cursor=self._adapter._outbox.latest_cursor(), channels=self._adapter.channel_list(), + # M5: lets the app dedupe sync-replayed notifications that + # already woke this device via push (docs/08 §8.7). + last_pushed_cursor=self._adapter._devices.last_pushed_cursor(device_id), ) try: await ws.send(ack.to_json())