diff --git a/apps/flipcash/shared/chat/build.gradle.kts b/apps/flipcash/shared/chat/build.gradle.kts index 8a52e3379d..189081537c 100644 --- a/apps/flipcash/shared/chat/build.gradle.kts +++ b/apps/flipcash/shared/chat/build.gradle.kts @@ -24,6 +24,10 @@ dependencies { implementation(libs.androidx.paging.runtime) + implementation(libs.androidx.work) + implementation(libs.hilt.worker) + testImplementation(libs.androidx.work.testing) + implementation(project(":apps:flipcash:shared:persistence:sources")) implementation(project(":apps:flipcash:shared:persistence:db")) implementation(project(":apps:flipcash:shared:contacts")) diff --git a/apps/flipcash/shared/chat/src/main/kotlin/com/flipcash/shared/chat/RosterSearchSource.kt b/apps/flipcash/shared/chat/src/main/kotlin/com/flipcash/shared/chat/RosterSearchSource.kt new file mode 100644 index 0000000000..70134105eb --- /dev/null +++ b/apps/flipcash/shared/chat/src/main/kotlin/com/flipcash/shared/chat/RosterSearchSource.kt @@ -0,0 +1,45 @@ +package com.flipcash.shared.chat + +import com.flipcash.services.models.chat.ChatId +import com.flipcash.services.models.chat.MediaItem +import com.getcode.opencode.model.core.ID + +/** + * Finds members of a chat by what the user has typed after `@`. + * + * The seam between the mention picker and where members are searched. The only implementation + * today searches the roster held on the device; a server-side search can replace or back it up + * behind this interface without its callers changing. + */ +interface RosterSearchSource { + + /** + * Up to [limit] members of [chatId] whose display name or handle has a word starting with each + * word of [query], ignoring case and diacritics. Never the current user. + * + * Ordered: members who spoke recently in the chat, most recent first; then a member whose + * handle is exactly [query]; then everyone else by display name. An empty [query] returns only + * the recent speakers. + */ + suspend fun search(chatId: ChatId, query: String, limit: Int = DEFAULT_LIMIT): List + + /** + * Brings the members a search of [chatId] will read up to date. Called when the picker opens: + * a profile change does not move the roster version, so held names can be stale without the + * event stream saying so. + */ + suspend fun refresh(chatId: ChatId) + + companion object { + const val DEFAULT_LIMIT = 20 + } +} + +/** A member a search found. */ +data class MemberMatch( + val userId: ID, + val displayName: String, + // Bare, without the `@`. Null when the member has not claimed one. + val username: String?, + val profilePicture: MediaItem?, +) diff --git a/apps/flipcash/shared/chat/src/main/kotlin/com/flipcash/shared/chat/inject/ChatModule.kt b/apps/flipcash/shared/chat/src/main/kotlin/com/flipcash/shared/chat/inject/ChatModule.kt index 89c8682e77..3f00741e3b 100644 --- a/apps/flipcash/shared/chat/src/main/kotlin/com/flipcash/shared/chat/inject/ChatModule.kt +++ b/apps/flipcash/shared/chat/src/main/kotlin/com/flipcash/shared/chat/inject/ChatModule.kt @@ -2,10 +2,16 @@ package com.flipcash.shared.chat.inject import com.flipcash.shared.chat.ChatCoordinator import com.flipcash.shared.chat.ChatDraftStore +import com.flipcash.shared.chat.RosterSearchSource import com.flipcash.shared.chat.internal.DmOutgoingEncryption +import com.flipcash.shared.chat.internal.LocalRosterSearchSource import com.flipcash.shared.chat.internal.OutgoingEncryption import com.flipcash.shared.chat.internal.RealChatCoordinator import com.flipcash.shared.chat.internal.RealChatDraftStore +import com.flipcash.shared.chat.internal.RosterReconcileScheduler +import com.flipcash.shared.chat.internal.RosterSync +import com.flipcash.shared.chat.internal.RosterSyncTrigger +import com.flipcash.shared.chat.internal.WorkManagerRosterReconcileScheduler import com.getcode.opencode.providers.SessionListener import dagger.Binds import dagger.Module @@ -35,6 +41,21 @@ abstract class ChatModule { impl: DmOutgoingEncryption ): OutgoingEncryption + @Binds + internal abstract fun bindRosterSearchSource( + impl: LocalRosterSearchSource + ): RosterSearchSource + + @Binds + abstract fun bindRosterReconcileScheduler( + impl: WorkManagerRosterReconcileScheduler + ): RosterReconcileScheduler + + @Binds + abstract fun bindRosterSyncTrigger( + impl: RosterSync + ): RosterSyncTrigger + @Binds @IntoSet abstract fun bindSessionListener( diff --git a/apps/flipcash/shared/chat/src/main/kotlin/com/flipcash/shared/chat/internal/LocalRosterSearchSource.kt b/apps/flipcash/shared/chat/src/main/kotlin/com/flipcash/shared/chat/internal/LocalRosterSearchSource.kt new file mode 100644 index 0000000000..1ba977b3c8 --- /dev/null +++ b/apps/flipcash/shared/chat/src/main/kotlin/com/flipcash/shared/chat/internal/LocalRosterSearchSource.kt @@ -0,0 +1,105 @@ +package com.flipcash.shared.chat.internal + +import com.flipcash.app.persistence.sources.ChatRosterDataSource +import com.flipcash.app.persistence.sources.RosterSearchCandidate +import com.flipcash.app.persistence.sources.search.MemberSearchText +import com.flipcash.services.models.chat.ChatId +import com.flipcash.services.user.UserManager +import com.flipcash.shared.chat.MemberMatch +import com.flipcash.shared.chat.RosterSearchSource +import com.getcode.utils.hexEncodedString +import javax.inject.Inject + +/** + * Searches the members held on the device, through the token index [ChatRosterDataSource] reads. + * + * A query of several words narrows: every word must start one of the member's words. The index + * answers the longest word with a range scan and the rest are checked against the few members it + * returns. + */ +internal class LocalRosterSearchSource @Inject constructor( + private val rosterDataSource: ChatRosterDataSource, + private val rosterSync: RosterSync, + private val userManager: UserManager, +) : RosterSearchSource { + + override suspend fun search(chatId: ChatId, query: String, limit: Int): List { + val selfId = userManager.accountId ?: return emptyList() + val words = MemberSearchText.queryWords(query) + + val candidates = if (words.isEmpty()) { + rosterDataSource.recentSpeakers(chatId, selfId, RECENT_MESSAGE_WINDOW) + } else { + rosterDataSource.searchByPrefix(chatId, selfId, words.maxBy { it.length }, RECENT_MESSAGE_WINDOW) + .filter { candidate -> matchesEveryWord(candidate, words) } + } + + return rankRosterMatches(candidates, words) + .take(limit) + .map { MemberMatch(it.userId, it.displayName, it.username, it.profilePicture) } + } + + override suspend fun refresh(chatId: ChatId) { + rosterSync.refreshFirstPage(chatId) + } + + private fun matchesEveryWord(candidate: RosterSearchCandidate, words: List): Boolean { + if (words.size == 1) return true // the range scan already matched it + val tokens = MemberSearchText.tokens(candidate.displayName, candidate.username) + return words.all { word -> tokens.any { it.startsWith(word) } } + } + + companion object { + /** How many of a chat's newest held messages decide who counts as a recent speaker. */ + const val RECENT_MESSAGE_WINDOW = 50 + } +} + +/** + * Orders search results: + * 1. members who sent one of the chat's newest held messages, most recent first; + * 2. a member whose handle equals the whole query; + * 3. everyone else, alphabetically by display name. + * + * Names compare in their normalized form, so "Érica" sorts with the e's. Ties fall back to the raw + * name and then the user id, so an order never depends on the order rows came back in. + * + * Names compare by code point, as Swift's String ordering does, so both platforms list the same + * members in the same order. Kotlin's [String.compareTo] compares UTF-16 units instead, which puts + * an emoji (a surrogate pair, D800–DFFF) below a full-width letter (FF00–FFEF). The user id breaks + * the last tie as lowercase hex, which orders the same as iOS's lowercase hyphenated UUID string. + */ +internal fun rankRosterMatches( + candidates: List, + queryWords: List, +): List { + val wholeQuery = queryWords.joinToString(" ").takeIf { it.isNotEmpty() } + fun isExactHandle(candidate: RosterSearchCandidate): Boolean = + wholeQuery != null && candidate.username?.removePrefix("@")?.let(MemberSearchText::normalize) == wholeQuery + + val keyed = candidates.map { it to MemberSearchText.normalize(it.displayName) } + return keyed.sortedWith( + compareBy> { (candidate, _) -> candidate.lastSpokeEpochMs == null } + .thenByDescending { (candidate, _) -> candidate.lastSpokeEpochMs ?: 0L } + .thenBy { (candidate, _) -> !isExactHandle(candidate) } + .thenComparing({ (_, name) -> name }, CodePointOrder) + .thenComparing({ (candidate, _) -> candidate.displayName }, CodePointOrder) + .thenBy { (candidate, _) -> candidate.userId.hexEncodedString() } + ).map { it.first } +} + +/** Orders strings by Unicode code point, not by UTF-16 unit. */ +internal object CodePointOrder : Comparator { + override fun compare(a: String, b: String): Int { + var i = 0 + var j = 0 + while (i < a.length && j < b.length) { + val x = a.codePointAt(i) + val y = b.codePointAt(j) + if (x != y) return x.compareTo(y) + i += Character.charCount(x) + j += Character.charCount(y) + } + return (a.length - i).compareTo(b.length - j) + } +} diff --git a/apps/flipcash/shared/chat/src/main/kotlin/com/flipcash/shared/chat/internal/RosterReconcile.kt b/apps/flipcash/shared/chat/src/main/kotlin/com/flipcash/shared/chat/internal/RosterReconcile.kt new file mode 100644 index 0000000000..90fa85e04c --- /dev/null +++ b/apps/flipcash/shared/chat/src/main/kotlin/com/flipcash/shared/chat/internal/RosterReconcile.kt @@ -0,0 +1,85 @@ +package com.flipcash.shared.chat.internal + +import android.content.Context +import androidx.hilt.work.HiltWorker +import androidx.work.BackoffPolicy +import androidx.work.Constraints +import androidx.work.CoroutineWorker +import androidx.work.ExistingWorkPolicy +import androidx.work.NetworkType +import androidx.work.OneTimeWorkRequestBuilder +import androidx.work.WorkManager +import androidx.work.WorkerParameters +import androidx.work.workDataOf +import com.flipcash.services.models.chat.ChatId +import com.getcode.utils.TraceType +import com.getcode.utils.trace +import dagger.assisted.Assisted +import dagger.assisted.AssistedInject +import dagger.hilt.android.qualifiers.ApplicationContext +import java.util.concurrent.TimeUnit +import javax.inject.Inject +import javax.inject.Singleton + +/** Queues a full read of a group's roster to run when it can. */ +interface RosterReconcileScheduler { + /** At most one read per chat is queued; asking again while one is waits on that one. */ + fun schedule(chatId: ChatId) +} + +/** + * Runs [RosterReconcileWorker] as unique work per chat, on any network, backing off on failure. + * Not expedited: a roster read is never what the user is waiting on. + */ +@Singleton +class WorkManagerRosterReconcileScheduler @Inject constructor( + @ApplicationContext private val context: Context, +) : RosterReconcileScheduler { + + override fun schedule(chatId: ChatId) { + val request = OneTimeWorkRequestBuilder() + .setInputData(workDataOf(RosterReconcileWorker.KEY_CHAT_ID to chatId.hex)) + .setConstraints(Constraints.Builder().setRequiredNetworkType(NetworkType.CONNECTED).build()) + .setBackoffCriteria(BackoffPolicy.EXPONENTIAL, BACKOFF_SECONDS, TimeUnit.SECONDS) + .build() + WorkManager.getInstance(context) + .enqueueUniqueWork(uniqueName(chatId), ExistingWorkPolicy.KEEP, request) + } + + companion object { + private const val BACKOFF_SECONDS = 30L + + fun uniqueName(chatId: ChatId) = "roster-reconcile-${chatId.hex}" + } +} + +/** Reads a group's whole roster and drops who has left: [RosterSync.reconcileNow]. */ +@HiltWorker +internal class RosterReconcileWorker @AssistedInject constructor( + @Assisted appContext: Context, + @Assisted private val params: WorkerParameters, + private val rosterSync: RosterSync, +) : CoroutineWorker(appContext, params) { + + override suspend fun doWork(): Result { + val chatId = params.inputData.getString(KEY_CHAT_ID)?.let(::ChatId) + ?: return Result.failure() + if (rosterSync.reconcileNow(chatId)) return Result.success() + if (params.runAttemptCount + 1 >= MAX_ATTEMPTS) { + // The pending flag stays set, so the chat's next open queues it again. + trace(tag = TAG, message = "Roster reconcile for $chatId gave up", type = TraceType.Error) + return Result.failure() + } + return Result.retry() + } + + companion object { + private const val TAG = "RosterReconcileWorker" + const val KEY_CHAT_ID = "chat_id" + private const val MAX_ATTEMPTS = 5 + } +} + +@OptIn(ExperimentalStdlibApi::class) +private val ChatId.hex: String + get() = bytes.toHexString() diff --git a/apps/flipcash/shared/chat/src/main/kotlin/com/flipcash/shared/chat/internal/RosterStateHolder.kt b/apps/flipcash/shared/chat/src/main/kotlin/com/flipcash/shared/chat/internal/RosterStateHolder.kt index 8cd27dbf6e..15edfa8567 100644 --- a/apps/flipcash/shared/chat/src/main/kotlin/com/flipcash/shared/chat/internal/RosterStateHolder.kt +++ b/apps/flipcash/shared/chat/src/main/kotlin/com/flipcash/shared/chat/internal/RosterStateHolder.kt @@ -2,6 +2,7 @@ package com.flipcash.shared.chat.internal import com.flipcash.app.persistence.sources.ChatMemberDataSource import com.flipcash.app.persistence.sources.ChatMetadataDataSource +import com.flipcash.app.persistence.sources.ChatRosterDataSource import com.flipcash.services.controllers.ChatController import com.flipcash.services.models.chat.ChatId import com.flipcash.services.models.chat.RosterChange @@ -33,6 +34,8 @@ class RosterStateHolder @Inject constructor( private val chatController: ChatController, private val metadataDataSource: ChatMetadataDataSource, private val memberDataSource: ChatMemberDataSource, + private val rosterDataSource: ChatRosterDataSource, + private val rosterSync: RosterSyncTrigger, ) { /** @@ -50,21 +53,25 @@ class RosterStateHolder @Inject constructor( when { incoming <= stored -> continue incoming > stored + 1 -> refetch(chatId, stored, incoming) - else -> applyChange(chatId, change) + else -> applyChange(chatId, change, stored) } } } - private suspend fun applyChange(chatId: ChatId, change: RosterChange) { + private suspend fun applyChange(chatId: ChatId, change: RosterChange, stored: Long) { when (change) { is RosterChange.MemberJoined -> memberDataSource.upsert(chatId, listOf(change.member)) - is RosterChange.MemberLeft -> memberDataSource.deleteMember(chatId, change.userId) + is RosterChange.MemberLeft -> + memberDataSource.markLeft(chatId, change.userId, version = change.rosterSummary.version) } metadataDataSource.updateRoster( chatId = chatId, memberCount = change.rosterSummary.memberCount, rosterVersion = change.rosterSummary.version, ) + // In sequence on a roster held whole keeps it whole. Behind, the watermark stays put and + // the next catch-up reads down to it. + rosterDataSource.advanceWatermark(chatId, from = stored, to = change.rosterSummary.version) } /** @@ -75,6 +82,9 @@ class RosterStateHolder @Inject constructor( * members the device legitimately holds — including the senders a transcript needs to name. * Departures come through [RosterChange.MemberLeft]; this is only here to get the count and * the version back in step with the server. + * + * The skipped changes may include joins and leaves the merge cannot see, so [RosterSync] also + * catches the roster up from its watermark. */ private suspend fun refetch(chatId: ChatId, stored: Long, incoming: Long) { trace( @@ -82,6 +92,7 @@ class RosterStateHolder @Inject constructor( message = "Roster version gap on $chatId: stored $stored, incoming $incoming", type = TraceType.Silent, ) + rosterSync.onRosterGap(chatId) val metadata = chatController.getChat(chatId).getOrElse { // Leaving the stored version alone is what makes this retryable: the next change on // this chat still reads as a gap, and tries again. diff --git a/apps/flipcash/shared/chat/src/main/kotlin/com/flipcash/shared/chat/internal/RosterSync.kt b/apps/flipcash/shared/chat/src/main/kotlin/com/flipcash/shared/chat/internal/RosterSync.kt new file mode 100644 index 0000000000..004d39296b --- /dev/null +++ b/apps/flipcash/shared/chat/src/main/kotlin/com/flipcash/shared/chat/internal/RosterSync.kt @@ -0,0 +1,289 @@ +package com.flipcash.shared.chat.internal + +import com.flipcash.app.persistence.sources.ChatMemberDataSource +import com.flipcash.app.persistence.sources.ChatMetadataDataSource +import com.flipcash.app.persistence.sources.ChatRosterDataSource +import com.flipcash.app.persistence.sources.RosterSyncState +import com.flipcash.libs.coroutines.DispatcherProvider +import com.flipcash.services.controllers.ChatController +import com.flipcash.services.models.PagingToken +import com.flipcash.services.models.QueryOptions +import com.flipcash.services.models.chat.ChatId +import com.flipcash.services.models.chat.ChatType +import com.flipcash.services.models.chat.RosterSummary +import com.getcode.opencode.model.core.ID +import com.getcode.utils.TraceType +import com.getcode.utils.trace +import kotlinx.coroutines.CoroutineScope +import kotlinx.coroutines.SupervisorJob +import kotlinx.coroutines.launch +import kotlinx.coroutines.sync.Mutex +import kotlinx.coroutines.sync.withLock +import java.util.concurrent.ConcurrentHashMap +import javax.inject.Inject +import javax.inject.Singleton + +/** + * Keeps the device's copy of a group's roster whole, so the mention picker can search every + * member rather than the ones the feed happened to carry. + * + * `GetRoster` pages at most [PAGE_SIZE] members, most recently joined first, and each member + * carries the roster version they joined at. There is no delta against a roster version, so + * repairs aim for the smallest read instead: + * + * - **Joins missed** since the watermark come back from the top of the roster: [catchUp] reads + * down to the first member at or below the watermark, usually one page. + * - **Leaves missed** are invisible from the top. A catch-up notices them only as more members + * held than `member_count`, and hands the chat to a [fullSync] in WorkManager, which is the one + * read that can say who is gone. + * + * Runs on a group's open, not on launch: reading every group's roster up front would cost a page + * per hundred members of each whether or not the user ever mentions anyone there. Versions are + * only ever compared, never subtracted: the proto calls them opaque. + */ +@Singleton +class RosterSync @Inject constructor( + private val chatController: ChatController, + private val metadataDataSource: ChatMetadataDataSource, + private val memberDataSource: ChatMemberDataSource, + private val rosterDataSource: ChatRosterDataSource, + private val reconcileScheduler: RosterReconcileScheduler, + dispatchers: DispatcherProvider, +) : RosterSyncTrigger { + + // Its own scope, like [SenderResolver]'s: a read outlives the screen that asked for it. + private val scope = CoroutineScope(dispatchers.IO + SupervisorJob()) + + // One read per chat at a time. A second open of the same chat waits, then finds nothing to do. + private val locks = ConcurrentHashMap() + + private suspend fun locked(chatId: ChatId, block: suspend () -> T): T = + locks.getOrPut(chatId) { Mutex() }.withLock { block() } + + override fun onChatOpened(chatId: ChatId) { + scope.launch { onOpen(chatId) } + } + + override fun onRosterGap(chatId: ChatId) { + scope.launch { onGap(chatId) } + } + + /** + * What a group's open does: the first full read if it never had one, the pending reconcile if + * one is owed, otherwise a catch-up if the roster has moved past the watermark. + */ + internal suspend fun onOpen(chatId: ChatId) { + if (metadataDataSource.getChatType(chatId) != ChatType.GROUP) return + locked(chatId) { + val state = rosterDataSource.getSyncState(chatId) + when { + state == null || !state.fullySynced -> firstSync(chatId) + // KEEP makes this a no-op while the work is queued, and re-queues it if it gave up. + state.reconcilePending -> reconcileScheduler.schedule(chatId) + metadataDataSource.getRosterVersion(chatId) > state.watermark -> catchUp(chatId, state) + } + } + } + + /** + * A version skipped on the stream. Only a chat already read in full is caught up here; one + * never read waits for its open, so a gap does not start a read of a group nobody opened. + */ + internal suspend fun onGap(chatId: ChatId) = locked(chatId) { + val state = rosterDataSource.getSyncState(chatId) + if (state == null || !state.fullySynced || state.reconcilePending) return@locked + catchUp(chatId, state) + } + + /** The full read WorkManager runs. False when it should be retried. */ + internal suspend fun reconcileNow(chatId: ChatId): Boolean = locked(chatId) { fullSync(chatId) } + + // A roster of one page is a single request, so it is read on the spot. Anything larger goes + // to WorkManager, where the walk survives the user leaving the chat or the process dying. + private suspend fun firstSync(chatId: ChatId) { + if (metadataDataSource.getMemberCount(chatId) > PAGE_SIZE) { + reconcileScheduler.schedule(chatId) + } else { + fullSync(chatId) + } + } + + /** + * Recovers the joins missed since [state]'s watermark by reading from the top of the roster, + * then checks the count for leaves the top cannot show. + */ + internal suspend fun catchUp(chatId: ChatId, state: RosterSyncState) { + val watermark = state.watermark + var token: PagingToken? = null + var pages = 0 + var summary: RosterSummary? = null + var readVersion = Long.MAX_VALUE + var recovered = 0 + var reachedStored = false + + while (true) { + val page = chatController.getRoster(chatId, QueryOptions(limit = PAGE_SIZE, token = token)) + .getOrElse { + trace(tag = TAG, message = "Roster catch-up failed for $chatId", type = TraceType.Error) + return + } + pages++ + if (summary == null) summary = page.rosterSummary + readVersion = minOf(readVersion, page.rosterSummary.version) + + // The stopping rule. Pages run most recently joined first, and Member.version is the + // version a member joined at, so the first member at or below the watermark is where + // the stored copy begins and everyone after them is already held. + // REVISIT: model.proto says Member.version will move on later member changes, such as + // a role change. Once it does, version order stops matching page order and this rule + // no longer holds: a long-standing member with a bumped version reads as a new join. + val newer = page.members.takeWhile { it.version > watermark } + memberDataSource.upsert(chatId, newer) + recovered += newer.size + + if (newer.size < page.members.size || !page.hasMore) { + reachedStored = true + break + } + val next = page.pagingToken + if (pages >= MAX_PAGES || next == null) break + token = next + } + + if (!reachedStored) { + // More joins missed than the cap reads: only a full read catches up now. + trace(tag = TAG, message = "Roster catch-up for $chatId ran past the cap", type = TraceType.Silent) + rosterDataSource.markReconcilePending(chatId) + reconcileScheduler.schedule(chatId) + return + } + + val held = memberDataSource.countMembers(chatId).toLong() + val memberCount = summary!!.memberCount + when { + held == memberCount -> { + // The page's version, not the stream's: a large group's pages can trail the stream, + // and the watermark may only claim what the pages showed. + rosterDataSource.setWatermark(chatId, readVersion) + } + held > memberCount -> { + // Someone left and the top of the roster cannot say who. Search keeps them until + // the full read settles it. + rosterDataSource.markReconcilePending(chatId) + reconcileScheduler.schedule(chatId) + } + // A roster cut off at the cap is always short, so short is no sign of a missed join. + state.truncated -> rosterDataSource.setWatermark(chatId, readVersion) + // Otherwise a join is still to come: the page trails the stream, and the join arrives + // as a roster update. (A chat never read in full does not get here: its open reads it.) + else -> Unit + } + + trace( + tag = TAG, + message = "Roster catch-up for $chatId: $recovered joins over $pages pages; " + + "holding $held of $memberCount", + type = TraceType.Silent, + ) + } + + /** + * Reads [chatId]'s whole roster, up to [MAX_PAGES], and drops the members it shows have left. + * False if a page failed or the database is not open, so nothing was recorded. + * + * The reconcile rule is the proto's: a held member absent from the read has left, unless they + * joined after the version the read described. With pages at different versions the lowest is + * used, since each page is only a promise about the roster as of its own version. + */ + internal suspend fun fullSync(chatId: ChatId): Boolean { + if (!rosterDataSource.isAvailable) return false + + var token: PagingToken? = null + var pages = 0 + var truncated = false + val seen = HashSet() + var readSummary: RosterSummary? = null + + while (true) { + val page = chatController.getRoster(chatId, QueryOptions(limit = PAGE_SIZE, token = token)) + .getOrElse { + trace(tag = TAG, message = "Roster page ${pages + 1} failed for $chatId", type = TraceType.Error) + return false + } + pages++ + memberDataSource.upsert(chatId, page.members) + page.members.mapTo(seen) { it.userId } + if (readSummary == null || page.rosterSummary.version < readSummary.version) { + readSummary = page.rosterSummary + } + + if (!page.hasMore) break + val next = page.pagingToken + if (pages >= MAX_PAGES || next == null) { + truncated = true + break + } + token = next + } + + val summary = readSummary!! + // A partial read cannot tell who left, so it drops no one. + if (!truncated) memberDataSource.reconcile(chatId, seen, readVersion = summary.version) + // No-op unless the read is ahead of what the stream has applied. + metadataDataSource.updateRoster(chatId, memberCount = summary.memberCount, rosterVersion = summary.version) + rosterDataSource.markFullySynced(chatId, watermark = summary.version, truncated = truncated) + + trace( + tag = TAG, + message = "Roster read for $chatId: ${seen.size} members over $pages pages at v${summary.version}" + + if (truncated) ", stopped at the cap" else "", + type = TraceType.Silent, + ) + return true + } + + /** + * Rewrites the first page of [chatId]'s roster: the members most recently joined, with their + * current profiles. Profile changes do not move the roster version, so this is what keeps held + * names fresh for a search that is about to start. + */ + suspend fun refreshFirstPage(chatId: ChatId) { + val page = chatController.getRoster(chatId, QueryOptions(limit = PAGE_SIZE)).getOrElse { + trace(tag = TAG, message = "Roster refresh failed for $chatId", type = TraceType.Error) + return + } + memberDataSource.upsert(chatId, page.members) + } + + companion object { + private const val TAG = "RosterSync" + + /** The most `GetRoster` returns in one page. */ + const val PAGE_SIZE = 100 + + /** Where a read stops: 2,000 members. A search covers whatever was read by then. */ + const val MAX_PAGES = 20 + } +} + +/** + * Told when a chat opens or its roster stream skips a version, so the roster can be repaired off + * the caller's path. + * + * A seam rather than [RosterSync] itself so the classes that report these can be built in tests + * without one. + */ +interface RosterSyncTrigger { + + /** Returns at once; any read it starts runs in the background. */ + fun onChatOpened(chatId: ChatId) + + /** Returns at once; any read it starts runs in the background. */ + fun onRosterGap(chatId: ChatId) + + /** Does nothing. What a class built without a roster sync -- a unit test -- is given. */ + object None : RosterSyncTrigger { + override fun onChatOpened(chatId: ChatId) = Unit + override fun onRosterGap(chatId: ChatId) = Unit + } +} 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..fdbb386279 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 @@ -43,6 +43,7 @@ import com.flipcash.shared.chat.PendingMutation import com.flipcash.shared.chat.UnreadBoundary import com.flipcash.shared.chat.internal.ChatStateHolder import com.flipcash.shared.chat.internal.OutgoingEncryption +import com.flipcash.shared.chat.internal.RosterSyncTrigger import com.flipcash.shared.chat.replacingText import com.flipcash.services.user.UserManager import com.flipcash.shared.chat.MessageLinkPrefetch @@ -96,6 +97,8 @@ class MessagingDelegate @Inject constructor( private val outgoing: OutgoingEncryption = OutgoingEncryption.None, /** Opens encrypted pushes; without it, every push keeps the server's body. */ private val incoming: IncomingMessageOpener? = null, + /** Reads a group's whole roster when it opens; without it, only what the feed carries is held. */ + private val rosterSync: RosterSyncTrigger = RosterSyncTrigger.None, ) : MessagingOperations { /** @@ -124,6 +127,7 @@ class MessagingDelegate @Inject constructor( override fun setActiveChatId(chatId: ChatId?) { stateHolder.update { it.copy(activeChat = chatId) } + if (chatId != null) rosterSync.onChatOpened(chatId) } override fun clearActiveChat(chatId: ChatId?) { diff --git a/apps/flipcash/shared/chat/src/test/kotlin/com/flipcash/shared/chat/LocalRosterSearchSourceTest.kt b/apps/flipcash/shared/chat/src/test/kotlin/com/flipcash/shared/chat/LocalRosterSearchSourceTest.kt new file mode 100644 index 0000000000..bf836773b8 --- /dev/null +++ b/apps/flipcash/shared/chat/src/test/kotlin/com/flipcash/shared/chat/LocalRosterSearchSourceTest.kt @@ -0,0 +1,72 @@ +package com.flipcash.shared.chat + +import com.flipcash.app.persistence.sources.ChatRosterDataSource +import com.flipcash.app.persistence.sources.RosterSearchCandidate +import com.flipcash.services.models.chat.ChatId +import com.flipcash.services.user.UserManager +import com.flipcash.shared.chat.internal.LocalRosterSearchSource +import com.flipcash.shared.chat.internal.RosterSync +import io.mockk.coEvery +import io.mockk.coVerify +import io.mockk.every +import io.mockk.mockk +import kotlinx.coroutines.test.runTest +import org.junit.Test +import kotlin.test.assertEquals + +class LocalRosterSearchSourceTest { + + private val chatId = ChatId("c0ffee") + private val self = listOf(1) + + private val rosterDataSource = mockk() + private val rosterSync = mockk(relaxed = true) + private val userManager = mockk { every { accountId } returns self } + + private val subject = LocalRosterSearchSource(rosterDataSource, rosterSync, userManager) + + private fun candidate(id: Int, name: String, username: String? = null, spokeAt: Long? = null) = + RosterSearchCandidate(listOf(id.toByte()), name, username, null, spokeAt) + + @Test + fun `an empty query returns the recent speakers, most recent first`() = runTest { + coEvery { rosterDataSource.recentSpeakers(chatId, self, any()) } returns + listOf(candidate(2, "Old", spokeAt = 1), candidate(3, "New", spokeAt = 2)) + + assertEquals(listOf("New", "Old"), subject.search(chatId, "@").map { it.displayName }) + coVerify(exactly = 0) { rosterDataSource.searchByPrefix(any(), any(), any(), any()) } + } + + @Test + fun `the index is asked for the query folded, by its longest word`() = runTest { + coEvery { rosterDataSource.searchByPrefix(chatId, self, "garc", any()) } returns + listOf(candidate(2, "María García"), candidate(3, "Luis García")) + + val names = subject.search(chatId, "@Ma Garc").map { it.displayName } + + // Every word must start one of the member's words. + assertEquals(listOf("María García"), names) + } + + @Test + fun `results are capped at the limit`() = runTest { + coEvery { rosterDataSource.searchByPrefix(chatId, self, "a", any()) } returns + (10..30).map { candidate(it, "A$it") } + + assertEquals(5, subject.search(chatId, "a", limit = 5).size) + } + + @Test + fun `nobody is signed in, nothing is searched`() = runTest { + every { userManager.accountId } returns null + + assertEquals(emptyList(), subject.search(chatId, "a")) + } + + @Test + fun `refresh rereads the first roster page`() = runTest { + subject.refresh(chatId) + + coVerify { rosterSync.refreshFirstPage(chatId) } + } +} diff --git a/apps/flipcash/shared/chat/src/test/kotlin/com/flipcash/shared/chat/RosterReconcileSchedulerTest.kt b/apps/flipcash/shared/chat/src/test/kotlin/com/flipcash/shared/chat/RosterReconcileSchedulerTest.kt new file mode 100644 index 0000000000..1d7ad1fa23 --- /dev/null +++ b/apps/flipcash/shared/chat/src/test/kotlin/com/flipcash/shared/chat/RosterReconcileSchedulerTest.kt @@ -0,0 +1,59 @@ +package com.flipcash.shared.chat + +import android.content.Context +import androidx.test.core.app.ApplicationProvider +import androidx.work.Configuration +import androidx.work.NetworkType +import androidx.work.WorkInfo +import androidx.work.WorkManager +import androidx.work.testing.SynchronousExecutor +import androidx.work.testing.WorkManagerTestInitHelper +import com.flipcash.services.models.chat.ChatId +import com.flipcash.shared.chat.internal.WorkManagerRosterReconcileScheduler +import org.junit.Before +import org.junit.Test +import org.junit.runner.RunWith +import org.robolectric.RobolectricTestRunner +import kotlin.test.assertEquals + +/** A chat has at most one full roster read queued, and it waits for a network. */ +@RunWith(RobolectricTestRunner::class) +class RosterReconcileSchedulerTest { + + private val context = ApplicationProvider.getApplicationContext() + private val subject = WorkManagerRosterReconcileScheduler(context) + + @Before + fun setUp() { + WorkManagerTestInitHelper.initializeTestWorkManager( + context, + Configuration.Builder().setExecutor(SynchronousExecutor()).build(), + ) + } + + private fun queued(chatId: ChatId): List = + WorkManager.getInstance(context) + .getWorkInfosForUniqueWork(WorkManagerRosterReconcileScheduler.uniqueName(chatId)) + .get() + + @Test + fun `scheduling twice queues one read`() { + val chatId = ChatId("c0ffee") + + subject.schedule(chatId) + subject.schedule(chatId) + + val work = queued(chatId).single() + assertEquals(WorkInfo.State.ENQUEUED, work.state) + assertEquals(NetworkType.CONNECTED, work.constraints.requiredNetworkType) + } + + @Test + fun `each chat gets its own read`() { + subject.schedule(ChatId("c0ffee")) + subject.schedule(ChatId("decade")) + + assertEquals(1, queued(ChatId("c0ffee")).size) + assertEquals(1, queued(ChatId("decade")).size) + } +} diff --git a/apps/flipcash/shared/chat/src/test/kotlin/com/flipcash/shared/chat/RosterSearchRankingTest.kt b/apps/flipcash/shared/chat/src/test/kotlin/com/flipcash/shared/chat/RosterSearchRankingTest.kt new file mode 100644 index 0000000000..659ae4a5ef --- /dev/null +++ b/apps/flipcash/shared/chat/src/test/kotlin/com/flipcash/shared/chat/RosterSearchRankingTest.kt @@ -0,0 +1,93 @@ +package com.flipcash.shared.chat + +import com.flipcash.app.persistence.sources.RosterSearchCandidate +import com.flipcash.app.persistence.sources.search.MemberSearchText +import com.flipcash.shared.chat.internal.CodePointOrder +import com.flipcash.shared.chat.internal.rankRosterMatches +import org.junit.Test +import kotlin.test.assertEquals +import kotlin.test.assertTrue + +class RosterSearchRankingTest { + + private fun candidate(id: Int, name: String, username: String? = null, spokeAt: Long? = null) = + RosterSearchCandidate( + userId = listOf(id.toByte()), + displayName = name, + username = username, + profilePicture = null, + lastSpokeEpochMs = spokeAt, + ) + + private fun rank(query: String, vararg candidates: RosterSearchCandidate) = + rankRosterMatches(candidates.toList(), MemberSearchText.queryWords(query)).map { it.displayName } + + @Test + fun `recent speakers first, then an exact handle, then by name`() { + assertEquals( + listOf("Sam Recent", "Sam Earlier", "Zed", "Adam", "Sammy"), + rank( + "sam", + candidate(1, "Sammy"), + candidate(2, "Adam", username = "samuel"), + candidate(3, "Zed", username = "sam"), + candidate(4, "Sam Earlier", spokeAt = 10), + candidate(5, "Sam Recent", spokeAt = 20), + ), + ) + } + + @Test + fun `a recent speaker outranks an exact handle`() { + assertEquals( + listOf("Talker", "Exact"), + rank("sam", candidate(1, "Exact", username = "sam"), candidate(2, "Talker", username = "samwise", spokeAt = 5)), + ) + } + + @Test + fun `the exact handle ignores case, diacritics and the mention trigger`() { + assertEquals( + listOf("Érica", "Alan"), + rank("@ERICA", candidate(1, "Alan", username = "ericaz"), candidate(2, "Érica", username = "érica")), + ) + } + + @Test + fun `names sort by code point, so an emoji sorts above U+FFxx`() { + // U+FFFD has no decomposition, so it survives folding. UTF-16 order would put the emoji's + // high surrogate (U+D83D) first; code point order, like Swift's, puts U+1F600 last. + assertEquals( + listOf("\uFFFD Box", "\uD83D\uDE00 Smile"), + rank("", candidate(1, "\uD83D\uDE00 Smile"), candidate(2, "\uFFFD Box")), + ) + } + + @Test + fun `the raw-name tiebreak compares code points too`() { + // A full-width letter (U+FF21) against an emoji, as the raw names reach the tiebreak. + assertTrue(CodePointOrder.compare("\uFF21", "\uD83D\uDE00") < 0) + assertTrue("\uFF21" > "\uD83D\uDE00") // what String.compareTo would have said + assertEquals(0, CodePointOrder.compare("abc", "abc")) + assertTrue(CodePointOrder.compare("ab", "abc") < 0) + } + + @Test + fun `the last tie breaks on the user id as lowercase hex`() { + // Signed bytes would put 0xAB (-85) before 0x0C (12); hex puts "0c" first. + val low = candidate(0x0C, "Same") + val high = candidate(0xAB, "Same") + assertEquals( + listOf(low.userId, high.userId), + rankRosterMatches(listOf(high, low), emptyList()).map { it.userId }, + ) + } + + @Test + fun `names sort with diacritics folded`() { + assertEquals( + listOf("Ana", "Érica", "Fran"), + rank("", candidate(1, "Fran"), candidate(2, "Érica"), candidate(3, "Ana")), + ) + } +} diff --git a/apps/flipcash/shared/chat/src/test/kotlin/com/flipcash/shared/chat/RosterStateHolderTest.kt b/apps/flipcash/shared/chat/src/test/kotlin/com/flipcash/shared/chat/RosterStateHolderTest.kt index 3cbc3cfd9c..19c88c099a 100644 --- a/apps/flipcash/shared/chat/src/test/kotlin/com/flipcash/shared/chat/RosterStateHolderTest.kt +++ b/apps/flipcash/shared/chat/src/test/kotlin/com/flipcash/shared/chat/RosterStateHolderTest.kt @@ -2,6 +2,8 @@ package com.flipcash.shared.chat import com.flipcash.app.persistence.sources.ChatMemberDataSource import com.flipcash.app.persistence.sources.ChatMetadataDataSource +import com.flipcash.app.persistence.sources.ChatRosterDataSource +import com.flipcash.shared.chat.internal.RosterSyncTrigger import com.flipcash.services.controllers.ChatController import com.flipcash.services.models.UserProfile import com.flipcash.services.models.chat.ChatId @@ -14,6 +16,7 @@ import com.flipcash.shared.chat.internal.RosterStateHolder import io.mockk.coEvery import io.mockk.coVerify import io.mockk.mockk +import io.mockk.verify import kotlinx.coroutines.test.runTest import org.junit.Test import kotlin.time.Instant @@ -32,11 +35,15 @@ class RosterStateHolderTest { private val controller = mockk(relaxed = true) private val metadataDataSource = mockk(relaxed = true) private val memberDataSource = mockk(relaxed = true) + private val rosterDataSource = mockk(relaxed = true) + private val rosterSync = mockk(relaxed = true) private val subject = RosterStateHolder( chatController = controller, metadataDataSource = metadataDataSource, memberDataSource = memberDataSource, + rosterDataSource = rosterDataSource, + rosterSync = rosterSync, ) private fun joiner(userId: List = joinerId) = ChatMember( @@ -71,12 +78,12 @@ class RosterStateHolderTest { } @Test - fun `a leave one version ahead removes the member`() = runTest { + fun `a leave one version ahead marks the member left at that version`() = runTest { storedVersion(4) subject.apply(chatId, listOf(left(version = 5))) - coVerify { memberDataSource.deleteMember(chatId, leaverId) } + coVerify { memberDataSource.markLeft(chatId, leaverId, version = 5) } coVerify { metadataDataSource.updateRoster(chatId, memberCount = 11, rosterVersion = 5) } } @@ -120,6 +127,18 @@ class RosterStateHolderTest { coVerify { memberDataSource.upsert(chatId, refetched.members) } // The skipped change is not applied on top of what the refetch returned. coVerify(exactly = 0) { metadataDataSource.updateRoster(chatId, 13, 7) } + // The skipped versions may hold joins and leaves the refetch cannot show, so it catches up. + verify { rosterSync.onRosterGap(chatId) } + } + + @Test + fun `a change in sequence moves the watermark along from where it was`() = runTest { + storedVersion(4) + + subject.apply(chatId, listOf(joined(version = 5))) + + coVerify { rosterDataSource.advanceWatermark(chatId, from = 4, to = 5) } + verify(exactly = 0) { rosterSync.onRosterGap(any()) } } @Test diff --git a/apps/flipcash/shared/chat/src/test/kotlin/com/flipcash/shared/chat/RosterSyncTest.kt b/apps/flipcash/shared/chat/src/test/kotlin/com/flipcash/shared/chat/RosterSyncTest.kt new file mode 100644 index 0000000000..ddc1ec37e3 --- /dev/null +++ b/apps/flipcash/shared/chat/src/test/kotlin/com/flipcash/shared/chat/RosterSyncTest.kt @@ -0,0 +1,340 @@ +package com.flipcash.shared.chat + +import com.flipcash.app.core.dispatchers.TestDispatchers +import com.flipcash.app.persistence.sources.ChatMemberDataSource +import com.flipcash.app.persistence.sources.ChatMetadataDataSource +import com.flipcash.app.persistence.sources.ChatRosterDataSource +import com.flipcash.app.persistence.sources.RosterSyncState +import com.flipcash.services.controllers.ChatController +import com.flipcash.services.models.QueryOptions +import com.flipcash.services.models.UserProfile +import com.flipcash.services.models.chat.ChatId +import com.flipcash.services.models.chat.ChatMember +import com.flipcash.services.models.chat.ChatType +import com.flipcash.services.models.chat.RosterPage +import com.flipcash.services.models.chat.RosterSummary +import com.flipcash.shared.chat.internal.RosterReconcileScheduler +import com.flipcash.shared.chat.internal.RosterSync +import io.mockk.coEvery +import io.mockk.coVerify +import io.mockk.every +import io.mockk.mockk +import io.mockk.slot +import io.mockk.verify +import kotlinx.coroutines.test.TestCoroutineScheduler +import kotlinx.coroutines.test.runTest +import org.junit.Test +import kotlin.test.assertEquals +import kotlin.test.assertFalse +import kotlin.test.assertTrue + +/** + * A group's roster is read in full once, then kept whole by catching up from the top of the roster + * to the watermark. Leaves the top cannot show go to a full read in WorkManager. + */ +class RosterSyncTest { + + private val chatId = ChatId("c0ffee") + + private val controller = mockk() + private val metadataDataSource = mockk(relaxed = true) + private val memberDataSource = mockk(relaxed = true) + private val rosterDataSource = mockk(relaxed = true) { + every { isAvailable } returns true + } + private val scheduler = mockk(relaxed = true) + + private val subject = RosterSync( + chatController = controller, + metadataDataSource = metadataDataSource, + memberDataSource = memberDataSource, + rosterDataSource = rosterDataSource, + reconcileScheduler = scheduler, + dispatchers = TestDispatchers(TestCoroutineScheduler()), + ) + + private val requests = mutableListOf() + + /** Member [n], who joined at roster version [n]. */ + private fun member(n: Int) = ChatMember( + userId = listOf(n.toByte(), (n shr 8).toByte()), + userProfile = UserProfile.Empty.copy(displayName = "m$n"), + pointers = emptyList(), + version = n.toLong(), + ) + + /** + * Serves a roster of members joined at versions [total] down to 1, most recent first, [perPage] + * to a page. Page i's token is `[i]`; each page reports [versionOf] of its index. + */ + private fun serve( + total: Int, + perPage: Int = 2, + versionOf: (Int) -> Long = { total.toLong() }, + memberCount: Long = total.toLong(), + ) { + val options = slot() + coEvery { controller.getRoster(chatId, capture(options)) } answers { + requests += options.captured + val index = options.captured.token?.single()?.toInt() ?: 0 + val from = total - index * perPage + val members = (from downTo maxOf(1, from - perPage + 1)).map(::member) + val hasMore = from - perPage > 0 + Result.success( + RosterPage( + members = members, + rosterSummary = RosterSummary(memberCount = memberCount, version = versionOf(index)), + pagingToken = if (hasMore) listOf((index + 1).toByte()) else null, + hasMore = hasMore, + ) + ) + } + } + + private fun state( + watermark: Long = 0, + fullySynced: Boolean = true, + truncated: Boolean = false, + reconcilePending: Boolean = false, + ) = RosterSyncState(watermark, fullySynced, truncated, reconcilePending) + + // Full read + + @Test + fun `a full read runs until has_more is false`() = runTest { + serve(total = 6) + + assertTrue(subject.fullSync(chatId)) + + assertEquals(3, requests.size) + assertTrue(requests.all { it.limit == RosterSync.PAGE_SIZE }) + assertEquals(listOf(null, listOf(1), listOf(2)), requests.map { it.token }) + coVerify { memberDataSource.reconcile(chatId, (1..6).map { member(it).userId }.toSet(), readVersion = 6) } + coVerify { rosterDataSource.markFullySynced(chatId, watermark = 6, truncated = false) } + } + + @Test + fun `a full read stops at the page cap and drops no one`() = runTest { + serve(total = (RosterSync.MAX_PAGES + 5) * 2) + + assertTrue(subject.fullSync(chatId)) + + assertEquals(RosterSync.MAX_PAGES, requests.size) + coVerify(exactly = 0) { memberDataSource.reconcile(any(), any(), any()) } + coVerify { rosterDataSource.markFullySynced(chatId, watermark = any(), truncated = true) } + } + + @Test + fun `a full read reconciles against its lowest page version`() = runTest { + // Later pages trail: each is only a promise about the roster as of its own version. + serve(total = 4, versionOf = { 9L - it }) + + subject.fullSync(chatId) + + coVerify { memberDataSource.reconcile(chatId, any(), readVersion = 8) } + coVerify { rosterDataSource.markFullySynced(chatId, watermark = 8, truncated = false) } + } + + @Test + fun `a failed page records nothing and asks to be retried`() = runTest { + serve(total = 6) + coEvery { controller.getRoster(chatId, match { it.token == listOf(1) }) } returns + Result.failure(Throwable("offline")) + + assertFalse(subject.fullSync(chatId)) + + coVerify(exactly = 0) { rosterDataSource.markFullySynced(any(), any(), any()) } + coVerify(exactly = 0) { memberDataSource.reconcile(any(), any(), any()) } + } + + @Test + fun `a full read with no database open asks to be retried`() = runTest { + every { rosterDataSource.isAvailable } returns false + + assertFalse(subject.fullSync(chatId)) + + coVerify(exactly = 0) { controller.getRoster(any(), any()) } + } + + // Catch-up + + @Test + fun `missed joins come back in one page and the read stops at the watermark`() = runTest { + // Held through version 10; members 13, 12 and 11 joined while the stream was away. + serve(total = 13, perPage = RosterSync.PAGE_SIZE) + coEvery { memberDataSource.countMembers(chatId) } returns 13 + + subject.catchUp(chatId, state(watermark = 10)) + + assertEquals(1, requests.size) + coVerify { memberDataSource.upsert(chatId, listOf(member(13), member(12), member(11))) } + } + + @Test + fun `a catch-up reads further pages only while every member is newer`() = runTest { + serve(total = 13, perPage = 2) + coEvery { memberDataSource.countMembers(chatId) } returns 13 + + subject.catchUp(chatId, state(watermark = 10)) + + // 13,12 | 11,10: the second page reaches the watermark. + assertEquals(2, requests.size) + coVerify { memberDataSource.upsert(chatId, listOf(member(11))) } + } + + @Test + fun `an equal count moves the watermark to the page's version, not the stream's`() = runTest { + serve(total = 13, perPage = RosterSync.PAGE_SIZE, versionOf = { 13 }) + coEvery { metadataDataSource.getRosterVersion(chatId) } returns 15 + coEvery { memberDataSource.countMembers(chatId) } returns 13 + + subject.catchUp(chatId, state(watermark = 10)) + + coVerify { rosterDataSource.setWatermark(chatId, 13) } + verify(exactly = 0) { scheduler.schedule(any()) } + } + + @Test + fun `more held than the count marks a reconcile and schedules it once`() = runTest { + serve(total = 13, perPage = RosterSync.PAGE_SIZE, memberCount = 12) + coEvery { memberDataSource.countMembers(chatId) } returns 13 + + subject.catchUp(chatId, state(watermark = 10)) + + coVerify { rosterDataSource.markReconcilePending(chatId) } + verify(exactly = 1) { scheduler.schedule(chatId) } + coVerify(exactly = 0) { rosterDataSource.setWatermark(any(), any()) } + } + + @Test + fun `fewer held than the count waits for the join on the stream`() = runTest { + serve(total = 13, perPage = RosterSync.PAGE_SIZE, memberCount = 14) + coEvery { memberDataSource.countMembers(chatId) } returns 13 + + subject.catchUp(chatId, state(watermark = 10)) + + coVerify(exactly = 0) { rosterDataSource.setWatermark(any(), any()) } + verify(exactly = 0) { scheduler.schedule(any()) } + } + + @Test + fun `a roster cut off at the cap still moves its watermark`() = runTest { + serve(total = 13, perPage = RosterSync.PAGE_SIZE, memberCount = 5_000) + coEvery { memberDataSource.countMembers(chatId) } returns 2_000 + + subject.catchUp(chatId, state(watermark = 10, truncated = true)) + + coVerify { rosterDataSource.setWatermark(chatId, 13) } + } + + @Test + fun `more missed joins than the cap reads hands over to a full read`() = runTest { + serve(total = (RosterSync.MAX_PAGES + 2) * 2, perPage = 2) + + subject.catchUp(chatId, state(watermark = 1)) + + assertEquals(RosterSync.MAX_PAGES, requests.size) + coVerify { rosterDataSource.markReconcilePending(chatId) } + verify(exactly = 1) { scheduler.schedule(chatId) } + } + + // Triggers + + @Test + fun `a DM is never read`() = runTest { + coEvery { metadataDataSource.getChatType(chatId) } returns ChatType.CONTACT_DM + + subject.onOpen(chatId) + + coVerify(exactly = 0) { controller.getRoster(any(), any()) } + verify(exactly = 0) { scheduler.schedule(any()) } + } + + @Test + fun `a small group never read is read on open`() = runTest { + group(memberCount = 40, syncState = null) + serve(total = 40, perPage = RosterSync.PAGE_SIZE) + + subject.onOpen(chatId) + + assertEquals(1, requests.size) + verify(exactly = 0) { scheduler.schedule(any()) } + } + + @Test + fun `a large group never read is handed to WorkManager`() = runTest { + group(memberCount = 250, syncState = state(fullySynced = false)) + + subject.onOpen(chatId) + + coVerify(exactly = 0) { controller.getRoster(any(), any()) } + verify(exactly = 1) { scheduler.schedule(chatId) } + } + + @Test + fun `a pending reconcile is re-queued on open`() = runTest { + group(memberCount = 40, syncState = state(watermark = 10, reconcilePending = true)) + + subject.onOpen(chatId) + + verify(exactly = 1) { scheduler.schedule(chatId) } + coVerify(exactly = 0) { controller.getRoster(any(), any()) } + } + + @Test + fun `a group whose roster moved past the watermark catches up on open`() = runTest { + group(memberCount = 13, syncState = state(watermark = 10)) + coEvery { metadataDataSource.getRosterVersion(chatId) } returns 13 + serve(total = 13, perPage = RosterSync.PAGE_SIZE) + + subject.onOpen(chatId) + + assertEquals(1, requests.size) + } + + @Test + fun `a group at its watermark is not read`() = runTest { + group(memberCount = 13, syncState = state(watermark = 13)) + coEvery { metadataDataSource.getRosterVersion(chatId) } returns 13 + + subject.onOpen(chatId) + + coVerify(exactly = 0) { controller.getRoster(any(), any()) } + } + + @Test + fun `a gap in a group never read waits for its open`() = runTest { + coEvery { rosterDataSource.getSyncState(chatId) } returns null + + subject.onGap(chatId) + + coVerify(exactly = 0) { controller.getRoster(any(), any()) } + } + + @Test + fun `a gap in a group read in full catches up`() = runTest { + coEvery { rosterDataSource.getSyncState(chatId) } returns state(watermark = 10) + serve(total = 13, perPage = RosterSync.PAGE_SIZE) + + subject.onGap(chatId) + + assertEquals(1, requests.size) + } + + @Test + fun `refreshing rewrites only the first page`() = runTest { + serve(total = 6) + + subject.refreshFirstPage(chatId) + + assertEquals(listOf?>(null), requests.map { it.token }) + coVerify(exactly = 1) { memberDataSource.upsert(chatId, any()) } + coVerify(exactly = 0) { rosterDataSource.markFullySynced(any(), any(), any()) } + } + + private fun group(memberCount: Long, syncState: RosterSyncState?) { + coEvery { metadataDataSource.getChatType(chatId) } returns ChatType.GROUP + coEvery { metadataDataSource.getMemberCount(chatId) } returns memberCount + coEvery { rosterDataSource.getSyncState(chatId) } returns syncState + } +} diff --git a/apps/flipcash/shared/persistence/db/build.gradle.kts b/apps/flipcash/shared/persistence/db/build.gradle.kts index 75f2a790f1..2ec22ac1e8 100644 --- a/apps/flipcash/shared/persistence/db/build.gradle.kts +++ b/apps/flipcash/shared/persistence/db/build.gradle.kts @@ -11,6 +11,14 @@ android { } } +// MigrationTestHelper reads exported schemas from the unit test's assets. +androidComponents { + onVariants { variant -> + variant.hostTests[com.android.build.api.variant.HostTestBuilder.UNIT_TEST_TYPE] + ?.sources?.assets?.addStaticSourceDirectory("$projectDir/schemas") + } +} + dependencies { testImplementation(kotlin("test")) testImplementation(libs.bundles.unit.testing) diff --git a/apps/flipcash/shared/persistence/db/schemas/com.flipcash.app.persistence.FlipcashDatabase/41.json b/apps/flipcash/shared/persistence/db/schemas/com.flipcash.app.persistence.FlipcashDatabase/41.json new file mode 100644 index 0000000000..b84f9eade4 --- /dev/null +++ b/apps/flipcash/shared/persistence/db/schemas/com.flipcash.app.persistence.FlipcashDatabase/41.json @@ -0,0 +1,1039 @@ +{ + "formatVersion": 1, + "database": { + "version": 41, + "identityHash": "d49e6cd616d1c1631dbb3e78901e2b00", + "entities": [ + { + "tableName": "messages", + "createSql": "CREATE TABLE IF NOT EXISTS `${TABLE_NAME}` (`idBase58` TEXT NOT NULL, `text` TEXT NOT NULL, `amountUsdc` INTEGER, `amountNative` INTEGER, `nativeCurrency` TEXT, `rate` REAL, `state` TEXT NOT NULL, `timestamp` INTEGER NOT NULL, `metadata` TEXT, `mintBase58` TEXT DEFAULT 'EPjFWdd5AufqSSqeM2qN1xzybapC8G4wEGGkZwyTDt1v', `textSubstitutions` TEXT, PRIMARY KEY(`idBase58`))", + "fields": [ + { + "fieldPath": "idBase58", + "columnName": "idBase58", + "affinity": "TEXT", + "notNull": true + }, + { + "fieldPath": "text", + "columnName": "text", + "affinity": "TEXT", + "notNull": true + }, + { + "fieldPath": "amountUsdc", + "columnName": "amountUsdc", + "affinity": "INTEGER" + }, + { + "fieldPath": "amountNative", + "columnName": "amountNative", + "affinity": "INTEGER" + }, + { + "fieldPath": "nativeCurrency", + "columnName": "nativeCurrency", + "affinity": "TEXT" + }, + { + "fieldPath": "rate", + "columnName": "rate", + "affinity": "REAL" + }, + { + "fieldPath": "state", + "columnName": "state", + "affinity": "TEXT", + "notNull": true + }, + { + "fieldPath": "timestamp", + "columnName": "timestamp", + "affinity": "INTEGER", + "notNull": true + }, + { + "fieldPath": "metadata", + "columnName": "metadata", + "affinity": "TEXT" + }, + { + "fieldPath": "mintBase58", + "columnName": "mintBase58", + "affinity": "TEXT", + "defaultValue": "'EPjFWdd5AufqSSqeM2qN1xzybapC8G4wEGGkZwyTDt1v'" + }, + { + "fieldPath": "textSubstitutions", + "columnName": "textSubstitutions", + "affinity": "TEXT" + } + ], + "primaryKey": { + "autoGenerate": false, + "columnNames": [ + "idBase58" + ] + } + }, + { + "tableName": "tokens", + "createSql": "CREATE TABLE IF NOT EXISTS `${TABLE_NAME}` (`address` TEXT NOT NULL, `decimals` INTEGER NOT NULL, `name` TEXT NOT NULL, `symbol` TEXT NOT NULL, `created_at` INTEGER, `description` TEXT NOT NULL, `image_url` TEXT NOT NULL, `social_links` TEXT, `bill_customizations` TEXT, `holder_metrics` TEXT, `market_cap_metrics` TEXT, `vm_vm` TEXT NOT NULL, `vm_authority` TEXT NOT NULL, `vm_lock_duration_days` INTEGER NOT NULL, `lp_currency_config` TEXT, `lp_liquidity_pool` TEXT, `lp_seed` TEXT, `lp_authority` TEXT, `lp_mint_vault` TEXT, `lp_core_mint_vault` TEXT, `lp_circulating_supply_quarks` INTEGER, `lp_sell_fee_bps` INTEGER, `lp_price_amount_usd` REAL, `lp_market_cap_amount_usd` REAL, PRIMARY KEY(`address`))", + "fields": [ + { + "fieldPath": "address", + "columnName": "address", + "affinity": "TEXT", + "notNull": true + }, + { + "fieldPath": "decimals", + "columnName": "decimals", + "affinity": "INTEGER", + "notNull": true + }, + { + "fieldPath": "name", + "columnName": "name", + "affinity": "TEXT", + "notNull": true + }, + { + "fieldPath": "symbol", + "columnName": "symbol", + "affinity": "TEXT", + "notNull": true + }, + { + "fieldPath": "createdAt", + "columnName": "created_at", + "affinity": "INTEGER" + }, + { + "fieldPath": "description", + "columnName": "description", + "affinity": "TEXT", + "notNull": true + }, + { + "fieldPath": "imageUrl", + "columnName": "image_url", + "affinity": "TEXT", + "notNull": true + }, + { + "fieldPath": "socialLinks", + "columnName": "social_links", + "affinity": "TEXT" + }, + { + "fieldPath": "billCustomizationsJson", + "columnName": "bill_customizations", + "affinity": "TEXT" + }, + { + "fieldPath": "holderMetricsJson", + "columnName": "holder_metrics", + "affinity": "TEXT" + }, + { + "fieldPath": "marketCapMetricsJson", + "columnName": "market_cap_metrics", + "affinity": "TEXT" + }, + { + "fieldPath": "vmMetadata.vm", + "columnName": "vm_vm", + "affinity": "TEXT", + "notNull": true + }, + { + "fieldPath": "vmMetadata.authority", + "columnName": "vm_authority", + "affinity": "TEXT", + "notNull": true + }, + { + "fieldPath": "vmMetadata.lockDurationInDays", + "columnName": "vm_lock_duration_days", + "affinity": "INTEGER", + "notNull": true + }, + { + "fieldPath": "launchpadMetadata.currencyConfig", + "columnName": "lp_currency_config", + "affinity": "TEXT" + }, + { + "fieldPath": "launchpadMetadata.liquidityPool", + "columnName": "lp_liquidity_pool", + "affinity": "TEXT" + }, + { + "fieldPath": "launchpadMetadata.seed", + "columnName": "lp_seed", + "affinity": "TEXT" + }, + { + "fieldPath": "launchpadMetadata.authority", + "columnName": "lp_authority", + "affinity": "TEXT" + }, + { + "fieldPath": "launchpadMetadata.mintVault", + "columnName": "lp_mint_vault", + "affinity": "TEXT" + }, + { + "fieldPath": "launchpadMetadata.coreMintVault", + "columnName": "lp_core_mint_vault", + "affinity": "TEXT" + }, + { + "fieldPath": "launchpadMetadata.currentCirculatingSupplyQuarks", + "columnName": "lp_circulating_supply_quarks", + "affinity": "INTEGER" + }, + { + "fieldPath": "launchpadMetadata.sellFeeBps", + "columnName": "lp_sell_fee_bps", + "affinity": "INTEGER" + }, + { + "fieldPath": "launchpadMetadata.priceAmount", + "columnName": "lp_price_amount_usd", + "affinity": "REAL" + }, + { + "fieldPath": "launchpadMetadata.marketCapAmount", + "columnName": "lp_market_cap_amount_usd", + "affinity": "REAL" + } + ], + "primaryKey": { + "autoGenerate": false, + "columnNames": [ + "address" + ] + } + }, + { + "tableName": "token_social_links", + "createSql": "CREATE TABLE IF NOT EXISTS `${TABLE_NAME}` (`id` INTEGER PRIMARY KEY AUTOINCREMENT NOT NULL, `token_address` TEXT NOT NULL, `type` TEXT NOT NULL, `value` TEXT NOT NULL, FOREIGN KEY(`token_address`) REFERENCES `tokens`(`address`) ON UPDATE NO ACTION ON DELETE CASCADE )", + "fields": [ + { + "fieldPath": "id", + "columnName": "id", + "affinity": "INTEGER", + "notNull": true + }, + { + "fieldPath": "tokenAddress", + "columnName": "token_address", + "affinity": "TEXT", + "notNull": true + }, + { + "fieldPath": "type", + "columnName": "type", + "affinity": "TEXT", + "notNull": true + }, + { + "fieldPath": "value", + "columnName": "value", + "affinity": "TEXT", + "notNull": true + } + ], + "primaryKey": { + "autoGenerate": true, + "columnNames": [ + "id" + ] + }, + "indices": [ + { + "name": "index_token_social_links_token_address", + "unique": false, + "columnNames": [ + "token_address" + ], + "orders": [], + "createSql": "CREATE INDEX IF NOT EXISTS `index_token_social_links_token_address` ON `${TABLE_NAME}` (`token_address`)" + } + ], + "foreignKeys": [ + { + "table": "tokens", + "onDelete": "CASCADE", + "onUpdate": "NO ACTION", + "columns": [ + "token_address" + ], + "referencedColumns": [ + "address" + ] + } + ] + }, + { + "tableName": "token_valuation", + "createSql": "CREATE TABLE IF NOT EXISTS `${TABLE_NAME}` (`token_address` TEXT NOT NULL, `balance_quarks` INTEGER NOT NULL, `cost_basis` REAL NOT NULL, PRIMARY KEY(`token_address`), FOREIGN KEY(`token_address`) REFERENCES `tokens`(`address`) ON UPDATE NO ACTION ON DELETE CASCADE )", + "fields": [ + { + "fieldPath": "tokenAddress", + "columnName": "token_address", + "affinity": "TEXT", + "notNull": true + }, + { + "fieldPath": "balanceQuarks", + "columnName": "balance_quarks", + "affinity": "INTEGER", + "notNull": true + }, + { + "fieldPath": "costBasis", + "columnName": "cost_basis", + "affinity": "REAL", + "notNull": true + } + ], + "primaryKey": { + "autoGenerate": false, + "columnNames": [ + "token_address" + ] + }, + "indices": [ + { + "name": "index_token_valuation_token_address", + "unique": false, + "columnNames": [ + "token_address" + ], + "orders": [], + "createSql": "CREATE INDEX IF NOT EXISTS `index_token_valuation_token_address` ON `${TABLE_NAME}` (`token_address`)" + } + ], + "foreignKeys": [ + { + "table": "tokens", + "onDelete": "CASCADE", + "onUpdate": "NO ACTION", + "columns": [ + "token_address" + ], + "referencedColumns": [ + "address" + ] + } + ] + }, + { + "tableName": "currency_creator_draft", + "createSql": "CREATE TABLE IF NOT EXISTS `${TABLE_NAME}` (`id` INTEGER PRIMARY KEY AUTOINCREMENT NOT NULL, `name` TEXT NOT NULL, `description` TEXT NOT NULL, `icon_uri` TEXT, `bill_customizations` TEXT, `attestations` TEXT, `current_step` TEXT NOT NULL, `created_mint` TEXT, `saved_at` INTEGER NOT NULL)", + "fields": [ + { + "fieldPath": "id", + "columnName": "id", + "affinity": "INTEGER", + "notNull": true + }, + { + "fieldPath": "name", + "columnName": "name", + "affinity": "TEXT", + "notNull": true + }, + { + "fieldPath": "description", + "columnName": "description", + "affinity": "TEXT", + "notNull": true + }, + { + "fieldPath": "iconUri", + "columnName": "icon_uri", + "affinity": "TEXT" + }, + { + "fieldPath": "billCustomizations", + "columnName": "bill_customizations", + "affinity": "TEXT" + }, + { + "fieldPath": "attestations", + "columnName": "attestations", + "affinity": "TEXT" + }, + { + "fieldPath": "currentStep", + "columnName": "current_step", + "affinity": "TEXT", + "notNull": true + }, + { + "fieldPath": "createdMint", + "columnName": "created_mint", + "affinity": "TEXT" + }, + { + "fieldPath": "savedAt", + "columnName": "saved_at", + "affinity": "INTEGER", + "notNull": true + } + ], + "primaryKey": { + "autoGenerate": true, + "columnNames": [ + "id" + ] + } + }, + { + "tableName": "contact_sync_state", + "createSql": "CREATE TABLE IF NOT EXISTS `${TABLE_NAME}` (`id` INTEGER NOT NULL, `checksumBytes` BLOB NOT NULL, `lastSyncTimestamp` INTEGER NOT NULL, `needsFullUpload` INTEGER NOT NULL, `hasDiscoveredFlipcashContacts` INTEGER NOT NULL DEFAULT 0, PRIMARY KEY(`id`))", + "fields": [ + { + "fieldPath": "id", + "columnName": "id", + "affinity": "INTEGER", + "notNull": true + }, + { + "fieldPath": "checksumBytes", + "columnName": "checksumBytes", + "affinity": "BLOB", + "notNull": true + }, + { + "fieldPath": "lastSyncTimestamp", + "columnName": "lastSyncTimestamp", + "affinity": "INTEGER", + "notNull": true + }, + { + "fieldPath": "needsFullUpload", + "columnName": "needsFullUpload", + "affinity": "INTEGER", + "notNull": true + }, + { + "fieldPath": "hasDiscoveredFlipcashContacts", + "columnName": "hasDiscoveredFlipcashContacts", + "affinity": "INTEGER", + "notNull": true, + "defaultValue": "0" + } + ], + "primaryKey": { + "autoGenerate": false, + "columnNames": [ + "id" + ] + } + }, + { + "tableName": "contact_mapping", + "createSql": "CREATE TABLE IF NOT EXISTS `${TABLE_NAME}` (`e164` TEXT NOT NULL, `androidContactId` INTEGER NOT NULL, `displayName` TEXT NOT NULL, `photoUri` TEXT, `isOnFlipcash` INTEGER NOT NULL, `displayNumber` TEXT NOT NULL DEFAULT '', `dmChatId` TEXT NOT NULL DEFAULT '', `joinedAtEpochSeconds` INTEGER NOT NULL DEFAULT 0, PRIMARY KEY(`e164`))", + "fields": [ + { + "fieldPath": "e164", + "columnName": "e164", + "affinity": "TEXT", + "notNull": true + }, + { + "fieldPath": "androidContactId", + "columnName": "androidContactId", + "affinity": "INTEGER", + "notNull": true + }, + { + "fieldPath": "displayName", + "columnName": "displayName", + "affinity": "TEXT", + "notNull": true + }, + { + "fieldPath": "photoUri", + "columnName": "photoUri", + "affinity": "TEXT" + }, + { + "fieldPath": "isOnFlipcash", + "columnName": "isOnFlipcash", + "affinity": "INTEGER", + "notNull": true + }, + { + "fieldPath": "displayNumber", + "columnName": "displayNumber", + "affinity": "TEXT", + "notNull": true, + "defaultValue": "''" + }, + { + "fieldPath": "dmChatId", + "columnName": "dmChatId", + "affinity": "TEXT", + "notNull": true, + "defaultValue": "''" + }, + { + "fieldPath": "joinedAtEpochSeconds", + "columnName": "joinedAtEpochSeconds", + "affinity": "INTEGER", + "notNull": true, + "defaultValue": "0" + } + ], + "primaryKey": { + "autoGenerate": false, + "columnNames": [ + "e164" + ] + } + }, + { + "tableName": "chat_metadata", + "createSql": "CREATE TABLE IF NOT EXISTS `${TABLE_NAME}` (`chat_id_hex` TEXT NOT NULL, `chat_type` TEXT NOT NULL, `last_activity_epoch_ms` INTEGER NOT NULL, `last_message_id` INTEGER, `latest_event_sequence` INTEGER NOT NULL DEFAULT 0, `is_hidden` INTEGER NOT NULL DEFAULT 0, `analytics_counted_through` INTEGER NOT NULL DEFAULT 0, `title` TEXT, `picture_json` TEXT, `member_count` INTEGER NOT NULL DEFAULT 0, `roster_version` INTEGER NOT NULL DEFAULT 0, `rules_json` TEXT, `is_member` INTEGER NOT NULL DEFAULT 1, `mute_until_epoch_ms` INTEGER, `mute_forever` INTEGER NOT NULL DEFAULT 0, `viewer_state_version` INTEGER NOT NULL DEFAULT 0, `can_edit` INTEGER NOT NULL DEFAULT 0, `creator_hex` TEXT, `use_e2ee` INTEGER NOT NULL DEFAULT 0, PRIMARY KEY(`chat_id_hex`))", + "fields": [ + { + "fieldPath": "chatIdHex", + "columnName": "chat_id_hex", + "affinity": "TEXT", + "notNull": true + }, + { + "fieldPath": "chatType", + "columnName": "chat_type", + "affinity": "TEXT", + "notNull": true + }, + { + "fieldPath": "lastActivityEpochMs", + "columnName": "last_activity_epoch_ms", + "affinity": "INTEGER", + "notNull": true + }, + { + "fieldPath": "lastMessageId", + "columnName": "last_message_id", + "affinity": "INTEGER" + }, + { + "fieldPath": "latestEventSequence", + "columnName": "latest_event_sequence", + "affinity": "INTEGER", + "notNull": true, + "defaultValue": "0" + }, + { + "fieldPath": "isHidden", + "columnName": "is_hidden", + "affinity": "INTEGER", + "notNull": true, + "defaultValue": "0" + }, + { + "fieldPath": "analyticsCountedThrough", + "columnName": "analytics_counted_through", + "affinity": "INTEGER", + "notNull": true, + "defaultValue": "0" + }, + { + "fieldPath": "title", + "columnName": "title", + "affinity": "TEXT" + }, + { + "fieldPath": "pictureJson", + "columnName": "picture_json", + "affinity": "TEXT" + }, + { + "fieldPath": "memberCount", + "columnName": "member_count", + "affinity": "INTEGER", + "notNull": true, + "defaultValue": "0" + }, + { + "fieldPath": "rosterVersion", + "columnName": "roster_version", + "affinity": "INTEGER", + "notNull": true, + "defaultValue": "0" + }, + { + "fieldPath": "rulesJson", + "columnName": "rules_json", + "affinity": "TEXT" + }, + { + "fieldPath": "isMember", + "columnName": "is_member", + "affinity": "INTEGER", + "notNull": true, + "defaultValue": "1" + }, + { + "fieldPath": "muteUntilEpochMs", + "columnName": "mute_until_epoch_ms", + "affinity": "INTEGER" + }, + { + "fieldPath": "muteForever", + "columnName": "mute_forever", + "affinity": "INTEGER", + "notNull": true, + "defaultValue": "0" + }, + { + "fieldPath": "viewerStateVersion", + "columnName": "viewer_state_version", + "affinity": "INTEGER", + "notNull": true, + "defaultValue": "0" + }, + { + "fieldPath": "canEdit", + "columnName": "can_edit", + "affinity": "INTEGER", + "notNull": true, + "defaultValue": "0" + }, + { + "fieldPath": "creatorHex", + "columnName": "creator_hex", + "affinity": "TEXT" + }, + { + "fieldPath": "useE2ee", + "columnName": "use_e2ee", + "affinity": "INTEGER", + "notNull": true, + "defaultValue": "0" + } + ], + "primaryKey": { + "autoGenerate": false, + "columnNames": [ + "chat_id_hex" + ] + }, + "indices": [ + { + "name": "index_chat_metadata_last_activity_epoch_ms", + "unique": false, + "columnNames": [ + "last_activity_epoch_ms" + ], + "orders": [], + "createSql": "CREATE INDEX IF NOT EXISTS `index_chat_metadata_last_activity_epoch_ms` ON `${TABLE_NAME}` (`last_activity_epoch_ms`)" + } + ] + }, + { + "tableName": "chat_messages", + "createSql": "CREATE TABLE IF NOT EXISTS `${TABLE_NAME}` (`chat_id_hex` TEXT NOT NULL, `message_id` INTEGER NOT NULL, `sender_id_hex` TEXT, `content_json` TEXT, `timestamp_epoch_ms` INTEGER NOT NULL, `unread_seq` INTEGER NOT NULL, `status` TEXT NOT NULL DEFAULT 'SENT', `pending_client_id_hex` TEXT, `event_sequence` INTEGER NOT NULL DEFAULT 0, `last_edited_ts_epoch_ms` INTEGER, `reactions_json` TEXT, `is_deleted` INTEGER NOT NULL DEFAULT 0, `ciphertext_json` TEXT, `encryption_state` TEXT, PRIMARY KEY(`chat_id_hex`, `message_id`))", + "fields": [ + { + "fieldPath": "chatIdHex", + "columnName": "chat_id_hex", + "affinity": "TEXT", + "notNull": true + }, + { + "fieldPath": "messageId", + "columnName": "message_id", + "affinity": "INTEGER", + "notNull": true + }, + { + "fieldPath": "senderIdHex", + "columnName": "sender_id_hex", + "affinity": "TEXT" + }, + { + "fieldPath": "contentJson", + "columnName": "content_json", + "affinity": "TEXT" + }, + { + "fieldPath": "timestampEpochMs", + "columnName": "timestamp_epoch_ms", + "affinity": "INTEGER", + "notNull": true + }, + { + "fieldPath": "unreadSeq", + "columnName": "unread_seq", + "affinity": "INTEGER", + "notNull": true + }, + { + "fieldPath": "status", + "columnName": "status", + "affinity": "TEXT", + "notNull": true, + "defaultValue": "'SENT'" + }, + { + "fieldPath": "pendingClientIdHex", + "columnName": "pending_client_id_hex", + "affinity": "TEXT" + }, + { + "fieldPath": "eventSequence", + "columnName": "event_sequence", + "affinity": "INTEGER", + "notNull": true, + "defaultValue": "0" + }, + { + "fieldPath": "lastEditedTsEpochMs", + "columnName": "last_edited_ts_epoch_ms", + "affinity": "INTEGER" + }, + { + "fieldPath": "reactionsJson", + "columnName": "reactions_json", + "affinity": "TEXT" + }, + { + "fieldPath": "isDeleted", + "columnName": "is_deleted", + "affinity": "INTEGER", + "notNull": true, + "defaultValue": "0" + }, + { + "fieldPath": "ciphertextJson", + "columnName": "ciphertext_json", + "affinity": "TEXT" + }, + { + "fieldPath": "encryptionState", + "columnName": "encryption_state", + "affinity": "TEXT" + } + ], + "primaryKey": { + "autoGenerate": false, + "columnNames": [ + "chat_id_hex", + "message_id" + ] + }, + "indices": [ + { + "name": "index_chat_messages_chat_id_hex_timestamp_epoch_ms", + "unique": false, + "columnNames": [ + "chat_id_hex", + "timestamp_epoch_ms" + ], + "orders": [], + "createSql": "CREATE INDEX IF NOT EXISTS `index_chat_messages_chat_id_hex_timestamp_epoch_ms` ON `${TABLE_NAME}` (`chat_id_hex`, `timestamp_epoch_ms`)" + } + ] + }, + { + "tableName": "chat_members", + "createSql": "CREATE TABLE IF NOT EXISTS `${TABLE_NAME}` (`chat_id_hex` TEXT NOT NULL, `user_id_hex` TEXT NOT NULL, `pointers_json` TEXT, `version` INTEGER NOT NULL DEFAULT 0, `is_member` INTEGER NOT NULL DEFAULT 1, PRIMARY KEY(`chat_id_hex`, `user_id_hex`))", + "fields": [ + { + "fieldPath": "chatIdHex", + "columnName": "chat_id_hex", + "affinity": "TEXT", + "notNull": true + }, + { + "fieldPath": "userIdHex", + "columnName": "user_id_hex", + "affinity": "TEXT", + "notNull": true + }, + { + "fieldPath": "pointersJson", + "columnName": "pointers_json", + "affinity": "TEXT" + }, + { + "fieldPath": "version", + "columnName": "version", + "affinity": "INTEGER", + "notNull": true, + "defaultValue": "0" + }, + { + "fieldPath": "isMember", + "columnName": "is_member", + "affinity": "INTEGER", + "notNull": true, + "defaultValue": "1" + } + ], + "primaryKey": { + "autoGenerate": false, + "columnNames": [ + "chat_id_hex", + "user_id_hex" + ] + } + }, + { + "tableName": "chat_draft", + "createSql": "CREATE TABLE IF NOT EXISTS `${TABLE_NAME}` (`chat_id_hex` TEXT NOT NULL, `text` TEXT NOT NULL, `reply_target_json` TEXT, `saved_at` INTEGER NOT NULL, PRIMARY KEY(`chat_id_hex`))", + "fields": [ + { + "fieldPath": "chatIdHex", + "columnName": "chat_id_hex", + "affinity": "TEXT", + "notNull": true + }, + { + "fieldPath": "text", + "columnName": "text", + "affinity": "TEXT", + "notNull": true + }, + { + "fieldPath": "replyTargetJson", + "columnName": "reply_target_json", + "affinity": "TEXT" + }, + { + "fieldPath": "savedAt", + "columnName": "saved_at", + "affinity": "INTEGER", + "notNull": true + } + ], + "primaryKey": { + "autoGenerate": false, + "columnNames": [ + "chat_id_hex" + ] + } + }, + { + "tableName": "blocked_users", + "createSql": "CREATE TABLE IF NOT EXISTS `${TABLE_NAME}` (`user_id_hex` TEXT NOT NULL, `blocked_at_epoch_ms` INTEGER NOT NULL, PRIMARY KEY(`user_id_hex`))", + "fields": [ + { + "fieldPath": "userIdHex", + "columnName": "user_id_hex", + "affinity": "TEXT", + "notNull": true + }, + { + "fieldPath": "blockedAtEpochMs", + "columnName": "blocked_at_epoch_ms", + "affinity": "INTEGER", + "notNull": true + } + ], + "primaryKey": { + "autoGenerate": false, + "columnNames": [ + "user_id_hex" + ] + } + }, + { + "tableName": "user_profiles", + "createSql": "CREATE TABLE IF NOT EXISTS `${TABLE_NAME}` (`user_id_hex` TEXT NOT NULL, `display_name` TEXT NOT NULL, `phone_value` TEXT, `phone_verified` INTEGER, `email_value` TEXT, `email_verified` INTEGER, `social_accounts_json` TEXT, `profile_picture_json` TEXT, `username` TEXT, `pending_migration_json` TEXT, PRIMARY KEY(`user_id_hex`))", + "fields": [ + { + "fieldPath": "userIdHex", + "columnName": "user_id_hex", + "affinity": "TEXT", + "notNull": true + }, + { + "fieldPath": "displayName", + "columnName": "display_name", + "affinity": "TEXT", + "notNull": true + }, + { + "fieldPath": "phoneValue", + "columnName": "phone_value", + "affinity": "TEXT" + }, + { + "fieldPath": "phoneVerified", + "columnName": "phone_verified", + "affinity": "INTEGER" + }, + { + "fieldPath": "emailValue", + "columnName": "email_value", + "affinity": "TEXT" + }, + { + "fieldPath": "emailVerified", + "columnName": "email_verified", + "affinity": "INTEGER" + }, + { + "fieldPath": "socialAccounts", + "columnName": "social_accounts_json", + "affinity": "TEXT" + }, + { + "fieldPath": "profilePicture", + "columnName": "profile_picture_json", + "affinity": "TEXT" + }, + { + "fieldPath": "username", + "columnName": "username", + "affinity": "TEXT" + }, + { + "fieldPath": "pendingMigrationJson", + "columnName": "pending_migration_json", + "affinity": "TEXT" + } + ], + "primaryKey": { + "autoGenerate": false, + "columnNames": [ + "user_id_hex" + ] + } + }, + { + "tableName": "link_previews", + "createSql": "CREATE TABLE IF NOT EXISTS `${TABLE_NAME}` (`key` TEXT NOT NULL, `json` TEXT NOT NULL, `updated_at` INTEGER NOT NULL, PRIMARY KEY(`key`))", + "fields": [ + { + "fieldPath": "key", + "columnName": "key", + "affinity": "TEXT", + "notNull": true + }, + { + "fieldPath": "json", + "columnName": "json", + "affinity": "TEXT", + "notNull": true + }, + { + "fieldPath": "updatedAt", + "columnName": "updated_at", + "affinity": "INTEGER", + "notNull": true + } + ], + "primaryKey": { + "autoGenerate": false, + "columnNames": [ + "key" + ] + } + }, + { + "tableName": "chat_member_search_tokens", + "createSql": "CREATE TABLE IF NOT EXISTS `${TABLE_NAME}` (`chat_id_hex` TEXT NOT NULL, `user_id_hex` TEXT NOT NULL, `token` TEXT NOT NULL, PRIMARY KEY(`chat_id_hex`, `user_id_hex`, `token`))", + "fields": [ + { + "fieldPath": "chatIdHex", + "columnName": "chat_id_hex", + "affinity": "TEXT", + "notNull": true + }, + { + "fieldPath": "userIdHex", + "columnName": "user_id_hex", + "affinity": "TEXT", + "notNull": true + }, + { + "fieldPath": "token", + "columnName": "token", + "affinity": "TEXT", + "notNull": true + } + ], + "primaryKey": { + "autoGenerate": false, + "columnNames": [ + "chat_id_hex", + "user_id_hex", + "token" + ] + }, + "indices": [ + { + "name": "index_chat_member_search_tokens_chat_id_hex_token", + "unique": false, + "columnNames": [ + "chat_id_hex", + "token" + ], + "orders": [], + "createSql": "CREATE INDEX IF NOT EXISTS `index_chat_member_search_tokens_chat_id_hex_token` ON `${TABLE_NAME}` (`chat_id_hex`, `token`)" + } + ] + }, + { + "tableName": "chat_roster_sync", + "createSql": "CREATE TABLE IF NOT EXISTS `${TABLE_NAME}` (`chat_id_hex` TEXT NOT NULL, `watermark` INTEGER NOT NULL, `fully_synced` INTEGER NOT NULL DEFAULT 0, `truncated` INTEGER NOT NULL DEFAULT 0, `reconcile_pending` INTEGER NOT NULL DEFAULT 0, PRIMARY KEY(`chat_id_hex`))", + "fields": [ + { + "fieldPath": "chatIdHex", + "columnName": "chat_id_hex", + "affinity": "TEXT", + "notNull": true + }, + { + "fieldPath": "watermark", + "columnName": "watermark", + "affinity": "INTEGER", + "notNull": true + }, + { + "fieldPath": "fullySynced", + "columnName": "fully_synced", + "affinity": "INTEGER", + "notNull": true, + "defaultValue": "0" + }, + { + "fieldPath": "truncated", + "columnName": "truncated", + "affinity": "INTEGER", + "notNull": true, + "defaultValue": "0" + }, + { + "fieldPath": "reconcilePending", + "columnName": "reconcile_pending", + "affinity": "INTEGER", + "notNull": true, + "defaultValue": "0" + } + ], + "primaryKey": { + "autoGenerate": false, + "columnNames": [ + "chat_id_hex" + ] + } + } + ], + "setupQueries": [ + "CREATE TABLE IF NOT EXISTS room_master_table (id INTEGER PRIMARY KEY,identity_hash TEXT)", + "INSERT OR REPLACE INTO room_master_table (id,identity_hash) VALUES(42, 'd49e6cd616d1c1631dbb3e78901e2b00')" + ] + } +} \ No newline at end of file diff --git a/apps/flipcash/shared/persistence/db/src/main/kotlin/com/flipcash/app/persistence/FlipcashDatabase.kt b/apps/flipcash/shared/persistence/db/src/main/kotlin/com/flipcash/app/persistence/FlipcashDatabase.kt index dc60110d2b..24e832e4e0 100644 --- a/apps/flipcash/shared/persistence/db/src/main/kotlin/com/flipcash/app/persistence/FlipcashDatabase.kt +++ b/apps/flipcash/shared/persistence/db/src/main/kotlin/com/flipcash/app/persistence/FlipcashDatabase.kt @@ -21,6 +21,7 @@ import com.flipcash.app.persistence.converters.TokenTypeConverters import com.flipcash.app.persistence.dao.BlockedUserDao import com.flipcash.app.persistence.dao.ChatDraftDao import com.flipcash.app.persistence.dao.ChatMemberDao +import com.flipcash.app.persistence.dao.ChatMemberSearchDao import com.flipcash.app.persistence.dao.ChatMessageDao import com.flipcash.app.persistence.dao.ChatMetadataDao import com.flipcash.app.persistence.dao.ContactDao @@ -32,8 +33,10 @@ import com.flipcash.app.persistence.dao.UserProfileDao import com.flipcash.app.persistence.entities.BlockedUserEntity import com.flipcash.app.persistence.entities.ChatDraftEntity import com.flipcash.app.persistence.entities.ChatMemberEntity +import com.flipcash.app.persistence.entities.ChatMemberSearchTokenEntity import com.flipcash.app.persistence.entities.ChatMessageEntity import com.flipcash.app.persistence.entities.ChatMetadataEntity +import com.flipcash.app.persistence.entities.ChatRosterSyncEntity import com.flipcash.app.persistence.entities.ContactMappingEntity import com.flipcash.app.persistence.entities.ContactSyncStateEntity import com.flipcash.app.persistence.entities.CurrencyCreatorDraftEntity @@ -64,6 +67,8 @@ import com.getcode.utils.subByteArray BlockedUserEntity::class, UserProfileEntity::class, LinkPreviewEntity::class, + ChatMemberSearchTokenEntity::class, + ChatRosterSyncEntity::class, ], autoMigrations = [ AutoMigration(from = 1, to = 2, spec = FlipcashDatabase.Migration1To2::class), @@ -111,8 +116,11 @@ import com.getcode.utils.subByteArray AutoMigration(from = 37, to = 38), // chat_metadata.creator_hex (nullable), use_e2ee (default 0) AutoMigration(from = 38, to = 39), // link_previews table AutoMigration(from = 39, to = 40, spec = FlipcashDatabase.Migration39To40::class), + // chat_member_search_tokens and chat_roster_sync. Both start empty: a group's first open + // after the upgrade finds no sync row and reads its roster, which fills the index. + AutoMigration(from = 40, to = 41), ], - version = 40, + version = 41, ) @TypeConverters(TokenTypeConverters::class, ChatTypeConverters::class) abstract class FlipcashDatabase : RoomDatabase() { @@ -124,6 +132,7 @@ abstract class FlipcashDatabase : RoomDatabase() { abstract fun chatMetadataDao(): ChatMetadataDao abstract fun chatMessageDao(): ChatMessageDao abstract fun chatMemberDao(): ChatMemberDao + abstract fun chatMemberSearchDao(): ChatMemberSearchDao abstract fun chatDraftDao(): ChatDraftDao abstract fun blockedUserDao(): BlockedUserDao abstract fun userProfileDao(): UserProfileDao diff --git a/apps/flipcash/shared/persistence/db/src/main/kotlin/com/flipcash/app/persistence/dao/ChatMemberDao.kt b/apps/flipcash/shared/persistence/db/src/main/kotlin/com/flipcash/app/persistence/dao/ChatMemberDao.kt index a2ca504711..ebcb6f8b06 100644 --- a/apps/flipcash/shared/persistence/db/src/main/kotlin/com/flipcash/app/persistence/dao/ChatMemberDao.kt +++ b/apps/flipcash/shared/persistence/db/src/main/kotlin/com/flipcash/app/persistence/dao/ChatMemberDao.kt @@ -14,15 +14,15 @@ import kotlinx.coroutines.flow.Flow interface ChatMemberDao { @Transaction - @Query("SELECT * FROM chat_members WHERE chat_id_hex = :chatIdHex") + @Query("SELECT * FROM chat_members WHERE chat_id_hex = :chatIdHex AND is_member = 1") suspend fun getMembersForChat(chatIdHex: String): List @Transaction - @Query("SELECT * FROM chat_members WHERE chat_id_hex = :chatIdHex") + @Query("SELECT * FROM chat_members WHERE chat_id_hex = :chatIdHex AND is_member = 1") fun observeMembersForChat(chatIdHex: String): Flow> @Transaction - @Query("SELECT * FROM chat_members") + @Query("SELECT * FROM chat_members WHERE is_member = 1") fun observeAll(): Flow> @Insert(onConflict = OnConflictStrategy.REPLACE) @@ -40,14 +40,42 @@ interface ChatMemberDao { suspend fun upsert(entity: ChatMemberEntity) { val existing = getMember(entity.chatIdHex, entity.userIdHex) ?: return insertOrReplace(entity) + // A member who left stays out until a rejoin, the only change with a version above the + // leave. A page that trails the stream still lists them, at their older join version. + if (!existing.isMember && entity.version <= existing.version) return insertOrReplace( entity.copy( pointersJson = mergePointers(existing.pointersJson, entity.pointersJson), + // Greater wins: a page that trails the stream must not wind a member back. + version = maxOf(existing.version, entity.version), + isMember = true, ) ) } + /** + * Applies a `MemberLeft` at roster [version]: [userIdHex] becomes a marker in [chatIdHex], so a + * trailing roster page cannot re-add them. A leave older than the member's join is stale and + * changes nothing. Pointers stay with the marker for a rejoin to carry on from. + * + * @return whether the leave applied, so the caller drops the member's search words only then. + */ + @Transaction + suspend fun markLeft(chatIdHex: String, userIdHex: String, version: Long): Boolean { + val existing = getMember(chatIdHex, userIdHex) + if (existing != null && existing.version > version) return false + insertOrReplace( + existing?.copy(version = version, isMember = false) + ?: ChatMemberEntity(chatIdHex, userIdHex, pointersJson = null, version = version, isMember = false) + ) + return true + } + + /** Clears [chatIdHex]'s leave markers at or below [version]: a complete read at that version outranks them. */ + @Query("DELETE FROM chat_members WHERE chat_id_hex = :chatIdHex AND is_member = 0 AND version <= :version") + suspend fun deleteMarkersUpTo(chatIdHex: String, version: Long) + @Transaction suspend fun upsert(entities: List) { for (entity in entities) upsert(entity) @@ -65,7 +93,7 @@ interface ChatMemberDao { """ SELECT m.chat_id_hex FROM chat_members m INNER JOIN chat_metadata c ON c.chat_id_hex = m.chat_id_hex - WHERE m.user_id_hex = :userIdHex AND c.chat_type = :chatType + WHERE m.user_id_hex = :userIdHex AND m.is_member = 1 AND c.chat_type = :chatType ORDER BY c.last_activity_epoch_ms DESC LIMIT 1 """ @@ -91,30 +119,33 @@ interface ChatMemberDao { userIdHex: String, pointer: MessagePointerSerialized, ) { - val existing = getMember(chatIdHex, userIdHex)?.pointersJson.orEmpty() + val row = getMember(chatIdHex, userIdHex) + val existing = row?.pointersJson.orEmpty() val current = existing.firstOrNull { it.type == pointer.type } if (current != null && current.value >= pointer.value) return + val pointers = existing.filterNot { it.type == pointer.type } + pointer + // Copied, not rebuilt: a fresh row would drop the member's version and re-add one who left. insertOrReplace( - ChatMemberEntity( - chatIdHex = chatIdHex, - userIdHex = userIdHex, - pointersJson = existing.filterNot { it.type == pointer.type } + pointer, - ) + row?.copy(pointersJson = pointers) + ?: ChatMemberEntity(chatIdHex = chatIdHex, userIdHex = userIdHex, pointersJson = pointers) ) } @Query("DELETE FROM chat_members WHERE chat_id_hex = :chatIdHex") suspend fun deleteForChat(chatIdHex: String) - /** Drops the members of [chatIdHex] that are no longer in [keepUserIdHexes]. */ + /** + * Drops the members of [chatIdHex] that are no longer in [keepUserIdHexes]. Leave markers stay: + * they guard against a trailing roster page, which a feed refresh does not change. + */ @Query( - "DELETE FROM chat_members WHERE chat_id_hex = :chatIdHex " + + "DELETE FROM chat_members WHERE chat_id_hex = :chatIdHex AND is_member = 1 " + "AND user_id_hex NOT IN (:keepUserIdHexes)" ) suspend fun deleteMembersNotIn(chatIdHex: String, keepUserIdHexes: List) - /** Drops one member of [chatIdHex]. What a `MemberLeft` roster change applies. */ + /** Drops one member row of [chatIdHex], marker or not. A `MemberLeft` uses [markLeft] instead. */ @Query("DELETE FROM chat_members WHERE chat_id_hex = :chatIdHex AND user_id_hex = :userIdHex") suspend fun deleteMember(chatIdHex: String, userIdHex: String) diff --git a/apps/flipcash/shared/persistence/db/src/main/kotlin/com/flipcash/app/persistence/dao/ChatMemberSearchDao.kt b/apps/flipcash/shared/persistence/db/src/main/kotlin/com/flipcash/app/persistence/dao/ChatMemberSearchDao.kt new file mode 100644 index 0000000000..aed4bcd913 --- /dev/null +++ b/apps/flipcash/shared/persistence/db/src/main/kotlin/com/flipcash/app/persistence/dao/ChatMemberSearchDao.kt @@ -0,0 +1,177 @@ +package com.flipcash.app.persistence.dao + +import androidx.room.ColumnInfo +import androidx.room.Dao +import androidx.room.Insert +import androidx.room.OnConflictStrategy +import androidx.room.Query +import androidx.room.Transaction +import com.flipcash.app.persistence.entities.ChatMemberSearchTokenEntity +import com.flipcash.app.persistence.entities.ChatRosterSyncEntity +import com.flipcash.services.models.chat.MediaItem + +/** + * The member search index of each chat, and how much of each group's roster has been read into it. + * + * Normalizing names into tokens is the caller's job. SQLite cannot fold diacritics, so the index + * holds already-normalized words and a query must be normalized the same way before it gets here. + */ +@Dao +interface ChatMemberSearchDao { + + // region tokens + + @Insert(onConflict = OnConflictStrategy.IGNORE) + suspend fun insertTokens(tokens: List) + + @Query("DELETE FROM chat_member_search_tokens WHERE chat_id_hex = :chatIdHex AND user_id_hex = :userIdHex") + suspend fun deleteTokensForMember(chatIdHex: String, userIdHex: String) + + /** Swaps [userIdHex]'s words in [chatIdHex] for [tokens], so a renamed member loses the old ones. */ + @Transaction + suspend fun replaceTokens(chatIdHex: String, userIdHex: String, tokens: Collection) { + deleteTokensForMember(chatIdHex, userIdHex) + insertTokens(tokens.map { ChatMemberSearchTokenEntity(chatIdHex, userIdHex, it) }) + } + + @Query("SELECT token FROM chat_member_search_tokens WHERE chat_id_hex = :chatIdHex AND user_id_hex = :userIdHex") + suspend fun getTokens(chatIdHex: String, userIdHex: String): List + + @Query("SELECT chat_id_hex FROM chat_members WHERE user_id_hex = :userIdHex AND is_member = 1") + suspend fun getChatIdsForMember(userIdHex: String): List + + /** + * Swaps [userIdHex]'s words for [tokens] in every chat they are a member of. For a profile + * written outside a roster write, which changes their name everywhere at once. + */ + @Transaction + suspend fun replaceTokensEverywhere(userIdHex: String, tokens: Collection) { + for (chatIdHex in getChatIdsForMember(userIdHex)) replaceTokens(chatIdHex, userIdHex, tokens) + } + + /** Drops the words of anyone in [chatIdHex]'s index who is no longer one of its members. */ + @Query( + "DELETE FROM chat_member_search_tokens WHERE chat_id_hex = :chatIdHex " + + "AND user_id_hex NOT IN (SELECT user_id_hex FROM chat_members WHERE chat_id_hex = :chatIdHex AND is_member = 1)" + ) + suspend fun deleteTokensOfFormerMembers(chatIdHex: String) + + @Query("DELETE FROM chat_member_search_tokens WHERE chat_id_hex = :chatIdHex") + suspend fun deleteTokensForChat(chatIdHex: String) + + @Query("DELETE FROM chat_member_search_tokens") + suspend fun deleteAllTokens() + + // endregion + + // region search + + /** + * Members of [chatIdHex] other than [selfIdHex] with a word in [lower, upper), each with when + * they last spoke among the chat's newest [recentWindow] held messages (null if they did not). + * + * Unordered: ranking folds diacritics the way the index does, which SQLite's collations cannot, + * so it happens in Kotlin. The token subquery is a range scan on + * `index_chat_member_search_tokens_chat_id_hex_token`; the recent-speaker subquery reads the + * newest rows of `index_chat_messages_chat_id_hex_timestamp_epoch_ms`. + */ + @Query( + """ + SELECT m.user_id_hex AS user_id_hex, p.display_name AS display_name, p.username AS username, + p.profile_picture_json AS profile_picture_json, r.last_spoke_epoch_ms AS last_spoke_epoch_ms + FROM chat_members m + LEFT JOIN user_profiles p ON p.user_id_hex = m.user_id_hex + LEFT JOIN ( + SELECT sender_id_hex, MAX(timestamp_epoch_ms) AS last_spoke_epoch_ms FROM ( + SELECT sender_id_hex, timestamp_epoch_ms FROM chat_messages + WHERE chat_id_hex = :chatIdHex AND sender_id_hex IS NOT NULL + ORDER BY timestamp_epoch_ms DESC LIMIT :recentWindow + ) GROUP BY sender_id_hex + ) r ON r.sender_id_hex = m.user_id_hex + WHERE m.chat_id_hex = :chatIdHex AND m.is_member = 1 AND m.user_id_hex != :selfIdHex + AND m.user_id_hex IN ( + SELECT user_id_hex FROM chat_member_search_tokens + WHERE chat_id_hex = :chatIdHex AND token >= :lower AND token < :upper + ) + """ + ) + suspend fun searchByTokenRange( + chatIdHex: String, + selfIdHex: String, + lower: String, + upper: String, + recentWindow: Int, + ): List + + /** Current members of [chatIdHex] other than [selfIdHex] who sent one of its newest [recentWindow] held messages. */ + @Query( + """ + SELECT m.user_id_hex AS user_id_hex, p.display_name AS display_name, p.username AS username, + p.profile_picture_json AS profile_picture_json, r.last_spoke_epoch_ms AS last_spoke_epoch_ms + FROM chat_members m + INNER JOIN ( + SELECT sender_id_hex, MAX(timestamp_epoch_ms) AS last_spoke_epoch_ms FROM ( + SELECT sender_id_hex, timestamp_epoch_ms FROM chat_messages + WHERE chat_id_hex = :chatIdHex AND sender_id_hex IS NOT NULL + ORDER BY timestamp_epoch_ms DESC LIMIT :recentWindow + ) GROUP BY sender_id_hex + ) r ON r.sender_id_hex = m.user_id_hex + LEFT JOIN user_profiles p ON p.user_id_hex = m.user_id_hex + WHERE m.chat_id_hex = :chatIdHex AND m.is_member = 1 AND m.user_id_hex != :selfIdHex + """ + ) + suspend fun recentSpeakers(chatIdHex: String, selfIdHex: String, recentWindow: Int): List + + // endregion + + // region roster sync + + @Query("SELECT COUNT(*) FROM chat_members WHERE chat_id_hex = :chatIdHex AND is_member = 1") + suspend fun countMembers(chatIdHex: String): Int + + /** Members of [chatIdHex] who joined at or before [version]: the ones a full read may drop. */ + @Query("SELECT user_id_hex FROM chat_members WHERE chat_id_hex = :chatIdHex AND is_member = 1 AND version <= :version") + suspend fun getMemberIdsJoinedBy(chatIdHex: String, version: Long): List + + @Query("SELECT * FROM chat_roster_sync WHERE chat_id_hex = :chatIdHex LIMIT 1") + suspend fun getSyncState(chatIdHex: String): ChatRosterSyncEntity? + + @Insert(onConflict = OnConflictStrategy.REPLACE) + suspend fun upsertSyncState(state: ChatRosterSyncEntity) + + @Insert(onConflict = OnConflictStrategy.IGNORE) + suspend fun insertSyncStateIfAbsent(state: ChatRosterSyncEntity) + + @Query("UPDATE chat_roster_sync SET watermark = :watermark WHERE chat_id_hex = :chatIdHex") + suspend fun setWatermark(chatIdHex: String, watermark: Long) + + /** + * Moves [chatIdHex]'s watermark from [from] to [to], and only from [from]: a stream change + * applied on top of a roster the device held in full keeps it in full. From anywhere else the + * watermark stays put, and the next catch-up reads down to it. + */ + @Query("UPDATE chat_roster_sync SET watermark = :to WHERE chat_id_hex = :chatIdHex AND watermark = :from") + suspend fun advanceWatermark(chatIdHex: String, from: Long, to: Long) + + @Query("UPDATE chat_roster_sync SET reconcile_pending = 1 WHERE chat_id_hex = :chatIdHex") + suspend fun setReconcilePending(chatIdHex: String) + + @Query("DELETE FROM chat_roster_sync WHERE chat_id_hex = :chatIdHex") + suspend fun deleteSyncState(chatIdHex: String) + + @Query("DELETE FROM chat_roster_sync") + suspend fun deleteAllSyncState() + + // endregion +} + +/** A member a search found, with what ranking needs to order it. */ +data class MemberSearchRow( + @ColumnInfo(name = "user_id_hex") val userIdHex: String, + // Null when the member's profile has not been written yet. + @ColumnInfo(name = "display_name") val displayName: String?, + @ColumnInfo(name = "username") val username: String?, + @ColumnInfo(name = "profile_picture_json") val profilePicture: MediaItem?, + // Null when the member sent none of the chat's newest held messages. + @ColumnInfo(name = "last_spoke_epoch_ms") val lastSpokeEpochMs: Long?, +) 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..5e56fc02d7 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 @@ -51,6 +51,10 @@ interface ChatMetadataDao { @Query("SELECT roster_version FROM chat_metadata WHERE chat_id_hex = :chatIdHex") suspend fun getRosterVersion(chatIdHex: String): Long? + /** The true size of [chatIdHex]'s roster, or null when the chat is not stored. */ + @Query("SELECT member_count FROM chat_metadata WHERE chat_id_hex = :chatIdHex") + suspend fun getMemberCount(chatIdHex: String): Long? + @Query("SELECT * FROM chat_metadata WHERE chat_id_hex = :chatIdHex") suspend fun getById(chatIdHex: String): ChatMetadataEntity? diff --git a/apps/flipcash/shared/persistence/db/src/main/kotlin/com/flipcash/app/persistence/entities/ChatMemberEntity.kt b/apps/flipcash/shared/persistence/db/src/main/kotlin/com/flipcash/app/persistence/entities/ChatMemberEntity.kt index 7d89127e69..94221accb9 100644 --- a/apps/flipcash/shared/persistence/db/src/main/kotlin/com/flipcash/app/persistence/entities/ChatMemberEntity.kt +++ b/apps/flipcash/shared/persistence/db/src/main/kotlin/com/flipcash/app/persistence/entities/ChatMemberEntity.kt @@ -14,6 +14,16 @@ data class ChatMemberEntity( @ColumnInfo(name = "chat_id_hex") val chatIdHex: String, @ColumnInfo(name = "user_id_hex") val userIdHex: String, @ColumnInfo(name = "pointers_json") val pointersJson: List?, + // The roster version the member joined at: `Member.version`. The merge key across roster pages + // and stream updates (greater wins), and what a full roster read checks before dropping a + // member it did not see. Zero for chat-creation joins, DM participants, and rows from before + // it was stored. + @ColumnInfo(name = "version", defaultValue = "0") val version: Long = 0, + // False for a marker left by a `MemberLeft`, with [version] set to the leave's roster version. + // Kept rather than deleted so a roster page that trails the stream cannot re-add the member: + // the greater version wins, and only a rejoin carries one above the leave. Every read of this + // table skips markers; a complete roster read clears those at or below its version. + @ColumnInfo(name = "is_member", defaultValue = "1") val isMember: Boolean = true, ) /** diff --git a/apps/flipcash/shared/persistence/db/src/main/kotlin/com/flipcash/app/persistence/entities/ChatMemberSearchTokenEntity.kt b/apps/flipcash/shared/persistence/db/src/main/kotlin/com/flipcash/app/persistence/entities/ChatMemberSearchTokenEntity.kt new file mode 100644 index 0000000000..8b56d0206e --- /dev/null +++ b/apps/flipcash/shared/persistence/db/src/main/kotlin/com/flipcash/app/persistence/entities/ChatMemberSearchTokenEntity.kt @@ -0,0 +1,31 @@ +package com.flipcash.app.persistence.entities + +import androidx.room.ColumnInfo +import androidx.room.Entity +import androidx.room.Index + +/** + * One searchable word of a chat member's name or handle, already normalized. + * + * Keyed per chat rather than per user, though a user's tokens are the same in every chat they are + * in: the lookup it serves is "members of this chat whose word starts with q", and with the chat + * as the leading column of [Index] that is a single range scan. Joining a per-user token table to + * `chat_members` would scan every user matching the prefix across all chats first. + * + * Queried as `token >= q AND token < q || U+10FFFF`, which the index answers without a table scan. + */ +@Entity( + tableName = "chat_member_search_tokens", + primaryKeys = ["chat_id_hex", "user_id_hex", "token"], + indices = [ + Index( + value = ["chat_id_hex", "token"], + name = "index_chat_member_search_tokens_chat_id_hex_token", + ), + ], +) +data class ChatMemberSearchTokenEntity( + @ColumnInfo(name = "chat_id_hex") val chatIdHex: String, + @ColumnInfo(name = "user_id_hex") val userIdHex: String, + @ColumnInfo(name = "token") val token: String, +) diff --git a/apps/flipcash/shared/persistence/db/src/main/kotlin/com/flipcash/app/persistence/entities/ChatRosterSyncEntity.kt b/apps/flipcash/shared/persistence/db/src/main/kotlin/com/flipcash/app/persistence/entities/ChatRosterSyncEntity.kt new file mode 100644 index 0000000000..f344de07f3 --- /dev/null +++ b/apps/flipcash/shared/persistence/db/src/main/kotlin/com/flipcash/app/persistence/entities/ChatRosterSyncEntity.kt @@ -0,0 +1,30 @@ +package com.flipcash.app.persistence.entities + +import androidx.room.ColumnInfo +import androidx.room.Entity +import androidx.room.PrimaryKey + +/** + * How far the device's copy of a group's roster can be trusted. No row means nothing is known: + * the roster has never been read to the end, and no watermark has been set. + * + * A table of its own rather than columns on `chat_metadata`, which is written whole from every + * feed payload: a column there would be reset by the next feed sync. For the same reason the + * watermark is not `chat_metadata.roster_version`, which a feed payload can move past changes the + * device never applied. + */ +@Entity(tableName = "chat_roster_sync") +data class ChatRosterSyncEntity( + @PrimaryKey @ColumnInfo(name = "chat_id_hex") val chatIdHex: String, + // The last roster version whose joins the device holds. A roster summary above it means joins + // were missed, and a read from the top of the roster down to it recovers them. + @ColumnInfo(name = "watermark") val watermark: Long, + // A read of the whole roster has finished at least once. + @ColumnInfo(name = "fully_synced", defaultValue = "0") val fullySynced: Boolean = false, + // The last full read stopped at the page cap, so the device holds fewer members than the + // roster has. Not a reason to read again: the next read would stop at the same place. + @ColumnInfo(name = "truncated", defaultValue = "0") val truncated: Boolean = false, + // Someone left while the device was not listening, and only a full read can say who. Set until + // that read finishes; search keeps the departed member until then. + @ColumnInfo(name = "reconcile_pending", defaultValue = "0") val reconcilePending: Boolean = false, +) diff --git a/apps/flipcash/shared/persistence/db/src/test/kotlin/com/flipcash/app/persistence/RosterSearchMigrationTest.kt b/apps/flipcash/shared/persistence/db/src/test/kotlin/com/flipcash/app/persistence/RosterSearchMigrationTest.kt new file mode 100644 index 0000000000..b58df84580 --- /dev/null +++ b/apps/flipcash/shared/persistence/db/src/test/kotlin/com/flipcash/app/persistence/RosterSearchMigrationTest.kt @@ -0,0 +1,152 @@ +package com.flipcash.app.persistence + +import android.content.Context +import androidx.room.Room +import androidx.room.testing.MigrationTestHelper +import androidx.sqlite.db.SupportSQLiteDatabase +import androidx.test.core.app.ApplicationProvider +import androidx.test.platform.app.InstrumentationRegistry +import kotlinx.coroutines.runBlocking +import org.junit.After +import org.junit.Rule +import org.junit.Test +import org.junit.runner.RunWith +import org.robolectric.RobolectricTestRunner +import kotlin.test.assertEquals +import kotlin.test.assertTrue + +/** + * 40 -> 41 is an AutoMigration: it adds `version` and `is_member` to chat_members and creates + * chat_member_search_tokens and chat_roster_sync. The schemas reach [MigrationTestHelper] as + * unit-test assets (see this module's build.gradle.kts). + */ +@RunWith(RobolectricTestRunner::class) +class RosterSearchMigrationTest { + + private val context = ApplicationProvider.getApplicationContext() + + @get:Rule + val helper = MigrationTestHelper( + InstrumentationRegistry.getInstrumentation(), + FlipcashDatabase::class.java, + ) + + private var room: FlipcashDatabase? = null + + @After + fun tearDown() { + room?.close() + context.deleteDatabase(DB) + } + + @Test + fun `migration keeps rows and matches the v41 schema`() { + helper.createDatabase(DB, 40).use { it.seedV40() } + + val db = helper.runMigrationsAndValidate(DB, 41, true) + + assertEquals(SEEDED_COUNTS, db.counts()) + db.query("SELECT text FROM messages").use { c -> + c.moveToFirst() + assertEquals("legacy", c.getString(0)) + } + db.query("SELECT pointers_json FROM chat_members WHERE user_id_hex = 'b0'").use { c -> + c.moveToFirst() + assertEquals("""[]""", c.getString(0)) + } + } + + @Test + fun `existing members take the column defaults`() { + helper.createDatabase(DB, 40).use { it.seedV40() } + + val db = helper.runMigrationsAndValidate(DB, 41, true) + + db.query("SELECT user_id_hex, version, is_member FROM chat_members ORDER BY user_id_hex").use { c -> + val rows = buildList { + while (c.moveToNext()) add(Triple(c.getString(0), c.getLong(1), c.getInt(2))) + } + assertEquals(listOf(Triple("a0", 0L, 1), Triple("b0", 0L, 1), Triple("c0", 0L, 1)), rows) + } + } + + @Test + fun `new tables exist and are empty`() { + helper.createDatabase(DB, 40).use { it.seedV40() } + + val db = helper.runMigrationsAndValidate(DB, 41, true) + + assertEquals(0, db.count("chat_member_search_tokens")) + // No sync row means the first open of each group does a full roster read. + assertEquals(0, db.count("chat_roster_sync")) + db.query("PRAGMA index_info(index_chat_member_search_tokens_chat_id_hex_token)").use { c -> + val columns = buildList { while (c.moveToNext()) add(c.getString(c.getColumnIndexOrThrow("name"))) } + assertEquals(listOf("chat_id_hex", "token"), columns) + } + } + + @Test + fun `the app's Room builder opens a v40 database without wiping it`() = runBlocking { + helper.createDatabase(DB, 40).use { it.seedV40() } + + // Same builder as FlipcashDatabase.init. It falls back to a destructive migration, + // so a broken AutoMigration would open an empty database instead of throwing. + val opened = Room.databaseBuilder(context, FlipcashDatabase::class.java, DB) + .addMigrations(FlipcashDatabase.MIGRATION_25_26) + .fallbackToDestructiveMigration() + .build() + .also { room = it } + + val sqlite = opened.openHelper.writableDatabase + assertEquals(41, sqlite.version) + assertEquals(SEEDED_COUNTS, sqlite.counts()) + + val search = opened.chatMemberSearchDao() + assertEquals(3, search.countMembers(CHAT)) + assertEquals(3, opened.chatMemberDao().getMembersForChat(CHAT).size) + + // Existing members become searchable once the roster read indexes them. + search.replaceTokens(CHAT, "b0", listOf("bob")) + val hits = search.searchByTokenRange(CHAT, selfIdHex = "a0", lower = "bo", upper = "bo􏿿", recentWindow = 50) + assertEquals(listOf("b0"), hits.map { it.userIdHex }) + assertTrue(hits.single().lastSpokeEpochMs != null, "b0 sent a held message") + } + + private fun SupportSQLiteDatabase.seedV40() { + execSQL( + "INSERT INTO chat_metadata (chat_id_hex, chat_type, last_activity_epoch_ms, last_message_id, " + + "title, member_count, roster_version) VALUES ('$CHAT', 'GROUP', 1000, 2, 'Group', 3, 9)" + ) + execSQL( + "INSERT INTO chat_metadata (chat_id_hex, chat_type, last_activity_epoch_ms) " + + "VALUES ('$OTHER_CHAT', 'CONTACT_DM', 500)" + ) + execSQL("INSERT INTO chat_members (chat_id_hex, user_id_hex, pointers_json) VALUES ('$CHAT', 'a0', NULL)") + execSQL("""INSERT INTO chat_members (chat_id_hex, user_id_hex, pointers_json) VALUES ('$CHAT', 'b0', '[]')""") + execSQL("INSERT INTO chat_members (chat_id_hex, user_id_hex, pointers_json) VALUES ('$CHAT', 'c0', NULL)") + execSQL( + "INSERT INTO chat_messages (chat_id_hex, message_id, sender_id_hex, content_json, " + + "timestamp_epoch_ms, unread_seq) VALUES ('$CHAT', 1, 'b0', NULL, 900, 1)" + ) + execSQL( + "INSERT INTO chat_messages (chat_id_hex, message_id, sender_id_hex, content_json, " + + "timestamp_epoch_ms, unread_seq) VALUES ('$CHAT', 2, 'a0', NULL, 1000, 2)" + ) + execSQL( + "INSERT INTO messages (idBase58, text, state, timestamp) VALUES ('m1', 'legacy', 'SENT', 100)" + ) + } + + private fun SupportSQLiteDatabase.count(table: String): Int = + query("SELECT COUNT(*) FROM $table").use { it.moveToFirst(); it.getInt(0) } + + private fun SupportSQLiteDatabase.counts(): Map = + SEEDED_COUNTS.keys.associateWith { count(it) } + + private companion object { + const val DB = "roster-search-migration-test" + const val CHAT = "c1" + const val OTHER_CHAT = "c2" + val SEEDED_COUNTS = mapOf("chat_metadata" to 2, "chat_members" to 3, "chat_messages" to 2, "messages" to 1) + } +} diff --git a/apps/flipcash/shared/persistence/db/src/test/kotlin/com/flipcash/app/persistence/dao/ChatMemberSearchDaoTest.kt b/apps/flipcash/shared/persistence/db/src/test/kotlin/com/flipcash/app/persistence/dao/ChatMemberSearchDaoTest.kt new file mode 100644 index 0000000000..70dab3542e --- /dev/null +++ b/apps/flipcash/shared/persistence/db/src/test/kotlin/com/flipcash/app/persistence/dao/ChatMemberSearchDaoTest.kt @@ -0,0 +1,191 @@ +package com.flipcash.app.persistence.dao + +import android.content.Context +import androidx.room.Room +import androidx.test.core.app.ApplicationProvider +import com.flipcash.app.persistence.FlipcashDatabase +import com.flipcash.app.persistence.entities.ChatMemberEntity +import com.flipcash.app.persistence.entities.ChatMessageEntity +import com.flipcash.app.persistence.entities.ChatRosterSyncEntity +import com.flipcash.app.persistence.entities.UserProfileEntity +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 kotlin.test.assertEquals +import kotlin.test.assertTrue + +/** + * The queries the member search runs. Tokens are written already normalized here; how names become + * tokens is `MemberSearchText`'s concern, tested in the sources module. + */ +@RunWith(RobolectricTestRunner::class) +class ChatMemberSearchDaoTest { + + private lateinit var db: FlipcashDatabase + private lateinit var dao: ChatMemberSearchDao + + @Before + fun setUp() { + val context = ApplicationProvider.getApplicationContext() + db = Room.inMemoryDatabaseBuilder(context, FlipcashDatabase::class.java) + .allowMainThreadQueries() + .build() + dao = db.chatMemberSearchDao() + } + + @After + fun tearDown() { + db.close() + } + + private suspend fun member(userIdHex: String, name: String, vararg tokens: String, chatIdHex: String = CHAT) { + db.userProfileDao().upsertFull(listOf(profile(userIdHex, name))) + db.chatMemberDao().upsert(ChatMemberEntity(chatIdHex = chatIdHex, userIdHex = userIdHex, pointersJson = null)) + dao.replaceTokens(chatIdHex, userIdHex, tokens.toList()) + } + + private suspend fun message(id: Long, senderHex: String, at: Long) { + db.chatMessageDao().insert( + ChatMessageEntity( + chatIdHex = CHAT, + messageId = id, + senderIdHex = senderHex, + contentJson = null, + timestampEpochMs = at, + unreadSeq = 0, + ) + ) + } + + private suspend fun search(prefix: String, window: Int = 50) = + dao.searchByTokenRange(CHAT, SELF, prefix, prefix + UPPER, window) + + @Test + fun `a prefix finds members with any word starting with it`() = runTest { + member("a1", "Érica Stone", "erica", "stone") + member("b2", "Mark Sterling", "mark", "sterling") + member("c3", "Alan", "alan") + + assertEquals(setOf("a1", "b2"), search("st").map { it.userIdHex }.toSet()) + assertEquals(listOf("a1"), search("eri").map { it.userIdHex }) + } + + @Test + fun `a member matching on two words comes back once`() = runTest { + member("a1", "Sam Samuels", "sam", "samuels") + + assertEquals(listOf("a1"), search("sam").map { it.userIdHex }) + } + + @Test + fun `the current user is never returned`() = runTest { + member(SELF, "Erin", "erin") + member("a1", "Eric", "eric") + + assertEquals(listOf("a1"), search("er").map { it.userIdHex }) + message(1, SELF, at = 10) + assertEquals(emptyList(), dao.recentSpeakers(CHAT, SELF, 50).map { it.userIdHex }) + } + + @Test + fun `other chats' members are not searched`() = runTest { + member("a1", "Eric", "eric", chatIdHex = "other") + + assertEquals(emptyList(), search("er")) + } + + @Test + fun `last spoke is the sender's newest message within the window`() = runTest { + member("a1", "Eric", "eric") + member("b2", "Erin", "erin") + message(1, "b2", at = 10) + message(2, "a1", at = 20) + message(3, "a1", at = 30) + + val byId = search("er").associate { it.userIdHex to it.lastSpokeEpochMs } + assertEquals(30L, byId["a1"]) + assertEquals(10L, byId["b2"]) + + // A window of the two newest messages leaves b2 out. + val narrow = search("er", window = 2).associate { it.userIdHex to it.lastSpokeEpochMs } + assertEquals(null, narrow["b2"]) + assertEquals(listOf("a1"), dao.recentSpeakers(CHAT, SELF, 2).map { it.userIdHex }) + } + + @Test + fun `a former member is not a recent speaker`() = runTest { + message(1, "gone", at = 10) + + assertEquals(emptyList(), dao.recentSpeakers(CHAT, SELF, 50)) + } + + @Test + fun `a renamed member's old words stop matching`() = runTest { + member("a1", "Eric", "eric") + dao.replaceTokens(CHAT, "a1", listOf("frank")) + + assertEquals(emptyList(), search("er")) + assertEquals(listOf("a1"), search("fr").map { it.userIdHex }) + } + + @Test + fun `a profile rewrite reaches every chat the member is in`() = runTest { + member("a1", "Eric", "eric") + member("a1", "Eric", "eric", chatIdHex = "other") + + dao.replaceTokensEverywhere("a1", listOf("frank")) + + assertEquals(listOf("frank"), dao.getTokens(CHAT, "a1")) + assertEquals(listOf("frank"), dao.getTokens("other", "a1")) + } + + @Test + fun `the prefix lookup is an index range scan`() { + val plan = db.openHelper.readableDatabase.query( + "EXPLAIN QUERY PLAN SELECT user_id_hex FROM chat_member_search_tokens " + + "WHERE chat_id_hex = ? AND token >= ? AND token < ?", + arrayOf(CHAT, "er", "er$UPPER"), + ).use { c -> + buildList { while (c.moveToNext()) add(c.getString(c.getColumnIndexOrThrow("detail"))) } + }.joinToString("\n") + + // A SEARCH bounded on both columns, not a SCAN of the table or the index. + assertTrue( + Regex("SEARCH .*INDEX index_chat_member_search_tokens_chat_id_hex_token \\(chat_id_hex=\\? AND token>\\? AND token<\\?\\)") + .containsMatchIn(plan), + plan, + ) + } + + @Test + fun `the stream moves the watermark only from where it stands`() = runTest { + dao.upsertSyncState(ChatRosterSyncEntity(chatIdHex = CHAT, watermark = 5, fullySynced = true)) + + dao.advanceWatermark(CHAT, from = 5, to = 6) + assertEquals(6, dao.getSyncState(CHAT)?.watermark) + + // Out of step: the device missed something between 6 and 8, so 9 does not vouch for it. + dao.advanceWatermark(CHAT, from = 8, to = 9) + assertEquals(6, dao.getSyncState(CHAT)?.watermark) + } + + private fun profile(userIdHex: String, name: String) = UserProfileEntity( + userIdHex = userIdHex, + displayName = name, + phoneValue = null, + phoneVerified = null, + emailValue = null, + emailVerified = null, + socialAccounts = null, + profilePicture = null, + ) + + private companion object { + const val CHAT = "c0ffee" + const val SELF = "5e1f" + const val UPPER = "􏿿" // U+10FFFF, as MemberSearchText.UPPER_BOUND + } +} diff --git a/apps/flipcash/shared/persistence/sources/build.gradle.kts b/apps/flipcash/shared/persistence/sources/build.gradle.kts index 744be06cbf..6e04a79dd7 100644 --- a/apps/flipcash/shared/persistence/sources/build.gradle.kts +++ b/apps/flipcash/shared/persistence/sources/build.gradle.kts @@ -17,6 +17,7 @@ dependencies { testImplementation(kotlin("test")) testImplementation(libs.bundles.unit.testing) testImplementation(testFixtures(project(":services:flipcash"))) + testImplementation(libs.robolectric) implementation(libs.bundles.kotlinx.serialization) diff --git a/apps/flipcash/shared/persistence/sources/src/main/kotlin/com/flipcash/app/persistence/sources/BlockedUserDataSource.kt b/apps/flipcash/shared/persistence/sources/src/main/kotlin/com/flipcash/app/persistence/sources/BlockedUserDataSource.kt index ef52898c90..5617f041d8 100644 --- a/apps/flipcash/shared/persistence/sources/src/main/kotlin/com/flipcash/app/persistence/sources/BlockedUserDataSource.kt +++ b/apps/flipcash/shared/persistence/sources/src/main/kotlin/com/flipcash/app/persistence/sources/BlockedUserDataSource.kt @@ -4,6 +4,7 @@ import androidx.paging.PagingSource import androidx.paging.PagingState import androidx.room.withTransaction import com.flipcash.app.persistence.FlipcashDatabase +import com.flipcash.app.persistence.sources.search.reindexMemberProfile import com.flipcash.app.persistence.entities.BlockedUserWithProfile import com.flipcash.app.persistence.sources.mapper.blocklist.BlockedUserEntityToProfileMapper import com.flipcash.app.persistence.sources.mapper.blocklist.BlockedUserToEntityMapper @@ -72,5 +73,6 @@ class BlockedUserDataSource @Inject constructor( profilePicture = resolved.profile?.profilePicture, username = resolved.profile?.username, ) + reindexMemberProfile(resolved.blocked.userId.hexEncodedString()) } } diff --git a/apps/flipcash/shared/persistence/sources/src/main/kotlin/com/flipcash/app/persistence/sources/ChatMemberDataSource.kt b/apps/flipcash/shared/persistence/sources/src/main/kotlin/com/flipcash/app/persistence/sources/ChatMemberDataSource.kt index b47bba7cb1..1411495c03 100644 --- a/apps/flipcash/shared/persistence/sources/src/main/kotlin/com/flipcash/app/persistence/sources/ChatMemberDataSource.kt +++ b/apps/flipcash/shared/persistence/sources/src/main/kotlin/com/flipcash/app/persistence/sources/ChatMemberDataSource.kt @@ -3,6 +3,7 @@ package com.flipcash.app.persistence.sources import androidx.room.withTransaction import com.flipcash.app.persistence.FlipcashDatabase import com.flipcash.app.persistence.sources.mapper.chat.ChatEntityMapper +import com.flipcash.app.persistence.sources.search.MemberSearchText import com.flipcash.services.models.chat.ChatId import com.flipcash.services.models.chat.ChatMember import com.flipcash.services.models.chat.ChatType @@ -71,6 +72,7 @@ class ChatMemberDataSource @Inject constructor( database.withTransaction { database.userProfileDao().upsertMembers(profileRows(members)) database.chatMemberDao().upsert(members.map { mapper.toEntity(hex, it) }) + index(database, hex, members) } } @@ -111,26 +113,92 @@ class ChatMemberDataSource @Inject constructor( chatIdHex = hex, keepUserIdHexes = members.map { mapper.userIdHex(it.userId) }, ) + index(database, hex, members) + database.chatMemberSearchDao().deleteTokensOfFormerMembers(hex) } } + /** How many of [chatId]'s members the device holds. For a group, compare with its `member_count`. */ + suspend fun countMembers(chatId: ChatId): Int = + db?.chatMemberSearchDao()?.countMembers(mapper.chatIdHex(chatId)) ?: 0 + /** - * Drops one member of [chatId]. What a `MemberLeft` roster change applies — [replaceMembers] - * cannot, because a roster change names who left rather than who remains. + * Drops the members of [chatId] a full roster read shows have left: held, not in [seen], and + * joined at or before [readVersion], the roster version the read described. A member who + * joined after it is newer than the read, not gone, and stays. + * + * Done a member at a time rather than as one `NOT IN` list: a large group's roster runs past + * the 999 bound variables the minSdk SQLite build allows in a statement. */ - suspend fun deleteMember(chatId: ChatId, userId: ID) { - db?.chatMemberDao()?.deleteMember( - chatIdHex = mapper.chatIdHex(chatId), - userIdHex = mapper.userIdHex(userId), - ) + suspend fun reconcile(chatId: ChatId, seen: Set, readVersion: Long) { + val database = db ?: return + val hex = mapper.chatIdHex(chatId) + val seenHexes = seen.mapTo(HashSet()) { mapper.userIdHex(it) } + database.withTransaction { + val departed = database.chatMemberSearchDao().getMemberIdsJoinedBy(hex, readVersion) + .filterNot { it in seenHexes } + for (userIdHex in departed) { + database.chatMemberDao().deleteMember(hex, userIdHex) + database.chatMemberSearchDao().deleteTokensForMember(hex, userIdHex) + } + // The read lists everyone in the roster as of [readVersion], so a leave at or below it + // no longer needs its marker. + database.chatMemberDao().deleteMarkersUpTo(hex, readVersion) + } + } + + /** + * Applies a `MemberLeft` at roster [version]: [userId] leaves search and the count, and a + * roster page that trails the stream cannot add them back (model.proto: a trailing page + * "cannot resurrect" a removed member). [replaceMembers] cannot do this, because a roster + * change names who left rather than who remains. + */ + suspend fun markLeft(chatId: ChatId, userId: ID, version: Long) { + val database = db ?: return + val hex = mapper.chatIdHex(chatId) + val userIdHex = mapper.userIdHex(userId) + database.withTransaction { + if (database.chatMemberDao().markLeft(chatIdHex = hex, userIdHex = userIdHex, version = version)) { + database.chatMemberSearchDao().deleteTokensForMember(chatIdHex = hex, userIdHex = userIdHex) + } + } } suspend fun deleteForChat(chatId: ChatId) { - db?.chatMemberDao()?.deleteForChat(mapper.chatIdHex(chatId)) + val database = db ?: return + val hex = mapper.chatIdHex(chatId) + database.withTransaction { + database.chatMemberDao().deleteForChat(hex) + database.chatMemberSearchDao().deleteTokensForChat(hex) + // The roster read is only complete for the members it wrote. + database.chatMemberSearchDao().deleteSyncState(hex) + } } suspend fun clear() { - db?.chatMemberDao()?.deleteAll() + val database = db ?: return + database.withTransaction { + database.chatMemberDao().deleteAll() + database.chatMemberSearchDao().deleteAllTokens() + database.chatMemberSearchDao().deleteAllSyncState() + } + } + + /** + * Rewrites the search tokens of [members] in [chatIdHex] from their stored profiles. + * + * Read back after the profile write rather than taken from [members]: a member can arrive with + * no profile (see [profileRows]), and is then indexed by the one already stored, if any. + */ + private suspend fun index(database: FlipcashDatabase, chatIdHex: String, members: List) { + val search = database.chatMemberSearchDao() + for (member in members) { + val userIdHex = mapper.userIdHex(member.userId) + // A member the upsert kept out, having left at a later version, gets no words. + if (database.chatMemberDao().getMember(chatIdHex, userIdHex)?.isMember != true) continue + val profile = database.userProfileDao().getByUserId(userIdHex) ?: continue + search.replaceTokens(chatIdHex, userIdHex, MemberSearchText.tokens(profile.displayName, profile.username)) + } } /** 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..99eeeb845c 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 @@ -59,6 +59,10 @@ class ChatMetadataDataSource @Inject constructor( suspend fun getRosterVersion(chatId: ChatId): Long = db?.chatMetadataDao()?.getRosterVersion(mapper.chatIdHex(chatId)) ?: 0L + /** [chatId]'s true roster size; 0 when the chat is not stored yet. */ + suspend fun getMemberCount(chatId: ChatId): Long = + db?.chatMetadataDao()?.getMemberCount(mapper.chatIdHex(chatId)) ?: 0L + suspend fun updateRoster(chatId: ChatId, memberCount: Long, rosterVersion: Long) { db?.chatMetadataDao()?.updateRosterIfNewer( chatIdHex = mapper.chatIdHex(chatId), diff --git a/apps/flipcash/shared/persistence/sources/src/main/kotlin/com/flipcash/app/persistence/sources/ChatRosterDataSource.kt b/apps/flipcash/shared/persistence/sources/src/main/kotlin/com/flipcash/app/persistence/sources/ChatRosterDataSource.kt new file mode 100644 index 0000000000..05b1ac8467 --- /dev/null +++ b/apps/flipcash/shared/persistence/sources/src/main/kotlin/com/flipcash/app/persistence/sources/ChatRosterDataSource.kt @@ -0,0 +1,128 @@ +package com.flipcash.app.persistence.sources + +import com.flipcash.app.persistence.FlipcashDatabase +import com.flipcash.app.persistence.dao.MemberSearchRow +import com.flipcash.app.persistence.entities.ChatRosterSyncEntity +import com.flipcash.app.persistence.sources.mapper.chat.ChatEntityMapper +import com.flipcash.app.persistence.sources.search.MemberSearchText +import com.flipcash.services.models.chat.ChatId +import com.flipcash.services.models.chat.MediaItem +import com.getcode.opencode.model.core.ID +import javax.inject.Inject +import javax.inject.Singleton + +/** + * Reads a chat's member search index, and records how much of each group's roster has been read + * into it. The index itself is written by [ChatMemberDataSource] alongside the members. + */ +@Singleton +class ChatRosterDataSource @Inject constructor( + private val mapper: ChatEntityMapper, +) { + + private val db: FlipcashDatabase? + get() = FlipcashDatabase.getInstance() + + /** + * Whether the signed-in user's database is open. A worker can start in a process that has not + * opened it yet, and a read with nowhere to go should be retried rather than recorded. + */ + val isAvailable: Boolean + get() = db != null + + /** + * Members of [chatId] other than [selfId] with a name or handle word starting with [prefix], + * which must already be normalized with [MemberSearchText.normalize]. Unordered. + */ + suspend fun searchByPrefix( + chatId: ChatId, + selfId: ID, + prefix: String, + recentWindow: Int, + ): List { + val dao = db?.chatMemberSearchDao() ?: return emptyList() + return dao.searchByTokenRange( + chatIdHex = mapper.chatIdHex(chatId), + selfIdHex = mapper.userIdHex(selfId), + lower = prefix, + upper = prefix + MemberSearchText.UPPER_BOUND, + recentWindow = recentWindow, + ).map { it.toCandidate() } + } + + /** Members of [chatId] other than [selfId] who sent one of its newest [recentWindow] held messages. */ + suspend fun recentSpeakers(chatId: ChatId, selfId: ID, recentWindow: Int): List { + val dao = db?.chatMemberSearchDao() ?: return emptyList() + return dao.recentSpeakers( + chatIdHex = mapper.chatIdHex(chatId), + selfIdHex = mapper.userIdHex(selfId), + recentWindow = recentWindow, + ).map { it.toCandidate() } + } + + suspend fun getSyncState(chatId: ChatId): RosterSyncState? = + db?.chatMemberSearchDao()?.getSyncState(mapper.chatIdHex(chatId))?.let { + RosterSyncState( + watermark = it.watermark, + fullySynced = it.fullySynced, + truncated = it.truncated, + reconcilePending = it.reconcilePending, + ) + } + + /** + * Records a full read of [chatId]'s roster that ended with [watermark] as the version it can + * vouch for, clearing any pending reconcile. + */ + suspend fun markFullySynced(chatId: ChatId, watermark: Long, truncated: Boolean) { + db?.chatMemberSearchDao()?.upsertSyncState( + ChatRosterSyncEntity( + chatIdHex = mapper.chatIdHex(chatId), + watermark = watermark, + fullySynced = true, + truncated = truncated, + reconcilePending = false, + ) + ) + } + + /** Sets [chatId]'s watermark after a catch-up read confirmed the roster at [watermark]. */ + suspend fun setWatermark(chatId: ChatId, watermark: Long) { + db?.chatMemberSearchDao()?.setWatermark(mapper.chatIdHex(chatId), watermark) + } + + /** Moves [chatId]'s watermark to [to] only if it is at [from]; see the DAO for why. */ + suspend fun advanceWatermark(chatId: ChatId, from: Long, to: Long) { + db?.chatMemberSearchDao()?.advanceWatermark(mapper.chatIdHex(chatId), from, to) + } + + /** Flags [chatId] as holding members who may have left, until a full read settles it. */ + suspend fun markReconcilePending(chatId: ChatId) { + db?.chatMemberSearchDao()?.setReconcilePending(mapper.chatIdHex(chatId)) + } + + private fun MemberSearchRow.toCandidate() = RosterSearchCandidate( + userId = mapper.userIdFromHex(userIdHex), + displayName = displayName.orEmpty(), + username = username, + profilePicture = profilePicture, + lastSpokeEpochMs = lastSpokeEpochMs, + ) +} + +/** A member a roster search found, before ranking. */ +data class RosterSearchCandidate( + val userId: ID, + val displayName: String, + val username: String?, + val profilePicture: MediaItem?, + // When they last sent one of the chat's newest held messages; null if they sent none of them. + val lastSpokeEpochMs: Long?, +) + +data class RosterSyncState( + val watermark: Long, + val fullySynced: Boolean, + val truncated: Boolean, + val reconcilePending: Boolean, +) diff --git a/apps/flipcash/shared/persistence/sources/src/main/kotlin/com/flipcash/app/persistence/sources/UserProfileDataSource.kt b/apps/flipcash/shared/persistence/sources/src/main/kotlin/com/flipcash/app/persistence/sources/UserProfileDataSource.kt index 1c2b07b3ff..df285a809f 100644 --- a/apps/flipcash/shared/persistence/sources/src/main/kotlin/com/flipcash/app/persistence/sources/UserProfileDataSource.kt +++ b/apps/flipcash/shared/persistence/sources/src/main/kotlin/com/flipcash/app/persistence/sources/UserProfileDataSource.kt @@ -1,7 +1,9 @@ package com.flipcash.app.persistence.sources +import androidx.room.withTransaction import com.flipcash.app.persistence.FlipcashDatabase import com.flipcash.app.persistence.entities.toSerialized +import com.flipcash.app.persistence.sources.search.reindexMemberProfile import com.flipcash.app.persistence.sources.mapper.toDomain import com.flipcash.services.models.UserProfile import com.getcode.opencode.model.core.ID @@ -44,14 +46,21 @@ class UserProfileDataSource @Inject constructor() { * Stores [profile]'s name, avatar and handle for [userId] (INSERT OR REPLACE, preserving any * existing phone/email/social columns). Used to back-fill the cache after a network resolve so * [observeProfiles] re-emits and consumers (e.g. the transaction list) resolve the row live. + * + * Also rewrites the user's member search tokens, so a renamed member is found by their new name. */ suspend fun store(userId: ID, profile: UserProfile) { - db?.userProfileDao()?.upsertNameAndAvatar( - userIdHex = userId.hexEncodedString(), - displayName = profile.displayName, - profilePicture = profile.profilePicture, - username = profile.username, - ) + val database = db ?: return + val userIdHex = userId.hexEncodedString() + database.withTransaction { + database.userProfileDao().upsertNameAndAvatar( + userIdHex = userIdHex, + displayName = profile.displayName, + profilePicture = profile.profilePicture, + username = profile.username, + ) + database.reindexMemberProfile(userIdHex) + } } /** diff --git a/apps/flipcash/shared/persistence/sources/src/main/kotlin/com/flipcash/app/persistence/sources/mapper/chat/ChatEntityMapper.kt b/apps/flipcash/shared/persistence/sources/src/main/kotlin/com/flipcash/app/persistence/sources/mapper/chat/ChatEntityMapper.kt index 212742203b..1ebb920184 100644 --- a/apps/flipcash/shared/persistence/sources/src/main/kotlin/com/flipcash/app/persistence/sources/mapper/chat/ChatEntityMapper.kt +++ b/apps/flipcash/shared/persistence/sources/src/main/kotlin/com/flipcash/app/persistence/sources/mapper/chat/ChatEntityMapper.kt @@ -268,6 +268,7 @@ class ChatEntityMapper @Inject constructor() { chatIdHex = chatIdHex, userIdHex = member.userId.hexEncodedString(), pointersJson = member.pointers.map { it.toSerialized() }, + version = member.version, ) } @@ -296,6 +297,7 @@ class ChatEntityMapper @Inject constructor() { userId = relation.member.userIdHex.hexToId(), userProfile = relation.profile?.toSerialized()?.toDomain() ?: UserProfile.Empty, pointers = relation.member.pointersJson?.map { it.toDomain() } ?: emptyList(), + version = relation.member.version, ) } diff --git a/apps/flipcash/shared/persistence/sources/src/main/kotlin/com/flipcash/app/persistence/sources/search/MemberSearchText.kt b/apps/flipcash/shared/persistence/sources/src/main/kotlin/com/flipcash/app/persistence/sources/search/MemberSearchText.kt new file mode 100644 index 0000000000..b771014855 --- /dev/null +++ b/apps/flipcash/shared/persistence/sources/src/main/kotlin/com/flipcash/app/persistence/sources/search/MemberSearchText.kt @@ -0,0 +1,64 @@ +package com.flipcash.app.persistence.sources.search + +import com.flipcash.app.persistence.FlipcashDatabase +import java.text.Normalizer +import java.util.Locale + +/** + * How member names and handles become search tokens, and how a query is folded to meet them. + * + * Both sides go through [normalize], so matching is case- and diacritic-insensitive: "eri" finds + * "Érica". iOS folds the same way, in the same order, so a query matches the same members on both + * platforms: + * + * 1. Unicode NFKD, which splits "É" into "E" plus a combining acute and folds compatibility forms + * such as full-width letters and ligatures. + * 2. Lowercase, locale-independent. After NFKD, not before: "İ" lowercases to "i" plus a combining + * dot, which step 3 then removes. + * 3. Drop every nonspacing combining mark (general category Mn). + */ +object MemberSearchText { + + /** + * Appended to a prefix to bound its range from above. The highest code point, so every token + * that starts with the prefix sorts below `prefix + UPPER_BOUND` under SQLite's byte-wise + * comparison, including tokens whose next character lies outside the Basic Multilingual Plane + * (an emoji), which a U+FFFF bound would leave out. + */ + const val UPPER_BOUND: String = "􏿿" + + private val combiningMarks = Regex("\\p{Mn}+") + private val whitespace = Regex("[\\s\\p{Z}]+") + + fun normalize(text: String): String = + combiningMarks.replace(Normalizer.normalize(text, Normalizer.Form.NFKD).lowercase(Locale.ROOT), "") + + /** [text]'s normalized words, split on whitespace. */ + fun words(text: String): List = + normalize(text).split(whitespace).filter { it.isNotEmpty() } + + /** The tokens a member is found by: each word of [displayName], and [username] without its `@`. */ + fun tokens(displayName: String?, username: String?): Set = buildSet { + displayName?.let { addAll(words(it)) } + username?.removePrefix("@")?.let { addAll(words(it)) } + } + + /** [query]'s words, each a prefix to match. A leading `@` on a word is the mention trigger, not part of it. */ + fun queryWords(query: String): List = + words(query).map { it.removePrefix("@") }.filter { it.isNotEmpty() } +} + +/** + * Re-derives [userIdHex]'s search tokens from their stored profile, in every chat they are in. + * + * For a profile write that does not go through a roster write, such as a sender resolved on its + * own. Read back from the row rather than taken from what was written, because a partial write + * keeps the fields it was not given. + */ +internal suspend fun FlipcashDatabase.reindexMemberProfile(userIdHex: String) { + val profile = userProfileDao().getByUserId(userIdHex) ?: return + chatMemberSearchDao().replaceTokensEverywhere( + userIdHex, + MemberSearchText.tokens(profile.displayName, profile.username), + ) +} diff --git a/apps/flipcash/shared/persistence/sources/src/test/kotlin/com/flipcash/app/persistence/sources/MemberSearchIndexTest.kt b/apps/flipcash/shared/persistence/sources/src/test/kotlin/com/flipcash/app/persistence/sources/MemberSearchIndexTest.kt new file mode 100644 index 0000000000..b37d4ade39 --- /dev/null +++ b/apps/flipcash/shared/persistence/sources/src/test/kotlin/com/flipcash/app/persistence/sources/MemberSearchIndexTest.kt @@ -0,0 +1,218 @@ +package com.flipcash.app.persistence.sources + +import android.content.Context +import androidx.test.core.app.ApplicationProvider +import com.flipcash.app.persistence.FlipcashDatabase +import com.flipcash.app.persistence.sources.mapper.chat.ChatEntityMapper +import com.flipcash.app.persistence.sources.search.MemberSearchText +import com.flipcash.services.models.UserProfile +import com.flipcash.services.models.chat.ChatId +import com.flipcash.services.models.chat.ChatMember +import com.flipcash.services.models.chat.MessagePointer +import com.flipcash.services.models.chat.PointerType +import com.getcode.opencode.model.core.ID +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 kotlin.test.assertEquals +import kotlin.test.assertNull +import kotlin.time.Instant + +/** + * The search index follows every member write: a roster read, a join or leave off the stream, and a + * profile resolved on its own. Run against the real database, since the index is kept in step inside + * the same transactions as the member rows. + */ +@RunWith(RobolectricTestRunner::class) +class MemberSearchIndexTest { + + private val context = ApplicationProvider.getApplicationContext() + private val mapper = ChatEntityMapper() + private val members = ChatMemberDataSource(mapper) + private val roster = ChatRosterDataSource(mapper) + private val profiles = UserProfileDataSource() + + @Before + fun setUp() { + FlipcashDatabase.init(context, ENTROPY) + } + + @After + fun tearDown() { + FlipcashDatabase.closeDb() + context.databaseList().forEach { context.deleteDatabase(it) } + } + + private suspend fun search(query: String): Set = + roster.searchByPrefix(CHAT, SELF, MemberSearchText.normalize(query), 50).map { it.userId }.toSet() + + @Test + fun `a joined member is searchable by any name word and by handle`() = runTest { + members.upsert(CHAT, listOf(member(ERICA, "Érica Stone", "estone"))) + + assertEquals(setOf(ERICA), search("eri")) + assertEquals(setOf(ERICA), search("STO")) + assertEquals(setOf(ERICA), search("est")) + assertEquals(emptySet(), search("ica")) + } + + @Test + fun `a member who left is no longer found`() = runTest { + members.upsert(CHAT, listOf(member(ERICA, "Erica"), member(ERIN, "Erin"))) + + members.markLeft(CHAT, ERICA, version = 5) + + assertEquals(setOf(ERIN), search("eri")) + assertEquals(1, members.countMembers(CHAT)) + } + + @Test + fun `a trailing roster page does not bring back a member who left`() = runTest { + members.upsert(CHAT, listOf(member(ERICA, "Erica", version = 3), member(ERIN, "Erin", version = 4))) + members.markLeft(CHAT, ERICA, version = 5) + + // A GetRoster page from before the leave still lists Erica, at the version she joined. + members.upsert(CHAT, listOf(member(ERICA, "Erica", version = 3))) + + assertEquals(setOf(ERIN), search("eri")) + assertEquals(1, members.countMembers(CHAT)) + assertEquals(listOf(ERIN), members.getMembersForChat(CHAT).map { it.userId }) + } + + @Test + fun `a rejoin after a leave brings the member back`() = runTest { + members.upsert(CHAT, listOf(member(ERICA, "Erica", version = 3))) + members.markLeft(CHAT, ERICA, version = 5) + + members.upsert(CHAT, listOf(member(ERICA, "Erica", version = 6))) + + assertEquals(setOf(ERICA), search("eri")) + assertEquals(1, members.countMembers(CHAT)) + } + + @Test + fun `a leave older than the member's join changes nothing`() = runTest { + members.upsert(CHAT, listOf(member(ERICA, "Erica", version = 6))) + + members.markLeft(CHAT, ERICA, version = 5) + + assertEquals(setOf(ERICA), search("eri")) + } + + @Test + fun `a read pointer for a member who left does not bring them back`() = runTest { + members.upsert(CHAT, listOf(member(ERICA, "Erica", version = 3))) + members.markLeft(CHAT, ERICA, version = 5) + + members.updatePointers(CHAT, MessagePointer(PointerType.READ, ERICA, value = 9, timestamp = Instant.fromEpochMilliseconds(0))) + members.upsert(CHAT, listOf(member(ERICA, "Erica", version = 3))) + + assertEquals(emptySet(), search("eri")) + assertEquals(0, members.countMembers(CHAT)) + } + + @Test + fun `a complete read clears leave markers at or below its version`() = runTest { + members.upsert(CHAT, listOf(member(ERICA, "Erica", version = 3), member(ERIN, "Erin", version = 4))) + members.markLeft(CHAT, ERICA, version = 5) + members.markLeft(CHAT, ERIN, version = 8) + + members.reconcile(CHAT, seen = emptySet(), readVersion = 7) + + val dao = FlipcashDatabase.getInstance()!!.chatMemberDao() + val chatHex = mapper.chatIdHex(CHAT) + assertNull(dao.getMember(chatHex, mapper.userIdHex(ERICA))) + // Erin's leave is newer than the read, so her marker still guards against older pages. + assertEquals(false, dao.getMember(chatHex, mapper.userIdHex(ERIN))?.isMember) + } + + @Test + fun `a full read drops members it did not return`() = runTest { + members.upsert(CHAT, listOf(member(ERICA, "Erica", version = 3), member(ERIN, "Erin", version = 4))) + + members.reconcile(CHAT, seen = setOf(ERIN), readVersion = 7) + + assertEquals(setOf(ERIN), search("eri")) + assertEquals(1, members.countMembers(CHAT)) + } + + @Test + fun `a full read keeps a member who joined after the version it read`() = runTest { + // Erica joined at 9 off the stream; the read described the roster at 7, before she was in it. + members.upsert(CHAT, listOf(member(ERICA, "Erica", version = 9), member(ERIN, "Erin", version = 4))) + + members.reconcile(CHAT, seen = setOf(ERIN), readVersion = 7) + + assertEquals(setOf(ERICA, ERIN), search("eri")) + } + + @Test + fun `a trailing page does not wind a member's version back`() = runTest { + members.upsert(CHAT, listOf(member(ERICA, "Erica", version = 9))) + members.upsert(CHAT, listOf(member(ERICA, "Erica", version = 2))) + + // Still above the read, so still kept. + members.reconcile(CHAT, seen = emptySet(), readVersion = 7) + + assertEquals(setOf(ERICA), search("eri")) + } + + @Test + fun `a renamed member is found by the new name only`() = runTest { + members.upsert(CHAT, listOf(member(ERICA, "Erica"))) + + profiles.store(ERICA, profile("Frankie", username = null)) + + assertEquals(emptySet(), search("eri")) + assertEquals(setOf(ERICA), search("fra")) + } + + @Test + fun `a member re-sent without a profile keeps the words already held`() = runTest { + members.upsert(CHAT, listOf(member(ERICA, "Erica"))) + + members.upsert(CHAT, listOf(member(ERICA, ""))) + + assertEquals(setOf(ERICA), search("eri")) + } + + @Test + fun `the current user is never a result`() = runTest { + members.upsert(CHAT, listOf(member(SELF, "Eric"), member(ERIN, "Erin"))) + + assertEquals(setOf(ERIN), search("eri")) + } + + @Test + fun `closing a chat clears its index and sync state`() = runTest { + members.upsert(CHAT, listOf(member(ERICA, "Erica"))) + roster.markFullySynced(CHAT, watermark = 4, truncated = false) + + members.deleteForChat(CHAT) + + assertEquals(emptySet(), search("eri")) + assertEquals(null, roster.getSyncState(CHAT)) + } + + private fun profile(name: String, username: String?) = UserProfile( + displayName = name, + socialAccounts = emptyList(), + phoneNumber = null, + email = null, + username = username, + ) + + private fun member(id: ID, name: String, username: String? = null, version: Long = 0) = + ChatMember(userId = id, userProfile = profile(name, username), pointers = emptyList(), version = version) + + private companion object { + const val ENTROPY = "bWVtYmVyLXNlYXJjaC1pbmRleC10ZXN0" + val CHAT = ChatId(listOf(0x0C, 0x0F)) + val SELF: ID = listOf(0x01) + val ERICA: ID = listOf(0x0E, 0x01) + val ERIN: ID = listOf(0x0E, 0x02) + } +} diff --git a/apps/flipcash/shared/persistence/sources/src/test/kotlin/com/flipcash/app/persistence/sources/search/MemberSearchTextTest.kt b/apps/flipcash/shared/persistence/sources/src/test/kotlin/com/flipcash/app/persistence/sources/search/MemberSearchTextTest.kt new file mode 100644 index 0000000000..1ebebe26bb --- /dev/null +++ b/apps/flipcash/shared/persistence/sources/src/test/kotlin/com/flipcash/app/persistence/sources/search/MemberSearchTextTest.kt @@ -0,0 +1,68 @@ +package com.flipcash.app.persistence.sources.search + +import org.junit.Test +import kotlin.test.assertEquals + +class MemberSearchTextTest { + + @Test + fun `normalizing folds case and diacritics`() { + assertEquals("erica", MemberSearchText.normalize("Érica")) + assertEquals("zoe", MemberSearchText.normalize("ZOË")) + assertEquals("francois", MemberSearchText.normalize("François")) + } + + @Test + fun `normalizing folds compatibility forms`() { + // Full-width letters and the "fi" ligature, via NFKD. + assertEquals("abc", MemberSearchText.normalize("ABC")) + assertEquals("fish", MemberSearchText.normalize("fish")) + } + + @Test + fun `a dotted capital I folds to a plain i`() { + assertEquals("istanbul", MemberSearchText.normalize("İstanbul")) + } + + @Test + fun `tokens are each display name word plus the bare handle`() { + assertEquals( + setOf("maria", "jose", "garcia", "mjg"), + MemberSearchText.tokens("María José García", "@mjg"), + ) + } + + @Test + fun `a member with no name is found by handle alone`() { + assertEquals(setOf("mjg"), MemberSearchText.tokens(null, "mjg")) + assertEquals(emptySet(), MemberSearchText.tokens("", null)) + } + + @Test + fun `a query drops the mention trigger and folds like the tokens`() { + assertEquals(listOf("eri"), MemberSearchText.queryWords("@Éri")) + assertEquals(listOf("mar", "ga"), MemberSearchText.queryWords(" Mar Ga ")) + assertEquals(emptyList(), MemberSearchText.queryWords("@")) + } + + @Test + fun `a query word loses exactly one leading mention trigger`() { + // The second @ is part of the word, so "@@eri" does not match a token "eri". + assertEquals(listOf("@eri"), MemberSearchText.queryWords("@@eri")) + // A full-width @ folds to @ under NFKD and is stripped the same way. + assertEquals(listOf("eri"), MemberSearchText.queryWords("@eri")) + assertEquals(listOf("a", "b"), MemberSearchText.queryWords("@a @b")) + } + + @Test + fun `the upper bound sorts above any continuation of a prefix`() { + // SQLite compares TEXT as UTF-8 bytes; so does this. + fun bytes(s: String) = s.toByteArray(Charsets.UTF_8).map { it.toInt() and 0xFF } + val bound = bytes("er" + MemberSearchText.UPPER_BOUND) + for (token in listOf("er", "erica", "er￿", "er😀")) { + val t = bytes(token) + val cmp = t.zip(bound).firstOrNull { (a, b) -> a != b }?.let { (a, b) -> a - b } ?: (t.size - bound.size) + assert(cmp < 0) { "$token should sort below the bound" } + } + } +}