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.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.ChatDao import core.database.data.MessageDao import core.database.data.MessageEntity import core.database.data.SyncStatus import core.network.ServerConfig 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 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: MessageDao, private val chatDao: ChatDao, private val database: ChatDatabase, private val hubClient: ChatHubClient, private val signalRHandler: MessageSignalRHandler, private val context: Context ) : ChatRepository { private val gson = Gson() private val baseUrl = serverConfig.getBaseUrl().removeSuffix("/api/") init { signalRHandler.startListening() Log.d(TAG, "ChatRepositoryImpl initialized") } override suspend fun getChats(): List { val currentUserId = tokenManager.getUserId() ?: "" return try { // Пробуем загрузить из сети val chats = api.getChats().map { it.toDomain(currentUserId, baseUrl) } // Кэшируем в Room val entities = chats.map { it.toEntity() } chatDao.insertChats(entities) // Кэшируем последние сообщения chats.forEach { chat -> chat.lastMessage?.let { msg -> messageDao.upsertMessage(msg.toEntity(gson)) } } android.util.Log.d(TAG, "Cached ${entities.size} chats with messages") chats } catch (e: Exception) { android.util.Log.d(TAG, "Network load failed, using cache") // При ошибке - возвращаем из кэша с загрузкой последних сообщений chatDao.getAllChats().map { entity -> val lastMessage = entity.lastMessageId?.let { messageId -> messageDao.getMessageById(messageId)?.toDomain(baseUrl, gson) } entity.toDomain(currentUserId, baseUrl, lastMessage) } } } override fun getChatsFlow(): Flow> { val currentUserId = tokenManager.getUserId() ?: "" return chatDao.getAllChatsFlow().map { entities -> entities.map { entity -> // Загружаем последнее сообщение из базы для каждого чата val lastMessage = entity.lastMessageId?.let { messageId -> messageDao.getMessageById(messageId)?.toDomain(baseUrl, gson) } entity.toDomain(currentUserId, baseUrl, lastMessage) } } } override fun getMessagesPaging(chatId: String): Flow> { val pagingConfig = PagingConfig( pageSize = 30, prefetchDistance = 10, initialLoadSize = 50, enablePlaceholders = false ) return Pager( config = pagingConfig, pagingSourceFactory = { messageDao.getMessagesPagingSource(chatId) }, remoteMediator = MessageRemoteMediator( chatId = chatId, api = api, database = database, dao = messageDao, serverConfig = serverConfig, tokenManager = tokenManager ) ).flow.map { pagingData -> pagingData.map { entity -> entity.toDomain(baseUrl, gson) } } } override fun getMessagesFlow(chatId: String): Flow> { return messageDao.getMessages(chatId).map { entities -> entities.map { it.toDomain(baseUrl, gson) } } } override suspend fun getMessages( chatId: String, cursor: String?, pivot: Long?, afterSequenceId: Long?, limit: Int? ): List { val currentUserId = tokenManager.getUserId() ?: "" return try { Log.d(TAG, "Fetching messages from API: chatId=$chatId, afterSequenceId=$afterSequenceId, limit=$limit") val messages = api.getMessages(chatId, cursor = cursor, pivot = pivot, afterSequenceId = afterSequenceId, limit = limit) if (messages.isNotEmpty()) { val entities = messages.map { it.toEntity(baseUrl, currentUserId, gson) } messageDao.upsertMessages(entities) Log.d(TAG, "Cached ${entities.size} messages") } messages.map { msg -> msg.toDomain(currentUserId, baseUrl).copy(isRead = true) } } catch (e: Exception) { Log.e(TAG, "Fetch messages failed", e) emptyList() } } override suspend fun getLastKnownSequenceId(chatId: String): Int? { return try { messageDao.getMaxSequenceId(chatId) } catch (e: Exception) { Log.e(TAG, "Failed to get last sequenceId", e) null } } override suspend fun sendMessage( chatId: String, content: String?, type: String, attachments: List?, replyToId: String?, forwardedFromId: String? ): Message { val userId = tokenManager.getUserId() ?: "" 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 ) messageDao.insertMessage(localMessage) Log.d(TAG, "Saved local message: $localId") // Пробуем отправить НЕМЕДЛЕННО через API try { Log.d(TAG, "Sending message immediately via API: $localId") val request = SendMessageRequest( content = content, type = type, attachments = attachments, replyToId = replyToId ) val response = api.sendMessage(chatId, request) Log.d(TAG, "Message sent successfully: ${response.id}") // Обновляем сообщение в базе с серверными данными val syncedMessage = localMessage.copy( id = response.id, sequenceId = response.sequenceId ?: 0, createdAt = response.createdAt ?: localMessage.createdAt, syncStatus = SyncStatus.SYNCED ) messageDao.insertMessage(syncedMessage) // Возвращаем доменную модель с серверными данными return response.toDomain(userId, baseUrl) } catch (e: Exception) { Log.e(TAG, "Failed to send message immediately, scheduling sync: ${e.message}") // Ошибка - планируем синхронизацию через WorkManager 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) { 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 { hubClient.readMessages(chats.data.remote.signalr.ReadMessagesRequest(chatId, lastMessageId, lastReadSequenceId)) } catch (e: Exception) { Log.e(TAG, "Error marking messages as read", e) } messageDao.markMessagesAsRead(chatId, lastReadSequenceId) } override suspend fun saveMessage(message: Message) { messageDao.insertMessage(message.toEntity(gson)) } override suspend fun deleteLocalMessage(messageId: String) { messageDao.markAsDeletedLocally(messageId) MessageSyncWorker.scheduleSync(context) } override suspend fun editLocalMessage(messageId: String, newContent: String) { messageDao.markAsEditedLocally(messageId, newContent) MessageSyncWorker.scheduleSync(context) } 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 = api.getTrendingGifs(page).data.data override suspend fun searchGifs(query: String, page: Int): List = api.searchGifs(query, page).data.data override suspend fun getGifCategories(): List = api.getGifCategories().data.categories override suspend fun createPersonalChat(userId: String): Chat { val currentUserId = tokenManager.getUserId() ?: "" 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) } override suspend fun editMessage(messageId: String, content: String): Message { val request = SendMessageRequest(content = content) val currentUserId = tokenManager.getUserId() ?: "" val response = api.editMessage(messageId, request) messageDao.insertMessage(response.toEntity(baseUrl, currentUserId, gson)) return response.toDomain(currentUserId, baseUrl) } companion object { private const val TAG = "ChatRepositoryImpl" } }