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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 4 additions & 0 deletions apps/flipcash/shared/chat/build.gradle.kts
Original file line number Diff line number Diff line change
Expand Up @@ -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"))
Expand Down
Original file line number Diff line number Diff line change
@@ -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<MemberMatch>

/**
* 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?,
)
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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(
Expand Down
Original file line number Diff line number Diff line change
@@ -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<MemberMatch> {
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<String>): 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<RosterSearchCandidate>,
queryWords: List<String>,
): List<RosterSearchCandidate> {
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<Pair<RosterSearchCandidate, String>> { (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<String> {
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)
}
}
Original file line number Diff line number Diff line change
@@ -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<RosterReconcileWorker>()
.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()
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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,
) {

/**
Expand All @@ -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)
}

/**
Expand All @@ -75,13 +82,17 @@ 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(
tag = TAG,
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.
Expand Down
Loading
Loading