417 lines
16 KiB
Kotlin
417 lines
16 KiB
Kotlin
package chats.data.repository
|
||
|
||
import android.util.Log
|
||
import androidx.paging.ExperimentalPagingApi
|
||
import androidx.paging.Pager
|
||
import androidx.paging.PagingConfig
|
||
import androidx.paging.PagingData
|
||
import androidx.paging.map
|
||
import androidx.work.ExistingPeriodicWorkPolicy
|
||
import androidx.work.ExistingWorkPolicy
|
||
import androidx.work.WorkManager
|
||
import chats.data.local.dao.ChatDao
|
||
import chats.data.local.dao.MessageDao
|
||
import chats.data.local.database.MessageEntity
|
||
import chats.data.local.mappers.createPendingMessageEntity
|
||
import chats.data.local.mappers.toDomain
|
||
import chats.data.local.mappers.toEntity
|
||
import chats.data.local.paging.MessagePagingSource
|
||
import chats.data.local.paging.MessageRemoteMediator
|
||
import chats.data.remote.api.ChatApi
|
||
import chats.data.remote.api.SendMessageRequest
|
||
import chats.data.remote.dto.MessageDto
|
||
import chats.data.workers.ChatSyncWorker
|
||
import chats.data.workers.SendMessageWorker
|
||
import chats.domain.model.Chat
|
||
import chats.domain.model.Message
|
||
import chats.domain.model.MessageStatus
|
||
import chats.domain.repository.ChatRepository
|
||
import core.network.ServerConfig
|
||
import core.security.TokenManager
|
||
import kotlinx.coroutines.flow.Flow
|
||
import kotlinx.coroutines.flow.map
|
||
import okhttp3.MediaType.Companion.toMediaTypeOrNull
|
||
import okhttp3.MultipartBody
|
||
import okhttp3.RequestBody.Companion.asRequestBody
|
||
import javax.inject.Inject
|
||
|
||
/**
|
||
* Реализация ChatRepository с поддержкой офлайн-режима.
|
||
* Использует паттерн Single Source of Truth: UI всегда берет данные из Room.
|
||
*/
|
||
class ChatRepositoryImpl @Inject constructor(
|
||
private val api: ChatApi,
|
||
private val tokenManager: TokenManager,
|
||
private val serverConfig: ServerConfig,
|
||
private val messageDao: MessageDao,
|
||
private val chatDao: ChatDao,
|
||
private val workManager: WorkManager,
|
||
private val appDatabase: chats.data.local.database.AppDatabase,
|
||
private val hubClient: chats.data.remote.signalr.ChatHubClient
|
||
) : ChatRepository {
|
||
|
||
private val gson = com.google.gson.Gson()
|
||
|
||
// ==================== Чаты ====================
|
||
|
||
override fun getChatsFlow(): Flow<List<Chat>> {
|
||
// Single Source of Truth - данные из Room
|
||
return chatDao.getAllChats().map { entities ->
|
||
entities.map { it.toDomain() }
|
||
}
|
||
}
|
||
|
||
override suspend fun getChats(): List<Chat> {
|
||
val currentUserId = tokenManager.getUserId() ?: ""
|
||
val baseUrl = serverConfig.getBaseUrl().removeSuffix("/api/")
|
||
val gson = com.google.gson.Gson()
|
||
|
||
return try {
|
||
// Пробуем получить с сервера
|
||
val remoteChats = api.getChats()
|
||
val domainChats = remoteChats.map { it.toDomain(currentUserId, baseUrl) }
|
||
|
||
// Сохраняем в локальную БД
|
||
val entities = domainChats.map { chat ->
|
||
val existing = chatDao.getChatByRemoteId(chat.id)
|
||
chat.toEntity().copy(
|
||
localId = existing?.localId ?: java.util.UUID.randomUUID().toString()
|
||
)
|
||
}
|
||
chatDao.insertChats(entities)
|
||
|
||
domainChats
|
||
} catch (e: Exception) {
|
||
Log.w("ChatRepo", "Failed to fetch chats from server, returning local", e)
|
||
// При ошибке возвращаем локальные данные
|
||
emptyList()
|
||
}
|
||
}
|
||
|
||
override suspend fun syncChats() {
|
||
try {
|
||
val currentUserId = tokenManager.getUserId() ?: ""
|
||
val baseUrl = serverConfig.getBaseUrl().removeSuffix("/api/")
|
||
val remoteChats = api.getChats()
|
||
val gson = com.google.gson.Gson()
|
||
|
||
val entities = remoteChats.map { dto ->
|
||
val existing = chatDao.getChatByRemoteId(dto.id)
|
||
val lastMessage = dto.messages.firstOrNull()
|
||
val timestamp = lastMessage?.createdAt?.let {
|
||
try { java.time.ZonedDateTime.parse(it).toInstant().toEpochMilli() }
|
||
catch (e: Exception) { null }
|
||
}
|
||
val lastMessageDomain = lastMessage?.toDomain(currentUserId, baseUrl)
|
||
dto.toDomain(currentUserId, baseUrl).toEntity().copy(
|
||
localId = existing?.localId ?: java.util.UUID.randomUUID().toString(),
|
||
lastMessageText = lastMessage?.content,
|
||
lastMessageTimestamp = timestamp,
|
||
lastMessageJson = lastMessageDomain?.let { gson.toJson(it) }
|
||
)
|
||
}
|
||
chatDao.insertChats(entities)
|
||
} catch (e: Exception) {
|
||
Log.e("ChatRepo", "Sync chats failed", e)
|
||
}
|
||
}
|
||
|
||
// ==================== Сообщения ====================
|
||
|
||
override fun getMessagesFlow(chatId: String): Flow<List<Message>> {
|
||
// Single Source of Truth - всегда из Room
|
||
return messageDao.getMessagesByChatId(chatId).map { entities ->
|
||
entities.map { it.toDomain() }
|
||
}
|
||
}
|
||
|
||
/**
|
||
* Получает сообщения с помощью Paging 3 и RemoteMediator
|
||
*/
|
||
@OptIn(ExperimentalPagingApi::class)
|
||
override fun getMessagesPagingSource(chatId: String): Flow<PagingData<Message>> {
|
||
return Pager(
|
||
config = PagingConfig(
|
||
pageSize = 30,
|
||
prefetchDistance = 10,
|
||
enablePlaceholders = false,
|
||
initialLoadSize = 60
|
||
),
|
||
pagingSourceFactory = {
|
||
MessagePagingSource(chatId, messageDao)
|
||
},
|
||
remoteMediator = MessageRemoteMediator(
|
||
chatId = chatId,
|
||
messageDao = messageDao,
|
||
chatDao = chatDao,
|
||
chatApi = api,
|
||
tokenManager = tokenManager,
|
||
appDatabase = appDatabase,
|
||
serverConfig = serverConfig
|
||
)
|
||
).flow.map { pagingData ->
|
||
pagingData.map { entity -> entity.toDomain() }
|
||
}
|
||
}
|
||
|
||
override suspend fun getMessages(chatId: String, cursor: String?, pivot: Long?, limit: Int?): List<Message> {
|
||
val baseUrl = serverConfig.getBaseUrl().removeSuffix("/api/")
|
||
val currentUserId = tokenManager.getUserId() ?: ""
|
||
|
||
return try {
|
||
Log.d("ChatRepo", "FETCH: chatId=$chatId, cursor=$cursor, limit=$limit")
|
||
val messages = api.getMessages(chatId, cursor = cursor, limit = limit)
|
||
|
||
if (messages.isNotEmpty()) {
|
||
Log.d("ChatRepo", "Received ${messages.size} messages")
|
||
// Сохраняем в локальную БД
|
||
saveMessagesToLocal(messages, chatId, currentUserId)
|
||
}
|
||
|
||
messages.map { msg ->
|
||
msg.toDomain(currentUserId, baseUrl).copy(isRead = true)
|
||
}
|
||
} catch (e: Exception) {
|
||
Log.e("ChatRepo", "Fetch messages failed, returning local", e)
|
||
emptyList()
|
||
}
|
||
}
|
||
|
||
override suspend fun sendMessage(
|
||
chatId: String,
|
||
content: String?,
|
||
type: String,
|
||
attachments: List<chats.data.remote.api.AttachmentRequest>?,
|
||
replyToId: String?,
|
||
forwardedFromId: String?
|
||
): Message {
|
||
val userId = tokenManager.getUserId() ?: ""
|
||
val userName = tokenManager.getUsername() ?: userId
|
||
|
||
// 1. Создаем локальное сообщение со статусом PENDING
|
||
val mediaType = when (type) {
|
||
"image" -> "IMAGE"
|
||
"video" -> "VIDEO"
|
||
"audio", "voice" -> "AUDIO"
|
||
else -> "TEXT"
|
||
}
|
||
|
||
val mediaJson = attachments?.map {
|
||
chats.domain.model.Media(
|
||
id = java.util.UUID.randomUUID().toString(),
|
||
type = it.type,
|
||
url = it.url,
|
||
filename = it.fileName,
|
||
size = it.fileSize
|
||
)
|
||
}?.let { gson.toJson(it) } ?: "[]"
|
||
|
||
val pendingEntity = createPendingMessageEntity(
|
||
chatId = chatId,
|
||
senderId = userId,
|
||
senderName = userName,
|
||
content = content,
|
||
mediaType = mediaType,
|
||
mediaJson = mediaJson,
|
||
replyToId = replyToId
|
||
)
|
||
|
||
// 2. Сохраняем в Room
|
||
messageDao.insertMessage(pendingEntity)
|
||
|
||
// 3. Ставим задачу в WorkManager для отправки
|
||
val workRequest = SendMessageWorker.createWorkRequest(pendingEntity.localId)
|
||
workManager.enqueueUniqueWork(
|
||
"send_${pendingEntity.localId}",
|
||
ExistingWorkPolicy.REPLACE,
|
||
workRequest
|
||
)
|
||
|
||
Log.d("ChatRepo", "Queued message for sending: ${pendingEntity.localId}")
|
||
|
||
// 4. Возвращаем доменную модель для немедленного отображения в UI
|
||
return pendingEntity.toDomain()
|
||
}
|
||
|
||
override suspend fun retryFailedMessage(localId: String) {
|
||
val message = messageDao.getMessageByLocalId(localId)
|
||
?: return
|
||
|
||
if (message.status != MessageStatus.FAILED) return
|
||
|
||
// Сбрасываем статус и ставим в очередь
|
||
messageDao.updateMessageStatus(
|
||
localId = localId,
|
||
status = MessageStatus.PENDING,
|
||
updatedAtMillis = System.currentTimeMillis()
|
||
)
|
||
|
||
val workRequest = SendMessageWorker.createWorkRequest(localId)
|
||
workManager.enqueueUniqueWork(
|
||
"send_$localId",
|
||
ExistingWorkPolicy.REPLACE,
|
||
workRequest
|
||
)
|
||
}
|
||
|
||
override suspend fun deleteLocalMessage(messageId: String) {
|
||
messageDao.deleteMessageByLocalId(messageId)
|
||
}
|
||
|
||
override suspend fun saveMessage(message: Message) {
|
||
// Сохраняем входящее сообщение из SignalR
|
||
val entity = message.toEntity(MessageStatus.DELIVERED)
|
||
messageDao.insertMessage(entity)
|
||
}
|
||
|
||
// ==================== Синхронизация ====================
|
||
|
||
override suspend fun syncMessagesForChat(chatId: String) {
|
||
val workRequest = ChatSyncWorker.createOneTimeWorkRequest(chatId)
|
||
workManager.enqueue(workRequest)
|
||
}
|
||
|
||
override suspend fun schedulePeriodicSync() {
|
||
val workRequest = ChatSyncWorker.createPeriodicWorkRequest()
|
||
workManager.enqueueUniquePeriodicWork(
|
||
"periodic_chat_sync",
|
||
ExistingPeriodicWorkPolicy.KEEP,
|
||
workRequest
|
||
)
|
||
}
|
||
|
||
// ==================== Вспомогательные методы ====================
|
||
|
||
private suspend fun saveMessagesToLocal(messages: List<chats.data.remote.dto.MessageDto>, chatId: String, currentUserId: String) {
|
||
val baseUrl = serverConfig.getBaseUrl().removeSuffix("/api/")
|
||
messages.forEach { dto ->
|
||
val existing = messageDao.getMessageByServerId(dto.id)
|
||
if (existing == null) {
|
||
val entity = MessageEntity(
|
||
localId = java.util.UUID.randomUUID().toString(),
|
||
serverId = dto.id,
|
||
idempotencyKey = dto.id,
|
||
chatId = chatId,
|
||
senderId = dto.senderId ?: dto.sender?.id ?: currentUserId,
|
||
senderName = dto.sender?.displayName ?: dto.sender?.username ?: "",
|
||
senderAvatar = dto.sender?.avatarUrl?.let { if (it.startsWith("http")) it else "$baseUrl$it" },
|
||
content = dto.content,
|
||
sequenceId = dto.sequenceId?.toLong() ?: 0L,
|
||
createdAt = dto.createdAt ?: java.time.ZonedDateTime.now().toString(),
|
||
mediaType = dto.type ?: "TEXT",
|
||
mediaJson = gson.toJson(dto.media.map { mediaItem ->
|
||
chats.domain.model.Media(
|
||
id = mediaItem.id,
|
||
type = mediaItem.type,
|
||
url = (mediaItem.url as? String ?: "").let { url -> if (url.startsWith("http")) url else "$baseUrl$url" },
|
||
filename = mediaItem.filename,
|
||
size = mediaItem.size,
|
||
duration = mediaItem.duration
|
||
)
|
||
}),
|
||
reactionsJson = gson.toJson(
|
||
dto.reactions?.associate { it.emoji to it.count } ?: emptyMap<String, Int>()
|
||
),
|
||
status = if (dto.senderId == currentUserId) MessageStatus.SENT else MessageStatus.DELIVERED,
|
||
createdAtMillis = System.currentTimeMillis(),
|
||
updatedAtMillis = System.currentTimeMillis()
|
||
)
|
||
messageDao.insertMessage(entity)
|
||
}
|
||
}
|
||
}
|
||
|
||
// ==================== Остальные методы ====================
|
||
|
||
override suspend fun addReaction(messageId: String, emoji: String) {
|
||
try {
|
||
api.addReaction(messageId, emoji)
|
||
} catch (e: Exception) {
|
||
Log.e("ChatRepo", "Add reaction failed", e)
|
||
throw e
|
||
}
|
||
}
|
||
|
||
override suspend fun sendTypingStatus(chatId: String) {
|
||
try {
|
||
api.sendTypingStatus(chatId)
|
||
} catch (e: Exception) {
|
||
// Игнорируем ошибки typing status
|
||
}
|
||
}
|
||
|
||
override suspend fun resetUnreadCount(chatId: String) {
|
||
chatDao.updateUnreadCount(chatId, 0)
|
||
}
|
||
|
||
override suspend fun markMessagesAsRead(chatId: String, lastMessageId: String, lastReadSequenceId: Int) {
|
||
try {
|
||
hubClient.readMessages(chats.data.remote.signalr.ReadMessagesRequest(chatId, lastMessageId, lastReadSequenceId))
|
||
} catch (e: Exception) {
|
||
Log.e("ChatRepo", "Error marking messages as read", e)
|
||
}
|
||
}
|
||
|
||
override suspend fun uploadMedia(file: java.io.File): String {
|
||
val mimeType = when (file.extension.lowercase()) {
|
||
"jpg", "jpeg" -> "image/jpeg"
|
||
"png" -> "image/png"
|
||
"webp" -> "image/webp"
|
||
"mp4" -> "video/mp4"
|
||
"mp3", "m4a", "wav" -> "audio/mpeg"
|
||
else -> "application/octet-stream"
|
||
}
|
||
val requestFile = file.asRequestBody(mimeType.toMediaTypeOrNull())
|
||
val body = MultipartBody.Part.createFormData("file", file.name, requestFile)
|
||
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 searchGifs(query: String, page: Int): List<chats.data.remote.api.KlipyGifDto> {
|
||
return api.searchGifs(query, page).data.data
|
||
}
|
||
|
||
override suspend fun getGifCategories(): List<chats.data.remote.api.GifCategoryDto> {
|
||
return api.getGifCategories().data.categories
|
||
}
|
||
|
||
override suspend fun getOrCreateFavorites(): Chat {
|
||
val currentUserId = tokenManager.getUserId() ?: ""
|
||
val baseUrl = serverConfig.getBaseUrl().removeSuffix("/api/")
|
||
val dto = api.getOrCreateFavorites()
|
||
val chat = dto.toDomain(currentUserId, baseUrl)
|
||
chatDao.insertChat(chat.toEntity())
|
||
return chat
|
||
}
|
||
|
||
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)
|
||
}
|
||
|
||
override suspend fun editMessage(messageId: String, content: String): Message {
|
||
val request = SendMessageRequest(content = content)
|
||
val currentUserId = tokenManager.getUserId() ?: ""
|
||
val returnedId = api.editMessage(messageId, request)
|
||
|
||
return Message(
|
||
id = returnedId,
|
||
chatId = "",
|
||
senderId = currentUserId,
|
||
senderName = "",
|
||
content = content,
|
||
sequenceId = 0,
|
||
createdAt = java.time.ZonedDateTime.now().toString()
|
||
)
|
||
}
|
||
}
|