diff --git a/app/src/main/java/com/nextcloud/talk/chat/ChatActivity.kt b/app/src/main/java/com/nextcloud/talk/chat/ChatActivity.kt index c8713fd2a1..7ff8182d0b 100644 --- a/app/src/main/java/com/nextcloud/talk/chat/ChatActivity.kt +++ b/app/src/main/java/com/nextcloud/talk/chat/ChatActivity.kt @@ -1030,11 +1030,7 @@ class ChatActivity : showEdit = sendingFailed || !isOnline, showDelete = sendingFailed || !isOnline, onResend = { - chatViewModel.resendMessage( - conversationUser!!.getCredentials(), - ApiUtils.getUrlForChat(chatApiVersion, conversationUser!!.baseUrl!!, roomToken), - msg - ) + chatViewModel.resendMessage(msg) }, onEdit = { messageInputViewModel.edit(msg) }, onDelete = { chatViewModel.deleteTempMessage(msg) }, diff --git a/app/src/main/java/com/nextcloud/talk/chat/MessageInputFragment.kt b/app/src/main/java/com/nextcloud/talk/chat/MessageInputFragment.kt index 5827bb6471..2fcdb14f5f 100644 --- a/app/src/main/java/com/nextcloud/talk/chat/MessageInputFragment.kt +++ b/app/src/main/java/com/nextcloud/talk/chat/MessageInputFragment.kt @@ -1049,12 +1049,8 @@ class MessageInputFragment : Fragment() { chatActivity.chatViewModel.onMessageSent() messageInputViewModel.sendChatMessage( - credentials = chatActivity.conversationUser!!.getCredentials(), - url = ApiUtils.getUrlForChat( - chatActivity.chatApiVersion, - chatActivity.conversationUser!!.baseUrl!!, - chatActivity.roomToken - ), + userId = chatActivity.conversationUser!!.id!!, + roomToken = chatActivity.roomToken, message = message, displayName = chatActivity.conversationUser!!.displayName ?: "", replyTo = chatActivity.getReplyToMessageId(), diff --git a/app/src/main/java/com/nextcloud/talk/chat/data/ChatMessageRepository.kt b/app/src/main/java/com/nextcloud/talk/chat/data/ChatMessageRepository.kt index ad9830c654..9ae1256ada 100644 --- a/app/src/main/java/com/nextcloud/talk/chat/data/ChatMessageRepository.kt +++ b/app/src/main/java/com/nextcloud/talk/chat/data/ChatMessageRepository.kt @@ -138,16 +138,11 @@ interface ChatMessageRepository : LifecycleAwareManager { threadTitle: String? ): Flow> - @Suppress("LongParameterList") - suspend fun resendChatMessage( - credentials: String, - url: String, - message: String, - displayName: String, - replyTo: Int, - sendWithoutNotification: Boolean, - referenceId: String - ): Flow> + /** + * Resets a previously failed temporary message back to PENDING so it can be handed to + * SendMessageWorker for another send attempt. Does not itself send anything. + */ + suspend fun markMessageForResend(referenceId: String): Flow> suspend fun addTemporaryMessage( message: CharSequence, diff --git a/app/src/main/java/com/nextcloud/talk/chat/data/network/OfflineFirstChatRepository.kt b/app/src/main/java/com/nextcloud/talk/chat/data/network/OfflineFirstChatRepository.kt index a55885699b..656a07289c 100644 --- a/app/src/main/java/com/nextcloud/talk/chat/data/network/OfflineFirstChatRepository.kt +++ b/app/src/main/java/com/nextcloud/talk/chat/data/network/OfflineFirstChatRepository.kt @@ -32,6 +32,7 @@ import com.nextcloud.talk.models.json.generic.GenericOverall import com.nextcloud.talk.models.json.participants.Participant import com.nextcloud.talk.utils.bundle.BundleKeys import com.nextcloud.talk.utils.message.SendMessageUtils +import kotlinx.coroutines.CancellationException import kotlinx.coroutines.delay import kotlinx.coroutines.flow.Flow import kotlinx.coroutines.flow.MutableSharedFlow @@ -656,16 +657,7 @@ class OfflineFirstChatRepository @Inject constructor( } } - @Suppress("LongParameterList") - override suspend fun resendChatMessage( - credentials: String, - url: String, - message: String, - displayName: String, - replyTo: Int, - sendWithoutNotification: Boolean, - referenceId: String - ): Flow> { + override suspend fun markMessageForResend(referenceId: String): Flow> { val messageToResend = chatDao.getTempMessageForConversation( internalConversationId, referenceId, @@ -678,16 +670,9 @@ class OfflineFirstChatRepository @Inject constructor( val messageToResendModel = messageToResend.toDomainModel() _updateMessageFlow.emit(messageToResendModel) - sendChatMessage( - credentials = credentials, - url = url, - message = message, - displayName = displayName, - replyTo = replyTo, - sendWithoutNotification = sendWithoutNotification, - referenceId = referenceId, - threadTitle = null - ) + flow { + emit(Result.success(messageToResendModel)) + } } else { flow { emit(Result.failure(IllegalStateException("No temporary message found to resend"))) @@ -713,6 +698,12 @@ class OfflineFirstChatRepository @Inject constructor( referenceId ) chatDao.upsertChatMessage(tempChatMessageEntity) + emit(Result.success(tempChatMessageEntity.toDomainModel())) + } catch (e: CancellationException) { + // a collector (e.g. first()/take(1)) is done with the flow, not a real failure - + // rethrow instead of turning it into a Result.failure emission, which would violate + // flow exception transparency since the collector already stopped listening + throw e } catch (e: Exception) { Log.e(TAG, "Something went wrong when adding temporary message", e) emit(Result.failure(e)) diff --git a/app/src/main/java/com/nextcloud/talk/chat/viewmodels/ChatViewModel.kt b/app/src/main/java/com/nextcloud/talk/chat/viewmodels/ChatViewModel.kt index 30ab6b1e4c..944ca5cc8d 100644 --- a/app/src/main/java/com/nextcloud/talk/chat/viewmodels/ChatViewModel.kt +++ b/app/src/main/java/com/nextcloud/talk/chat/viewmodels/ChatViewModel.kt @@ -41,8 +41,8 @@ import com.nextcloud.talk.dagger.modules.ApplicationScope import com.nextcloud.talk.data.database.mappers.toDomainModel import com.nextcloud.talk.data.database.model.ChatMessageEntity import com.nextcloud.talk.data.user.model.User -import com.nextcloud.talk.extensions.toIntOrZero import com.nextcloud.talk.jobs.ReadMarkerSyncWorker +import com.nextcloud.talk.jobs.SendMessageWorker import com.nextcloud.talk.jobs.ShareOperationWorker import androidx.lifecycle.asFlow import androidx.work.WorkManager @@ -2559,19 +2559,17 @@ class ChatViewModel @AssistedInject constructor( } } - fun resendMessage(credentials: String, urlForChat: String, message: ChatMessage) { + fun resendMessage(message: ChatMessage) { + val referenceId = message.referenceId.orEmpty() viewModelScope.launch { - chatRepository.resendChatMessage( - credentials, - urlForChat, - message.message.orEmpty(), - message.actorDisplayName.orEmpty(), - message.parentMessageId?.toIntOrZero() ?: 0, - false, - message.referenceId.orEmpty() - ).collect { result -> + chatRepository.markMessageForResend(referenceId).collect { result -> if (result.isSuccess) { - Log.d(TAG, "resend successful") + Log.d(TAG, "message marked pending for resend") + SendMessageWorker.enqueue( + internalConversationId = "${currentUser.id}@$chatRoomToken", + referenceId = referenceId, + threadTitle = null + ) } else { Log.e(TAG, "resend failed") } diff --git a/app/src/main/java/com/nextcloud/talk/chat/viewmodels/MessageInputViewModel.kt b/app/src/main/java/com/nextcloud/talk/chat/viewmodels/MessageInputViewModel.kt index 5f38d1b79f..a6b79ed97c 100644 --- a/app/src/main/java/com/nextcloud/talk/chat/viewmodels/MessageInputViewModel.kt +++ b/app/src/main/java/com/nextcloud/talk/chat/viewmodels/MessageInputViewModel.kt @@ -23,6 +23,7 @@ import com.nextcloud.talk.chat.data.io.AudioRecorderManager import com.nextcloud.talk.chat.data.io.MediaPlayerManager import com.nextcloud.talk.chat.data.model.ChatMessage import com.nextcloud.talk.chat.data.network.ChatNetworkDataSource +import com.nextcloud.talk.jobs.SendMessageWorker import com.nextcloud.talk.models.MessageDraft import com.nextcloud.talk.models.json.chat.ChatOverallSingleMessage import com.nextcloud.talk.models.json.chat.ChatUtils @@ -33,6 +34,8 @@ import kotlinx.coroutines.flow.MutableStateFlow import kotlinx.coroutines.flow.StateFlow import kotlinx.coroutines.flow.update import kotlinx.coroutines.launch +import kotlinx.coroutines.sync.Mutex +import kotlinx.coroutines.sync.withLock import javax.inject.Inject @Suppress("Detekt.TooManyFunctions") @@ -77,6 +80,13 @@ class MessageInputViewModel : lateinit var currentLifeCycleFlag: LifeCycleFlag val disposableSet = mutableSetOf() + // Serializes addTemporaryMessage()+SendMessageWorker.enqueue() across sendChatMessage() calls. + // Without it, several messages sent in quick succession each await their own independent local + // DB write before enqueueing, and those writes aren't guaranteed to finish in the order they + // were started - so the worker chain (and therefore the order messages actually reach the + // server) could end up scrambled relative to the order the user sent them. + private val sendMessageMutex = Mutex() + fun setData(chatMessageRepository: ChatMessageRepository) { chatRepository = chatMessageRepository } @@ -164,8 +174,8 @@ class MessageInputViewModel : @Suppress("LongParameterList") fun sendChatMessage( - credentials: String, - url: String, + userId: Long, + roomToken: String, message: String, displayName: String, replyTo: Int, @@ -176,40 +186,28 @@ class MessageInputViewModel : Log.d(TAG, "Random SHA-256 Hash: $referenceId") viewModelScope.launch { - chatRepository.addTemporaryMessage( - message, - displayName, - replyTo, - sendWithoutNotification, - referenceId - ).collect { result -> - if (result.isSuccess) { - Log.d(TAG, "temp message ref id: " + (result.getOrNull()?.referenceId ?: "none")) - - _sendChatMessageViewState.value = SendChatMessageSuccessState(message) - } else { - _sendChatMessageViewState.value = SendChatMessageErrorState(message) - } - } - } - - viewModelScope.launch { - chatRepository.sendChatMessage( - credentials, - url, - message, - displayName, - replyTo, - sendWithoutNotification, - referenceId, - threadTitle - ).collect { result -> - if (result.isSuccess) { - Log.d(TAG, "received ref id: " + (result.getOrNull()?.referenceId ?: "none")) - - _sendChatMessageViewState.value = SendChatMessageSuccessState(message) - } else { - _sendChatMessageViewState.value = SendChatMessageErrorState(message) + // Holding the lock across the DB write and the enqueue call, rather than just around + // the enqueue call, is what actually guarantees ordering: it forces this whole + // insert-then-enqueue step for one message to finish before the next queued send is + // even allowed to start its own DB write. + sendMessageMutex.withLock { + chatRepository.addTemporaryMessage( + message, + displayName, + replyTo, + sendWithoutNotification, + referenceId + ).collect { result -> + if (result.isSuccess) { + _sendChatMessageViewState.value = SendChatMessageSuccessState(message) + SendMessageWorker.enqueue( + internalConversationId = "$userId@$roomToken", + referenceId = referenceId, + threadTitle = threadTitle + ) + } else { + _sendChatMessageViewState.value = SendChatMessageErrorState(message) + } } } } diff --git a/app/src/main/java/com/nextcloud/talk/jobs/SendMessageWorker.kt b/app/src/main/java/com/nextcloud/talk/jobs/SendMessageWorker.kt new file mode 100644 index 0000000000..e8a3ecab50 --- /dev/null +++ b/app/src/main/java/com/nextcloud/talk/jobs/SendMessageWorker.kt @@ -0,0 +1,182 @@ +/* + * Nextcloud Talk - Android Client + * + * SPDX-FileCopyrightText: 2026 Nextcloud GmbH and Nextcloud contributors + * SPDX-License-Identifier: GPL-3.0-or-later + */ +package com.nextcloud.talk.jobs + +import android.content.Context +import android.util.Log +import androidx.work.BackoffPolicy +import androidx.work.Constraints +import androidx.work.CoroutineWorker +import androidx.work.Data +import androidx.work.ExistingWorkPolicy +import androidx.work.NetworkType +import androidx.work.OneTimeWorkRequest +import androidx.work.WorkManager +import androidx.work.WorkRequest +import androidx.work.WorkerParameters +import autodagger.AutoInjector +import com.nextcloud.talk.application.NextcloudTalkApplication +import com.nextcloud.talk.application.NextcloudTalkApplication.Companion.sharedApplication +import com.nextcloud.talk.chat.data.network.ChatNetworkDataSource +import com.nextcloud.talk.data.database.dao.ChatMessagesDao +import com.nextcloud.talk.data.database.model.ChatMessageEntity +import com.nextcloud.talk.data.database.model.SendStatus +import com.nextcloud.talk.users.UserManager +import com.nextcloud.talk.utils.ApiUtils +import kotlinx.coroutines.flow.firstOrNull +import java.io.IOException +import java.util.concurrent.TimeUnit +import javax.inject.Inject + +/** + * Sends a chat text message to the server, retrying transient failures with backoff. + * + * Enqueued via WorkManager instead of running in a viewModelScope coroutine, so the send survives + * leaving the conversation, the app going to background, or process death - unlike a plain + * viewModelScope coroutine, which is cancelled the moment ChatActivity's ViewModel is cleared, + * potentially leaving a message stuck as PENDING with nothing left running to retry it. Work is + * chained per conversation with [ExistingWorkPolicy.APPEND_OR_REPLACE] (mirroring + * [UploadAndShareFilesWorker]) so messages typed in a row still arrive in the order they were + * sent, and a failed/cancelled message ahead in the queue doesn't block the ones behind it. + */ +@AutoInjector(NextcloudTalkApplication::class) +class SendMessageWorker(context: Context, workerParams: WorkerParameters) : CoroutineWorker(context, workerParams) { + + @Inject + lateinit var userManager: UserManager + + @Inject + lateinit var chatNetworkDataSource: ChatNetworkDataSource + + @Inject + lateinit var chatDao: ChatMessagesDao + + override suspend fun doWork(): Result { + sharedApplication!!.componentApplication.inject(this) + + val internalConversationId = inputData.getString(KEY_INTERNAL_CONVERSATION_ID) + val referenceId = inputData.getString(KEY_REFERENCE_ID) + val threadTitle = inputData.getString(KEY_THREAD_TITLE) + + if (internalConversationId.isNullOrEmpty() || referenceId.isNullOrEmpty()) { + Log.e(TAG, "Missing required input data, dropping message send") + return Result.failure() + } + + return sendMessage(internalConversationId, referenceId, threadTitle) + } + + @Suppress("Detekt.TooGenericExceptionCaught") + private suspend fun sendMessage(internalConversationId: String, referenceId: String, threadTitle: String?): Result { + // Re-checked on every attempt (including retries): the user may have deleted this + // temporary message - e.g. from the message actions sheet while offline - after it was + // enqueued but before a connection was available to actually send it. Without this check + // the message would still be posted to the server even though it no longer exists locally. + // The message text and its other send parameters are also read fresh from this row rather + // than passed in via WorkManager's input data, so an edit made to a still-queued message + // while offline (editTempChatMessage only updates this row) is picked up instead of sending + // stale data. + val tempMessage = chatDao.getTempMessageForConversation(internalConversationId, referenceId, null) + .firstOrNull() + if (tempMessage == null) { + Log.d(TAG, "Temporary message $referenceId no longer exists, skipping send") + return Result.success() + } + + val user = userManager.getUserWithId(tempMessage.accountId).blockingGet() + val credentials = user?.let { ApiUtils.getCredentials(it.username, it.token) } + if (user == null || credentials == null) { + Log.e(TAG, "No user or credentials found for account id ${tempMessage.accountId}, failing message send") + return failMessage(tempMessage) + } + + val apiVersion = ApiUtils.getChatApiVersion(user.capabilities!!.spreedCapability!!, intArrayOf(ApiUtils.API_V1)) + val url = ApiUtils.getUrlForChat(apiVersion, user.baseUrl!!, tempMessage.token) + + return try { + chatNetworkDataSource.sendChatMessage( + credentials, + url, + tempMessage.message, + tempMessage.actorDisplayName, + tempMessage.parentMessageId?.toInt() ?: 0, + tempMessage.silent, + referenceId, + threadTitle + ) + updateStatus(tempMessage, SendStatus.SENT_PENDING_ACK) + Log.d(TAG, "sending chat message succeeded: ${tempMessage.message}") + Result.success() + } catch (e: IOException) { + Log.w(TAG, "Network error while sending message (attempt ${runAttemptCount + 1}/$MAX_SEND_ATTEMPTS)", e) + retryOrFail(tempMessage) + } catch (e: Exception) { + Log.e(TAG, "Something went wrong when sending message", e) + failMessage(tempMessage) + } + } + + private fun retryOrFail(tempMessage: ChatMessageEntity): Result = + if (runAttemptCount < MAX_SEND_ATTEMPTS - 1) { + Result.retry() + } else { + failMessage(tempMessage) + } + + private fun failMessage(tempMessage: ChatMessageEntity): Result { + updateStatus(tempMessage, SendStatus.FAILED) + return Result.failure() + } + + private fun updateStatus(tempMessage: ChatMessageEntity, status: SendStatus) { + tempMessage.sendStatus = status + chatDao.updateChatMessage(tempMessage) + } + + companion object { + private val TAG = SendMessageWorker::class.simpleName + private const val KEY_INTERNAL_CONVERSATION_ID = "INTERNAL_CONVERSATION_ID" + private const val KEY_REFERENCE_ID = "REFERENCE_ID" + private const val KEY_THREAD_TITLE = "THREAD_TITLE" + + // Total attempts allowed for a single message (1 initial run + retries) before giving up on + // a transient network failure and marking it FAILED so the user can resend manually. + private const val MAX_SEND_ATTEMPTS = 4 + + // userId, roomToken, message text, displayName, replyTo and sendWithoutNotification are + // deliberately not passed in here: they're all already persisted on the temporary message + // row (accountId, token, actorDisplayName, parentMessageId, silent), which doWork() reads + // fresh instead, so a later edit or deletion of that row is always reflected. threadTitle + // isn't persisted on that row, so it still has to travel through the work request. + fun enqueue(internalConversationId: String, referenceId: String, threadTitle: String?) { + val data = Data.Builder() + .putString(KEY_INTERNAL_CONVERSATION_ID, internalConversationId) + .putString(KEY_REFERENCE_ID, referenceId) + .putString(KEY_THREAD_TITLE, threadTitle) + .build() + + val sendWork = OneTimeWorkRequest.Builder(SendMessageWorker::class.java) + .setInputData(data) + .setConstraints(Constraints.Builder().setRequiredNetworkType(NetworkType.CONNECTED).build()) + .setBackoffCriteria(BackoffPolicy.EXPONENTIAL, WorkRequest.MIN_BACKOFF_MILLIS, TimeUnit.MILLISECONDS) + .build() + + // Chained per conversation (not enqueueUniqueWork(referenceId, ...), which would run every + // message fully independently) so several messages typed in a row still arrive in the + // order they were sent. APPEND_OR_REPLACE rather than APPEND: if the message ahead in the + // queue permanently failed, this starts a fresh chain instead of cascading that failure + // onto every message queued behind it. + WorkManager.getInstance().enqueueUniqueWork( + sendQueueName(internalConversationId), + ExistingWorkPolicy.APPEND_OR_REPLACE, + sendWork + ) + } + + private fun sendQueueName(internalConversationId: String) = "send_message_queue_$internalConversationId" + } +} diff --git a/app/src/test/java/com/nextcloud/talk/chat/data/network/OfflineFirstChatRepositoryTest.kt b/app/src/test/java/com/nextcloud/talk/chat/data/network/OfflineFirstChatRepositoryTest.kt index 00343b1930..e01a43f256 100644 --- a/app/src/test/java/com/nextcloud/talk/chat/data/network/OfflineFirstChatRepositoryTest.kt +++ b/app/src/test/java/com/nextcloud/talk/chat/data/network/OfflineFirstChatRepositoryTest.kt @@ -14,6 +14,8 @@ import com.nextcloud.talk.data.database.dao.ChatBlocksDao import com.nextcloud.talk.data.database.dao.ChatMessagesDao 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.SendStatus import com.nextcloud.talk.data.network.NetworkMonitor import com.nextcloud.talk.data.user.model.User import com.nextcloud.talk.logger.Logger @@ -25,10 +27,13 @@ import com.nextcloud.talk.models.json.chat.ChatOCS import com.nextcloud.talk.models.json.chat.ChatOverall import com.nextcloud.talk.models.json.conversations.Conversation import kotlinx.coroutines.flow.MutableStateFlow +import kotlinx.coroutines.flow.first import kotlinx.coroutines.flow.flowOf +import kotlinx.coroutines.flow.toList import kotlinx.coroutines.test.runTest import org.junit.Assert.assertEquals import org.junit.Assert.assertFalse +import org.junit.Assert.assertTrue import org.junit.Before import org.junit.Test import org.mockito.kotlin.any @@ -36,6 +41,7 @@ import org.mockito.kotlin.argumentCaptor import org.mockito.kotlin.eq import org.mockito.kotlin.mock import org.mockito.kotlin.never +import org.mockito.kotlin.verify import org.mockito.kotlin.verifyBlocking import org.mockito.kotlin.whenever import org.mockito.kotlin.wheneverBlocking @@ -190,6 +196,96 @@ class OfflineFirstChatRepositoryTest { verifyBlocking(network, never()) { pullChatMessages(any(), any(), any()) } } + @Test + fun `addTemporaryMessage upserts a pending local echo and emits it on success`() = + runTest { + repository.updateConversation(conversation(lastReadMessage = 0, unreadMessages = 0)) + + val result = repository.addTemporaryMessage( + message = "hello", + displayName = "Me", + replyTo = 0, + sendWithoutNotification = false, + referenceId = "ref-1" + ).toList().single() + + assertTrue(result.isSuccess) + assertEquals("hello", result.getOrNull()?.message) + assertEquals(SendStatus.PENDING, result.getOrNull()?.sendStatus) + assertEquals(true, result.getOrNull()?.isTemporary) + + val entityCaptor = argumentCaptor() + verifyBlocking(chatDao) { upsertChatMessage(entityCaptor.capture()) } + assertEquals(INTERNAL_CONVERSATION_ID, entityCaptor.firstValue.internalConversationId) + assertEquals("ref-1", entityCaptor.firstValue.referenceId) + } + + @Test + fun `addTemporaryMessage tolerates a collector that stops after the first value`() = + runTest { + // first()/take(1) cancel the flow right after receiving one value, which used to + // surface as "Flow exception transparency is violated" because that cancellation was + // caught by addTemporaryMessage's own catch(Exception) block and turned into a second, + // illegal emit() call. + repository.updateConversation(conversation(lastReadMessage = 0, unreadMessages = 0)) + + val result = repository.addTemporaryMessage( + message = "hello", + displayName = "Me", + replyTo = 0, + sendWithoutNotification = false, + referenceId = "ref-1b" + ).first() + + assertTrue(result.isSuccess) + } + + @Test + fun `markMessageForResend resets a failed temp message back to PENDING without sending anything`() = + runTest { + val failedMessage = tempMessageEntity(referenceId = "ref-2", sendStatus = SendStatus.FAILED) + whenever(chatDao.getTempMessageForConversation(INTERNAL_CONVERSATION_ID, "ref-2", null)) + .thenReturn(flowOf(failedMessage)) + + val result = repository.markMessageForResend("ref-2").toList().single() + + assertTrue(result.isSuccess) + assertEquals(SendStatus.PENDING, result.getOrNull()?.sendStatus) + assertEquals(SendStatus.PENDING, failedMessage.sendStatus) + verify(chatDao).updateChatMessage(failedMessage) + // resending is only a local status reset - the actual send is left to SendMessageWorker + verifyBlocking(network, never()) { sendChatMessage(any(), any(), any(), any(), any(), any(), any(), any()) } + } + + @Test + fun `markMessageForResend fails when no temp message exists for the reference id`() = + runTest { + whenever(chatDao.getTempMessageForConversation(INTERNAL_CONVERSATION_ID, "missing", null)) + .thenReturn(flowOf(null)) + + val result = repository.markMessageForResend("missing").toList().single() + + assertTrue(result.isFailure) + verify(chatDao, never()).updateChatMessage(any()) + } + + private fun tempMessageEntity(referenceId: String, sendStatus: SendStatus): ChatMessageEntity = + ChatMessageEntity( + internalId = "$INTERNAL_CONVERSATION_ID@_temp_$referenceId", + accountId = ACCOUNT_ID, + token = ROOM_TOKEN, + internalConversationId = INTERNAL_CONVERSATION_ID, + actorDisplayName = "Me", + message = "hello", + actorId = "me", + actorType = "users", + messageType = "comment", + systemMessageType = ChatMessage.SystemMessageType.DUMMY, + referenceId = referenceId, + isTemporary = true, + sendStatus = sendStatus + ) + private fun givenLatestBlock(block: ChatBlockEntity?) { whenever(chatBlocksDao.getLatestChatBlock(INTERNAL_CONVERSATION_ID, null)) .thenReturn(flowOf(block)) @@ -206,6 +302,7 @@ class OfflineFirstChatRepositoryTest { id = ACCOUNT_ID, userId = "me", username = "me", + displayName = "Me", baseUrl = "https://server.example.com", capabilities = Capabilities().apply { spreedCapability = SpreedCapability().apply { features = listOf("chat-keep-notifications") }