diff --git a/app/shared/src/commonMain/kotlin/iris/state/IrisController.kt b/app/shared/src/commonMain/kotlin/iris/state/IrisController.kt index e01944f..d5bc8ec 100644 --- a/app/shared/src/commonMain/kotlin/iris/state/IrisController.kt +++ b/app/shared/src/commonMain/kotlin/iris/state/IrisController.kt @@ -82,6 +82,7 @@ import iris.ui.theme.UserTheme import iris.util.IrisLog import kotlinx.coroutines.CoroutineScope import kotlinx.coroutines.Dispatchers +import kotlinx.coroutines.Job import kotlinx.coroutines.SupervisorJob import kotlinx.coroutines.delay import kotlinx.coroutines.flow.MutableStateFlow @@ -209,6 +210,37 @@ class IrisController( * 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 + /** M5: highest outbox cursor this device has CONSUMED (applied from the + * event stream). The stream resumes from [store.syncCursor] on every + * (re)connect, so this high-water mark is persisted there (debounced + + * on backgrounding): without it, a restart re-delivers every frame + * consumed since the last `sync.done` — and replayed `notification` + * frames re-showed their banner on every app reopen (docs/08 §8.4). + * @Volatile: written on the frame-handler coroutine, read on the + * foreground/state collectors. */ + @Volatile + private var consumedCursor: Long = store.syncCursor + + /** Debounced persist of [consumedCursor] (one store write per burst, + * not per frame). */ + private var cursorSaveJob: Job? = null + + private fun persistCursorDebounced() { + cursorSaveJob?.cancel() + cursorSaveJob = + scope.launch { + delay(CACHE_SAVE_DEBOUNCE_MS) + // maxOf: sync.done may have advanced the store past us. + store.syncCursor = maxOf(store.syncCursor, consumedCursor) + } + } + + /** Flush [consumedCursor] to the store now (app backgrounding / exit). */ + private fun persistCursorNow() { + cursorSaveJob?.cancel() + store.syncCursor = maxOf(store.syncCursor, consumedCursor) + } + // ── M3: threads toggle (per-app for now; per-channel lands later) ───── // Persisted (Settings → "Threads"). private val _threadsEnabled = MutableStateFlow(store.threadsEnabled) @@ -507,12 +539,23 @@ class IrisController( // clears on scroll-to-bottom while already focused.) scope.launch { foreground.collect { fg -> + if (!fg) persistCursorNow() if (fg && currentLaneAtBottom) markCurrentLaneRead() } } scope.launch { client.events.collect { frame -> try { + // Consume the frame's outbox cursor BEFORE dispatching: + // a frame that arrives a second time (SSE catch-up and + // the explicit sync replay both carry the cursor) must + // not re-notify — only its first consumption may. + val cursor = frame.cursor + val alreadyConsumed = cursor != null && cursor <= consumedCursor + if (cursor != null && cursor > consumedCursor) { + consumedCursor = cursor + persistCursorDebounced() + } when (frame.type) { TYPE_TOOL_START, TYPE_TOOL_PROGRESS, @@ -562,7 +605,7 @@ class IrisController( frame.chatId?.let { cid -> noteIncomingAssistantMessage(chat.laneKey(cid, frame.threadId)) } - if (!isPushedReplay(frame)) { + if (!alreadyConsumed && !isPushedReplay(frame)) { notifyMessageIfBackgrounded(frame.chatId, frame.threadId, it.finalText) } } @@ -575,7 +618,7 @@ class IrisController( frame.chatId?.let { cid -> noteIncomingAssistantMessage(chat.laneKey(cid, frame.threadId)) } - if (!isPushedReplay(frame)) { + if (!alreadyConsumed && !isPushedReplay(frame)) { notifyMessageIfBackgrounded(frame.chatId, frame.threadId, it.text) } } @@ -649,8 +692,13 @@ class IrisController( TYPE_SYNC_DONE -> { // Replayed frames already flowed through [events]; the - // cursor is authoritative server-side (outbox). - frame.payloadAs()?.let { store.syncCursor = it.cursor } + // cursor is authoritative server-side (outbox). Keep the + // consumed high-water mark in step (it may be ahead of + // us — live frames consumed after the replay started). + frame.payloadAs()?.let { + consumedCursor = maxOf(consumedCursor, it.cursor) + store.syncCursor = maxOf(store.syncCursor, it.cursor) + } } TYPE_HISTORY -> { @@ -689,15 +737,21 @@ class IrisController( TYPE_NOTIFICATION -> { frame.payloadAs()?.let { p -> - pushBanner(p.kind, p.title, p.body, p.chatId, p.threadId) - // 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). 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) + // alreadyConsumed: this frame was applied earlier + // (its cursor is at/below the high-water mark) — + // a duplicate delivery (restart catch-up or the + // sync replay) must not re-show the banner. + if (!alreadyConsumed) { + pushBanner(p.kind, p.title, p.body, p.chatId, p.threadId) + // 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). 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) + } } } }