From 4df6ce412eca2a240dfe3aba71dba54ab909d700 Mon Sep 17 00:00:00 2001 From: Brandon McAnsh Date: Fri, 2 Oct 2026 11:34:48 -0400 Subject: [PATCH 1/4] fix(chat): land the launch feed sync as one list update Share the in-flight feed sync with onStart and push refreshes, write DM and group feeds (rows, members, previews) in one transaction, make catch-up a single forward-only metadata write, rebuild the list when a chat's newest message changes, and stage read stamps into the reconcile with bounded concurrency. Catch-up now runs the open chat first. --- .../chat/internal/RealChatCoordinator.kt | 70 ++++- .../internal/delegates/EventStreamDelegate.kt | 27 +- .../internal/delegates/FeedSyncDelegate.kt | 289 +++++++++++++----- .../internal/delegates/GroupFeedDelegate.kt | 32 +- .../internal/delegates/MessagingDelegate.kt | 15 +- .../flipcash/shared/chat/FeedReconcileTest.kt | 216 +++++++++++++ .../shared/chat/FeedSyncPreviewWriteTest.kt | 2 +- .../shared/chat/GroupChatRoutingTest.kt | 6 +- .../shared/chat/MessagingLoadCursorTest.kt | 10 +- .../shared/chat/MessagingPushedMessageTest.kt | 10 +- .../app/persistence/dao/ChatMessageDao.kt | 15 + .../app/persistence/dao/ChatMetadataDao.kt | 33 ++ .../persistence/ChatMetadataCatchUpTest.kt | 99 ++++++ .../app/persistence/sources/ChatFeedWriter.kt | 37 +++ .../sources/ChatMessageDataSource.kt | 44 ++- .../sources/ChatMetadataDataSource.kt | 15 + 16 files changed, 784 insertions(+), 136 deletions(-) create mode 100644 apps/flipcash/shared/chat/src/test/kotlin/com/flipcash/shared/chat/FeedReconcileTest.kt create mode 100644 apps/flipcash/shared/persistence/db/src/test/kotlin/com/flipcash/app/persistence/ChatMetadataCatchUpTest.kt create mode 100644 apps/flipcash/shared/persistence/sources/src/main/kotlin/com/flipcash/app/persistence/sources/ChatFeedWriter.kt diff --git a/apps/flipcash/shared/chat/src/main/kotlin/com/flipcash/shared/chat/internal/RealChatCoordinator.kt b/apps/flipcash/shared/chat/src/main/kotlin/com/flipcash/shared/chat/internal/RealChatCoordinator.kt index 8462c33412..73eebdc760 100644 --- a/apps/flipcash/shared/chat/src/main/kotlin/com/flipcash/shared/chat/internal/RealChatCoordinator.kt +++ b/apps/flipcash/shared/chat/src/main/kotlin/com/flipcash/shared/chat/internal/RealChatCoordinator.kt @@ -19,6 +19,7 @@ import com.flipcash.shared.chat.FeedOperations import com.flipcash.shared.chat.GroupOperations import com.flipcash.shared.chat.MessagingOperations import com.flipcash.shared.chat.ReactionOperations +import com.flipcash.services.models.chat.ChatMetadata import com.flipcash.shared.chat.internal.delegates.EventStreamDelegate import com.flipcash.shared.chat.internal.delegates.FeedSyncDelegate import com.flipcash.shared.chat.internal.delegates.GroupFeedDelegate @@ -34,6 +35,8 @@ import kotlinx.coroutines.CoroutineScope import kotlinx.coroutines.ExperimentalCoroutinesApi import kotlinx.coroutines.FlowPreview import kotlinx.coroutines.Job +import kotlinx.coroutines.async +import kotlinx.coroutines.coroutineScope import kotlinx.coroutines.SupervisorJob import kotlinx.coroutines.flow.MutableStateFlow import kotlinx.coroutines.flow.StateFlow @@ -107,6 +110,7 @@ class RealChatCoordinator @Inject constructor( companion object { private const val TAG = "ChatCoordinator" private val FEED_READ_WAIT = 3.seconds + private val FEED_PAIR_WAIT = 2.seconds } // Recreated on re-login: [reset] cancels [supervisorJob] on logout, which would @@ -204,7 +208,7 @@ class RealChatCoordinator @Inject constructor( .onEach { event -> when (event) { is EventStreamDelegate.Event.SyncFeedRequested -> - syncFeeds() + syncFeeds(fresh = true) is EventStreamDelegate.Event.LoadMessages -> messagingDelegate.loadMessages(event.chatId) is EventStreamDelegate.Event.RosterChanged -> @@ -254,6 +258,11 @@ class RealChatCoordinator @Inject constructor( messagingDelegate.clearActiveChat(chatId) } + private var feedSyncJob: Job? = null + + @Volatile + private var feedSyncRerun = false + override fun onStart(owner: LifecycleOwner) { foregrounded.value = true backgroundedActiveChat?.let { @@ -325,14 +334,67 @@ class RealChatCoordinator @Inject constructor( * overlap rather than queue. A group feed that fails is traced and dropped by * [GroupFeedDelegate.performGroupFeedSync]; it cannot take the DM list down with it. */ - private fun syncFeeds() { - feedDelegate.syncFeed() - groupFeedDelegate.syncGroupFeed() + private fun syncFeeds(fresh: Boolean = false) { + if (feedSyncJob?.isActive == true) { + // Shared with the one in flight, so a foreground, a reconnect or a heartbeat landing + // mid-launch does not cancel and restart the fetch the list is waiting on. `fresh` + // is for a trigger that says the server has something this fetch may have missed. + if (fresh) feedSyncRerun = true + return + } + feedSyncJob = scope.launch { + do { + feedSyncRerun = false + reconcileFeeds() + } while (feedSyncRerun) + } + } + + /** + * Fetches the DM and group feeds concurrently and lands them in one transaction, so the list + * rebuilds once instead of once per feed. A feed that is slow past [FEED_PAIR_WAIT] does not + * hold the other back: the one in hand is written, and the slow one follows when it arrives. + */ + private suspend fun reconcileFeeds() = coroutineScope { + val dm = async { feedDelegate.fetchFeed() } + val groups = async { groupFeedDelegate.fetchGroupFeed() } + withTimeoutOrNull(FEED_PAIR_WAIT) { + dm.await() + groups.await() + } + + var dmChats: List? = null + var groupChats: List? = null + if (dm.isCompleted) { + dm.await().onSuccess { dmChats = it }.onFailure { feedDelegate.onFeedFailed(it) } + } + if (groups.isCompleted) groupChats = groups.await() + val first = (dmChats.orEmpty()) + (groupChats.orEmpty()) + // Written together even when one half is missing or empty; a failed half contributes + // nothing and cannot take the other down. + feedDelegate.writeFeed(first) + dmChats?.let { feedDelegate.onFeedCommitted(it) } + groupChats?.let { groupFeedDelegate.requestCatchUp(it) } + + if (!dm.isCompleted) { + dm.await().onSuccess { + feedDelegate.writeFeed(it) + feedDelegate.onFeedCommitted(it) + }.onFailure { feedDelegate.onFeedFailed(it) } + } + if (!groups.isCompleted) { + groups.await()?.let { + feedDelegate.writeFeed(it) + groupFeedDelegate.requestCatchUp(it) + } + } } override suspend fun teardown() { eventStreamDelegate.stopHeartbeat() eventStreamDelegate.close() + feedSyncJob?.cancel() + feedSyncJob = null feedDelegate.cancelJobs() groupFeedDelegate.cancelJobs() networkObserverJob?.cancel() diff --git a/apps/flipcash/shared/chat/src/main/kotlin/com/flipcash/shared/chat/internal/delegates/EventStreamDelegate.kt b/apps/flipcash/shared/chat/src/main/kotlin/com/flipcash/shared/chat/internal/delegates/EventStreamDelegate.kt index 53c746ae83..3d8c3da0b4 100644 --- a/apps/flipcash/shared/chat/src/main/kotlin/com/flipcash/shared/chat/internal/delegates/EventStreamDelegate.kt +++ b/apps/flipcash/shared/chat/src/main/kotlin/com/flipcash/shared/chat/internal/delegates/EventStreamDelegate.kt @@ -295,21 +295,14 @@ class EventStreamDelegate @Inject constructor( if (delta.messages.isNotEmpty()) { linkPrefetch.prefetch(delta.messages, MessageLinkPrefetch.LIVE_WAIT) messageDataSource.upsert(chatId, delta.messages) - val latest = delta.messages.maxByOrNull { it.messageId } - latest?.let { msg -> - metadataDataSource.updateLastMessageId(chatId, msg.messageId) - // Only advance lastActivity — a partial page from a - // delta sync must not regress it to an older timestamp. - val existing = metadataDataSource.getLastActivity(chatId) - val incoming = msg.timestamp.toEpochMilliseconds() - if (existing == null || incoming > existing) { - metadataDataSource.updateLastActivity(chatId, incoming) - } - } - } - if (delta.latestSequence > afterSequence) { - metadataDataSource.updateLatestEventSequence(chatId, delta.latestSequence) } + // One write that only moves forward: the cursor, and the newest message only + // when it is newer than the stored one. See ChatMetadataDao.applyCatchUp. + metadataDataSource.applyCatchUp( + chatId, + if (delta.latestSequence > afterSequence) delta.latestSequence else 0L, + delta.messages.maxByOrNull { it.messageId }, + ) sequenceTracker.resetTo(chatId, delta.latestSequence) trace(tag = TAG, message = "Delta sync complete: ${delta.messages.size} messages, sequence ${delta.latestSequence}", type = TraceType.Process) } @@ -442,10 +435,8 @@ class EventStreamDelegate @Inject constructor( // Update lastMessageId and lastActivity AFTER metadata updates so that // incoming message timestamps always take precedence over a potentially // stale FullRefresh. - lastMsg?.let { msg -> - metadataDataSource.updateLastMessageId(chatId, msg.messageId) - metadataDataSource.updateLastActivity(chatId, msg.timestamp.toEpochMilliseconds()) - } + // One write, forward-only: a message that is not newer than the stored one changes nothing. + lastMsg?.let { msg -> metadataDataSource.applyCatchUp(chatId, 0L, msg) } // --- Roster changes --- diff --git a/apps/flipcash/shared/chat/src/main/kotlin/com/flipcash/shared/chat/internal/delegates/FeedSyncDelegate.kt b/apps/flipcash/shared/chat/src/main/kotlin/com/flipcash/shared/chat/internal/delegates/FeedSyncDelegate.kt index 7010fa0c6e..7b77ebeb61 100644 --- a/apps/flipcash/shared/chat/src/main/kotlin/com/flipcash/shared/chat/internal/delegates/FeedSyncDelegate.kt +++ b/apps/flipcash/shared/chat/src/main/kotlin/com/flipcash/shared/chat/internal/delegates/FeedSyncDelegate.kt @@ -9,6 +9,7 @@ import androidx.paging.PagingData import androidx.paging.filter import androidx.paging.map import com.flipcash.app.persistence.entities.ChatMetadataEntity +import com.flipcash.app.persistence.sources.ChatFeedWriter import com.flipcash.app.persistence.sources.ChatMemberDataSource import com.flipcash.app.persistence.sources.ChatMessageDataSource import com.flipcash.app.persistence.sources.ChatMetadataDataSource @@ -43,12 +44,19 @@ import kotlinx.coroutines.channels.Channel import kotlinx.coroutines.coroutineScope import kotlinx.coroutines.flow.Flow import kotlinx.coroutines.flow.combine +import kotlinx.coroutines.flow.distinctUntilChanged import kotlinx.coroutines.flow.launchIn import kotlinx.coroutines.flow.map import kotlinx.coroutines.flow.mapNotNull import kotlinx.coroutines.flow.onEach +import kotlinx.coroutines.flow.onStart import kotlinx.coroutines.flow.receiveAsFlow import kotlinx.coroutines.launch +import kotlinx.coroutines.sync.Semaphore +import kotlinx.coroutines.sync.withPermit +import kotlinx.coroutines.awaitAll +import kotlinx.coroutines.withTimeoutOrNull +import kotlin.time.Duration.Companion.milliseconds import java.util.concurrent.ConcurrentHashMap import javax.inject.Inject import javax.inject.Singleton @@ -77,6 +85,8 @@ class FeedSyncDelegate @Inject constructor( private val userManager: UserManager, private val messagingController: ChatMessagingController, private val linkPrefetch: MessageLinkPrefetch = MessageLinkPrefetch.None, + private val feedWriter: ChatFeedWriter = + ChatFeedWriter(metadataDataSource, memberDataSource, messageDataSource), ) : FeedOperations { companion object { @@ -85,6 +95,12 @@ class FeedSyncDelegate @Inject constructor( // The feed's rows are cheap and the merge in the mediator costs a round trip per source // per page, so this is larger than it would be for a transcript. private const val FEED_PAGE_SIZE = 20 + + // The read-stamp lookups a reconcile folds in. Bounded in number and in time: the list is + // already on screen from the cache, so a slow lookup costs one dot staying a dot for a + // while longer, never a held-up list. + private const val READ_STAMP_CONCURRENCY = 4 + private val READ_STAMP_WAIT = 1500.milliseconds } sealed interface Event { @@ -125,11 +141,27 @@ class FeedSyncDelegate @Inject constructor( private var syncJob: Job? = null private var feedObserverJob: Job? = null + // Stamps a reconcile fetched, parked until the feed rebuild that follows its write picks them + // up. Applying them to state straight away would repaint the rows' counts a beat before the + // rows themselves, as two emissions. + private val stagedStamps = ConcurrentHashMap, Long>() + + private fun FeedSyncState.isSettled() = this == FeedSyncState.Synced || this == FeedSyncState.Error + + // What the reconcile could not stage in time is fetched now, one lookup per stamp. + private fun fetchLeftoverReadStamps() { + stateHolder.current.feed?.let(::fetchMissingReadStamps) + } + // region FeedOperations override fun feed(vararg chatTypes: ChatType): Flow> { val requested = chatTypes.toSet() - return stateHolder.state.mapNotNull { state -> summaries(state, requested) } + // Distinct: state also moves for reasons the list does not show (sync status, typing, + // overlays), and an equal list repaints nothing but still re-runs every collector. + return stateHolder.state + .mapNotNull { state -> summaries(state, requested) } + .distinctUntilChanged() } override fun currentFeed(vararg chatTypes: ChatType): List? = @@ -226,17 +258,47 @@ class FeedSyncDelegate @Inject constructor( feedObserverJob = combine( metadataDataSource.observeAll(), memberDataSource.observeAll(), - ) { metadataEntities, membersByChat -> + // A preview comes from `chat_messages`, which the two above do not touch. Without this + // a row whose newest message changed without its metadata row changing (an edit, a + // delete, a preview written a beat after the row) kept its old preview until some + // unrelated write. Distinct, so only a change to some chat's newest visible message + // rebuilds, not every message write; the initial value keeps a slow first read from + // holding the list up. + messageDataSource.observeLatestVisibleChanges() + .onStart { emit(emptyList()) } + .distinctUntilChanged(), + ) { metadataEntities, membersByChat, _ -> buildFeedFromDb(metadataEntities, membersByChat) }.onEach { (feed, readStamps) -> - stateHolder.update { it.copy(feed = feed, readStamps = readStamps) } - fetchMissingReadStamps(feed) + val staged = drainStagedStamps() + stateHolder.update { + it.copy( + feed = feed, + readStamps = readStamps, + fetchedReadStamps = if (staged.isEmpty()) it.fetchedReadStamps else it.fetchedReadStamps + staged, + ) + } + // Until the first sync settles, the reconcile looks the stamps up itself, bounded and + // folded into its write; doing it here as well would be the unbounded pop-in it avoids. + if (stateHolder.current.feedSyncState.isSettled()) fetchMissingReadStamps(feed) }.launchIn(scope) } + private fun drainStagedStamps(): Map, Long> { + if (stagedStamps.isEmpty()) return emptyMap() + val drained = HashMap(stagedStamps) + drained.keys.forEach { stagedStamps.remove(it) } + return drained + } + + /** + * Standalone sync of the DM half. The coordinator syncs the DM and group halves together + * through [fetchFeed] / [writeFeed] so the list is written once; this is the single-delegate + * path. A sync already in flight is shared, not cancelled and restarted. + */ internal fun syncFeed() { val scope = scope ?: return - syncJob?.cancel() + if (syncJob?.isActive == true) return syncJob = scope.launch { performFeedSync() } } @@ -256,6 +318,7 @@ class FeedSyncDelegate @Inject constructor( feedObserverJob = null readStampsInFlight.clear() readStampsNotFound.clear() + stagedStamps.clear() } /** @@ -282,28 +345,73 @@ class FeedSyncDelegate @Inject constructor( if (!readStampsInFlight.add(key)) continue scope.launch { - messagingController.getMessage(metadata.chatId, readPointer) - .onSuccess { message -> - stateHolder.update { - it.copy(fetchedReadStamps = it.fetchedReadStamps + (key to message.unreadSeq)) - } - } - .onFailure { cause -> - if (cause is GetMessageError.NotFound) readStampsNotFound.add(key) - val expected = cause is GetMessageError.NotFound || cause is GetMessageError.Denied - trace( - tag = TAG, - message = "Fetching the read pointer's message $readPointer failed", - error = cause, - // Error routes through ErrorUtils, which drops transport failures. - type = if (expected) TraceType.Log else TraceType.Error, - ) - } + lookUpReadStamp(key)?.let { stamp -> + stateHolder.update { it.copy(fetchedReadStamps = it.fetchedReadStamps + (key to stamp)) } + } readStampsInFlight.remove(key) } } } + private suspend fun lookUpReadStamp(key: Pair): Long? { + val (chatId, readPointer) = key + return messagingController.getMessage(chatId, readPointer) + .onFailure { cause -> + if (cause is GetMessageError.NotFound) readStampsNotFound.add(key) + val expected = cause is GetMessageError.NotFound || cause is GetMessageError.Denied + trace( + tag = TAG, + message = "Fetching the read pointer's message $readPointer failed", + error = cause, + // Error routes through ErrorUtils, which drops transport failures. + type = if (expected) TraceType.Log else TraceType.Error, + ) + } + .getOrNull() + ?.unreadSeq + } + + /** + * The reconcile's share of [fetchMissingReadStamps]: looks up, concurrently but bounded, the + * stamps the incoming [chats] will need, and parks them in [stagedStamps] so the rebuild the + * write triggers shows exact counts on its first paint. + * + * Not derivable locally: a count is `lastMessage.unreadSeq - stamp(readPointer)`, and the + * stamp lives on the message the pointer names, which the device holds only if it ever loaded + * that far back. The event sequence counts every event (edits, reactions), not unread + * messages, so it cannot stand in. + */ + private suspend fun stageReadStamps(chats: List) { + val selfId = userManager.accountId ?: return + val selfPhone = userManager.profile?.verifiedPhoneNumber + val state = stateHolder.current + val wanted = chats.mapNotNull { chat -> + if (!isRenderable(chat, selfId, selfPhone)) return@mapNotNull null + val readPointer = selfReadPointer(chat, selfId) ?: return@mapNotNull null + val key = chat.chatId to readPointer + when { + state.readStampAt(chat.chatId, readPointer) != null -> null + key in readStampsNotFound -> null + // Read, or counted without a stamp: nothing to look up. + unreadCount(chat, selfId) { null } != null -> null + messageDataSource.getUnreadSeq(metadataDataSource.chatIdHex(chat.chatId), readPointer) != null -> null + else -> key + } + }.distinct() + if (wanted.isEmpty()) return + + val gate = Semaphore(READ_STAMP_CONCURRENCY) + withTimeoutOrNull(READ_STAMP_WAIT) { + coroutineScope { + wanted.map { key -> + async { + gate.withPermit { lookUpReadStamp(key)?.let { stagedStamps[key] = it } } + } + }.awaitAll() + } + } + } + /** * Emits [Event.ReadPointerUnreported] when the stored READ pointer for [chat] is ahead of the * one the feed just returned. @@ -373,64 +481,95 @@ class FeedSyncDelegate @Inject constructor( Result.success(contact.chats + (tip?.chats ?: emptyList())) } + /** Standalone DM sync: fetch, write, then the post-write steps. See [syncFeed]. */ internal suspend fun performFeedSync() { - stateHolder.update { it.copy(feedSyncState = FeedSyncState.Syncing) } - fetchCombinedFeed() + fetchFeed() .onSuccess { chats -> - metadataDataSource.upsert(chats) + writeFeed(chats) + onFeedCommitted(chats) + } + .onFailure { onFeedFailed(it) } + } - for (chat in chats) { - memberDataSource.upsert(chat.chatId, chat.members) - } - val previews = chats.lastMessagesByChat() - // Not waited on: a preview is one message per chat, and the feed is not held for it. - linkPrefetch.prefetch(previews.values.flatten()) - messageDataSource.upsertAll(previews) - - stateHolder.update { it.copy(feedSyncState = FeedSyncState.Synced) } - trace(tag = TAG, message = "Feed synced: ${chats.size} chats", type = TraceType.Process) - - val selfId = userManager.accountId - - for (chat in chats) { - if (selfId != null) reportUnreportedRead(chat, selfId) - - // The applied cursor, not the presence of messages, is what says whether a - // transcript was ever pulled: the loop above persists each chat's last-message - // preview, so "has messages" is true for nearly every chat in a feed the client - // has otherwise never fetched. Only a message load or an applied delta seats a - // cursor. - val cursor = metadataDataSource.getLatestEventSequence(chat.chatId) - when { - // Never fetched: take the newest page. Resuming a delta from 0 would instead - // re-pull the entire history as a "gap". - cursor <= 0L -> _events.send(Event.LoadMessages(chat.chatId)) - // Fetched, but the server has moved on: stream the missed window. - chat.latestEventSequence > cursor -> - _events.send(Event.DeltaSyncNeeded(chat.chatId)) - } - } + /** The network half of a DM sync. Nothing is written. */ + internal suspend fun fetchFeed(): Result> { + stateHolder.update { it.copy(feedSyncState = FeedSyncState.Syncing) } + return fetchCombinedFeed() + } - _events.send(Event.CatchUpComplete) - } - .onFailure { error -> - stateHolder.update { state -> - state.copy( - feedSyncState = FeedSyncState.Error, - // Don't downgrade a hydration that already succeeded: a later failure means - // this sync missed, not that the cache stopped being trustworthy. Moving off - // Unknown at all matters though — callers waiting on hydration have to stop - // waiting when the server is unreachable, or the wallet spins forever - // offline. - historyHydration = if (state.historyHydration == ChatHydrationState.Unknown) { - ChatHydrationState.Unavailable - } else { - state.historyHydration - }, - ) - } - trace(tag = TAG, message = "Feed sync failed: ${error.message}", type = TraceType.Error) + /** + * Writes [chats] — whatever mix of DM and group chats the caller has in hand — in one + * transaction, so the list rebuilds once. Link prefetch and read stamps go first: neither + * blocks the list, and the stamps ride in with the rebuild. + */ + internal suspend fun writeFeed(chats: List) { + if (chats.isEmpty()) return + // Not waited on: a preview is one message per chat, and the feed is not held for it. + linkPrefetch.prefetch(chats.lastMessagesByChat().values.flatten()) + stageReadStamps(chats) + feedWriter.write(chats) + } + + /** What follows a successful write of [chats]: the synced state and the catch-up work. */ + internal suspend fun onFeedCommitted(chats: List) { + stateHolder.update { it.copy(feedSyncState = FeedSyncState.Synced) } + fetchLeftoverReadStamps() + trace(tag = TAG, message = "Feed synced: ${chats.size} chats", type = TraceType.Process) + + val selfId = userManager.accountId + + for (chat in catchUpOrder(chats)) { + if (selfId != null) reportUnreportedRead(chat, selfId) + + // The applied cursor, not the presence of messages, is what says whether a + // transcript was ever pulled: the write above persists each chat's last-message + // preview, so "has messages" is true for nearly every chat in a feed the client + // has otherwise never fetched. Only a message load or an applied delta seats a + // cursor. + val cursor = metadataDataSource.getLatestEventSequence(chat.chatId) + when { + // Never fetched: take the newest page. Resuming a delta from 0 would instead + // re-pull the entire history as a "gap". + cursor <= 0L -> _events.send(Event.LoadMessages(chat.chatId)) + // Fetched, but the server has moved on: stream the missed window. + chat.latestEventSequence > cursor -> + _events.send(Event.DeltaSyncNeeded(chat.chatId)) } + } + + _events.send(Event.CatchUpComplete) + } + + internal fun onFeedFailed(error: Throwable) { + stateHolder.update { state -> + state.copy( + feedSyncState = FeedSyncState.Error, + // Don't downgrade a hydration that already succeeded: a later failure means + // this sync missed, not that the cache stopped being trustworthy. Moving off + // Unknown at all matters though — callers waiting on hydration have to stop + // waiting when the server is unreachable, or the wallet spins forever + // offline. + historyHydration = if (state.historyHydration == ChatHydrationState.Unknown) { + ChatHydrationState.Unavailable + } else { + state.historyHydration + }, + ) + } + fetchLeftoverReadStamps() + trace(tag = TAG, message = "Feed sync failed: ${error.message}", type = TraceType.Error) + } + + /** + * Catch-up runs one chat at a time, so the order is the order the user sees things fill in: + * the chat on screen first, then the list from the top down. + */ + internal fun catchUpOrder(chats: List): List { + val active = stateHolder.current.activeChat + return chats.sortedWith( + compareByDescending { it.chatId == active } + .thenByDescending { it.lastActivity }, + ) } // endregion diff --git a/apps/flipcash/shared/chat/src/main/kotlin/com/flipcash/shared/chat/internal/delegates/GroupFeedDelegate.kt b/apps/flipcash/shared/chat/src/main/kotlin/com/flipcash/shared/chat/internal/delegates/GroupFeedDelegate.kt index 7892ae60ec..283789557b 100644 --- a/apps/flipcash/shared/chat/src/main/kotlin/com/flipcash/shared/chat/internal/delegates/GroupFeedDelegate.kt +++ b/apps/flipcash/shared/chat/src/main/kotlin/com/flipcash/shared/chat/internal/delegates/GroupFeedDelegate.kt @@ -1,5 +1,6 @@ package com.flipcash.shared.chat.internal.delegates +import com.flipcash.app.persistence.sources.ChatFeedWriter import com.flipcash.app.persistence.sources.ChatMemberDataSource import com.flipcash.app.persistence.sources.ChatMessageDataSource import com.flipcash.app.persistence.sources.ChatMetadataDataSource @@ -51,6 +52,8 @@ class GroupFeedDelegate @Inject constructor( private val rosterStateHolder: RosterStateHolder, private val userManager: UserManager, private val linkPrefetch: MessageLinkPrefetch = MessageLinkPrefetch.None, + private val feedWriter: ChatFeedWriter = + ChatFeedWriter(metadataDataSource, memberDataSource, messageDataSource), ) : GroupOperations { sealed interface Event { @@ -75,9 +78,10 @@ class GroupFeedDelegate @Inject constructor( this.scope = scope } + /** Standalone group sync. A sync already in flight is shared, not cancelled and restarted. */ internal fun syncGroupFeed() { val scope = scope ?: return - syncJob?.cancel() + if (syncJob?.isActive == true) return syncJob = scope.launch { performGroupFeedSync() } } @@ -95,17 +99,28 @@ class GroupFeedDelegate @Inject constructor( * does not. */ internal suspend fun performGroupFeedSync() { + val chats = fetchGroupFeed() ?: return + persist(chats) + requestCatchUp(chats) + } + + /** The network half of a group sync: the first page, or null when the fetch failed. */ + internal suspend fun fetchGroupFeed(): List? { val page = chatController.getGroupChatFeed().getOrElse { trace(tag = TAG, message = "Group feed sync failed", type = TraceType.Error) - return + return null } - persist(page.chats) + return page.chats + } + + /** What follows a committed write of [chats]: transcripts to fetch or gaps to fill. */ + internal suspend fun requestCatchUp(chats: List) { // Caching the feed is not the same as having the transcript: [persist] writes each group's // last-message preview and nothing else, so without this a synced group opens to a single // bubble. The applied cursor is what says whether the transcript was ever pulled — a feed // sync never writes it. Same rule the DM feed applies in [FeedSyncDelegate]. - for (chat in page.chats) { + for (chat in chats) { val cursor = metadataDataSource.getLatestEventSequence(chat.chatId) when { cursor <= 0L -> _events.send(Event.LoadMessages(chat.chatId)) @@ -226,14 +241,9 @@ class GroupFeedDelegate @Inject constructor( } private suspend fun persist(chats: List) { - metadataDataSource.upsert(chats) - for (chat in chats) { - memberDataSource.upsert(chat.chatId, chat.members) - } - val previews = chats.lastMessagesByChat() // Not waited on: a preview is one message per chat, and the feed is not held for it. - linkPrefetch.prefetch(previews.values.flatten()) - messageDataSource.upsertAll(previews) + linkPrefetch.prefetch(chats.lastMessagesByChat().values.flatten()) + feedWriter.write(chats) } private companion object { diff --git a/apps/flipcash/shared/chat/src/main/kotlin/com/flipcash/shared/chat/internal/delegates/MessagingDelegate.kt b/apps/flipcash/shared/chat/src/main/kotlin/com/flipcash/shared/chat/internal/delegates/MessagingDelegate.kt index 245e1f4739..6d8cb8f77e 100644 --- a/apps/flipcash/shared/chat/src/main/kotlin/com/flipcash/shared/chat/internal/delegates/MessagingDelegate.kt +++ b/apps/flipcash/shared/chat/src/main/kotlin/com/flipcash/shared/chat/internal/delegates/MessagingDelegate.kt @@ -264,14 +264,10 @@ class MessagingDelegate @Inject constructor( // transcript as fetched, and it lets a following catch-up resume from head and // append genuinely newer messages instead of re-pulling the whole history from // sequence 0. Only ever advanced — a page older than the cursor must not rewind it. + // One write, and one that changes nothing when this page has nothing newer than + // what the feed already stored: the list rebuilds only for a real change. val head = messages.maxOfOrNull { it.eventSequence } ?: 0L - if (head > metadataDataSource.getLatestEventSequence(chatId)) { - metadataDataSource.updateLatestEventSequence(chatId, head) - } - - val latest = messages.maxByOrNull { it.messageId } ?: return@onSuccess - metadataDataSource.updateLastMessageId(chatId, latest.messageId) - metadataDataSource.updateLastActivity(chatId, latest.timestamp.toEpochMilliseconds()) + metadataDataSource.applyCatchUp(chatId, head, messages.maxByOrNull { it.messageId }) } } @@ -284,9 +280,10 @@ class MessagingDelegate @Inject constructor( // Deliberately not advancing the event-log cursor. A push carries one message, not // a page, so seating the cursor at its sequence would let a later catch-up resume // from a frontier it never actually fetched and skip whatever it missed in between. + // The row the list reads is updated here, so a message that arrived by push is already + // on the list at the next cold launch. if (message.messageId > (metadataDataSource.getLastMessageId(chatId) ?: 0L)) { - metadataDataSource.updateLastMessageId(chatId, message.messageId) - metadataDataSource.updateLastActivity(chatId, message.timestamp.toEpochMilliseconds()) + metadataDataSource.applyCatchUp(chatId, 0L, message) } } diff --git a/apps/flipcash/shared/chat/src/test/kotlin/com/flipcash/shared/chat/FeedReconcileTest.kt b/apps/flipcash/shared/chat/src/test/kotlin/com/flipcash/shared/chat/FeedReconcileTest.kt new file mode 100644 index 0000000000..a7b1a331ab --- /dev/null +++ b/apps/flipcash/shared/chat/src/test/kotlin/com/flipcash/shared/chat/FeedReconcileTest.kt @@ -0,0 +1,216 @@ +package com.flipcash.shared.chat + +import androidx.lifecycle.LifecycleOwner +import com.flipcash.app.core.dispatchers.TestDispatchers +import com.flipcash.app.persistence.sources.ChatFeedWriter +import com.flipcash.app.persistence.sources.ChatMemberDataSource +import com.flipcash.app.persistence.sources.ChatMessageDataSource +import com.flipcash.app.persistence.sources.ChatMetadataDataSource +import com.flipcash.app.persistence.sources.ContactDataSource +import com.flipcash.app.tokens.TokenCoordinator +import com.flipcash.services.controllers.ChatController +import com.flipcash.services.controllers.ChatMessagingController +import com.flipcash.services.controllers.EventStreamingController +import com.flipcash.services.models.chat.ChatId +import com.flipcash.services.models.chat.ChatFeedPage +import com.flipcash.services.models.chat.ChatMetadata +import com.flipcash.services.models.chat.ChatType +import com.flipcash.services.models.chat.ChatUpdate +import com.flipcash.services.models.chat.RosterChange +import com.flipcash.services.models.chat.RosterSummary +import com.flipcash.services.user.UserManager +import com.flipcash.shared.chat.internal.ChatIdGenerator +import com.flipcash.shared.chat.internal.ChatStateHolder +import com.flipcash.shared.chat.internal.RealChatCoordinator +import com.flipcash.shared.chat.internal.delegates.DmChatResolverDelegate +import com.flipcash.shared.chat.internal.delegates.EventStreamDelegate +import com.flipcash.shared.chat.internal.delegates.FeedSyncDelegate +import com.flipcash.shared.chat.internal.delegates.GroupFeedDelegate +import com.flipcash.shared.chat.internal.delegates.MessagingDelegate +import com.getcode.utils.network.NetworkConnectivityListener +import com.getcode.opencode.model.accounts.AccountCluster +import io.mockk.coEvery +import io.mockk.coVerify +import io.mockk.every +import io.mockk.mockk +import kotlinx.coroutines.ExperimentalCoroutinesApi +import kotlinx.coroutines.CompletableDeferred +import kotlinx.coroutines.channels.Channel +import kotlinx.coroutines.flow.receiveAsFlow +import kotlinx.coroutines.test.TestCoroutineScheduler +import kotlinx.coroutines.test.TestScope +import kotlinx.coroutines.test.runCurrent +import kotlinx.coroutines.test.runTest +import org.junit.Test +import org.junit.runner.RunWith +import org.robolectric.RobolectricTestRunner +import kotlin.time.Instant + +/** + * Launch fetches the DM and group feeds once and writes them once. A foreground that lands while + * that fetch is in flight shares it rather than cancelling and restarting it, which is what used to + * double the launch traffic and let the list change twice. + */ +@OptIn(ExperimentalCoroutinesApi::class) +@RunWith(RobolectricTestRunner::class) +class FeedReconcileTest { + + private val selfId = listOf(1, 2, 3) + private val dmChat = chat("aa", ChatType.CONTACT_DM) + private val groupChat = chat("bb", ChatType.GROUP) + + private fun chat(hex: String, type: ChatType) = ChatMetadata( + chatId = ChatId(hex), + type = type, + members = emptyList(), + lastMessage = null, + lastActivity = Instant.fromEpochSeconds(1000), + ) + + private val groupGate = CompletableDeferred() + private val chatController = mockk(relaxed = true).also { + coEvery { it.getDmChatFeed(ChatType.CONTACT_DM, any()) } returns + Result.success(ChatFeedPage(listOf(dmChat), null, false)) + coEvery { it.getDmChatFeed(ChatType.TIP_DM, any()) } returns + Result.success(ChatFeedPage(emptyList(), null, false)) + } + private val groupFeedDelegate = mockk(relaxed = true).also { + coEvery { it.fetchGroupFeed() } coAnswers { + groupGate.await() + listOf(groupChat) + } + } + private val feedWriter = mockk(relaxed = true) + private val testDispatchers = TestDispatchers(TestCoroutineScheduler()) + + private fun coordinator(): RealChatCoordinator { + val userManager = mockk(relaxed = true).also { + every { it.accountId } returns selfId + } + val eventStreamingController = mockk(relaxed = true).also { + every { it.chatUpdates } returns Channel().receiveAsFlow() + every { it.isConnected } returns true + every { it.isStreamActive } returns true + } + val metadataDataSource = mockk(relaxed = true) + val messageDataSource = mockk(relaxed = true) + val memberDataSource = mockk(relaxed = true) + val stateHolder = ChatStateHolder() + val messagingController = mockk(relaxed = true).also { + coEvery { it.getMessages(any(), any()) } returns Result.success(emptyList()) + } + + return RealChatCoordinator( + feedDelegate = FeedSyncDelegate( + messagingController = mockk(relaxed = true), + chatController = chatController, + metadataDataSource = metadataDataSource, + messageDataSource = messageDataSource, + memberDataSource = memberDataSource, + stateHolder = stateHolder, + userManager = userManager, + feedWriter = feedWriter, + ), + eventStreamDelegate = EventStreamDelegate( + eventStreamingController = eventStreamingController, + messagingController = mockk(relaxed = true), + metadataDataSource = metadataDataSource, + messageDataSource = messageDataSource, + memberDataSource = memberDataSource, + tokenCoordinator = mockk(relaxed = true), + userManager = userManager, + stateHolder = stateHolder, + analytics = mockk(relaxed = true), + exchange = mockk(relaxed = true), + ), + dmChatResolverDelegate = DmChatResolverDelegate( + chatIdGenerator = ChatIdGenerator(), + userManager = userManager, + contactDataSource = mockk(relaxed = true), + memberDataSource = memberDataSource, + ), + messagingDelegate = MessagingDelegate( + chatController = chatController, + messagingController = messagingController, + metadataDataSource = metadataDataSource, + messageDataSource = messageDataSource, + memberDataSource = memberDataSource, + notificationManager = mockk(relaxed = true), + userManager = userManager, + stateHolder = stateHolder, + analytics = mockk(relaxed = true), + senderResolver = mockk(relaxed = true), + ), + groupFeedDelegate = groupFeedDelegate, + reactionsDelegate = mockk(relaxed = true), + stateHolder = stateHolder, + draftStore = mockk(relaxed = true), + userManager = userManager, + networkObserver = mockk(relaxed = true), + dispatchers = testDispatchers, + ) + } + + private suspend fun TestScope.loggedIn(block: suspend (RealChatCoordinator) -> Unit) { + val subject = coordinator() + subject.onUserLoggedIn(mockk(relaxed = true)) + runCurrent() + try { + block(subject) + } finally { + groupGate.complete(Unit) + subject.teardown() + } + } + + @Test + fun `the DM and group feeds land in one write`() = runTest(testDispatchers.dispatcher) { + loggedIn { + groupGate.complete(Unit) + runCurrent() + + coVerify(exactly = 1) { feedWriter.write(match { it.toSet() == setOf(dmChat, groupChat) }) } + } + } + + @Test + fun `a foreground during the launch sync shares it instead of restarting it`() = + runTest(testDispatchers.dispatcher) { + loggedIn { subject -> + // The group fetch is still in flight: the launch sync has not finished. + subject.onStart(mockk(relaxed = true)) + runCurrent() + groupGate.complete(Unit) + runCurrent() + + coVerify(exactly = 1) { chatController.getDmChatFeed(ChatType.CONTACT_DM, any()) } + coVerify(exactly = 1) { groupFeedDelegate.fetchGroupFeed() } + coVerify(exactly = 1) { feedWriter.write(any()) } + } + } + + @Test + fun `a push refresh during the launch sync shares it too`() = runTest(testDispatchers.dispatcher) { + loggedIn { subject -> + subject.refreshFeed() + runCurrent() + groupGate.complete(Unit) + runCurrent() + + coVerify(exactly = 1) { chatController.getDmChatFeed(ChatType.CONTACT_DM, any()) } + } + } + + @Test + fun `a refresh after the sync has finished fetches again`() = runTest(testDispatchers.dispatcher) { + loggedIn { subject -> + groupGate.complete(Unit) + runCurrent() + + subject.refreshFeed() + runCurrent() + + coVerify(exactly = 2) { chatController.getDmChatFeed(ChatType.CONTACT_DM, any()) } + } + } +} diff --git a/apps/flipcash/shared/chat/src/test/kotlin/com/flipcash/shared/chat/FeedSyncPreviewWriteTest.kt b/apps/flipcash/shared/chat/src/test/kotlin/com/flipcash/shared/chat/FeedSyncPreviewWriteTest.kt index 1779ea9fdf..d07d02b825 100644 --- a/apps/flipcash/shared/chat/src/test/kotlin/com/flipcash/shared/chat/FeedSyncPreviewWriteTest.kt +++ b/apps/flipcash/shared/chat/src/test/kotlin/com/flipcash/shared/chat/FeedSyncPreviewWriteTest.kt @@ -63,7 +63,7 @@ class FeedSyncPreviewWriteTest { ).performFeedSync() coVerify(exactly = 1) { - messageDataSource.upsertAll( + messageDataSource.prepare( mapOf( ChatId("aa") to listOf(message(1)), ChatId("bb") to listOf(message(2)), diff --git a/apps/flipcash/shared/chat/src/test/kotlin/com/flipcash/shared/chat/GroupChatRoutingTest.kt b/apps/flipcash/shared/chat/src/test/kotlin/com/flipcash/shared/chat/GroupChatRoutingTest.kt index b9a8d3b64e..5d3ad27ee7 100644 --- a/apps/flipcash/shared/chat/src/test/kotlin/com/flipcash/shared/chat/GroupChatRoutingTest.kt +++ b/apps/flipcash/shared/chat/src/test/kotlin/com/flipcash/shared/chat/GroupChatRoutingTest.kt @@ -185,7 +185,7 @@ class GroupChatRoutingTest { @Test fun `logging in syncs the group feed`() = runTest(testDispatchers.dispatcher) { loggedIn { - verify(exactly = 1) { groupFeedDelegate.syncGroupFeed() } + coVerify(exactly = 1) { groupFeedDelegate.fetchGroupFeed() } } } @@ -198,7 +198,7 @@ class GroupChatRoutingTest { subject.refreshFeed() runCurrent() - verify(exactly = 1) { groupFeedDelegate.syncGroupFeed() } + coVerify(exactly = 1) { groupFeedDelegate.fetchGroupFeed() } } } @@ -226,7 +226,7 @@ class GroupChatRoutingTest { subject.onStart(mockk(relaxed = true)) runCurrent() - verify(exactly = 1) { groupFeedDelegate.syncGroupFeed() } + coVerify(exactly = 1) { groupFeedDelegate.fetchGroupFeed() } } } diff --git a/apps/flipcash/shared/chat/src/test/kotlin/com/flipcash/shared/chat/MessagingLoadCursorTest.kt b/apps/flipcash/shared/chat/src/test/kotlin/com/flipcash/shared/chat/MessagingLoadCursorTest.kt index 28066303a5..f07f338d34 100644 --- a/apps/flipcash/shared/chat/src/test/kotlin/com/flipcash/shared/chat/MessagingLoadCursorTest.kt +++ b/apps/flipcash/shared/chat/src/test/kotlin/com/flipcash/shared/chat/MessagingLoadCursorTest.kt @@ -57,11 +57,11 @@ class MessagingLoadCursorTest { delegateWith(messagingController, metadataDataSource).loadMessages(chatId) - coVerify(exactly = 1) { metadataDataSource.updateLatestEventSequence(chatId, 9) } + coVerify(exactly = 1) { metadataDataSource.applyCatchUp(chatId, 9, match { it.messageId == 2L }) } } @Test - fun `page older than the cursor does not rewind it`() = runTest { + fun `page older than the cursor goes through the forward-only write`() = runTest { val messagingController = mockk(relaxed = true) coEvery { messagingController.getMessages(chatId, any()) } returns Result.success(listOf(message(id = 1, eventSequence = 4))) @@ -70,7 +70,12 @@ class MessagingLoadCursorTest { delegateWith(messagingController, metadataDataSource).loadMessages(chatId) + // The rewind guard is the write's own (ChatMetadataCatchUpTest); what matters here is that + // the catch-up is one write and never the three separate ones that each invalidated the list. + coVerify(exactly = 1) { metadataDataSource.applyCatchUp(chatId, 4, any()) } coVerify(exactly = 0) { metadataDataSource.updateLatestEventSequence(chatId, any()) } + coVerify(exactly = 0) { metadataDataSource.updateLastMessageId(chatId, any()) } + coVerify(exactly = 0) { metadataDataSource.updateLastActivity(chatId, any()) } } @Test @@ -82,6 +87,7 @@ class MessagingLoadCursorTest { delegateWith(messagingController, metadataDataSource).loadMessages(chatId) + coVerify(exactly = 1) { metadataDataSource.applyCatchUp(chatId, 0, null) } coVerify(exactly = 0) { metadataDataSource.updateLatestEventSequence(chatId, any()) } } } diff --git a/apps/flipcash/shared/chat/src/test/kotlin/com/flipcash/shared/chat/MessagingPushedMessageTest.kt b/apps/flipcash/shared/chat/src/test/kotlin/com/flipcash/shared/chat/MessagingPushedMessageTest.kt index 252f9661d8..d5d60940cf 100644 --- a/apps/flipcash/shared/chat/src/test/kotlin/com/flipcash/shared/chat/MessagingPushedMessageTest.kt +++ b/apps/flipcash/shared/chat/src/test/kotlin/com/flipcash/shared/chat/MessagingPushedMessageTest.kt @@ -96,8 +96,7 @@ class MessagingPushedMessageTest { delegateWith(metadataDataSource) .applyPushedMessage(chatId, message(id = 12, eventSequence = 9, epochSeconds = 1_757_000_042)) - coVerify(exactly = 1) { metadataDataSource.updateLastMessageId(chatId, 12) } - coVerify(exactly = 1) { metadataDataSource.updateLastActivity(chatId, 1_757_000_042_000) } + coVerify(exactly = 1) { metadataDataSource.applyCatchUp(chatId, 0L, match { it.messageId == 12L }) } } @Test @@ -107,8 +106,7 @@ class MessagingPushedMessageTest { delegateWith(metadataDataSource).applyPushedMessage(chatId, message(id = 12, eventSequence = 9)) - coVerify(exactly = 0) { metadataDataSource.updateLastMessageId(chatId, any()) } - coVerify(exactly = 0) { metadataDataSource.updateLastActivity(chatId, any()) } + coVerify(exactly = 0) { metadataDataSource.applyCatchUp(chatId, any(), any()) } } @Test @@ -121,7 +119,7 @@ class MessagingPushedMessageTest { delegate.applyPushedMessage(chatId, pushed) delegate.applyPushedMessage(chatId, pushed) - coVerify(exactly = 1) { metadataDataSource.updateLastMessageId(chatId, 12) } + coVerify(exactly = 1) { metadataDataSource.applyCatchUp(chatId, 0L, match { it.messageId == 12L }) } } @Test @@ -131,6 +129,6 @@ class MessagingPushedMessageTest { delegateWith(metadataDataSource).applyPushedMessage(chatId, message(id = 1, eventSequence = 3)) - coVerify(exactly = 1) { metadataDataSource.updateLastMessageId(chatId, 1) } + coVerify(exactly = 1) { metadataDataSource.applyCatchUp(chatId, 0L, match { it.messageId == 1L }) } } } diff --git a/apps/flipcash/shared/persistence/db/src/main/kotlin/com/flipcash/app/persistence/dao/ChatMessageDao.kt b/apps/flipcash/shared/persistence/db/src/main/kotlin/com/flipcash/app/persistence/dao/ChatMessageDao.kt index 9d30d931a2..56c58df1bc 100644 --- a/apps/flipcash/shared/persistence/db/src/main/kotlin/com/flipcash/app/persistence/dao/ChatMessageDao.kt +++ b/apps/flipcash/shared/persistence/db/src/main/kotlin/com/flipcash/app/persistence/dao/ChatMessageDao.kt @@ -65,6 +65,21 @@ interface ChatMessageDao { ) suspend fun getLatestVisibleForAllChats(): List + /** + * [getLatestVisibleForAllChats] as a stream. The conversation list previews these rows, so it + * rebuilds when they change; a caller that de-duplicates the emissions rebuilds only when a + * chat's preview did, not on every message write. + */ + @Query( + "SELECT * FROM chat_messages WHERE rowid IN (" + + "SELECT (SELECT m.rowid FROM chat_messages m " + + "WHERE m.chat_id_hex = c.chat_id_hex AND m.is_deleted = 0 " + + "AND (m.encryption_state IS NULL OR m.encryption_state != 'KEY_PENDING') " + + "ORDER BY m.timestamp_epoch_ms DESC, m.message_id DESC LIMIT 1) " + + "FROM (SELECT DISTINCT chat_id_hex FROM chat_messages) c)" + ) + fun observeLatestVisibleForAllChats(): Flow> + @Query( "SELECT * FROM chat_messages " + "WHERE chat_id_hex = :chatIdHex " + diff --git a/apps/flipcash/shared/persistence/db/src/main/kotlin/com/flipcash/app/persistence/dao/ChatMetadataDao.kt b/apps/flipcash/shared/persistence/db/src/main/kotlin/com/flipcash/app/persistence/dao/ChatMetadataDao.kt index 8518a1e520..95d99cb35a 100644 --- a/apps/flipcash/shared/persistence/db/src/main/kotlin/com/flipcash/app/persistence/dao/ChatMetadataDao.kt +++ b/apps/flipcash/shared/persistence/db/src/main/kotlin/com/flipcash/app/persistence/dao/ChatMetadataDao.kt @@ -239,6 +239,39 @@ interface ChatMetadataDao { @Query("SELECT latest_event_sequence FROM chat_metadata WHERE chat_id_hex = :chatIdHex") suspend fun getLatestEventSequence(chatIdHex: String): Long? + /** + * The catch-up write: one statement, so a catch-up costs the list one invalidation instead of + * three, and idempotent, so one that learns nothing costs it none. + * + * - `latest_event_sequence` only moves forward. + * - `last_message_id` and `last_activity_epoch_ms` move only when [messageId] is strictly newer + * than the stored `last_message_id`. A feed sync has usually written the same message + * already, and re-stamping its activity from the message's own timestamp (which can differ + * from the server's `last_activity`) reordered rows that had nothing new. + * - Even when the message is newer, `last_activity_epoch_ms` never moves backwards: `MAX`. + * + * The `WHERE` guard turns the no-op case into zero changed rows, which SQLite does not report + * to Room's invalidation tracker. Pass `0` for [latestEventSequence] or [messageId] to leave + * that half alone. + */ + @Query( + "UPDATE chat_metadata SET " + + "latest_event_sequence = MAX(latest_event_sequence, :latestEventSequence), " + + "last_activity_epoch_ms = CASE WHEN :messageId > COALESCE(last_message_id, 0) " + + "THEN MAX(last_activity_epoch_ms, :timestampEpochMs) ELSE last_activity_epoch_ms END, " + + "last_message_id = CASE WHEN :messageId > COALESCE(last_message_id, 0) " + + "THEN :messageId ELSE last_message_id END " + + "WHERE chat_id_hex = :chatIdHex AND (" + + ":latestEventSequence > latest_event_sequence " + + "OR :messageId > COALESCE(last_message_id, 0))" + ) + suspend fun applyCatchUp( + chatIdHex: String, + latestEventSequence: Long, + messageId: Long, + timestampEpochMs: Long, + ) + @Query("SELECT chat_type FROM chat_metadata WHERE chat_id_hex = :chatIdHex") suspend fun getChatType(chatIdHex: String): String? diff --git a/apps/flipcash/shared/persistence/db/src/test/kotlin/com/flipcash/app/persistence/ChatMetadataCatchUpTest.kt b/apps/flipcash/shared/persistence/db/src/test/kotlin/com/flipcash/app/persistence/ChatMetadataCatchUpTest.kt new file mode 100644 index 0000000000..5571949a37 --- /dev/null +++ b/apps/flipcash/shared/persistence/db/src/test/kotlin/com/flipcash/app/persistence/ChatMetadataCatchUpTest.kt @@ -0,0 +1,99 @@ +package com.flipcash.app.persistence + +import android.content.Context +import androidx.test.core.app.ApplicationProvider +import com.flipcash.app.persistence.entities.ChatMetadataEntity +import kotlinx.coroutines.flow.first +import kotlinx.coroutines.runBlocking +import org.junit.After +import org.junit.Before +import org.junit.Test +import org.junit.runner.RunWith +import org.robolectric.RobolectricTestRunner +import kotlin.test.assertEquals + +/** + * A catch-up that learns nothing must leave the row byte-for-byte alone, because the chat list + * reorders on `last_activity_epoch_ms` and rebuilds on any change to the row. + */ +@RunWith(RobolectricTestRunner::class) +class ChatMetadataCatchUpTest { + + private val context = ApplicationProvider.getApplicationContext() + private lateinit var dao: com.flipcash.app.persistence.dao.ChatMetadataDao + + @Before + fun setUp() { + FlipcashDatabase.init(context, "cccccccccccccccccccccccc") + dao = FlipcashDatabase.requireInstance().chatMetadataDao() + runBlocking { + dao.upsert( + ChatMetadataEntity( + chatIdHex = CHAT, + chatType = "CONTACT_DM", + lastActivityEpochMs = 5_000, + lastMessageId = 10, + latestEventSequence = 20, + ), + ) + } + } + + @After + fun tearDown() { + FlipcashDatabase.closeDb() + } + + private fun row() = runBlocking { dao.observeAll().first().single() } + + @Test + fun `the same newest message changes nothing, even with a later timestamp`() = runBlocking { + val before = row() + + dao.applyCatchUp(CHAT, latestEventSequence = 20, messageId = 10, timestampEpochMs = 9_999) + + assertEquals(before, row()) + } + + @Test + fun `an older message and an older cursor change nothing`() = runBlocking { + val before = row() + + dao.applyCatchUp(CHAT, latestEventSequence = 3, messageId = 4, timestampEpochMs = 9_999) + + assertEquals(before, row()) + } + + @Test + fun `a newer message moves the id and activity together with the cursor`() = runBlocking { + dao.applyCatchUp(CHAT, latestEventSequence = 25, messageId = 11, timestampEpochMs = 6_000) + + val after = row() + assertEquals(11L, after.lastMessageId) + assertEquals(6_000L, after.lastActivityEpochMs) + assertEquals(25L, after.latestEventSequence) + } + + @Test + fun `a newer message with an older timestamp never moves activity backwards`() = runBlocking { + dao.applyCatchUp(CHAT, latestEventSequence = 0, messageId = 11, timestampEpochMs = 1_000) + + val after = row() + assertEquals(11L, after.lastMessageId) + assertEquals(5_000L, after.lastActivityEpochMs) + } + + @Test + fun `a cursor-only catch-up leaves the message fields alone`() = runBlocking { + dao.applyCatchUp(CHAT, latestEventSequence = 30, messageId = 0, timestampEpochMs = 0) + + val after = row() + assertEquals(30L, after.latestEventSequence) + assertEquals(10L, after.lastMessageId) + assertEquals(5_000L, after.lastActivityEpochMs) + } + + private companion object { + const val CHAT = "aabb" + } +} diff --git a/apps/flipcash/shared/persistence/sources/src/main/kotlin/com/flipcash/app/persistence/sources/ChatFeedWriter.kt b/apps/flipcash/shared/persistence/sources/src/main/kotlin/com/flipcash/app/persistence/sources/ChatFeedWriter.kt new file mode 100644 index 0000000000..1c6b907771 --- /dev/null +++ b/apps/flipcash/shared/persistence/sources/src/main/kotlin/com/flipcash/app/persistence/sources/ChatFeedWriter.kt @@ -0,0 +1,37 @@ +package com.flipcash.app.persistence.sources + +import androidx.room.withTransaction +import com.flipcash.app.persistence.FlipcashDatabase +import com.flipcash.services.models.chat.ChatMetadata +import javax.inject.Inject +import javax.inject.Singleton + +/** + * Lands a batch of feed chats — rows, members and preview messages, DM and group alike — in one + * Room transaction. + * + * Room tells observers about a transaction once, when it commits. Written as separate calls, the + * chat list saw the rows reorder first, then each chat's members, then the previews, and rebuilt + * for every one of them. + */ +@Singleton +class ChatFeedWriter @Inject constructor( + private val metadataDataSource: ChatMetadataDataSource, + private val memberDataSource: ChatMemberDataSource, + private val messageDataSource: ChatMessageDataSource, +) { + suspend fun write(chats: List) { + if (chats.isEmpty()) return + // Opened before the transaction: opening can fetch a key over the network, and a + // transaction holds the database's one writer. + val prepared = messageDataSource.prepare(chats.lastMessagesByChat()) + val body: suspend () -> Unit = { + metadataDataSource.upsert(chats) + for (chat in chats) memberDataSource.upsert(chat.chatId, chat.members) + messageDataSource.writePrepared(prepared) + } + val database = FlipcashDatabase.getInstance() + if (database != null) database.withTransaction { body() } else body() + messageDataSource.afterWrite(prepared) + } +} diff --git a/apps/flipcash/shared/persistence/sources/src/main/kotlin/com/flipcash/app/persistence/sources/ChatMessageDataSource.kt b/apps/flipcash/shared/persistence/sources/src/main/kotlin/com/flipcash/app/persistence/sources/ChatMessageDataSource.kt index a1453d65ae..3810eec0ac 100644 --- a/apps/flipcash/shared/persistence/sources/src/main/kotlin/com/flipcash/app/persistence/sources/ChatMessageDataSource.kt +++ b/apps/flipcash/shared/persistence/sources/src/main/kotlin/com/flipcash/app/persistence/sources/ChatMessageDataSource.kt @@ -36,6 +36,9 @@ data class PendingMessage( val clientMessageId: ClientMessageId, ) +/** Messages already opened by [ChatMessageDataSource.prepare], ready for one write. */ +class PreparedMessages internal constructor(internal val byChat: Map>) + @Singleton class ChatMessageDataSource @Inject constructor( private val mapper: ChatEntityMapper, @@ -134,6 +137,14 @@ class ChatMessageDataSource @Inject constructor( ?.associate { it.chatIdHex to toChatMessage(it) } .orEmpty() + /** + * Emits whenever the set of newest-visible messages changes (one per chat). Carries the rows + * only so callers can de-duplicate on them; the list reads previews with + * [getLatestVisibleByChat]. + */ + fun observeLatestVisibleChanges(): Flow> = + db?.chatMessageDao()?.observeLatestVisibleForAllChats()?.distinctUntilChanged() ?: emptyFlow() + /** * The server's unread stamp on one stored message, or null when that message isn't stored — * the running count of unread-eligible messages, so two stamps subtract to the count between. @@ -293,16 +304,35 @@ class ChatMessageDataSource @Inject constructor( * chat in the feed. */ suspend fun upsertAll(messagesByChat: Map>) { + if (db == null || messagesByChat.isEmpty()) return + val prepared = prepare(messagesByChat) + writePrepared(prepared) + afterWrite(prepared) + } + + /** + * The network half of [upsertAll]: opens what can be opened. Split out so a caller can run it + * before it takes a transaction, since opening can fetch a key. + */ + suspend fun prepare(messagesByChat: Map>): PreparedMessages = + PreparedMessages( + messagesByChat.mapValues { (chatId, messages) -> + openEncrypted(chatId, mapper.chatIdHex(chatId), messages) + }, + ) + + /** The write half of [upsertAll]. Joins a transaction the caller already holds. */ + suspend fun writePrepared(prepared: PreparedMessages) { val database = db ?: return - if (messagesByChat.isEmpty()) return - // Opened before the transaction: opening can fetch a key over the network. - val opened = messagesByChat.mapValues { (chatId, messages) -> - openEncrypted(chatId, mapper.chatIdHex(chatId), messages) - } + if (prepared.byChat.isEmpty()) return database.withTransaction { - for ((chatId, messages) in opened) write(mapper.chatIdHex(chatId), messages) + for ((chatId, messages) in prepared.byChat) write(mapper.chatIdHex(chatId), messages) } - for ((chatId, messages) in opened) { + } + + /** Runs once the transaction that carried [prepared] has committed. */ + suspend fun afterWrite(prepared: PreparedMessages) { + for ((chatId, messages) in prepared.byChat) { if (messages.none { it.encryption == MessageEncryption.KeyPending }) reopenKeyPending(chatId) } } diff --git a/apps/flipcash/shared/persistence/sources/src/main/kotlin/com/flipcash/app/persistence/sources/ChatMetadataDataSource.kt b/apps/flipcash/shared/persistence/sources/src/main/kotlin/com/flipcash/app/persistence/sources/ChatMetadataDataSource.kt index 3dc964ab34..fdf71e5def 100644 --- a/apps/flipcash/shared/persistence/sources/src/main/kotlin/com/flipcash/app/persistence/sources/ChatMetadataDataSource.kt +++ b/apps/flipcash/shared/persistence/sources/src/main/kotlin/com/flipcash/app/persistence/sources/ChatMetadataDataSource.kt @@ -116,6 +116,21 @@ class ChatMetadataDataSource @Inject constructor( db?.chatMetadataDao()?.updateLastMessageId(mapper.chatIdHex(chatId), messageId) } + /** + * Records what a catch-up learned about [chatId] in one write: the event cursor [latestEventSequence] + * (0 for none) and the [newest] message it carried (null for none). Only advances; see + * [ChatMetadataDao.applyCatchUp][com.flipcash.app.persistence.dao.ChatMetadataDao.applyCatchUp]. + */ + suspend fun applyCatchUp(chatId: ChatId, latestEventSequence: Long, newest: ChatMessage?) { + if (latestEventSequence <= 0L && newest == null) return + db?.chatMetadataDao()?.applyCatchUp( + chatIdHex = mapper.chatIdHex(chatId), + latestEventSequence = latestEventSequence, + messageId = newest?.messageId ?: 0L, + timestampEpochMs = newest?.timestamp?.toEpochMilliseconds() ?: 0L, + ) + } + suspend fun updateLatestEventSequence(chatId: ChatId, sequence: Long) { db?.chatMetadataDao()?.updateLatestEventSequence(mapper.chatIdHex(chatId), sequence) } From b3f329faa5cfd24dd91f2ec7b9a6a11f30b8f266 Mon Sep 17 00:00:00 2001 From: Brandon McAnsh Date: Fri, 2 Oct 2026 12:31:03 -0400 Subject: [PATCH 2/4] fix(chat): make the feed upsert newer-wins on activity and last message A feed response in flight while a stream event advanced a row overwrote last_activity and last_message_id with older server values, rolling the row back. The feed upsert now takes MAX on last_activity and replaces last_message_id only when strictly newer. Other feed-owned fields are unchanged. --- .../app/persistence/dao/ChatMetadataDao.kt | 11 ++- .../persistence/ChatFeedUpsertClampTest.kt | 83 +++++++++++++++++++ 2 files changed, 91 insertions(+), 3 deletions(-) create mode 100644 apps/flipcash/shared/persistence/db/src/test/kotlin/com/flipcash/app/persistence/ChatFeedUpsertClampTest.kt diff --git a/apps/flipcash/shared/persistence/db/src/main/kotlin/com/flipcash/app/persistence/dao/ChatMetadataDao.kt b/apps/flipcash/shared/persistence/db/src/main/kotlin/com/flipcash/app/persistence/dao/ChatMetadataDao.kt index 95d99cb35a..574b5d24f9 100644 --- a/apps/flipcash/shared/persistence/db/src/main/kotlin/com/flipcash/app/persistence/dao/ChatMetadataDao.kt +++ b/apps/flipcash/shared/persistence/db/src/main/kotlin/com/flipcash/app/persistence/dao/ChatMetadataDao.kt @@ -72,7 +72,10 @@ interface ChatMetadataDao { suspend fun insertIfAbsent(entity: ChatMetadataEntity): Long /** - * Overwrites only the columns the server owns. `latest_event_sequence` and + * Overwrites only the columns the server owns, except that `last_activity_epoch_ms` and + * `last_message_id` are newer-wins: a feed response that was in flight while a stream event + * advanced the row is older than the row, and must not roll it back. Activity only moves + * forward (`MAX`); the message id is replaced only when strictly newer. `latest_event_sequence` and * `analytics_counted_through` are client-owned watermarks that no server payload * carries, so they are deliberately absent here. The roster and viewer-state columns * are absent too — they are versioned, and go through [updateRosterIfNewer] and @@ -80,8 +83,10 @@ interface ChatMetadataDao { */ @Query( "UPDATE chat_metadata SET chat_type = :chatType, " + - "last_activity_epoch_ms = :lastActivityEpochMs, " + - "last_message_id = :lastMessageId, " + + "last_activity_epoch_ms = MAX(last_activity_epoch_ms, :lastActivityEpochMs), " + + "last_message_id = CASE WHEN :lastMessageId IS NOT NULL " + + "AND :lastMessageId > COALESCE(last_message_id, 0) " + + "THEN :lastMessageId ELSE last_message_id END, " + "is_hidden = :isHidden, " + "title = :title, " + "picture_json = :pictureJson, " + diff --git a/apps/flipcash/shared/persistence/db/src/test/kotlin/com/flipcash/app/persistence/ChatFeedUpsertClampTest.kt b/apps/flipcash/shared/persistence/db/src/test/kotlin/com/flipcash/app/persistence/ChatFeedUpsertClampTest.kt new file mode 100644 index 0000000000..07d948fc85 --- /dev/null +++ b/apps/flipcash/shared/persistence/db/src/test/kotlin/com/flipcash/app/persistence/ChatFeedUpsertClampTest.kt @@ -0,0 +1,83 @@ +package com.flipcash.app.persistence + +import android.content.Context +import androidx.test.core.app.ApplicationProvider +import com.flipcash.app.persistence.dao.ChatMetadataDao +import com.flipcash.app.persistence.entities.ChatMetadataEntity +import kotlinx.coroutines.flow.first +import kotlinx.coroutines.runBlocking +import org.junit.After +import org.junit.Before +import org.junit.Test +import org.junit.runner.RunWith +import org.robolectric.RobolectricTestRunner +import kotlin.test.assertEquals + +/** A delayed feed response must not roll back a row a newer stream event already advanced. */ +@RunWith(RobolectricTestRunner::class) +class ChatFeedUpsertClampTest { + + private val context = ApplicationProvider.getApplicationContext() + private lateinit var dao: ChatMetadataDao + + private fun feedRow(activity: Long, messageId: Long?, title: String? = null) = ChatMetadataEntity( + chatIdHex = CHAT, + chatType = "CONTACT_DM", + lastActivityEpochMs = activity, + lastMessageId = messageId, + title = title, + ) + + @Before + fun setUp() { + FlipcashDatabase.init(context, "dddddddddddddddddddddddd") + dao = FlipcashDatabase.requireInstance().chatMetadataDao() + runBlocking { dao.upsert(feedRow(5_000, 10)) } + } + + @After + fun tearDown() = FlipcashDatabase.closeDb() + + private fun row() = runBlocking { dao.observeAll().first().single() } + + @Test + fun `an older feed response after a stream event leaves the row unchanged`() = runBlocking { + dao.applyCatchUp(CHAT, latestEventSequence = 30, messageId = 12, timestampEpochMs = 8_000) + val advanced = row() + + dao.upsert(feedRow(activity = 5_000, messageId = 10)) + + assertEquals(advanced, row()) + } + + @Test + fun `a newer feed response advances activity and message id`() = runBlocking { + dao.upsert(feedRow(activity = 9_000, messageId = 14)) + + assertEquals(9_000L, row().lastActivityEpochMs) + assertEquals(14L, row().lastMessageId) + } + + @Test + fun `other server-owned fields still take the feed's value and the cursor is untouched`() = runBlocking { + dao.applyCatchUp(CHAT, latestEventSequence = 30, messageId = 12, timestampEpochMs = 8_000) + + dao.upsert(feedRow(activity = 5_000, messageId = 10, title = "renamed")) + + assertEquals("renamed", row().title) + assertEquals(30L, row().latestEventSequence) + assertEquals(12L, row().lastMessageId) + assertEquals(8_000L, row().lastActivityEpochMs) + } + + @Test + fun `a feed row with no last message keeps the stored one`() = runBlocking { + dao.upsert(feedRow(activity = 5_000, messageId = null)) + + assertEquals(10L, row().lastMessageId) + } + + private companion object { + const val CHAT = "aabb" + } +} From ec8da060f39e36cc73ac3a935a112de4793e121b Mon Sep 17 00:00:00 2001 From: Brandon McAnsh Date: Fri, 2 Oct 2026 14:25:03 -0400 Subject: [PATCH 3/4] fix(chat): queue one trailing feed sync for a push during a sync A push refresh that lands mid-sync joined that sync, which may have fetched before the change. It now also queues a single trailing run, coalesced across pushes. onStart, reconnect and heartbeat still only share the in-flight sync. --- .../shared/chat/internal/RealChatCoordinator.kt | 4 +++- .../flipcash/shared/chat/FeedReconcileTest.kt | 16 ++++++++++++++-- 2 files changed, 17 insertions(+), 3 deletions(-) diff --git a/apps/flipcash/shared/chat/src/main/kotlin/com/flipcash/shared/chat/internal/RealChatCoordinator.kt b/apps/flipcash/shared/chat/src/main/kotlin/com/flipcash/shared/chat/internal/RealChatCoordinator.kt index 73eebdc760..13cd3e1664 100644 --- a/apps/flipcash/shared/chat/src/main/kotlin/com/flipcash/shared/chat/internal/RealChatCoordinator.kt +++ b/apps/flipcash/shared/chat/src/main/kotlin/com/flipcash/shared/chat/internal/RealChatCoordinator.kt @@ -298,7 +298,9 @@ class RealChatCoordinator @Inject constructor( * conversation list, and none of them know which half a chat belongs to. */ override fun refreshFeed() { - syncFeeds() + // A push says the server has something new: join a sync in flight, but queue one trailing + // run behind it, since that sync may have fetched before the change. Coalesced. + syncFeeds(fresh = true) } /** diff --git a/apps/flipcash/shared/chat/src/test/kotlin/com/flipcash/shared/chat/FeedReconcileTest.kt b/apps/flipcash/shared/chat/src/test/kotlin/com/flipcash/shared/chat/FeedReconcileTest.kt index a7b1a331ab..ecbb733560 100644 --- a/apps/flipcash/shared/chat/src/test/kotlin/com/flipcash/shared/chat/FeedReconcileTest.kt +++ b/apps/flipcash/shared/chat/src/test/kotlin/com/flipcash/shared/chat/FeedReconcileTest.kt @@ -190,14 +190,26 @@ class FeedReconcileTest { } @Test - fun `a push refresh during the launch sync shares it too`() = runTest(testDispatchers.dispatcher) { + fun `a push during an in-flight sync queues exactly one trailing fetch`() = runTest(testDispatchers.dispatcher) { loggedIn { subject -> subject.refreshFeed() runCurrent() groupGate.complete(Unit) runCurrent() - coVerify(exactly = 1) { chatController.getDmChatFeed(ChatType.CONTACT_DM, any()) } + coVerify(exactly = 2) { chatController.getDmChatFeed(ChatType.CONTACT_DM, any()) } + } + } + + @Test + fun `several pushes during one sync still produce one trailing fetch`() = runTest(testDispatchers.dispatcher) { + loggedIn { subject -> + repeat(5) { subject.refreshFeed() } + runCurrent() + groupGate.complete(Unit) + runCurrent() + + coVerify(exactly = 2) { chatController.getDmChatFeed(ChatType.CONTACT_DM, any()) } } } From a698d35dee321a1001a09f00f3850a2d47de96c4 Mon Sep 17 00:00:00 2001 From: Brandon McAnsh Date: Fri, 2 Oct 2026 15:03:08 -0400 Subject: [PATCH 4/4] fix(chat): name a group preview's sender on the chat list's first frame The first-frame draw passed no sender profiles, so a cold launch drew group previews without the sender prefix until the user_profiles read landed. Names are already persisted there; expose the last read synchronously and use it for that draw. Adds tests that a stored name is in the next session's first emission and that re-storing an unchanged name does not re-emit. --- .../app/tipping/internal/ChatsViewModel.kt | 10 ++- .../flipcash/shared/chat/ChatCoordinator.kt | 7 ++ .../shared/chat/internal/SenderResolver.kt | 17 +++++ .../internal/delegates/MessagingDelegate.kt | 2 + .../shared/chat/SenderResolverTest.kt | 13 ++++ .../persistence/sources/build.gradle.kts | 1 + .../sources/UserProfileDataSourceTest.kt | 71 +++++++++++++++++++ 7 files changed, 119 insertions(+), 2 deletions(-) create mode 100644 apps/flipcash/shared/persistence/sources/src/test/kotlin/com/flipcash/app/persistence/sources/UserProfileDataSourceTest.kt diff --git a/apps/flipcash/features/tipping/src/main/kotlin/com/flipcash/app/tipping/internal/ChatsViewModel.kt b/apps/flipcash/features/tipping/src/main/kotlin/com/flipcash/app/tipping/internal/ChatsViewModel.kt index f6765f4ce5..be2f8eaa5b 100644 --- a/apps/flipcash/features/tipping/src/main/kotlin/com/flipcash/app/tipping/internal/ChatsViewModel.kt +++ b/apps/flipcash/features/tipping/src/main/kotlin/com/flipcash/app/tipping/internal/ChatsViewModel.kt @@ -50,7 +50,7 @@ internal class ChatsViewModel @Inject constructor( fun conversations( summaries: List, tokens: List, - // Null on the first-frame draw, before the table has been read: asking then would + // Null on the first-frame draw if the table has not been read yet: asking then would // re-fetch every sender already on disk. senderProfiles: Map?, ): List { @@ -72,7 +72,13 @@ internal class ChatsViewModel @Inject constructor( chatCoordinator.currentChatListFeed()?.let { summaries -> dispatchEvent( Event.ChatsUpdated( - Loadable.Loaded(conversations(summaries, tokenCoordinator.cachedTokens(), null)) + Loadable.Loaded( + conversations( + summaries, + tokenCoordinator.cachedTokens(), + chatCoordinator.currentSenderProfiles(), + ) + ) ) ) } diff --git a/apps/flipcash/shared/chat/src/main/kotlin/com/flipcash/shared/chat/ChatCoordinator.kt b/apps/flipcash/shared/chat/src/main/kotlin/com/flipcash/shared/chat/ChatCoordinator.kt index a07750b4ff..d2e70089c6 100644 --- a/apps/flipcash/shared/chat/src/main/kotlin/com/flipcash/shared/chat/ChatCoordinator.kt +++ b/apps/flipcash/shared/chat/src/main/kotlin/com/flipcash/shared/chat/ChatCoordinator.kt @@ -351,6 +351,13 @@ interface MessagingOperations { */ fun observeSenderProfiles(): Flow> + /** + * The last read of the profiles [observeSenderProfiles] emits, or null before the first read + * lands. For a caller that must draw synchronously and would otherwise leave a group preview + * without its sender's name until the first emission. + */ + fun currentSenderProfiles(): Map? + /** * Asks for [userId]'s profile if nothing has asked already, for a sender the roster subset * does not cover or whose cached profile has no name. Returns immediately; the answer arrives diff --git a/apps/flipcash/shared/chat/src/main/kotlin/com/flipcash/shared/chat/internal/SenderResolver.kt b/apps/flipcash/shared/chat/src/main/kotlin/com/flipcash/shared/chat/internal/SenderResolver.kt index b5a2dae50b..e751c58e83 100644 --- a/apps/flipcash/shared/chat/src/main/kotlin/com/flipcash/shared/chat/internal/SenderResolver.kt +++ b/apps/flipcash/shared/chat/src/main/kotlin/com/flipcash/shared/chat/internal/SenderResolver.kt @@ -99,6 +99,23 @@ class SenderResolver @Inject constructor( /** Every profile this device holds, keyed by user-id hex. The transcript indexes into it. */ val profiles: Flow> = userProfileDataSource.observeProfiles() + @Volatile + private var snapshot: Map? = null + + /** + * The last `user_profiles` read, or null before the first one lands. Lets a caller that must + * draw synchronously (the chat list's first frame) name a group's last sender from names + * persisted by an earlier session, instead of drawing it unattributed until [profiles] emits. + * Kept in its own scope so [clear] does not stop it; [profiles] follows the open database. + */ + val cachedProfiles: Map? get() = snapshot + + init { + CoroutineScope(dispatchers.IO + SupervisorJob()).launch { + profiles.collect { snapshot = it } + } + } + /** * Asks for [userId]'s profile if nothing has asked already. Returns immediately. */ diff --git a/apps/flipcash/shared/chat/src/main/kotlin/com/flipcash/shared/chat/internal/delegates/MessagingDelegate.kt b/apps/flipcash/shared/chat/src/main/kotlin/com/flipcash/shared/chat/internal/delegates/MessagingDelegate.kt index 6d8cb8f77e..4c14928351 100644 --- a/apps/flipcash/shared/chat/src/main/kotlin/com/flipcash/shared/chat/internal/delegates/MessagingDelegate.kt +++ b/apps/flipcash/shared/chat/src/main/kotlin/com/flipcash/shared/chat/internal/delegates/MessagingDelegate.kt @@ -237,6 +237,8 @@ class MessagingDelegate @Inject constructor( override fun observeSenderProfiles(): Flow> = senderResolver.profiles + override fun currentSenderProfiles(): Map? = senderResolver.cachedProfiles + override fun requestSenderProfile(userId: ID) = senderResolver.request(userId) override fun observeOldestEncryptedMessageId(chatId: ChatId): Flow = diff --git a/apps/flipcash/shared/chat/src/test/kotlin/com/flipcash/shared/chat/SenderResolverTest.kt b/apps/flipcash/shared/chat/src/test/kotlin/com/flipcash/shared/chat/SenderResolverTest.kt index ab14321ec2..f591dd0fc3 100644 --- a/apps/flipcash/shared/chat/src/test/kotlin/com/flipcash/shared/chat/SenderResolverTest.kt +++ b/apps/flipcash/shared/chat/src/test/kotlin/com/flipcash/shared/chat/SenderResolverTest.kt @@ -55,6 +55,19 @@ class SenderResolverTest { dispatchers = dispatchers, ) + @Test + fun `cachedProfiles holds the persisted names as soon as the table is read`() = + runTest(dispatchers.dispatcher) { + every { userProfileDataSource.observeProfiles() } returns + MutableStateFlow(mapOf(userIdHex to profile)) + val resolver = subject() + assertEquals(null, resolver.cachedProfiles) + + runCurrent() + + assertEquals("Ada", resolver.cachedProfiles?.get(userIdHex)?.displayName) + } + @Test fun `a miss triggers one fetch and is written to user_profiles`() = runTest(dispatchers.dispatcher) { coEvery { profileController.getProfileForUser(userId) } returns Result.success(profile) diff --git a/apps/flipcash/shared/persistence/sources/build.gradle.kts b/apps/flipcash/shared/persistence/sources/build.gradle.kts index 744be06cbf..5fa1f607cf 100644 --- a/apps/flipcash/shared/persistence/sources/build.gradle.kts +++ b/apps/flipcash/shared/persistence/sources/build.gradle.kts @@ -16,6 +16,7 @@ android { dependencies { testImplementation(kotlin("test")) testImplementation(libs.bundles.unit.testing) + testImplementation(libs.robolectric) testImplementation(testFixtures(project(":services:flipcash"))) implementation(libs.bundles.kotlinx.serialization) diff --git a/apps/flipcash/shared/persistence/sources/src/test/kotlin/com/flipcash/app/persistence/sources/UserProfileDataSourceTest.kt b/apps/flipcash/shared/persistence/sources/src/test/kotlin/com/flipcash/app/persistence/sources/UserProfileDataSourceTest.kt new file mode 100644 index 0000000000..ee494c58c0 --- /dev/null +++ b/apps/flipcash/shared/persistence/sources/src/test/kotlin/com/flipcash/app/persistence/sources/UserProfileDataSourceTest.kt @@ -0,0 +1,71 @@ +package com.flipcash.app.persistence.sources + +import app.cash.turbine.test +import com.flipcash.app.persistence.FlipcashDatabase +import com.flipcash.services.models.UserProfile +import com.getcode.utils.hexEncodedString +import kotlinx.coroutines.test.runTest +import org.junit.After +import org.junit.Before +import org.junit.Test +import org.junit.runner.RunWith +import org.robolectric.RobolectricTestRunner +import org.robolectric.RuntimeEnvironment +import kotlin.test.assertEquals + +/** + * A group preview is attributed from `user_profiles`. These pin the two properties the chat list's + * launch behavior rests on: a name stored by one session is in the next session's first emission, + * and re-storing an unchanged name does not make the list rebuild. + */ +@RunWith(RobolectricTestRunner::class) +class UserProfileDataSourceTest { + + private val dataSource = UserProfileDataSource() + private val userId = ByteArray(16) { 3 }.toList() + private val profile = UserProfile.Empty.copy(displayName = "Ada", userId = userId) + + @Before + fun setUp() { + FlipcashDatabase.init(RuntimeEnvironment.getApplication(), ENTROPY) + } + + @After + fun tearDown() { + FlipcashDatabase.closeDb() + } + + @Test + fun `a name stored in one session is in the first emission of the next`() = runTest { + dataSource.store(userId, profile) + + // Process restart: the database closes and the same account's file is reopened. + FlipcashDatabase.closeDb() + FlipcashDatabase.init(RuntimeEnvironment.getApplication(), ENTROPY) + + dataSource.observeProfiles().test { + assertEquals("Ada", awaitItem()[userId.hexEncodedString()]?.displayName) + cancelAndIgnoreRemainingEvents() + } + } + + @Test + fun `storing the same name again does not re-emit`() = runTest { + dataSource.store(userId, profile) + + dataSource.observeProfiles().test { + assertEquals("Ada", awaitItem()[userId.hexEncodedString()]?.displayName) + + dataSource.store(userId, profile) + expectNoEvents() + + dataSource.store(userId, profile.copy(displayName = "Ada L")) + assertEquals("Ada L", awaitItem()[userId.hexEncodedString()]?.displayName) + cancelAndIgnoreRemainingEvents() + } + } + + private companion object { + const val ENTROPY = "dGVzdC1lbnRyb3B5LWZvci11c2VyLXByb2ZpbGVz" + } +}