Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -50,7 +50,7 @@ internal class ChatsViewModel @Inject constructor(
fun conversations(
summaries: List<ChatSummary>,
tokens: List<Token>,
// 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<String, UserProfile>?,
): List<ConversationReference> {
Expand All @@ -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(),
)
)
)
)
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -351,6 +351,13 @@ interface MessagingOperations {
*/
fun observeSenderProfiles(): Flow<Map<String, UserProfile>>

/**
* 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<String, UserProfile>?

/**
* 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
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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 ->
Expand Down Expand Up @@ -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 {
Expand Down Expand Up @@ -289,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)
}

/**
Expand Down Expand Up @@ -325,14 +336,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<ChatMetadata>? = null
var groupChats: List<ChatMetadata>? = 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()
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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<Map<String, UserProfile>> = userProfileDataSource.observeProfiles()

@Volatile
private var snapshot: Map<String, UserProfile>? = 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<String, UserProfile>? get() = snapshot

init {
CoroutineScope(dispatchers.IO + SupervisorJob()).launch {
profiles.collect { snapshot = it }
}
}

/**
* Asks for [userId]'s profile if nothing has asked already. Returns immediately.
*/
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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)
}
Expand Down Expand Up @@ -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 ---

Expand Down
Loading
Loading