Кэш чатов

This commit is contained in:
Халимов Рустам
2026-04-20 00:25:40 +03:00
parent 945134f029
commit 7c66e1c0c0
35 changed files with 2131 additions and 92 deletions
@@ -0,0 +1,80 @@
package chats.data.paging
import androidx.paging.PagingSource
import androidx.paging.PagingState
import core.database.data.MessageDao
import core.database.data.MessageEntity
import android.util.Log
/**
* PagingSource для загрузки сообщений из локальной базы Room
* Данные сортируются по sequenceId DESC (новые сообщения первыми)
*/
class MessagePagingSource(
private val dao: MessageDao,
private val chatId: String
) : PagingSource<Int, MessageEntity>() {
companion object {
private const val TAG = "MessagePagingSource"
}
override fun getRefreshKey(state: PagingState<Int, MessageEntity>): Int? {
Log.d(TAG, "getRefreshKey() called")
return state.anchorPosition?.let { anchorPosition ->
val anchorPage = state.closestPageToPosition(anchorPosition)
anchorPage?.prevKey?.plus(1) ?: anchorPage?.nextKey?.minus(1)
}
}
override suspend fun load(params: LoadParams<Int>): LoadResult<Int, MessageEntity> {
Log.d(TAG, ">>> load() called with: position=${params.key}, loadSize=${params.loadSize}")
return try {
// При DESC порядке:
// - key = sequenceId последнего сообщения в текущей странице
// - APPEND = загружаем сообщения С МЕНЬШИМ sequenceId (старые)
// - PREPEND = загружаем сообщения С БОЛЬШИМ sequenceId (новые)
val position = params.key // sequenceId граничного сообщения
val loadSize = params.loadSize
Log.d(TAG, "Loading: chatId=$chatId, position=$position, loadSize=$loadSize")
val messages = if (position == null) {
// Первая загрузка - получаем самые последние сообщения (включая maxSeq)
val maxSeq = dao.getMaxSequenceId(chatId)
Log.d(TAG, "First load: maxSeq=$maxSeq")
if (maxSeq == null) {
Log.d(TAG, "No messages in database")
emptyList()
} else {
// Используем <= чтобы включить самое последнее сообщение
dao.getMessagesUpToAndIncluding(chatId, maxSeq, loadSize)
}
} else {
// Загружаем сообщения с sequenceId < position (старые)
Log.d(TAG, "Append: loading messages before sequenceId=$position")
dao.getMessagesBefore(chatId, position, loadSize)
}
Log.d(TAG, "Loaded ${messages.size} messages, first=${messages.firstOrNull()?.sequenceId}, last=${messages.lastOrNull()?.sequenceId}")
// Для DESC порядка:
// - prevKey = максимальный sequenceId в странице (для загрузки новых)
// - nextKey = минимальный sequenceId в странице (для загрузки старых)
val prevKey = messages.firstOrNull()?.sequenceId?.plus(1)
val nextKey = messages.lastOrNull()?.sequenceId?.minus(1)
Log.d(TAG, "prevKey=$prevKey, nextKey=$nextKey")
LoadResult.Page(
data = messages,
prevKey = prevKey,
nextKey = nextKey
)
} catch (e: Exception) {
Log.e(TAG, "Error loading messages", e)
LoadResult.Error(e)
}
}
}
@@ -0,0 +1,203 @@
package chats.data.paging
import android.util.Log
import androidx.paging.*
import chats.data.remote.api.ChatApi
import chats.data.remote.dto.MessageDto
import chats.domain.model.MediaType
import com.google.gson.Gson
import core.database.data.ChatDatabase
import core.database.data.MessageDao
import core.database.data.MessageEntity
import core.database.data.SyncStatus
import core.network.ServerConfig
import core.security.TokenManager
/**
* RemoteMediator для Paging 3
* Управляет загрузкой сообщений из сети и кэшированием в Room
*
* Логика работы:
* 1. При первой загрузке (REFRESH) - загружаем последние сообщения
* 2. При прокрутке вниз (APPEND) - загружаем более старые сообщения
* 3. При прокрутке вверх (PREPEND) - загружаем более новые сообщения
* 4. Данные сохраняются в Room, Paging читает из базы
*/
@OptIn(ExperimentalPagingApi::class)
class MessageRemoteMediator(
private val chatId: String,
private val api: ChatApi,
private val database: ChatDatabase,
private val dao: core.database.data.MessageDao,
private val serverConfig: ServerConfig,
private val tokenManager: TokenManager
) : RemoteMediator<Int, MessageEntity>() {
private val gson = Gson()
private val baseUrl = serverConfig.getBaseUrl().removeSuffix("/api/")
/**
* Состояние пагинации
*/
data class MessageRemoteMediatorState(
val lastSequenceId: Int?,
val firstSequenceId: Int?
)
override suspend fun initialize(): InitializeAction {
// Всегда запускаем refresh для проверки кэша
return InitializeAction.LAUNCH_INITIAL_REFRESH
}
override suspend fun load(
loadType: LoadType,
state: PagingState<Int, MessageEntity>
): MediatorResult {
Log.d(TAG, ">>> load() loadType=$loadType")
return try {
// Проверяем наличие данных в локальной БД
val localMinSeq = dao.getMinSequenceId(chatId)
val localMaxSeq = dao.getMaxSequenceId(chatId)
Log.d(TAG, "Local DB: minSeq=$localMinSeq, maxSeq=$localMaxSeq")
when (loadType) {
LoadType.REFRESH -> {
// Если есть локальные данные - не делаем API запрос
if (localMaxSeq != null) {
Log.d(TAG, "REFRESH: Using cached data (maxSeq=$localMaxSeq)")
return MediatorResult.Success(endOfPaginationReached = false)
}
// Нет данных - загружаем последние сообщения
Log.d(TAG, "REFRESH: No cached data, fetching from API")
}
LoadType.APPEND -> {
Log.d(TAG, "APPEND: Loading older messages")
}
LoadType.PREPEND -> {
Log.d(TAG, "PREPEND: Loading newer messages")
}
}
// Определяем pivot для API запроса
val pivotValue: Long? = when (loadType) {
LoadType.REFRESH -> null // Последние сообщения
LoadType.APPEND -> {
// Старые сообщения - берём минимальный sequenceId из текущей страницы
(state.pages.lastOrNull()?.data?.lastOrNull()?.sequenceId
?: localMinSeq)?.toLong()?.minus(1)
}
LoadType.PREPEND -> {
// Новые сообщения - берём максимальный sequenceId
(state.pages.firstOrNull()?.data?.firstOrNull()?.sequenceId
?: localMaxSeq)?.toLong()?.plus(1)
}
}
val limit = when (loadType) {
LoadType.REFRESH -> state.config.initialLoadSize
else -> state.config.pageSize
}
Log.d(TAG, "API call: chatId=$chatId, pivot=$pivotValue, limit=$limit")
// Выполняем запрос к API
val messages = try {
api.getMessages(
chatId = chatId,
cursor = null,
pivot = pivotValue,
limit = limit
)
} catch (e: Exception) {
Log.e(TAG, "API call failed - offline?", e)
// При ошибке сети и наличии локальных данных - успех
if ((localMinSeq != null || localMaxSeq != null)) {
Log.d(TAG, "Offline mode: returning success with cached data")
return MediatorResult.Success(endOfPaginationReached = true)
}
throw e
}
Log.d(TAG, "Received ${messages.size} messages from API")
if (messages.isEmpty()) {
Log.d(TAG, "End of pagination - no more messages")
return MediatorResult.Success(endOfPaginationReached = true)
}
// Сохраняем в базу
val currentUserId = tokenManager.getUserId() ?: ""
val entities = messages.map { dto ->
dto.toEntity(currentUserId, baseUrl, gson)
}
Log.d(TAG, "Saving ${entities.size} messages to DB, first seq=${entities.firstOrNull()?.sequenceId}, last seq=${entities.lastOrNull()?.sequenceId}")
dao.upsertMessages(entities)
// Проверяем что данные сохранились
val savedCount = dao.getMessagesCount(chatId)
Log.d(TAG, "Total messages in DB after save: $savedCount")
val endOfPaginationReached = messages.size < limit
Log.d(TAG, "endOfPaginationReached=$endOfPaginationReached")
MediatorResult.Success(endOfPaginationReached = endOfPaginationReached)
} catch (e: Exception) {
Log.e(TAG, "Error loading messages", e)
MediatorResult.Error(e)
}
}
companion object {
private const val TAG = "MessageRemoteMediator"
}
}
/**
* Extension function для преобразования DTO в Entity
*/
private fun MessageDto.toEntity(
currentUserId: String,
baseUrl: String,
gson: Gson
): MessageEntity {
// Определяем тип медиа
val mediaType = when {
media.isEmpty() -> MediaType.TEXT.name
media.any { it.type.startsWith("image") || it.url.endsWith(".gif") } -> MediaType.GIF.name
media.any { it.type.startsWith("image") } -> MediaType.IMAGE.name
media.any { it.type.startsWith("video") } -> MediaType.VIDEO.name
media.any { it.type.startsWith("audio") } -> MediaType.AUDIO.name
else -> MediaType.FILE.name
}
// Сериализуем медиа в JSON
val mediaJson = gson.toJson(media)
// Сериализуем реакции в JSON
val reactionsMap = reactions?.associate { it.emoji to it.count } ?: emptyMap()
val reactionsJson = gson.toJson(reactionsMap)
// Определяем, прочитано ли сообщение текущим пользователем
val isRead = senderId == currentUserId
return MessageEntity(
id = id,
chatId = chatId ?: "",
senderId = senderId ?: "",
senderName = sender?.displayName ?: sender?.username ?: "Unknown",
senderAvatar = sender?.avatarUrl,
content = content,
sequenceId = sequenceId ?: 0,
createdAt = createdAt ?: "",
mediaType = mediaType,
mediaJson = mediaJson,
reactionsJson = reactionsJson,
isRead = isRead,
replyToId = replyTo?.id,
syncStatus = SyncStatus.SYNCED,
isDeletedLocally = false,
isEditedLocally = false,
editedContent = null,
lastUpdated = System.currentTimeMillis()
)
}
@@ -1,113 +1,211 @@
package chats.data.repository
import android.content.Context
import android.util.Log
import androidx.paging.*
import chats.data.paging.MessageRemoteMediator
import chats.data.remote.api.ChatApi
import chats.data.remote.api.SendMessageRequest
import chats.data.remote.dto.ChatDto
import chats.data.remote.dto.MessageDto
import chats.data.remote.dto.MediaItemDto
import chats.data.remote.dto.ReactionDto
import chats.data.signalr.MessageSignalRHandler
import chats.data.sync.MessageSyncWorker
import chats.domain.model.Chat
import chats.domain.model.Message
import chats.domain.model.MediaType
import chats.domain.repository.ChatRepository
import core.database.data.ChatDatabase
import core.database.data.MessageDao
import core.database.data.MessageEntity
import core.database.data.SyncStatus
import core.network.ServerConfig
import chats.data.remote.signalr.ReadMessagesRequest
import chats.data.remote.signalr.ChatHubClient
import core.security.TokenManager
import com.google.gson.Gson
import kotlinx.coroutines.flow.Flow
import kotlinx.coroutines.flow.map
import okhttp3.MediaType.Companion.toMediaTypeOrNull
import okhttp3.MultipartBody
import okhttp3.RequestBody.Companion.asRequestBody
import java.text.SimpleDateFormat
import java.util.*
import javax.inject.Inject
import kotlinx.coroutines.flow.map
import javax.inject.Singleton
/**
* Основная реализация репозитория чатов
*
* Архитектура Offline-first:
* 1. Все данные читаются из локальной базы Room
* 2. При изменении данных - сначала запись в БД, потом синхронизация с сервером
* 3. SignalR обновления сразу записываются в БД
* 4. WorkManager обрабатывает фоновую синхронизацию
*
* Conflict Resolution:
* - Серверные данные имеют приоритет над локальными
* - Исключение: сообщения в процессе отправки (SYNCING) или редактирования
*/
@OptIn(ExperimentalPagingApi::class)
@Singleton
class ChatRepositoryImpl @Inject constructor(
private val api: ChatApi,
private val tokenManager: TokenManager,
private val serverConfig: ServerConfig,
private val messageDao: core.database.data.MessageDao,
private val hubClient: chats.data.remote.signalr.ChatHubClient
private val dao: MessageDao,
private val database: ChatDatabase,
private val hubClient: ChatHubClient,
private val signalRHandler: MessageSignalRHandler,
private val context: Context
) : ChatRepository {
private val gson = com.google.gson.Gson()
private val gson = Gson()
private val baseUrl = serverConfig.getBaseUrl().removeSuffix("/api/")
init {
signalRHandler.startListening()
Log.d(TAG, "ChatRepositoryImpl initialized")
}
override suspend fun getChats(): List<Chat> {
val currentUserId = tokenManager.getUserId() ?: ""
val baseUrl = serverConfig.getBaseUrl().removeSuffix("/api/")
return api.getChats().map { it.toDomain(currentUserId, baseUrl) }
}
override fun getMessagesFlow(chatId: String): kotlinx.coroutines.flow.Flow<List<Message>> {
val baseUrl = serverConfig.getBaseUrl().removeSuffix("/api/")
return messageDao.getMessages(chatId).map { entities: List<core.database.data.MessageEntity> ->
override fun getMessagesPaging(chatId: String): Flow<PagingData<Message>> {
val pagingConfig = PagingConfig(
pageSize = 30,
prefetchDistance = 10,
initialLoadSize = 50,
enablePlaceholders = false
)
return Pager(
config = pagingConfig,
pagingSourceFactory = { dao.getMessagesPagingSource(chatId) },
remoteMediator = MessageRemoteMediator(
chatId = chatId,
api = api,
database = database,
dao = dao,
serverConfig = serverConfig,
tokenManager = tokenManager
)
).flow.map { pagingData ->
pagingData.map { entity -> entity.toDomain(baseUrl, gson) }
}
}
override fun getMessagesFlow(chatId: String): Flow<List<Message>> {
return dao.getMessages(chatId).map { entities ->
entities.map { it.toDomain(baseUrl, gson) }
}
}
override suspend fun getMessages(chatId: String, cursor: String?, pivot: Long?, limit: Int?): List<Message> {
val baseUrl = serverConfig.getBaseUrl().removeSuffix("/api/")
override suspend fun getMessages(
chatId: String, cursor: String?, pivot: Long?, limit: Int?
): List<Message> {
val currentUserId = tokenManager.getUserId() ?: ""
return try {
android.util.Log.d("ChatRepo", "FETCH: chatId=$chatId, cursor=$cursor, limit=$limit")
Log.d(TAG, "Fetching messages from API: chatId=$chatId")
val messages = api.getMessages(chatId, cursor = cursor, limit = limit)
if (messages.isNotEmpty()) {
android.util.Log.d("ChatRepo", "Received ${messages.size} messages. TopSeq: ${messages.first().sequenceId}, BottomSeq: ${messages.last().sequenceId}")
val entities = messages.map { it.toEntity(baseUrl, currentUserId, gson) }
dao.upsertMessages(entities)
Log.d(TAG, "Cached ${entities.size} messages")
}
// Мапим в доменные модели. По умолчанию считаем прочитанными,
// так как unreadCount нам тут не критичен для истории.
messages.map { msg ->
msg.toDomain(currentUserId, baseUrl).copy(isRead = true)
}
messages.map { msg -> msg.toDomain(currentUserId, baseUrl).copy(isRead = true) }
} catch (e: Exception) {
android.util.Log.e("ChatRepo", "Fetch messages failed", e)
Log.e(TAG, "Fetch messages failed", e)
emptyList()
}
}
override suspend fun sendMessage(
chatId: String,
content: String?,
type: String,
chatId: String, content: String?, type: String,
attachments: List<chats.data.remote.api.AttachmentRequest>?,
replyToId: String?,
forwardedFromId: String?
replyToId: String?, forwardedFromId: String?
): Message {
val request = SendMessageRequest(
content = content,
type = type,
attachments = attachments,
replyToId = replyToId,
forwardedFromId = forwardedFromId
)
val baseUrl = serverConfig.getBaseUrl().removeSuffix("/api/")
val userId = tokenManager.getUserId() ?: ""
return api.sendMessage(chatId, request).toDomain(userId, baseUrl)
val currentTime = System.currentTimeMillis()
val localId = "local_${currentTime}_${chatId}"
val localMessage = MessageEntity(
id = localId, chatId = chatId, senderId = userId,
senderName = "Вы", senderAvatar = null, content = content,
sequenceId = 0,
createdAt = SimpleDateFormat("yyyy-MM-dd'T'HH:mm:ss.SSS'Z'", Locale.US).format(Date(currentTime)),
mediaType = type.uppercase(),
mediaJson = "[]",
reactionsJson = "{}",
isRead = true, replyToId = replyToId,
syncStatus = SyncStatus.SYNCING,
isDeletedLocally = false, isEditedLocally = false,
editedContent = null, lastUpdated = currentTime
)
dao.insertMessage(localMessage)
Log.d(TAG, "Saved local message: $localId")
MessageSyncWorker.scheduleSync(context)
// Создаём доменную модель вручную для локального сообщения
return Message(
id = localId,
chatId = chatId,
senderId = userId,
senderName = "Вы",
senderAvatar = null,
content = content,
sequenceId = 0,
createdAt = SimpleDateFormat("yyyy-MM-dd'T'HH:mm:ss.SSS'Z'", Locale.US).format(Date(currentTime)),
media = emptyList(),
mediaType = when (type) {
"image" -> MediaType.IMAGE
"video" -> MediaType.VIDEO
"audio" -> MediaType.AUDIO
"gif" -> MediaType.GIF
else -> MediaType.TEXT
},
reactions = emptyMap(),
isRead = true,
isPinned = false,
isForwarded = false,
forwardedFromName = null,
replyTo = null
)
}
override suspend fun addReaction(messageId: String, emoji: String) {
api.addReaction(messageId, emoji)
hubClient.addReaction(messageId, "", emoji)
}
override suspend fun sendTypingStatus(chatId: String) {
api.sendTypingStatus(chatId)
hubClient.sendTypingIndicator(chatId)
}
override suspend fun markMessagesAsRead(chatId: String, lastMessageId: String, lastReadSequenceId: Int) {
try {
android.util.Log.d("ChatRepoImpl", "markMessagesAsRead CALLED FOR $chatId. Caller stack: ${android.util.Log.getStackTraceString(Throwable())}")
hubClient.readMessages(ReadMessagesRequest(chatId, lastMessageId, lastReadSequenceId))
hubClient.readMessages(chats.data.remote.signalr.ReadMessagesRequest(chatId, lastMessageId, lastReadSequenceId))
} catch (e: Exception) {
android.util.Log.e("ChatRepo", "Error marking messages as read", e)
Log.e(TAG, "Error marking messages as read", e)
}
// Обновляем локальную БД
messageDao.markMessagesAsRead(chatId, lastReadSequenceId)
dao.markMessagesAsRead(chatId, lastReadSequenceId)
}
override suspend fun saveMessage(message: Message) {
android.util.Log.d("ChatRepo", "DB cache disabled, skipping save: ${message.id}")
dao.insertMessage(message.toEntity(gson))
}
override suspend fun deleteLocalMessage(messageId: String) {
messageDao.deleteMessage(messageId)
dao.markAsDeletedLocally(messageId)
MessageSyncWorker.scheduleSync(context)
}
override suspend fun editLocalMessage(messageId: String, newContent: String) {
dao.markAsEditedLocally(messageId, newContent)
MessageSyncWorker.scheduleSync(context)
}
override suspend fun uploadMedia(file: java.io.File): String {
@@ -124,34 +222,36 @@ class ChatRepositoryImpl @Inject constructor(
return api.uploadFile(body).url
}
override suspend fun getTrendingGifs(page: Int): List<chats.data.remote.api.KlipyGifDto> {
return api.getTrendingGifs(page).data.data
}
override suspend fun getTrendingGifs(page: Int): List<chats.data.remote.api.KlipyGifDto> =
api.getTrendingGifs(page).data.data
override suspend fun searchGifs(query: String, page: Int): List<chats.data.remote.api.KlipyGifDto> {
return api.searchGifs(query, page).data.data
}
override suspend fun searchGifs(query: String, page: Int): List<chats.data.remote.api.KlipyGifDto> =
api.searchGifs(query, page).data.data
override suspend fun getGifCategories(): List<chats.data.remote.api.GifCategoryDto> {
return api.getGifCategories().data.categories
}
override suspend fun getGifCategories(): List<chats.data.remote.api.GifCategoryDto> =
api.getGifCategories().data.categories
override suspend fun createPersonalChat(userId: String): Chat {
val currentUserId = tokenManager.getUserId() ?: ""
val baseUrl = serverConfig.getBaseUrl().removeSuffix("/api/")
val request = chats.data.remote.api.CreatePersonalChatRequest(userId)
return api.createPersonalChat(request).toDomain(currentUserId, baseUrl)
}
override suspend fun deleteMessage(messageId: String, forEveryone: Boolean) {
api.deleteMessage(messageId, forEveryone)
messageDao.deleteMessage(messageId)
dao.deleteMessage(messageId)
}
override suspend fun editMessage(messageId: String, content: String): Message {
val request = SendMessageRequest(content = content)
val currentUserId = tokenManager.getUserId() ?: ""
val baseUrl = serverConfig.getBaseUrl().removeSuffix("/api/")
return api.editMessage(messageId, request).toDomain(currentUserId, baseUrl)
val response = api.editMessage(messageId, request)
dao.insertMessage(response.toEntity(baseUrl, currentUserId, gson))
return response.toDomain(currentUserId, baseUrl)
}
companion object {
private const val TAG = "ChatRepositoryImpl"
}
}
@@ -0,0 +1,194 @@
package chats.data.signalr
import android.util.Log
import chats.data.remote.dto.MessageDto
import chats.data.remote.signalr.ChatEvent
import chats.data.remote.signalr.ChatHubClient
import chats.data.remote.signalr.MessagesReadEvent
import chats.data.remote.signalr.ReactionEvent
import chats.data.repository.toEntity
import com.google.gson.Gson
import core.database.data.MessageDao
import core.database.data.MessageEntity
import core.database.data.SyncStatus
import core.network.ServerConfig
import core.security.TokenManager
import kotlinx.coroutines.*
import kotlinx.coroutines.flow.collectLatest
import javax.inject.Inject
import javax.inject.Singleton
/**
* Обработчик SignalR событий для обновления локального кэша
*
* Обрабатывает события:
* - new_message: новое сообщение в чате
* - message_edited: сообщение отредактировано
* - message_deleted: сообщение удалено
* - messages_read: сообщения прочитаны
* - reaction_added/removed: реакция добавлена/удалена
*
* Все изменения сразу записываются в Room, UI обновляется через Flow
*/
@Singleton
class MessageSignalRHandler @Inject constructor(
private val hubClient: ChatHubClient,
private val dao: MessageDao,
private val serverConfig: ServerConfig,
private val tokenManager: TokenManager
) {
private val gson = Gson()
private val scope = CoroutineScope(Dispatchers.IO + SupervisorJob())
private val baseUrl = serverConfig.getBaseUrl().removeSuffix("/api/")
private var isListening = false
/**
* Начинает прослушивание SignalR событий
* Вызывается один раз при инициализации приложения
*/
fun startListening() {
if (isListening) {
Log.w(TAG, "Already listening to SignalR events")
return
}
isListening = true
Log.d(TAG, "Started listening to SignalR events")
scope.launch {
hubClient.events.collectLatest { event ->
handleEvent(event)
}
}
}
/**
* Останавливает прослушивание событий
*/
fun stopListening() {
isListening = false
scope.cancel()
Log.d(TAG, "Stopped listening to SignalR events")
}
private suspend fun handleEvent(event: ChatEvent) {
try {
when (event) {
is ChatEvent.NewMessage -> handleNewMessage(event.message)
is ChatEvent.MessageEdited -> handleMessageEdited(event.messageId, event.chatId, event.content)
is ChatEvent.MessageDeleted -> handleMessageDeleted(event.messageId, event.chatId)
is ChatEvent.MessagesRead -> handleMessagesRead(event.chatId, event.lastReadSequenceId)
is ChatEvent.ReactionUpdated -> handleReactionUpdated(
event.messageId,
event.emoji,
event.isRemoved
)
else -> Log.d(TAG, "Unhandled event: ${event::class.simpleName}")
}
} catch (e: Exception) {
Log.e(TAG, "Error handling SignalR event: ${event::class.simpleName}", e)
}
}
private suspend fun handleNewMessage(dto: MessageDto) {
Log.d(TAG, "New message received: ${dto.id} in chat ${dto.chatId}")
val currentUserId = tokenManager.getUserId() ?: ""
val entity = dto.toEntity(baseUrl, currentUserId, gson).copy(
syncStatus = SyncStatus.SYNCED
)
val existing = dao.getMessageById(dto.id)
if (existing != null && existing.syncStatus == SyncStatus.SYNCING) {
val syncedEntity = entity.copy(
syncStatus = SyncStatus.SYNCED,
isDeletedLocally = false,
isEditedLocally = false,
editedContent = null
)
dao.insertMessage(syncedEntity)
Log.d(TAG, "Merged local message with server response: ${dto.id}")
} else {
dao.insertMessage(entity)
Log.d(TAG, "Saved new message: ${dto.id}")
}
}
private suspend fun handleMessageEdited(messageId: String, chatId: String, content: String) {
Log.d(TAG, "Message edited: $messageId")
val existing = dao.getMessageById(messageId)
if (existing != null) {
if (existing.isEditedLocally) {
Log.d(TAG, "Message is being edited locally, skipping server update: $messageId")
return
}
val updated = existing.copy(
content = content,
lastUpdated = System.currentTimeMillis()
)
dao.insertMessage(updated)
Log.d(TAG, "Updated edited message: $messageId")
} else {
Log.w(TAG, "Edited message not found in cache: $messageId")
}
}
private suspend fun handleMessageDeleted(messageId: String, chatId: String) {
Log.d(TAG, "Message deleted: $messageId")
val existing = dao.getMessageById(messageId)
if (existing?.isDeletedLocally == true) {
dao.deleteMessage(messageId)
Log.d(TAG, "Completed local deletion: $messageId")
} else {
dao.deleteMessage(messageId)
Log.d(TAG, "Removed deleted message from cache: $messageId")
}
}
private suspend fun handleMessagesRead(chatId: String, lastReadSequenceId: Int?) {
if (lastReadSequenceId == null) {
Log.w(TAG, "Messages read event with null sequenceId")
return
}
Log.d(TAG, "Messages read in chat $chatId up to sequenceId $lastReadSequenceId")
dao.markMessagesAsRead(chatId, lastReadSequenceId)
}
private suspend fun handleReactionUpdated(messageId: String, emoji: String, isRemoved: Boolean) {
Log.d(TAG, "Reaction ${if (isRemoved) "removed" else "added"}: $emoji on message $messageId")
val existing = dao.getMessageById(messageId)
if (existing != null) {
val reactions = gson.fromJson<Map<String, Int>>(existing.reactionsJson, Map::class.java) ?: emptyMap()
val updatedReactions = reactions.toMutableMap()
if (isRemoved) {
val currentCount = updatedReactions[emoji] ?: 0
if (currentCount > 1) {
updatedReactions[emoji] = currentCount - 1
} else {
updatedReactions.remove(emoji)
}
} else {
updatedReactions[emoji] = (updatedReactions[emoji] ?: 0) + 1
}
val updated = existing.copy(
reactionsJson = gson.toJson(updatedReactions),
lastUpdated = System.currentTimeMillis()
)
dao.insertMessage(updated)
Log.d(TAG, "Updated reactions for message: $messageId")
}
}
companion object {
private const val TAG = "MessageSignalRHandler"
}
}
@@ -0,0 +1,225 @@
package chats.data.sync
import android.content.Context
import android.util.Log
import androidx.hilt.work.HiltWorker
import androidx.work.*
import chats.data.remote.api.ChatApi
import chats.data.remote.api.SendMessageRequest
import core.database.data.MessageDao
import core.database.data.MessageEntity
import core.database.data.SyncStatus
import core.network.ServerConfig
import core.security.TokenManager
import com.google.gson.Gson
import dagger.assisted.Assisted
import dagger.assisted.AssistedInject
import java.util.concurrent.TimeUnit
/**
* WorkManager Worker для фоновой синхронизации сообщений
*
* Обрабатывает:
* 1. Отправку новых сообщений (SYNCING)
* 2. Обновление отредактированных сообщений (isEditedLocally = true)
* 3. Удаление сообщений (isDeletedLocally = true)
* 4. Повторную отправку при ошибках (FAILED)
*
* Политика повторных попыток:
* - Экспоненциальная задержка (backoff)
* - Максимум 3 попытки
* - Таймаут 10 минут на выполнение
*/
@HiltWorker
class MessageSyncWorker @AssistedInject constructor(
@Assisted appContext: Context,
@Assisted params: WorkerParameters,
private val dao: MessageDao,
private val api: ChatApi,
private val serverConfig: ServerConfig,
private val tokenManager: TokenManager
) : CoroutineWorker(appContext, params) {
private val gson = Gson()
companion object {
const val WORK_NAME = "message_sync_worker"
private const val TAG = "MessageSyncWorker"
/**
* Планирует синхронизацию
* Вызывается при изменении сообщений в базе
*/
fun scheduleSync(context: Context) {
Log.d(TAG, "Scheduling sync")
val constraints = Constraints.Builder()
.setRequiredNetworkType(NetworkType.CONNECTED)
.setRequiresBatteryNotLow(false)
.build()
val workRequest = OneTimeWorkRequestBuilder<MessageSyncWorker>()
.setConstraints(constraints)
.setBackoffCriteria(
BackoffPolicy.EXPONENTIAL,
WorkRequest.MIN_BACKOFF_MILLIS,
TimeUnit.MILLISECONDS
)
.addTag(WORK_NAME)
.build()
WorkManager.getInstance(context).enqueueUniqueWork(
WORK_NAME,
ExistingWorkPolicy.REPLACE,
workRequest
)
Log.d(TAG, "Sync scheduled")
}
/**
* Отменяет запланированную синхронизацию
*/
fun cancelSync(context: Context) {
WorkManager.getInstance(context).cancelUniqueWork(WORK_NAME)
Log.d(TAG, "Sync cancelled")
}
}
override suspend fun doWork(): Result {
Log.d(TAG, "Starting message sync. Attempt: ${runAttemptCount + 1}")
if (runAttemptCount >= 3) {
Log.e(TAG, "Max retry attempts reached")
return Result.failure()
}
return try {
val pendingMessages = dao.getPendingSyncMessages()
Log.d(TAG, "Found ${pendingMessages.size} messages to sync")
if (pendingMessages.isEmpty()) {
Log.d(TAG, "No pending messages. Sync complete.")
return Result.success()
}
var successCount = 0
var failureCount = 0
for (message in pendingMessages) {
try {
when {
message.isDeletedLocally -> {
handleDeleteMessage(message)
successCount++
}
message.isEditedLocally -> {
handleEditMessage(message)
successCount++
}
message.syncStatus == SyncStatus.SYNCING -> {
handleSendMessage(message)
successCount++
}
}
} catch (e: Exception) {
Log.e(TAG, "Failed to sync message ${message.id}", e)
dao.markAsSyncFailed(message.id)
failureCount++
}
}
Log.d(TAG, "Sync completed. Success: $successCount, Failed: $failureCount")
if (failureCount > 0 && successCount == 0) {
Result.retry()
} else {
Result.success()
}
} catch (e: Exception) {
Log.e(TAG, "Sync failed with exception", e)
Result.retry()
}
}
private suspend fun handleSendMessage(message: MessageEntity) {
Log.d(TAG, "Sending message: ${message.id}")
val attachments = parseAttachments(message.mediaJson)
val request = SendMessageRequest(
content = message.content,
type = message.mediaType.lowercase(),
attachments = attachments,
replyToId = message.replyToId
)
val response = api.sendMessage(message.chatId, request)
val syncedMessage = message.copy(
id = response.id,
sequenceId = response.sequenceId ?: message.sequenceId,
createdAt = response.createdAt ?: message.createdAt,
syncStatus = SyncStatus.SYNCED,
isDeletedLocally = false,
isEditedLocally = false,
editedContent = null,
lastUpdated = System.currentTimeMillis()
)
dao.insertMessage(syncedMessage)
Log.d(TAG, "Message sent successfully: ${response.id}")
}
private suspend fun handleEditMessage(message: MessageEntity) {
Log.d(TAG, "Editing message: ${message.id}")
val newContent = message.editedContent ?: message.content
val request = SendMessageRequest(content = newContent)
val response = api.editMessage(message.id, request)
val syncedMessage = message.copy(
content = response.content,
syncStatus = SyncStatus.SYNCED,
isEditedLocally = false,
editedContent = null,
lastUpdated = System.currentTimeMillis()
)
dao.insertMessage(syncedMessage)
Log.d(TAG, "Message edited successfully: ${message.id}")
}
private suspend fun handleDeleteMessage(message: MessageEntity) {
Log.d(TAG, "Deleting message: ${message.id}")
val response = api.deleteMessage(message.id, forEveryone = false)
if (response.isSuccessful || response.code() == 404) {
dao.deleteMessage(message.id)
Log.d(TAG, "Message deleted successfully: ${message.id}")
} else {
throw Exception("Delete failed with code: ${response.code()}")
}
}
private fun parseAttachments(mediaJson: String): List<chats.data.remote.api.AttachmentRequest>? {
return try {
val mediaList = gson.fromJson(mediaJson, Array::class.java)
?.map { elem ->
val map = elem as Map<*, *>
chats.data.remote.api.AttachmentRequest(
type = map["type"] as? String ?: "file",
url = map["url"] as? String ?: "",
fileName = map["filename"] as? String ?: "file",
fileSize = (map["size"] as? Number)?.toLong() ?: 0L
)
}
mediaList?.takeIf { it.isNotEmpty() }
} catch (e: Exception) {
Log.e(TAG, "Failed to parse attachments", e)
null
}
}
}