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>(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" } }