Skip to content
Draft
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
7 changes: 4 additions & 3 deletions app/src/main/java/com/nextcloud/talk/api/NcApi.java
Original file line number Diff line number Diff line change
Expand Up @@ -95,9 +95,10 @@ Observable<ResponseBody> getContactsWithSearchParam(@Header("Authorization") Str
Server URL is: baseUrl + ocsApiVersion + spreedApiVersion + /room
*/
@GET
Observable<RoomsOverall> getRooms(@Header("Authorization") String authorization,
@Url String url,
@Nullable @Query("includeStatus") Boolean includeStatus);
Observable<Response<RoomsOverall>> 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
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -279,7 +279,7 @@ class ConversationsListActivity : BaseActivity() {
isMaintenanceModeState.value = false
isRefreshingState.value = true
appPreferences.setConversationListPositionAndOffset(0, 0)
fetchRooms()
fetchRooms(forceFullSync = true)
fetchPendingInvitations()
},
onFabClick = {
Expand Down Expand Up @@ -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() {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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]
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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<ConversationDto>, val modifiedBefore: Long?, val wasDelta: Boolean)

interface ConversationsNetworkDataSource {
fun getRooms(user: User, url: String, includeStatus: Boolean): Observable<List<ConversationDto>>
/**
* 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<RoomListResult>
}
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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
Expand All @@ -49,13 +51,15 @@ 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,
private val chatNetworkDataSource: ChatNetworkDataSource,
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 {
Expand Down Expand Up @@ -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)
}
}

Expand Down Expand Up @@ -163,56 +168,114 @@ class OfflineFirstConversationsRepository @Inject constructor(
}

@Suppress("Detekt.TooGenericExceptionCaught")
private suspend fun getRoomsFromServer(user: User): List<ConversationEntity>? {
private suspend fun getRoomsFromServer(user: User, forceFullSync: Boolean = false): List<ConversationEntity>? {
var conversationsFromSync: List<ConversationEntity>? = null

if (!networkMonitor.isOnline.value) {
Log.d(TAG, "Device is offline, can't load conversations from server")
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)
}
}
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
Expand Down Expand Up @@ -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)
}
}
}
}
Expand Down Expand Up @@ -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"
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -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<List<ConversationDto>> {
override fun getRooms(
user: User,
url: String,
includeStatus: Boolean,
modifiedSince: Long?
): Observable<RoomListResult> {
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"
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -204,6 +205,7 @@ class RepositoryModule {
networkMonitor: NetworkMonitor,
chatMessageSyncer: ChatMessageSyncer,
conversationListUpdater: ConversationListUpdater,
arbitraryStorageManager: ArbitraryStorageManager,
context: Context,
logger: Logger
): OfflineConversationsRepository =
Expand All @@ -214,6 +216,7 @@ class RepositoryModule {
networkMonitor,
chatMessageSyncer,
conversationListUpdater,
arbitraryStorageManager,
context,
logger
)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -26,9 +26,11 @@ interface ConversationsDao {
fun getConversationForUser(accountId: Long, token: String): Flow<ConversationEntity?>

/**
* 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(
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -193,6 +194,7 @@ class ComposePreviewUtils private constructor(context: Context) {
networkMonitor,
chatMessageSyncer,
conversationListUpdater,
ArbitraryStorageManager(DummyArbitraryStoragesRepositoryImpl()),
mContext,
logger
)
Expand Down
Loading
Loading