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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -300,6 +300,7 @@ class ChooseAccountDialogCompose {
if (userManager.setUserAsActive(userItem.user)) {
cookieManager.cookieStore.removeAll()
val intent = Intent(activity, ConversationsListActivity::class.java)
intent.putExtra(BundleKeys.KEY_INTERNAL_USER_ID, userItem.user.id)
intent.addFlags(Intent.FLAG_ACTIVITY_CLEAR_TOP)
activity.startActivity(intent)
onSelected()
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -198,14 +198,10 @@ class ConversationsListActivity : BaseActivity() {
NextcloudTalkApplication.sharedApplication!!.componentApplication.inject(this)
ecosystemManager = EcosystemManager(this@ConversationsListActivity)

val targetUserId = intent.getLongExtra(KEY_INTERNAL_USER_ID, 0L)
currentUser = if (targetUserId != 0L) {
runBlocking { userManager.getUserWithId(targetUserId) }!!
} else {
currentUserProviderOld.currentUser.blockingGet()
}
currentUser = resolveUser()

conversationsListViewModel = ViewModelProvider(this, viewModelFactory)[ConversationsListViewModel::class.java]
currentUser?.let { conversationsListViewModel.setUser(it) }
conversationTagsViewModel = ViewModelProvider(this, viewModelFactory)[ConversationTagsViewModel::class.java]

setSupportActionBar(null)
Expand All @@ -230,6 +226,26 @@ class ConversationsListActivity : BaseActivity() {
initObservers()
}

/**
* The account this screen shows: the one passed via [KEY_INTERNAL_USER_ID], otherwise the
* active one. [UserManager.currentUserFlow] is preferred over [currentUserProviderOld], since
* it is updated synchronously by [UserManager.setUserAsActive].
*/
private fun resolveUser(): User? {
val targetUserId = intent.getLongExtra(KEY_INTERNAL_USER_ID, 0L)
return if (targetUserId != 0L) {
runBlocking { userManager.getUserWithId(targetUserId) }!!
} else {
userManager.currentUserFlow.value ?: currentUserProviderOld.currentUser.blockingGet()
}
}

/** True when this screen follows the active account, and the active account has changed. */
private fun isActiveUserChanged(): Boolean {
val activeUser = userManager.currentUserFlow.value
return !intent.hasExtra(KEY_INTERNAL_USER_ID) && activeUser != null && activeUser.id != currentUser?.id
}

override fun onSaveInstanceState(outState: Bundle) {
super.onSaveInstanceState(outState)
outState.putBoolean(KEY_ACCOUNT_DIALOG_VISIBLE, showAccountDialogState.value)
Expand Down Expand Up @@ -353,6 +369,7 @@ class ConversationsListActivity : BaseActivity() {
if (user != null) {
userManager.setUserAsActive(user)
val intent = Intent(context, ConversationsListActivity::class.java)
intent.putExtra(KEY_INTERNAL_USER_ID, user.id)
startActivity(intent)
} else {
showSnackbar(getString(R.string.nc_no_account_found))
Expand Down Expand Up @@ -415,6 +432,11 @@ class ConversationsListActivity : BaseActivity() {
eventBus.register(this)
}

if (isActiveUserChanged()) {
recreate()
return
}

if (currentUser != null) {
if (isServerEOL(currentUser!!.serverVersion?.major)) {
showServerEOLDialog()
Expand Down Expand Up @@ -671,7 +693,7 @@ class ConversationsListActivity : BaseActivity() {
}

fun fetchRooms() {
conversationsListViewModel.getRooms(currentUser!!)
conversationsListViewModel.getRooms()
}

private fun fetchPendingInvitations() {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -16,20 +16,21 @@ import kotlinx.coroutines.flow.Flow
interface OfflineConversationsRepository {

/**
* Live stream of the observed account's conversations, for use in the conversation list.
* Backed by the local database: it re-emits whenever conversation rows change (room list
* sync, background catch-up, optimistic updates), with unchanged lists deduplicated.
* Live stream of the conversations of the account [accountId], for use in the conversation
* list. Backed by the local database: it re-emits whenever conversation rows change (room
* list sync, background catch-up, optimistic updates), with unchanged lists deduplicated.
*/
val roomListFlow: Flow<List<ConversationModel>>
fun observeRooms(accountId: Long): Flow<List<ConversationModel>>

/**
* Emits when [getRooms] fails to sync with the server (e.g. a dropped/reset connection on a
* slow network) while there are no locally cached conversations to fall back on for that
* account, so the UI can tell the user why the list is empty instead of failing silently.
* A failed sync while conversations are already cached does not emit here, since
* [roomListFlow] already has data to show and the sync is a best-effort background refresh.
* [observeRooms] already has data to show and the sync is a best-effort background refresh.
* Collectors must ignore errors whose [SyncError.accountId] is not the account they show.
*/
val syncErrorFlow: Flow<Throwable>
val syncErrorFlow: Flow<SyncError>

/**
* Stream of a single conversation, for use in each conversations settings.
Expand All @@ -38,9 +39,8 @@ interface OfflineConversationsRepository {
val conversationFlow: Flow<ConversationModel>

/**
* Selects the account observed by [roomListFlow] and synchronizes its conversations with
* the server (when online). The synced changes surface through [roomListFlow], which
* observes the database.
* Synchronizes the conversations of [user] with the server (when online). The synced changes
* surface through [observeRooms], which observes the database.
*/
@Deprecated("use observeConversation")
fun getRooms(user: User): Job
Expand All @@ -53,7 +53,7 @@ interface OfflineConversationsRepository {
fun getRoom(user: User, roomToken: String): Job

/**
* Updates a single conversation in the local database. [roomListFlow] observes the database
* Updates a single conversation in the local database. [observeRooms] observes the database
* and re-emits the updated list on its own.
*/
suspend fun updateConversation(conversationModel: ConversationModel)
Expand All @@ -62,4 +62,6 @@ interface OfflineConversationsRepository {
suspend fun getLocallyStoredConversation(user: User, roomToken: String): ConversationModel?

fun observeConversation(accountId: Long, roomToken: String): Flow<ConversationResult>

data class SyncError(val accountId: Long, val throwable: Throwable)
}
Original file line number Diff line number Diff line change
Expand Up @@ -32,15 +32,11 @@ import io.reactivex.android.schedulers.AndroidSchedulers
import io.reactivex.schedulers.Schedulers
import kotlinx.coroutines.CoroutineScope
import kotlinx.coroutines.Dispatchers
import kotlinx.coroutines.ExperimentalCoroutinesApi
import kotlinx.coroutines.Job
import kotlinx.coroutines.flow.Flow
import kotlinx.coroutines.flow.MutableSharedFlow
import kotlinx.coroutines.flow.MutableStateFlow
import kotlinx.coroutines.flow.distinctUntilChanged
import kotlinx.coroutines.flow.filterNotNull
import kotlinx.coroutines.flow.first
import kotlinx.coroutines.flow.flatMapLatest
import kotlinx.coroutines.flow.map
import kotlinx.coroutines.coroutineScope
import kotlinx.coroutines.launch
Expand All @@ -59,33 +55,24 @@ class OfflineFirstConversationsRepository @Inject constructor(
private val context: Context,
private val logger: Logger
) : OfflineConversationsRepository {
private val observedAccountId = MutableStateFlow<Long?>(null)

/**
* The conversation list as a live view of the local database — the single source of truth.
* Every write to the conversations table (room list sync, background message catch-up,
* optimistic read state, drafts) reaches collectors reactively; [getRooms] only selects the
* account to observe and triggers the background sync, which stays in place as the authority
* and self-healing safeguard.
* optimistic read state, drafts) reaches collectors reactively; [getRooms] only triggers the
* background sync, which stays in place as the authority and self-healing safeguard.
*/
@OptIn(ExperimentalCoroutinesApi::class)
override val roomListFlow: Flow<List<ConversationModel>> =
observedAccountId
.filterNotNull()
.distinctUntilChanged()
.flatMapLatest { accountId ->
dao.getConversationsForUser(accountId)
.map { entities -> entities.map(ConversationEntity::toDomainModel) }
}
override fun observeRooms(accountId: Long): Flow<List<ConversationModel>> =
dao.getConversationsForUser(accountId)
.map { entities -> entities.map(ConversationEntity::toDomainModel) }
.distinctUntilChanged()

override val conversationFlow: Flow<ConversationModel>
get() = _conversationFlow
private val _conversationFlow: MutableSharedFlow<ConversationModel> = MutableSharedFlow()

override val syncErrorFlow: Flow<Throwable>
override val syncErrorFlow: Flow<OfflineConversationsRepository.SyncError>
get() = _syncErrorFlow
private val _syncErrorFlow: MutableSharedFlow<Throwable> = MutableSharedFlow()
private val _syncErrorFlow: MutableSharedFlow<OfflineConversationsRepository.SyncError> = MutableSharedFlow()

private val scope = CoroutineScope(Dispatchers.IO)

Expand All @@ -109,8 +96,6 @@ class OfflineFirstConversationsRepository @Inject constructor(

override fun getRooms(user: User): Job =
scope.launch {
observedAccountId.value = user.id!!

if (networkMonitor.isOnline.value) {
getRoomsFromServer(user)
}
Expand Down Expand Up @@ -207,7 +192,7 @@ class OfflineFirstConversationsRepository @Inject constructor(
Log.e(TAG, "Something went wrong when fetching conversations", e)
val hasCachedConversations = dao.getConversationsForUser(user.id!!).first().isNotEmpty()
if (!hasCachedConversations) {
_syncErrorFlow.emit(e)
_syncErrorFlow.emit(OfflineConversationsRepository.SyncError(user.id!!, e))
}
}
return conversationsFromSync
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -42,13 +42,13 @@ import com.nextcloud.talk.utils.ApiUtils
import com.nextcloud.talk.utils.CapabilitiesUtil.hasSpreedFeatureCapability
import com.nextcloud.talk.utils.SpreedFeatures
import com.nextcloud.talk.utils.UserIdUtils
import com.nextcloud.talk.utils.database.user.CurrentUserProviderOld
import com.nextcloud.talk.utils.withRetry
import io.reactivex.Observer
import io.reactivex.android.schedulers.AndroidSchedulers
import io.reactivex.disposables.Disposable
import io.reactivex.schedulers.Schedulers
import kotlinx.coroutines.Dispatchers
import kotlinx.coroutines.ExperimentalCoroutinesApi
import kotlinx.coroutines.Job
import kotlinx.coroutines.flow.Flow
import kotlinx.coroutines.flow.MutableStateFlow
Expand All @@ -57,10 +57,16 @@ import kotlinx.coroutines.flow.StateFlow
import kotlinx.coroutines.flow.asStateFlow
import kotlinx.coroutines.flow.catch
import kotlinx.coroutines.flow.combine
import kotlinx.coroutines.flow.distinctUntilChanged
import kotlinx.coroutines.flow.filter
import kotlinx.coroutines.flow.filterNotNull
import kotlinx.coroutines.flow.flatMapLatest
import kotlinx.coroutines.flow.flow
import kotlinx.coroutines.flow.flowOn
import kotlinx.coroutines.flow.launchIn
import kotlinx.coroutines.flow.map
import kotlinx.coroutines.flow.onEach
import kotlinx.coroutines.flow.onStart
import kotlinx.coroutines.flow.stateIn
import kotlinx.coroutines.flow.update
import kotlinx.coroutines.launch
Expand All @@ -72,7 +78,6 @@ import javax.inject.Inject
class ConversationsListViewModel @Inject constructor(
private val repository: OfflineConversationsRepository,
private val threadsRepository: ThreadsRepository,
private val currentUserProvider: CurrentUserProviderOld,
private val openConversationsRepository: OpenConversationsRepository,
private val contactsRepository: ContactsRepository,
private val unifiedSearchRepository: UnifiedSearchRepository,
Expand All @@ -84,11 +89,37 @@ class ConversationsListViewModel @Inject constructor(
private val logger: Logger
) : ViewModel() {

private val _currentUser = currentUserProvider.currentUser.blockingGet()
val currentUser: User = _currentUser
val credentials = ApiUtils.getCredentials(_currentUser.username, _currentUser.token) ?: ""
private val userFlow = MutableStateFlow<User?>(null)

private val searchHelper = MessageSearchHelper(unifiedSearchRepository, currentUser)
/** The account this list shows. Must be set via [setUser] before the view model is used. */
val currentUser: User
get() = checkNotNull(userFlow.value) { "setUser must be called before using the view model" }

val credentials: String
get() = ApiUtils.getCredentials(currentUser.username, currentUser.token) ?: ""

private var searchHelper: MessageSearchHelper? = null

private val accountIdFlow = userFlow
.map { it?.id }
.filterNotNull()
.distinctUntilChanged()

/**
* Binds the view model to [user]. Switching to another account drops the previous account's
* search state; the room list follows via [accountIdFlow].
*/
fun setUser(user: User) {
val previous = userFlow.value
userFlow.value = user
if (previous?.id != user.id) {
searchHelper?.cancelSearch()
searchHelper = MessageSearchHelper(unifiedSearchRepository, user)
if (previous != null) {
cancelSearch()
}
}
}

sealed interface ViewState

Expand Down Expand Up @@ -145,7 +176,9 @@ class ConversationsListViewModel @Inject constructor(
val getRoomsViewState: LiveData<ViewState>
get() = _getRoomsViewState

val getRoomsFlow = repository.roomListFlow
@OptIn(ExperimentalCoroutinesApi::class)
val getRoomsFlow = accountIdFlow
.flatMapLatest { accountId -> repository.observeRooms(accountId) }
.onEach { list ->
_getRoomsViewState.value = GetRoomsSuccessState(list.isNotEmpty())
}.catch {
Expand All @@ -154,7 +187,8 @@ class ConversationsListViewModel @Inject constructor(

init {
repository.syncErrorFlow
.onEach { throwable -> _getRoomsViewState.value = GetRoomsErrorState(throwable) }
.filter { error -> error.accountId == userFlow.value?.id }
.onEach { error -> _getRoomsViewState.value = GetRoomsErrorState(error.throwable) }
.launchIn(viewModelScope)
}

Expand All @@ -166,8 +200,12 @@ class ConversationsListViewModel @Inject constructor(
*/
val isLoadingRooms: StateFlow<Boolean> = _isLoadingRooms.asStateFlow()

val getRoomsStateFlow = repository
.roomListFlow
@OptIn(ExperimentalCoroutinesApi::class)
val getRoomsStateFlow = accountIdFlow
.flatMapLatest { accountId ->
// Clear the previous account's rooms until the new account's first emission.
repository.observeRooms(accountId).onStart { emit(emptyList()) }
}
.catch { throwable ->
Log.e(TAG, "Error observing the conversation list", throwable)
_getRoomsViewState.value = GetRoomsErrorState(throwable)
Expand Down Expand Up @@ -348,7 +386,7 @@ class ConversationsListViewModel @Inject constructor(
_isSearchLoadingFlow.value = false
searchJob?.cancel()
searchJob = null
searchHelper.cancelSearch()
searchHelper?.cancelSearch()
searchResultEntries.value = emptyList()
_currentSearchQueryFlow.value = ""
}
Expand Down Expand Up @@ -536,13 +574,13 @@ class ConversationsListViewModel @Inject constructor(

private fun getMessagesFlow(search: String): Flow<MessageSearchResults> =
flow {
emit(searchHelper.startMessageSearch(search))
emit(checkNotNull(searchHelper).startMessageSearch(search))
}.flowOn(Dispatchers.IO)

fun loadMoreMessages(context: Context) {
viewModelScope.launch {
val result = withContext(Dispatchers.IO) {
searchHelper.loadMore()
searchHelper?.loadMore()
} ?: return@launch

val newEntries: List<ConversationListEntry> =
Expand All @@ -563,11 +601,11 @@ class ConversationsListViewModel @Inject constructor(
}
}

fun getRooms(user: User) {
fun getRooms() {
val startNanoTime = System.nanoTime()
Log.d(TAG, "fetchData - getRooms - calling: $startNanoTime")
_isLoadingRooms.value = true
val job = repository.getRooms(user)
val job = repository.getRooms(currentUser)
viewModelScope.launch {
job.join()
_isLoadingRooms.value = false
Expand All @@ -576,7 +614,7 @@ class ConversationsListViewModel @Inject constructor(

fun checkIfThreadsExist() {
val limitForFollowedThreadsExistenceCheck = 1
val accountId = UserIdUtils.getIdForUser(currentUserProvider.currentUser.blockingGet())
val accountId = UserIdUtils.getIdForUser(currentUser)

fun isLastCheckTooOld(lastCheckDate: Long): Boolean {
val currentTimeMillis = System.currentTimeMillis()
Expand Down Expand Up @@ -944,8 +982,6 @@ class ConversationsListViewModel @Inject constructor(
}

override fun onNext(invitationsModel: InvitationsModel) {
val currentUser = currentUserProvider.currentUser.blockingGet()

if (invitationsModel.user.userId?.equals(currentUser.userId) == true &&
invitationsModel.user.baseUrl?.equals(currentUser.baseUrl) == true
) {
Expand Down
Loading
Loading