diff --git a/app/src/main/java/com/nextcloud/talk/api/NcApi.java b/app/src/main/java/com/nextcloud/talk/api/NcApi.java index b43918cde7..5f04f4c558 100644 --- a/app/src/main/java/com/nextcloud/talk/api/NcApi.java +++ b/app/src/main/java/com/nextcloud/talk/api/NcApi.java @@ -95,9 +95,10 @@ Observable getContactsWithSearchParam(@Header("Authorization") Str Server URL is: baseUrl + ocsApiVersion + spreedApiVersion + /room */ @GET - Observable getRooms(@Header("Authorization") String authorization, - @Url String url, - @Nullable @Query("includeStatus") Boolean includeStatus); + Observable> getRooms(@Header("Authorization") String authorization, + @Url String url, + @Nullable @Query("includeStatus") Boolean includeStatus, + @Nullable @Query("modifiedSince") Long modifiedSince); /* Server URL is: baseUrl + ocsApiVersion + spreedApiVersion + /room/roomToken diff --git a/app/src/main/java/com/nextcloud/talk/conversationlist/ConversationsListActivity.kt b/app/src/main/java/com/nextcloud/talk/conversationlist/ConversationsListActivity.kt index d27efe8209..246d3df227 100644 --- a/app/src/main/java/com/nextcloud/talk/conversationlist/ConversationsListActivity.kt +++ b/app/src/main/java/com/nextcloud/talk/conversationlist/ConversationsListActivity.kt @@ -279,7 +279,7 @@ class ConversationsListActivity : BaseActivity() { isMaintenanceModeState.value = false isRefreshingState.value = true appPreferences.setConversationListPositionAndOffset(0, 0) - fetchRooms() + fetchRooms(forceFullSync = true) fetchPendingInvitations() }, onFabClick = { @@ -670,8 +670,8 @@ class ConversationsListActivity : BaseActivity() { lifecycleScope.launch { snackbarHostState.showSnackbar(text) } } - fun fetchRooms() { - conversationsListViewModel.getRooms(currentUser!!) + fun fetchRooms(forceFullSync: Boolean = false) { + conversationsListViewModel.getRooms(currentUser!!, forceFullSync) } private fun fetchPendingInvitations() { diff --git a/app/src/main/java/com/nextcloud/talk/conversationlist/data/OfflineConversationsRepository.kt b/app/src/main/java/com/nextcloud/talk/conversationlist/data/OfflineConversationsRepository.kt index 72d9efe613..27c0637258 100644 --- a/app/src/main/java/com/nextcloud/talk/conversationlist/data/OfflineConversationsRepository.kt +++ b/app/src/main/java/com/nextcloud/talk/conversationlist/data/OfflineConversationsRepository.kt @@ -41,9 +41,14 @@ interface OfflineConversationsRepository { * 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. + * + * The sync asks the server only for what changed since the last one where it can. Set + * [forceFullSync] for the cases where that is not good enough and the whole list is wanted + * back - a pull to refresh, where a conversation the user left elsewhere should be gone by the + * time the indicator stops spinning, rather than within the next five minutes. */ @Deprecated("use observeConversation") - fun getRooms(user: User): Job + fun getRooms(user: User, forceFullSync: Boolean = false): Job /** * Called once onStart to emit a conversation to [conversationFlow] diff --git a/app/src/main/java/com/nextcloud/talk/conversationlist/data/network/ConversationsNetworkDataSource.kt b/app/src/main/java/com/nextcloud/talk/conversationlist/data/network/ConversationsNetworkDataSource.kt index 370b162eb3..8018c0dc51 100644 --- a/app/src/main/java/com/nextcloud/talk/conversationlist/data/network/ConversationsNetworkDataSource.kt +++ b/app/src/main/java/com/nextcloud/talk/conversationlist/data/network/ConversationsNetworkDataSource.kt @@ -11,6 +11,28 @@ import com.nextcloud.talk.data.user.model.User import com.nextcloud.talk.models.json.conversations.ConversationDto import io.reactivex.Observable +/** + * The outcome of a conversation list request. + * + * [conversations] are the rooms the server returned, [modifiedBefore] the timestamp it reported in + * `X-Nextcloud-Talk-Modified-Before`, or null when it sent none. [wasDelta] is true when the + * response covers only the conversations changed since a given timestamp, and false when it covers + * all of them - including when a `modifiedSince` was sent but the server answered in full anyway. + */ +data class RoomListResult(val conversations: List, val modifiedBefore: Long?, val wasDelta: Boolean) + interface ConversationsNetworkDataSource { - fun getRooms(user: User, url: String, includeStatus: Boolean): Observable> + /** + * Loads [user]'s conversation list. + * + * With [includeStatus] the response carries the user status of one-to-one rooms. With + * [modifiedSince] set, and on a server that supports it, the response covers only the + * conversations changed since that timestamp. + */ + fun getRooms( + user: User, + url: String, + includeStatus: Boolean, + modifiedSince: Long? = null + ): Observable } diff --git a/app/src/main/java/com/nextcloud/talk/conversationlist/data/network/OfflineFirstConversationsRepository.kt b/app/src/main/java/com/nextcloud/talk/conversationlist/data/network/OfflineFirstConversationsRepository.kt index 91738cb1a9..5535a96438 100644 --- a/app/src/main/java/com/nextcloud/talk/conversationlist/data/network/OfflineFirstConversationsRepository.kt +++ b/app/src/main/java/com/nextcloud/talk/conversationlist/data/network/OfflineFirstConversationsRepository.kt @@ -13,6 +13,7 @@ import android.database.sqlite.SQLiteConstraintException import android.net.ConnectivityManager import android.os.PowerManager import android.util.Log +import com.nextcloud.talk.arbitrarystorage.ArbitraryStorageManager import com.nextcloud.talk.chat.data.network.ChatMessageSyncer import com.nextcloud.talk.chat.data.network.ChatNetworkDataSource import com.nextcloud.talk.conversationlist.data.OfflineConversationsRepository @@ -30,6 +31,7 @@ import com.nextcloud.talk.utils.SpreedFeatures import com.nextcloud.talk.utils.withRetry import io.reactivex.android.schedulers.AndroidSchedulers import io.reactivex.schedulers.Schedulers +import kotlinx.coroutines.CancellationException import kotlinx.coroutines.CoroutineScope import kotlinx.coroutines.Dispatchers import kotlinx.coroutines.ExperimentalCoroutinesApi @@ -49,6 +51,7 @@ import kotlinx.coroutines.sync.withPermit import javax.inject.Inject import kotlin.collections.map +@Suppress("LongParameterList") class OfflineFirstConversationsRepository @Inject constructor( private val dao: ConversationsDao, private val network: ConversationsNetworkDataSource, @@ -56,6 +59,7 @@ class OfflineFirstConversationsRepository @Inject constructor( private val networkMonitor: NetworkMonitor, private val chatMessageSyncer: ChatMessageSyncer, private val conversationListUpdater: ConversationListUpdater, + private val arbitraryStorageManager: ArbitraryStorageManager, private val context: Context, private val logger: Logger ) : OfflineConversationsRepository { @@ -107,12 +111,13 @@ class OfflineFirstConversationsRepository @Inject constructor( } } - override fun getRooms(user: User): Job = + override fun getRooms(user: User, forceFullSync: Boolean): Job = scope.launch { + val accountChanged = observedAccountId.value != user.id observedAccountId.value = user.id!! if (networkMonitor.isOnline.value) { - getRoomsFromServer(user) + getRoomsFromServer(user, forceFullSync = forceFullSync || accountChanged) } } @@ -163,7 +168,7 @@ class OfflineFirstConversationsRepository @Inject constructor( } @Suppress("Detekt.TooGenericExceptionCaught") - private suspend fun getRoomsFromServer(user: User): List? { + private suspend fun getRoomsFromServer(user: User, forceFullSync: Boolean = false): List? { var conversationsFromSync: List? = null if (!networkMonitor.isOnline.value) { @@ -171,41 +176,53 @@ class OfflineFirstConversationsRepository @Inject constructor( return null } - val includeStatus = isUserStatusAvailable(user) + val accountId = user.id!! + val modifiedSince = modifiedSinceFor(user, forceFullSync) + + val includeStatus = modifiedSince == null && isUserStatusAvailable(user) try { - val conversationsList = withRetry( + val roomList = withRetry( retries = NETWORK_FETCH_RETRIES, initialDelayMillis = NETWORK_FETCH_RETRY_INITIAL_DELAY_MS, maxDelayMillis = NETWORK_FETCH_RETRY_MAX_DELAY_MS ) { - network.getRooms(user, user.baseUrl!!, includeStatus) + network.getRooms(user, user.baseUrl!!, includeStatus, modifiedSince) .subscribeOn(Schedulers.io()) .observeOn(AndroidSchedulers.mainThread()) .blockingSingle() } - conversationsFromSync = conversationsList.map { - it.asEntity(user.id!!) + conversationsFromSync = roomList.conversations.map { + it.asEntity(accountId) } - val previousConversations = dao.getConversationsForUser(user.id!!).first() + val previousConversations = dao.getConversationsForUser(accountId).first() .associateBy { it.internalId } dao.syncConversationsForUser( - accountId = user.id!!, + accountId = accountId, serverItems = conversationListUpdater.preservePendingLocalState( previousConversations, conversationsFromSync ), - conversationIdsToDelete = determineLeftConversationIds(previousConversations, conversationsFromSync) + conversationIdsToDelete = if (roomList.wasDelta) { + emptyList() + } else { + determineLeftConversationIds(previousConversations, conversationsFromSync) + } ) + rememberSyncedState(accountId, roomList) + val roomsWithNewMessages = getRoomsWithNewMessages(conversationsFromSync, previousConversations) scope.launch { catchUpRoomsWithNewMessages(user, roomsWithNewMessages) } + } catch (e: CancellationException) { + throw e } catch (e: Exception) { Log.e(TAG, "Something went wrong when fetching conversations", e) - val hasCachedConversations = dao.getConversationsForUser(user.id!!).first().isNotEmpty() + storeTimestamp(accountId, KEY_MODIFIED_SINCE, null) + val hasCachedConversations = dao.getConversationsForUser(accountId).first().isNotEmpty() if (!hasCachedConversations) { _syncErrorFlow.emit(e) } @@ -213,6 +230,52 @@ class OfflineFirstConversationsRepository @Inject constructor( return conversationsFromSync } + /** + * The value to send as `modifiedSince`, or null when this sync has to be a full one. + * + * A filtered response cannot express a removal, so the server asks clients to refresh in full + * regularly (`docs/conversation.md`): at least every five minutes, and always when the internal + * signaling backend is in use, since there is no signaling server to announce a change out of + * band. On top of that a caller can demand one, for a pull to refresh or an account switch. + */ + private fun modifiedSinceFor(user: User, forceFullSync: Boolean): Long? { + val accountId = user.id!! + val lastFullSyncAt = readTimestamp(accountId, KEY_LAST_FULL_SYNC_AT) + val fullSyncIsRecent = lastFullSyncAt != null && + System.currentTimeMillis() - lastFullSyncAt in 0 until FULL_SYNC_INTERVAL_MILLIS + val usesExternalSignaling = !user.externalSignalingServer?.externalSignalingServer.isNullOrEmpty() + + return if (!forceFullSync && usesExternalSignaling && fullSyncIsRecent) { + readTimestamp(accountId, KEY_MODIFIED_SINCE) + } else { + null + } + } + + /** + * Stores what the next sync needs, and only once the response is safely in the database: a + * timestamp kept ahead of a write that then failed would permanently skip the conversations + * that write was carrying. + */ + private fun rememberSyncedState(accountId: Long, roomList: RoomListResult) { + storeTimestamp(accountId, KEY_MODIFIED_SINCE, roomList.modifiedBefore) + if (!roomList.wasDelta) { + storeTimestamp(accountId, KEY_LAST_FULL_SYNC_AT, System.currentTimeMillis()) + } + } + + /** Null for anything that is not a plausible timestamp, so a bad value asks for a full sync. */ + private fun readTimestamp(accountId: Long, key: String): Long? = + arbitraryStorageManager.getStorageSetting(accountId, key, "") + .blockingGet() + ?.value + ?.toLongOrNull() + ?.takeIf { it > 0 } + + private fun storeTimestamp(accountId: Long, key: String, value: Long?) { + arbitraryStorageManager.storeStorageSetting(accountId, key, value?.toString(), "") + } + /** * Determines the rooms whose messages should be caught up in the background: rooms with * activity newer than the last synced state (matching the iOS behavior) plus unread rooms that @@ -273,7 +336,10 @@ class OfflineFirstConversationsRepository @Inject constructor( unreadMessages = room.unreadMessages ) } - .onFailure { Log.e(TAG, "Message catch-up failed for room ${room.token}", it) } + .onFailure { + if (it is CancellationException) throw it + Log.e(TAG, "Message catch-up failed for room ${room.token}", it) + } } } } @@ -346,5 +412,8 @@ class OfflineFirstConversationsRepository @Inject constructor( private const val NETWORK_FETCH_RETRIES = 3 private const val NETWORK_FETCH_RETRY_INITIAL_DELAY_MS = 1000L private const val NETWORK_FETCH_RETRY_MAX_DELAY_MS = 8000L + private const val FULL_SYNC_INTERVAL_MILLIS = 5 * 60 * 1000L + private const val KEY_MODIFIED_SINCE = "conversation_list_modified_since" + private const val KEY_LAST_FULL_SYNC_AT = "conversation_list_last_full_sync_at" } } diff --git a/app/src/main/java/com/nextcloud/talk/conversationlist/data/network/RetrofitConversationsNetwork.kt b/app/src/main/java/com/nextcloud/talk/conversationlist/data/network/RetrofitConversationsNetwork.kt index 434378eebf..cb11c8677a 100644 --- a/app/src/main/java/com/nextcloud/talk/conversationlist/data/network/RetrofitConversationsNetwork.kt +++ b/app/src/main/java/com/nextcloud/talk/conversationlist/data/network/RetrofitConversationsNetwork.kt @@ -8,21 +8,40 @@ package com.nextcloud.talk.conversationlist.data.network import com.nextcloud.talk.api.NcApi import com.nextcloud.talk.data.user.model.User -import com.nextcloud.talk.models.json.conversations.ConversationDto import com.nextcloud.talk.utils.ApiUtils import io.reactivex.Observable +import retrofit2.HttpException class RetrofitConversationsNetwork(private val ncApi: NcApi) : ConversationsNetworkDataSource { - override fun getRooms(user: User, url: String, includeStatus: Boolean): Observable> { + override fun getRooms( + user: User, + url: String, + includeStatus: Boolean, + modifiedSince: Long? + ): Observable { val credentials: String = ApiUtils.getCredentials(user.username, user.token)!! val apiVersion = ApiUtils.getConversationApiVersion(user, intArrayOf(ApiUtils.API_V4, ApiUtils.API_V3, 1)) + val delta = modifiedSince?.takeIf { apiVersion >= ApiUtils.API_V4 && it > 0 } + return ncApi.getRooms( credentials, ApiUtils.getUrlForRooms(apiVersion, user.baseUrl!!), - includeStatus - ).map { it -> - it.ocs?.data?.map { it } ?: listOf() + includeStatus, + delta + ).map { response -> + if (!response.isSuccessful) throw HttpException(response) + + RoomListResult( + conversations = response.body()?.ocs?.data ?: emptyList(), + modifiedBefore = response.headers()[HEADER_MODIFIED_BEFORE]?.toLongOrNull()?.takeIf { it > 0 }, + wasDelta = delta != null + ) } } + + companion object { + /** Taken by the server before it queries, to be sent as `modifiedSince` on the next request. */ + private const val HEADER_MODIFIED_BEFORE = "X-Nextcloud-Talk-Modified-Before" + } } diff --git a/app/src/main/java/com/nextcloud/talk/conversationlist/viewmodels/ConversationsListViewModel.kt b/app/src/main/java/com/nextcloud/talk/conversationlist/viewmodels/ConversationsListViewModel.kt index ecb3ab613d..6508a1697e 100644 --- a/app/src/main/java/com/nextcloud/talk/conversationlist/viewmodels/ConversationsListViewModel.kt +++ b/app/src/main/java/com/nextcloud/talk/conversationlist/viewmodels/ConversationsListViewModel.kt @@ -563,11 +563,11 @@ class ConversationsListViewModel @Inject constructor( } } - fun getRooms(user: User) { + fun getRooms(user: User, forceFullSync: Boolean = false) { val startNanoTime = System.nanoTime() Log.d(TAG, "fetchData - getRooms - calling: $startNanoTime") _isLoadingRooms.value = true - val job = repository.getRooms(user) + val job = repository.getRooms(user, forceFullSync) viewModelScope.launch { job.join() _isLoadingRooms.value = false diff --git a/app/src/main/java/com/nextcloud/talk/dagger/modules/RepositoryModule.kt b/app/src/main/java/com/nextcloud/talk/dagger/modules/RepositoryModule.kt index 4b2ce8e3d8..91f2477d31 100644 --- a/app/src/main/java/com/nextcloud/talk/dagger/modules/RepositoryModule.kt +++ b/app/src/main/java/com/nextcloud/talk/dagger/modules/RepositoryModule.kt @@ -17,6 +17,7 @@ import com.nextcloud.talk.api.NcApi import com.nextcloud.talk.api.NcApiCoroutines import com.nextcloud.talk.chat.data.ChatMessageRepository import com.nextcloud.talk.chat.data.network.ChatMessageSyncer +import com.nextcloud.talk.arbitrarystorage.ArbitraryStorageManager import com.nextcloud.talk.chat.data.network.ChatNetworkDataSource import com.nextcloud.talk.chat.data.network.OfflineFirstChatRepository import com.nextcloud.talk.logger.Logger @@ -204,6 +205,7 @@ class RepositoryModule { networkMonitor: NetworkMonitor, chatMessageSyncer: ChatMessageSyncer, conversationListUpdater: ConversationListUpdater, + arbitraryStorageManager: ArbitraryStorageManager, context: Context, logger: Logger ): OfflineConversationsRepository = @@ -214,6 +216,7 @@ class RepositoryModule { networkMonitor, chatMessageSyncer, conversationListUpdater, + arbitraryStorageManager, context, logger ) diff --git a/app/src/main/java/com/nextcloud/talk/data/database/dao/ConversationsDao.kt b/app/src/main/java/com/nextcloud/talk/data/database/dao/ConversationsDao.kt index 1008111556..dde86001f1 100644 --- a/app/src/main/java/com/nextcloud/talk/data/database/dao/ConversationsDao.kt +++ b/app/src/main/java/com/nextcloud/talk/data/database/dao/ConversationsDao.kt @@ -26,9 +26,11 @@ interface ConversationsDao { fun getConversationForUser(accountId: Long, token: String): Flow /** - * Applies a full room list sync atomically: left conversations are deleted and the server - * items are upserted in one transaction, so observers of the conversations table see a single - * consistent update per sync instead of intermediate states. + * Applies a room list sync atomically: the conversations named in [conversationIdsToDelete] are + * deleted and [serverItems] are upserted in one transaction, so observers of the conversations + * table see a single consistent update per sync instead of intermediate states. + * + * Conversations that are named by neither argument are left untouched. */ @Transaction suspend fun syncConversationsForUser( diff --git a/app/src/main/java/com/nextcloud/talk/utils/preview/ComposePreviewUtils.kt b/app/src/main/java/com/nextcloud/talk/utils/preview/ComposePreviewUtils.kt index 9b7da27ae0..350478f4a8 100644 --- a/app/src/main/java/com/nextcloud/talk/utils/preview/ComposePreviewUtils.kt +++ b/app/src/main/java/com/nextcloud/talk/utils/preview/ComposePreviewUtils.kt @@ -17,6 +17,7 @@ import com.nextcloud.android.common.ui.theme.utils.DialogViewThemeUtils import com.nextcloud.android.common.ui.theme.utils.MaterialViewThemeUtils import com.nextcloud.talk.api.NcApi import com.nextcloud.talk.api.NcApiCoroutines +import com.nextcloud.talk.arbitrarystorage.ArbitraryStorageManager import com.nextcloud.talk.chat.data.ChatMessageRepository import com.nextcloud.talk.chat.data.io.AudioFocusRequestManager import com.nextcloud.talk.chat.data.io.MediaRecorderManager @@ -193,6 +194,7 @@ class ComposePreviewUtils private constructor(context: Context) { networkMonitor, chatMessageSyncer, conversationListUpdater, + ArbitraryStorageManager(DummyArbitraryStoragesRepositoryImpl()), mContext, logger ) diff --git a/app/src/main/java/com/nextcloud/talk/utils/preview/ComposePreviewUtilsDaos.kt b/app/src/main/java/com/nextcloud/talk/utils/preview/ComposePreviewUtilsDaos.kt index 1dc27b7788..1b55ec88f5 100644 --- a/app/src/main/java/com/nextcloud/talk/utils/preview/ComposePreviewUtilsDaos.kt +++ b/app/src/main/java/com/nextcloud/talk/utils/preview/ComposePreviewUtilsDaos.kt @@ -13,9 +13,13 @@ import com.nextcloud.talk.data.database.dao.ConversationsDao import com.nextcloud.talk.data.database.model.ChatBlockEntity import com.nextcloud.talk.data.database.model.ChatMessageEntity import com.nextcloud.talk.data.database.model.ConversationEntity +import com.nextcloud.talk.data.storage.ArbitraryStoragesRepository +import com.nextcloud.talk.data.storage.model.ArbitraryStorage +import com.nextcloud.talk.data.storage.model.ArbitraryStorageEntity import com.nextcloud.talk.data.user.UsersDao import com.nextcloud.talk.data.user.model.UserEntity import com.nextcloud.talk.models.json.push.PushConfigurationState +import io.reactivex.Maybe import kotlinx.coroutines.flow.Flow import kotlinx.coroutines.flow.flowOf @@ -309,3 +313,17 @@ class DummyChatBlocksDaoImpl : ChatBlocksDao { override suspend fun getChatBlocksForConversation(internalConversationId: String): List = emptyList() } + +class DummyArbitraryStoragesRepositoryImpl : ArbitraryStoragesRepository { + override fun getStorageSetting( + accountIdentifier: Long, + key: String, + objectString: String + ): Maybe = Maybe.empty() + + override fun deleteArbitraryStorage(accountIdentifier: Long): Int = 0 + + override fun saveArbitraryStorage(arbitraryStorage: ArbitraryStorage): Long = 0 + + override fun getAll(): Maybe> = Maybe.empty() +} diff --git a/app/src/test/java/com/nextcloud/talk/conversationlist/data/network/ConversationListDeltaSyncIntegrationTest.kt b/app/src/test/java/com/nextcloud/talk/conversationlist/data/network/ConversationListDeltaSyncIntegrationTest.kt new file mode 100644 index 0000000000..c4e0b161b5 --- /dev/null +++ b/app/src/test/java/com/nextcloud/talk/conversationlist/data/network/ConversationListDeltaSyncIntegrationTest.kt @@ -0,0 +1,299 @@ +/* + * Nextcloud Talk - Android Client + * + * SPDX-FileCopyrightText: 2026 Andy Scherzinger + * SPDX-License-Identifier: GPL-3.0-or-later + */ + +package com.nextcloud.talk.conversationlist.data.network + +import android.app.Application +import android.content.Context +import androidx.room.Room +import androidx.test.core.app.ApplicationProvider +import com.nextcloud.talk.arbitrarystorage.ArbitraryStorageManager +import com.nextcloud.talk.chat.data.model.ChatMessage +import com.nextcloud.talk.chat.data.network.ChatMessageSyncer +import com.nextcloud.talk.chat.data.network.ChatNetworkDataSource +import com.nextcloud.talk.data.database.model.ChatBlockEntity +import com.nextcloud.talk.data.database.model.ChatMessageEntity +import com.nextcloud.talk.data.network.NetworkMonitor +import com.nextcloud.talk.data.source.local.TalkDatabase +import com.nextcloud.talk.data.storage.ArbitraryStoragesRepositoryImpl +import com.nextcloud.talk.data.user.model.User +import com.nextcloud.talk.data.user.model.UserEntity +import com.nextcloud.talk.logger.Logger +import com.nextcloud.talk.models.ExternalSignalingServer +import com.nextcloud.talk.models.json.capabilities.CapabilitiesDto +import com.nextcloud.talk.models.json.capabilities.SpreedCapabilityDto +import com.nextcloud.talk.models.json.capabilities.UserStatusCapabilityDto +import com.nextcloud.talk.models.json.chat.ChatOCS +import com.nextcloud.talk.models.json.chat.ChatOverall +import com.nextcloud.talk.models.json.conversations.ConversationDto +import io.reactivex.Observable +import io.reactivex.android.plugins.RxAndroidPlugins +import io.reactivex.schedulers.Schedulers +import kotlinx.coroutines.flow.MutableStateFlow +import kotlinx.coroutines.flow.first +import kotlinx.coroutines.runBlocking +import org.junit.After +import org.junit.Assert.assertEquals +import org.junit.Assert.assertNull +import org.junit.Assert.assertTrue +import org.junit.Before +import org.junit.Test +import org.junit.runner.RunWith +import org.mockito.kotlin.any +import org.mockito.kotlin.anyOrNull +import org.mockito.kotlin.argumentCaptor +import org.mockito.kotlin.mock +import org.mockito.kotlin.times +import org.mockito.kotlin.verify +import org.mockito.kotlin.whenever +import org.mockito.kotlin.wheneverBlocking +import org.robolectric.RobolectricTestRunner +import org.robolectric.annotation.Config +import retrofit2.Response + +/** + * Integration tests for the `modifiedSince` delta sync against an in-memory database. + * + * Covers that a delta response leaves the conversations it omits, and their cached messages and + * chat blocks, in place; that a full response still removes them; what a delta request sends; that + * an internal-signaling account never gets one; and that a failed sync drops the stored timestamp. + */ +private const val ACCOUNT_ID = 1L +private const val BASE_URL = "https://server.example.com" + +@RunWith(RobolectricTestRunner::class) +@Config(application = Application::class, sdk = [33]) +class ConversationListDeltaSyncIntegrationTest { + + private lateinit var db: TalkDatabase + private lateinit var arbitraryStorageManager: ArbitraryStorageManager + private lateinit var repository: OfflineFirstConversationsRepository + + private val conversationsNetwork: ConversationsNetworkDataSource = mock() + private val chatNetwork: ChatNetworkDataSource = mock() + private val networkMonitor: NetworkMonitor = mock() + + @Before + fun setUp() { + RxAndroidPlugins.setInitMainThreadSchedulerHandler { Schedulers.trampoline() } + RxAndroidPlugins.setMainThreadSchedulerHandler { Schedulers.trampoline() } + + val context = ApplicationProvider.getApplicationContext() + db = Room.inMemoryDatabaseBuilder(context, TalkDatabase::class.java) + .allowMainThreadQueries() + .build() + runBlocking { + db.usersDao().saveUser(UserEntity(id = ACCOUNT_ID, userId = "me", username = "me", baseUrl = BASE_URL)) + } + + whenever(networkMonitor.isOnline).thenReturn(MutableStateFlow(true)) + wheneverBlocking { chatNetwork.pullChatMessages(any(), any(), any()) } + .thenReturn(Response.success(ChatOverall(ocs = ChatOCS(meta = null, data = emptyList())))) + + val conversationListUpdater = + ConversationListUpdater(db.chatMessagesDao(), db.chatBlocksDao(), db.conversationsDao()) + val syncer = ChatMessageSyncer( + db.chatMessagesDao(), + db.chatBlocksDao(), + chatNetwork, + networkMonitor, + conversationListUpdater + ) + arbitraryStorageManager = ArbitraryStorageManager(ArbitraryStoragesRepositoryImpl(db.arbitraryStoragesDao())) + repository = OfflineFirstConversationsRepository( + db.conversationsDao(), + conversationsNetwork, + chatNetwork, + networkMonitor, + syncer, + conversationListUpdater, + arbitraryStorageManager, + context, + mock() + ) + } + + @After + fun tearDown() { + db.close() + RxAndroidPlugins.reset() + } + + @Test + fun `a delta sync keeps the conversations it does not mention, with their messages and blocks`() { + val allRooms = (1..ROOM_COUNT).map { conversation("room$it", lastActivity = 10) } + whenever(conversationsNetwork.getRooms(any(), any(), any(), anyOrNull())).thenReturn( + fullResponse(allRooms, modifiedBefore = 1000), + deltaResponse(listOf(conversation("room1", lastActivity = 20)), modifiedBefore = 2000) + ) + + runBlocking { + repository.getRooms(user()).join() + assertEquals(ROOM_COUNT, storedConversationCount()) + allRooms.forEach { seedCachedChat(it.token!!) } + + repository.getRooms(user()).join() + + assertEquals( + "a delta response must not reconcile away the conversations it left out", + ROOM_COUNT, + storedConversationCount() + ) + assertTrue("cached messages must survive a delta sync", cachedMessageIds("room7").isNotEmpty()) + assertTrue("chat blocks must survive a delta sync", cachedBlocks("room7").isNotEmpty()) + } + } + + @Test + fun `a full sync reconciles away a conversation the response left out`() { + whenever(conversationsNetwork.getRooms(any(), any(), any(), anyOrNull())).thenReturn( + fullResponse(listOf(conversation("room1"), conversation("room2")), modifiedBefore = 1000), + fullResponse(listOf(conversation("room1")), modifiedBefore = 2000) + ) + + runBlocking { + repository.getRooms(user()).join() + assertEquals(2, storedConversationCount()) + + repository.getRooms(user(), forceFullSync = true).join() + + assertEquals("a full response is the one that can say a conversation is gone", 1, storedConversationCount()) + } + } + + @Test + fun `a delta sync sends the stored timestamp and asks without includeStatus`() { + whenever(conversationsNetwork.getRooms(any(), any(), any(), anyOrNull())).thenReturn( + fullResponse(listOf(conversation("room1")), modifiedBefore = 1000), + deltaResponse(listOf(conversation("room1")), modifiedBefore = 2000) + ) + + runBlocking { + repository.getRooms(user()).join() + repository.getRooms(user()).join() + } + + val includeStatus = argumentCaptor() + val modifiedSince = argumentCaptor() + verify(conversationsNetwork, times(2)) + .getRooms(any(), any(), includeStatus.capture(), modifiedSince.capture()) + + assertEquals("the first sync has nothing to anchor a delta on", null, modifiedSince.firstValue) + assertTrue("a full sync keeps user status fresh", includeStatus.firstValue) + assertEquals("the second sync echoes the header of the first", 1000L, modifiedSince.secondValue) + assertEquals("includeStatus would return every one-to-one room anyway", false, includeStatus.secondValue) + } + + @Test + fun `the internal signaling backend never gets a delta sync`() { + whenever(conversationsNetwork.getRooms(any(), any(), any(), anyOrNull())) + .thenReturn(fullResponse(listOf(conversation("room1")), modifiedBefore = 1000)) + + runBlocking { + repository.getRooms(userWithInternalSignaling()).join() + repository.getRooms(userWithInternalSignaling()).join() + } + + val modifiedSince = argumentCaptor() + verify(conversationsNetwork, times(2)) + .getRooms(any(), any(), any(), modifiedSince.capture()) + assertNull("no signaling server exists to announce a removal out of band", modifiedSince.secondValue) + } + + @Test + fun `a failed sync drops the stored timestamp so the next one asks for everything`() { + whenever(conversationsNetwork.getRooms(any(), any(), any(), anyOrNull())).thenReturn( + fullResponse(listOf(conversation("room1")), modifiedBefore = 1000), + Observable.error(IllegalStateException("server said no")) + ) + + runBlocking { + repository.getRooms(user()).join() + assertEquals("1000", storedValue(KEY_MODIFIED_SINCE)) + + repository.getRooms(user()).join() + + assertNull( + "a delta anchored on a sync that did not land would skip what it carried", + storedValue(KEY_MODIFIED_SINCE) + ) + } + } + + private fun storedValue(key: String): String? = + arbitraryStorageManager.getStorageSetting(ACCOUNT_ID, key, "").blockingGet()?.value + + private suspend fun storedConversationCount(): Int = + db.conversationsDao().getConversationsForUser(ACCOUNT_ID).first().size + + private suspend fun seedCachedChat(roomToken: String) { + val internalConversationId = "$ACCOUNT_ID@$roomToken" + db.chatMessagesDao().upsertChatMessages( + listOf( + ChatMessageEntity( + internalId = "$internalConversationId@1", + accountId = ACCOUNT_ID, + token = roomToken, + id = 1, + internalConversationId = internalConversationId, + actorDisplayName = "Other User", + message = "cached message", + actorId = "other", + actorType = "users", + messageType = "comment", + systemMessageType = ChatMessage.SystemMessageType.DUMMY + ) + ) + ) + db.chatBlocksDao().upsertChatBlock( + ChatBlockEntity( + internalConversationId = internalConversationId, + accountId = ACCOUNT_ID, + token = roomToken, + oldestMessageId = 1, + newestMessageId = 1, + hasHistory = false + ) + ) + } + + private suspend fun cachedMessageIds(roomToken: String): List = + db.chatMessagesDao().getMessagesForConversation("$ACCOUNT_ID@$roomToken", null).first().map { it.id } + + private suspend fun cachedBlocks(roomToken: String): List = + db.chatBlocksDao().getChatBlocksForConversation("$ACCOUNT_ID@$roomToken") + + companion object { + private const val ROOM_COUNT = 10 + private const val KEY_MODIFIED_SINCE = "conversation_list_modified_since" + } +} + +private fun conversation(roomToken: String, lastActivity: Long = 10): ConversationDto = + ConversationDto(token = roomToken, lastActivity = lastActivity, unreadMessages = 0) + +private fun fullResponse(rooms: List, modifiedBefore: Long): Observable = + Observable.just(RoomListResult(rooms, modifiedBefore = modifiedBefore, wasDelta = false)) + +private fun deltaResponse(rooms: List, modifiedBefore: Long): Observable = + Observable.just(RoomListResult(rooms, modifiedBefore = modifiedBefore, wasDelta = true)) + +private fun user(): User = + User( + id = ACCOUNT_ID, + userId = "me", + username = "me", + baseUrl = BASE_URL, + token = "app-password", + externalSignalingServer = ExternalSignalingServer(externalSignalingServer = "https://hpb.example.com"), + capabilities = CapabilitiesDto().apply { + spreedCapability = SpreedCapabilityDto().apply { features = listOf("chat-keep-notifications") } + userStatusCapability = UserStatusCapabilityDto(enabled = true, restore = false, supportsEmoji = true) + } + ) + +private fun userWithInternalSignaling(): User = user().copy(externalSignalingServer = null) diff --git a/app/src/test/java/com/nextcloud/talk/conversationlist/data/network/ConversationListFreshnessIntegrationTest.kt b/app/src/test/java/com/nextcloud/talk/conversationlist/data/network/ConversationListFreshnessIntegrationTest.kt index 99fc651fae..10ed445f4a 100644 --- a/app/src/test/java/com/nextcloud/talk/conversationlist/data/network/ConversationListFreshnessIntegrationTest.kt +++ b/app/src/test/java/com/nextcloud/talk/conversationlist/data/network/ConversationListFreshnessIntegrationTest.kt @@ -12,6 +12,7 @@ import android.content.Context import androidx.room.Room import androidx.test.core.app.ApplicationProvider import com.bluelinelabs.logansquare.LoganSquare +import com.nextcloud.talk.arbitrarystorage.ArbitraryStorageManager import com.nextcloud.talk.chat.data.model.ChatMessage import com.nextcloud.talk.chat.data.network.ChatMessageSyncer import com.nextcloud.talk.chat.data.network.ChatNetworkDataSource @@ -19,6 +20,7 @@ import com.nextcloud.talk.data.database.mappers.asEntity import com.nextcloud.talk.data.database.model.ConversationEntity import com.nextcloud.talk.data.network.NetworkMonitor import com.nextcloud.talk.data.source.local.TalkDatabase +import com.nextcloud.talk.data.storage.ArbitraryStoragesRepositoryImpl import com.nextcloud.talk.data.user.model.User import com.nextcloud.talk.data.user.model.UserEntity import com.nextcloud.talk.logger.Logger @@ -47,6 +49,7 @@ import org.junit.Before import org.junit.Test import org.junit.runner.RunWith import org.mockito.kotlin.any +import org.mockito.kotlin.anyOrNull import org.mockito.kotlin.mock import org.mockito.kotlin.whenever import org.mockito.kotlin.wheneverBlocking @@ -187,6 +190,7 @@ class ConversationListFreshnessIntegrationTest { networkMonitor, syncer, conversationListUpdater, + ArbitraryStorageManager(ArbitraryStoragesRepositoryImpl(db.arbitraryStoragesDao())), ApplicationProvider.getApplicationContext(), mock() ) @@ -256,13 +260,14 @@ class ConversationListFreshnessIntegrationTest { networkMonitor, syncer, conversationListUpdater, + ArbitraryStorageManager(ArbitraryStoragesRepositoryImpl(db.arbitraryStoragesDao())), ApplicationProvider.getApplicationContext(), mock() ) - whenever(conversationsNetwork.getRooms(any(), any(), any())).thenReturn( - Observable.just(listOf(staleServerRoom(lastReadMessage = 10, unreadMessages = 2))), - Observable.just(listOf(staleServerRoom(lastReadMessage = 12, unreadMessages = 0))), - Observable.just(listOf(staleServerRoom(lastReadMessage = 8, unreadMessages = 4))) + whenever(conversationsNetwork.getRooms(any(), any(), any(), anyOrNull())).thenReturn( + roomList(listOf(staleServerRoom(lastReadMessage = 10, unreadMessages = 2))), + roomList(listOf(staleServerRoom(lastReadMessage = 12, unreadMessages = 0))), + roomList(listOf(staleServerRoom(lastReadMessage = 8, unreadMessages = 4))) ) runBlocking { @@ -289,10 +294,10 @@ class ConversationListFreshnessIntegrationTest { fun `a sync cannot revert a pending favorite until the server confirms it`() { val user = user(withKeepNotificationsCapability = false) val repository = repository() - whenever(conversationsNetwork.getRooms(any(), any(), any())).thenReturn( - Observable.just(listOf(staleServerRoom(lastReadMessage = 10, unreadMessages = 0))), - Observable.just(listOf(staleServerRoom(lastReadMessage = 10, unreadMessages = 0, favorite = true))), - Observable.just(listOf(staleServerRoom(lastReadMessage = 10, unreadMessages = 0))) + whenever(conversationsNetwork.getRooms(any(), any(), any(), anyOrNull())).thenReturn( + roomList(listOf(staleServerRoom(lastReadMessage = 10, unreadMessages = 0))), + roomList(listOf(staleServerRoom(lastReadMessage = 10, unreadMessages = 0, favorite = true))), + roomList(listOf(staleServerRoom(lastReadMessage = 10, unreadMessages = 0))) ) runBlocking { @@ -315,9 +320,9 @@ class ConversationListFreshnessIntegrationTest { fun `a sync cannot revert a pending mark as unread until the server confirms it`() { val user = user(withKeepNotificationsCapability = false) val repository = repository() - whenever(conversationsNetwork.getRooms(any(), any(), any())).thenReturn( - Observable.just(listOf(staleServerRoom(lastReadMessage = 12, unreadMessages = 0))), - Observable.just(listOf(staleServerRoom(lastReadMessage = 9, unreadMessages = 3))) + whenever(conversationsNetwork.getRooms(any(), any(), any(), anyOrNull())).thenReturn( + roomList(listOf(staleServerRoom(lastReadMessage = 12, unreadMessages = 0))), + roomList(listOf(staleServerRoom(lastReadMessage = 9, unreadMessages = 3))) ) runBlocking { @@ -337,10 +342,10 @@ class ConversationListFreshnessIntegrationTest { fun `a sync cannot revert a pending archive until the server confirms it`() { val user = user(withKeepNotificationsCapability = false) val repository = repository() - whenever(conversationsNetwork.getRooms(any(), any(), any())).thenReturn( - Observable.just(listOf(staleServerRoom(lastReadMessage = 10, unreadMessages = 0))), - Observable.just(listOf(staleServerRoom(lastReadMessage = 10, unreadMessages = 0, hasArchived = true))), - Observable.just(listOf(staleServerRoom(lastReadMessage = 10, unreadMessages = 0))) + whenever(conversationsNetwork.getRooms(any(), any(), any(), anyOrNull())).thenReturn( + roomList(listOf(staleServerRoom(lastReadMessage = 10, unreadMessages = 0))), + roomList(listOf(staleServerRoom(lastReadMessage = 10, unreadMessages = 0, hasArchived = true))), + roomList(listOf(staleServerRoom(lastReadMessage = 10, unreadMessages = 0))) ) runBlocking { @@ -369,11 +374,12 @@ class ConversationListFreshnessIntegrationTest { networkMonitor, syncer, conversationListUpdater, + ArbitraryStorageManager(ArbitraryStoragesRepositoryImpl(db.arbitraryStoragesDao())), ApplicationProvider.getApplicationContext(), mock() ) - whenever(conversationsNetwork.getRooms(any(), any(), any())).thenReturn( - Observable.just(listOf(ConversationDto(token = ROOM_TOKEN, lastActivity = 10, unreadMessages = 1))) + whenever(conversationsNetwork.getRooms(any(), any(), any(), anyOrNull())).thenReturn( + roomList(listOf(ConversationDto(token = ROOM_TOKEN, lastActivity = 10, unreadMessages = 1))) ) runBlocking { @@ -400,6 +406,7 @@ class ConversationListFreshnessIntegrationTest { networkMonitor, syncer, conversationListUpdater, + ArbitraryStorageManager(ArbitraryStoragesRepositoryImpl(db.arbitraryStoragesDao())), ApplicationProvider.getApplicationContext(), mock() ) @@ -500,3 +507,6 @@ class ConversationListFreshnessIntegrationTest { private const val POLL_INTERVAL_MILLIS = 50L } } + +private fun roomList(conversations: List): Observable = + Observable.just(RoomListResult(conversations, modifiedBefore = null, wasDelta = false)) diff --git a/app/src/test/java/com/nextcloud/talk/conversationlist/data/network/OfflineFirstConversationsRepositoryTest.kt b/app/src/test/java/com/nextcloud/talk/conversationlist/data/network/OfflineFirstConversationsRepositoryTest.kt index 73556891b8..9e00d419a6 100644 --- a/app/src/test/java/com/nextcloud/talk/conversationlist/data/network/OfflineFirstConversationsRepositoryTest.kt +++ b/app/src/test/java/com/nextcloud/talk/conversationlist/data/network/OfflineFirstConversationsRepositoryTest.kt @@ -10,6 +10,7 @@ package com.nextcloud.talk.conversationlist.data.network import android.content.Context import android.net.ConnectivityManager import android.os.PowerManager +import com.nextcloud.talk.arbitrarystorage.ArbitraryStorageManager import com.nextcloud.talk.chat.data.network.ChatMessageSyncer import com.nextcloud.talk.chat.data.network.ChatNetworkDataSource import com.nextcloud.talk.data.database.dao.ConversationsDao @@ -23,6 +24,7 @@ import com.nextcloud.talk.models.json.capabilities.CapabilitiesDto import com.nextcloud.talk.models.json.capabilities.SpreedCapabilityDto import com.nextcloud.talk.models.json.conversations.ConversationDto import com.nextcloud.talk.utils.SpreedFeatures +import io.reactivex.Maybe import io.reactivex.Observable import io.reactivex.android.plugins.RxAndroidPlugins import io.reactivex.schedulers.Schedulers @@ -78,6 +80,8 @@ class OfflineFirstConversationsRepositoryTest { private val powerManager: PowerManager = mock() private val connectivityManager: ConnectivityManager = mock() + private val arbitraryStorageManager: ArbitraryStorageManager = mock() + private lateinit var repository: OfflineFirstConversationsRepository @Before @@ -94,6 +98,8 @@ class OfflineFirstConversationsRepositoryTest { .thenReturn(ConnectivityManager.RESTRICT_BACKGROUND_STATUS_DISABLED) whenever(dao.getConversationsForUser(ACCOUNT_ID)).thenReturn(flowOf(emptyList())) + whenever(arbitraryStorageManager.getStorageSetting(any(), any(), any())) + .thenReturn(Maybe.empty()) whenever(conversationListUpdater.preservePendingLocalState(any(), any())) .thenAnswer { invocation -> invocation.getArgument>(1) } @@ -104,6 +110,7 @@ class OfflineFirstConversationsRepositoryTest { networkMonitor, chatMessageSyncer, conversationListUpdater, + arbitraryStorageManager, context, mock() ) @@ -121,14 +128,14 @@ class OfflineFirstConversationsRepositoryTest { repository.getRooms(user()).join() - verifyBlocking(network, never()) { getRooms(any(), any(), any()) } + verifyBlocking(network, never()) { getRooms(any(), any(), any(), anyOrNull()) } } @Test fun `getRooms fetches conversations from the server and syncs them locally when online`() = runBlocking { val room = conversation(token = ROOM_TOKEN, lastActivity = 5, unreadMessages = 0) - whenever(network.getRooms(any(), any(), any())).thenReturn(Observable.just(listOf(room))) + whenever(network.getRooms(any(), any(), any(), anyOrNull())).thenReturn(roomList(listOf(room))) repository.getRooms(user()).join() @@ -143,7 +150,7 @@ class OfflineFirstConversationsRepositoryTest { runBlocking { val previous = conversation(token = ROOM_TOKEN, lastActivity = 5, unreadMessages = 0).asEntity(ACCOUNT_ID) whenever(dao.getConversationsForUser(ACCOUNT_ID)).thenReturn(flowOf(listOf(previous))) - whenever(network.getRooms(any(), any(), any())).thenReturn(Observable.just(emptyList())) + whenever(network.getRooms(any(), any(), any(), anyOrNull())).thenReturn(roomList(listOf())) repository.getRooms(user()).join() @@ -159,7 +166,7 @@ class OfflineFirstConversationsRepositoryTest { val leaving = conversation(token = "roomB", lastActivity = 5, unreadMessages = 0).asEntity(ACCOUNT_ID) whenever(dao.getConversationsForUser(ACCOUNT_ID)).thenReturn(flowOf(listOf(staying, leaving))) val stayingRoom = conversation(token = "roomA", lastActivity = 5, unreadMessages = 0) - whenever(network.getRooms(any(), any(), any())).thenReturn(Observable.just(listOf(stayingRoom))) + whenever(network.getRooms(any(), any(), any(), anyOrNull())).thenReturn(roomList(listOf(stayingRoom))) repository.getRooms(user()).join() @@ -174,7 +181,7 @@ class OfflineFirstConversationsRepositoryTest { val previous = conversation(token = ROOM_TOKEN, lastActivity = 5, unreadMessages = 0).asEntity(ACCOUNT_ID) whenever(dao.getConversationsForUser(ACCOUNT_ID)).thenReturn(flowOf(listOf(previous))) val serverRoom = conversation(token = ROOM_TOKEN, lastActivity = 6, unreadMessages = 1) - whenever(network.getRooms(any(), any(), any())).thenReturn(Observable.just(listOf(serverRoom))) + whenever(network.getRooms(any(), any(), any(), anyOrNull())).thenReturn(roomList(listOf(serverRoom))) repository.getRooms(user()).join() @@ -189,7 +196,7 @@ class OfflineFirstConversationsRepositoryTest { fun `background message catch-up is skipped when the server lacks the chat-keep-notifications capability`() = runBlocking { val room = conversation(token = ROOM_TOKEN, lastActivity = 5, unreadMessages = 2) - whenever(network.getRooms(any(), any(), any())).thenReturn(Observable.just(listOf(room))) + whenever(network.getRooms(any(), any(), any(), anyOrNull())).thenReturn(roomList(listOf(room))) repository.getRooms(user(withKeepNotificationsCapability = false)).join() @@ -203,7 +210,7 @@ class OfflineFirstConversationsRepositoryTest { runBlocking { whenever(powerManager.isPowerSaveMode).thenReturn(true) val room = conversation(token = ROOM_TOKEN, lastActivity = 5, unreadMessages = 2) - whenever(network.getRooms(any(), any(), any())).thenReturn(Observable.just(listOf(room))) + whenever(network.getRooms(any(), any(), any(), anyOrNull())).thenReturn(roomList(listOf(room))) repository.getRooms(user()).join() @@ -219,7 +226,7 @@ class OfflineFirstConversationsRepositoryTest { whenever(connectivityManager.restrictBackgroundStatus) .thenReturn(ConnectivityManager.RESTRICT_BACKGROUND_STATUS_ENABLED) val room = conversation(token = ROOM_TOKEN, lastActivity = 5, unreadMessages = 2) - whenever(network.getRooms(any(), any(), any())).thenReturn(Observable.just(listOf(room))) + whenever(network.getRooms(any(), any(), any(), anyOrNull())).thenReturn(roomList(listOf(room))) repository.getRooms(user()).join() @@ -235,7 +242,7 @@ class OfflineFirstConversationsRepositoryTest { whenever(connectivityManager.restrictBackgroundStatus) .thenReturn(ConnectivityManager.RESTRICT_BACKGROUND_STATUS_DISABLED) val room = conversation(token = ROOM_TOKEN, lastActivity = 5, unreadMessages = 2) - whenever(network.getRooms(any(), any(), any())).thenReturn(Observable.just(listOf(room))) + whenever(network.getRooms(any(), any(), any(), anyOrNull())).thenReturn(roomList(listOf(room))) stubCatchUpRoom() repository.getRooms(user()).join() @@ -257,8 +264,8 @@ class OfflineFirstConversationsRepositoryTest { val unchangedRoom = conversation(token = "unchanged", lastActivity = 10, unreadMessages = 0) val noBlockRoom = conversation(token = "noBlockUnread", lastActivity = 10, unreadMessages = 3) - whenever(network.getRooms(any(), any(), any())) - .thenReturn(Observable.just(listOf(unchangedRoom, noBlockRoom))) + whenever(network.getRooms(any(), any(), any(), anyOrNull())) + .thenReturn(roomList(listOf(unchangedRoom, noBlockRoom))) whenever(chatMessageSyncer.hasLocalChatBlock("$ACCOUNT_ID@unchanged", null)).thenReturn(true) whenever(chatMessageSyncer.hasLocalChatBlock("$ACCOUNT_ID@noBlockUnread", null)).thenReturn(false) @@ -279,7 +286,7 @@ class OfflineFirstConversationsRepositoryTest { val rooms = (0 until 25).map { i -> conversation(token = "room$i", lastActivity = i.toLong(), unreadMessages = 0) } - whenever(network.getRooms(any(), any(), any())).thenReturn(Observable.just(rooms)) + whenever(network.getRooms(any(), any(), any(), anyOrNull())).thenReturn(roomList(rooms)) stubCatchUpRoom() repository.getRooms(user()).join() @@ -454,3 +461,6 @@ class OfflineFirstConversationsRepositoryTest { private const val COLLECTOR_STARTUP_MILLIS = 100L } } + +private fun roomList(conversations: List): Observable = + Observable.just(RoomListResult(conversations, modifiedBefore = null, wasDelta = false)) diff --git a/app/src/test/java/com/nextcloud/talk/conversationlist/data/network/RoomListMessagePrefetchIntegrationTest.kt b/app/src/test/java/com/nextcloud/talk/conversationlist/data/network/RoomListMessagePrefetchIntegrationTest.kt index 5133dfd0a0..47e2d3c706 100644 --- a/app/src/test/java/com/nextcloud/talk/conversationlist/data/network/RoomListMessagePrefetchIntegrationTest.kt +++ b/app/src/test/java/com/nextcloud/talk/conversationlist/data/network/RoomListMessagePrefetchIntegrationTest.kt @@ -11,11 +11,13 @@ import android.app.Application import android.content.Context import androidx.room.Room import androidx.test.core.app.ApplicationProvider +import com.nextcloud.talk.arbitrarystorage.ArbitraryStorageManager import com.nextcloud.talk.chat.data.model.ChatMessage import com.nextcloud.talk.chat.data.network.ChatMessageSyncer import com.nextcloud.talk.chat.data.network.ChatNetworkDataSource import com.nextcloud.talk.data.network.NetworkMonitor import com.nextcloud.talk.data.source.local.TalkDatabase +import com.nextcloud.talk.data.storage.ArbitraryStoragesRepositoryImpl import com.nextcloud.talk.data.user.model.User import com.nextcloud.talk.data.user.model.UserEntity import com.nextcloud.talk.logger.Logger @@ -38,6 +40,7 @@ import org.junit.Before import org.junit.Test import org.junit.runner.RunWith import org.mockito.kotlin.any +import org.mockito.kotlin.anyOrNull import org.mockito.kotlin.eq import org.mockito.kotlin.mock import org.mockito.kotlin.whenever @@ -94,6 +97,7 @@ class RoomListMessagePrefetchIntegrationTest { networkMonitor, syncer, conversationListUpdater, + ArbitraryStorageManager(ArbitraryStoragesRepositoryImpl(db.arbitraryStoragesDao())), context, mock() ) @@ -111,7 +115,8 @@ class RoomListMessagePrefetchIntegrationTest { conversation(ROOM_A, unreadMessages = 2, lastMessageId = 12), conversation(ROOM_B, unreadMessages = 1, lastMessageId = 7) ) - whenever(conversationsNetwork.getRooms(any(), any(), any())).thenReturn(Observable.just(rooms)) + whenever(conversationsNetwork.getRooms(any(), any(), any(), anyOrNull())) + .thenReturn(roomList(rooms)) wheneverBlocking { chatNetwork.pullChatMessages(any(), eq(chatUrl(ROOM_A)), any()) } .thenReturn(Response.success(overall(message(10, ROOM_A), message(11, ROOM_A), message(12, ROOM_A)))) wheneverBlocking { chatNetwork.pullChatMessages(any(), eq(chatUrl(ROOM_B)), any()) } @@ -207,3 +212,6 @@ class RoomListMessagePrefetchIntegrationTest { private const val POLL_INTERVAL_MILLIS = 50L } } + +private fun roomList(conversations: List): Observable = + Observable.just(RoomListResult(conversations, modifiedBefore = null, wasDelta = false))