package chats.data.remote.signalr import android.content.Context import core.network.NetworkManager import core.network.ServerConfig import core.notifications.data.ActiveChatTracker import core.notifications.data.NotificationHelper import core.security.TokenManager import chats.data.sync.MessageSyncWorker import dagger.hilt.android.qualifiers.ApplicationContext import kotlinx.coroutines.CoroutineScope import kotlinx.coroutines.Dispatchers import kotlinx.coroutines.SupervisorJob import kotlinx.coroutines.launch import kotlinx.coroutines.flow.filterIsInstance import kotlinx.coroutines.flow.launchIn import kotlinx.coroutines.flow.onEach import kotlinx.coroutines.flow.filter import kotlinx.coroutines.flow.distinctUntilChanged import kotlinx.coroutines.flow.collect import javax.inject.Inject import javax.inject.Singleton import chats.data.remote.signalr.ConnectionStatus import okhttp3.OkHttpClient import okhttp3.Request import java.util.concurrent.TimeUnit @Singleton class SignalRNotificationObserver @Inject constructor( private val signalrClient: ChatHubClient, private val activeChatTracker: ActiveChatTracker, private val tokenManager: TokenManager, private val chatRepository: chats.domain.repository.ChatRepository, private val serverConfig: ServerConfig, private val networkManager: NetworkManager, @ApplicationContext private val context: Context ) { private val scope = CoroutineScope(SupervisorJob() + Dispatchers.Main) private var isStarted = false private val processedMessageIds = mutableSetOf() // OkHttpClient для ping запроса private val pingClient = OkHttpClient.Builder() .connectTimeout(5, TimeUnit.SECONDS) .readTimeout(5, TimeUnit.SECONDS) .build() fun start() { if (isStarted) return isStarted = true // Запускаем мониторинг сети networkManager.startMonitoring() // Принудительно обновляем состояние сети при старте networkManager.refreshNetworkState() // Подключаемся к SignalR при старте приложения connectSignalR() // Initial count load refresh() // Запускаем периодическую проверку подключения SignalR startConnectionHealthCheck() // Слушаем восстановление сети и переподключаем SignalR networkManager.isOnline .filter { it } // Только переход в онлайн .distinctUntilChanged() .onEach { android.util.Log.d("SignalRNtfObserver", "Network restored! Reconnecting SignalR and syncing...") // Небольшая задержка чтобы сеть стабилизировалась kotlinx.coroutines.delay(1500) reconnectOnNetworkRestored() } .launchIn(scope) // Также отслеживаем состояние SignalR для переподключения signalrClient.status .filter { it == ConnectionStatus.DISCONNECTED } .onEach { android.util.Log.d("SignalRNtfObserver", "SignalR disconnected, checking network...") // Если мы offline и SignalR отключен - не делаем ничего // Подключимся когда сеть восстановится } .launchIn(scope) // При восстановлении соединения SignalR обновляем список чатов и вступаем в них signalrClient.status .filter { it == ConnectionStatus.CONNECTED } .onEach { android.util.Log.d("SignalRNtfObserver", "SignalR connected, refreshing chats and joining rooms") onSignalRConnected() } .launchIn(scope) // Слушаем события SignalR для уведомлений и обновления списка чатов android.util.Log.d("SignalRNtfObserver", "Starting to listen to signalrClient.events") signalrClient.events .onEach { event -> android.util.Log.d("SignalRNtfObserver", ">>> Received event: ${event::class.simpleName}") when (event) { is ChatEvent.NewMessage -> { val currentUserId = tokenManager.getUserId() val message = event.message android.util.Log.d("SignalRNtfObserver", "New message: ${message.id} from ${message.senderId}, current: $currentUserId") // Don't show if it's our message or already processed if (message.senderId == currentUserId) { android.util.Log.d("SignalRNtfObserver", "Skipping - own message") return@onEach } if (processedMessageIds.contains(message.id)) { android.util.Log.d("SignalRNtfObserver", "Skipping - already processed") return@onEach } // Mark as processed processedMessageIds.add(message.id) if (processedMessageIds.size > 200) { processedMessageIds.remove(processedMessageIds.first()) } // Increment immediately for UI feedback activeChatTracker.incrementUnreadCount() // Refresh total count from source of truth in background refresh() // Don't show if this chat is currently open if (activeChatTracker.currentChatId.value == message.chatId) { android.util.Log.d("SignalRNtfObserver", "Skipping - chat is open: ${message.chatId}") return@onEach } android.util.Log.d("SignalRNtfObserver", "Showing notification for chat: ${message.chatId}") NotificationHelper.showNotification( context = context, title = message.sender?.displayName ?: "Новое сообщение", body = message.content ?: "Вам прислали вложение", type = "chat", chatId = message.chatId, notificationId = message.id.hashCode(), totalCount = activeChatTracker.totalUnreadCount.value ) } is ChatEvent.MessagesRead -> { android.util.Log.d("SignalRNtfObserver", "Messages read event") // If anyone read messages, sync our total count refresh() } else -> { android.util.Log.d("SignalRNtfObserver", "Unhandled event: ${event::class.simpleName}") Unit } } } .launchIn(scope) } private fun startConnectionHealthCheck() { // Каждые 10 секунд проверяем подключение и при необходимости переподключаемся scope.launch { var consecutiveFailures = 0 while (true) { kotlinx.coroutines.delay(10000) val currentStatus = signalrClient.status.value val isOnline = networkManager.isOnline.value if (isOnline && currentStatus == ConnectionStatus.DISCONNECTED) { consecutiveFailures++ android.util.Log.d("SignalRNtfObserver", "Health check: Network online but SignalR disconnected (failures: $consecutiveFailures), reconnecting...") signalrClient.reconnect() } else if (isOnline && currentStatus == ConnectionStatus.CONNECTED) { // Сбрасываем счетчик ошибок при успешном подключении consecutiveFailures = 0 android.util.Log.d("SignalRNtfObserver", "Health check: Connection healthy") } } } } private fun onSignalRConnected() { android.util.Log.d("SignalRNtfObserver", "SignalR connected - syncing missed messages...") refresh() // Вступаем во все чаты для получения событий joinAllChats() // Синхронизируем пропущенные сообщения syncMissedMessages() } /** * Синхронизирует пропущенные сообщения после восстановления соединения * Запрашивает только НОВЫЕ сообщения с последнего известного sequenceId */ private fun syncMissedMessages() { scope.launch { try { android.util.Log.d("SignalRNtfObserver", "Starting missed messages sync...") val chats = chatRepository.getChats() android.util.Log.d("SignalRNtfObserver", "Syncing ${chats.size} chats for missed messages") // Запрашиваем только новые сообщения для каждого чата chats.forEach { chat -> try { // Получаем последний известный sequenceId из локальной базы val lastSequenceId = chatRepository.getLastKnownSequenceId(chat.id) if (lastSequenceId != null) { // Запрашиваем сообщения ПОСЛЕ lastSequenceId (только новые) android.util.Log.d("SignalRNtfObserver", "Chat ${chat.id}: last known seqId=$lastSequenceId, fetching newer...") chatRepository.getMessages( chatId = chat.id, afterSequenceId = lastSequenceId.toLong(), limit = 100 ) } else { // Нет локальных сообщений - загружаем последние 50 android.util.Log.d("SignalRNtfObserver", "Chat ${chat.id}: no local messages, fetching last 50") chatRepository.getMessages(chatId = chat.id, limit = 50) } android.util.Log.d("SignalRNtfObserver", "Synced chat ${chat.id}") // Небольшая пауза между чатами чтобы не перегружать сервер kotlinx.coroutines.delay(100) } catch (e: Exception) { android.util.Log.e("SignalRNtfObserver", "Failed to sync chat ${chat.id}", e) } } android.util.Log.d("SignalRNtfObserver", "Missed messages sync completed") } catch (e: Exception) { android.util.Log.e("SignalRNtfObserver", "Sync failed", e) } } } /** * Получает последний известный sequenceId для чата из локальной базы */ private suspend fun getLastKnownSequenceId(chatId: String): Int? { return try { chatRepository.getLastKnownSequenceId(chatId) } catch (e: Exception) { android.util.Log.e("SignalRNtfObserver", "Failed to get last sequenceId for $chatId", e) null } } private fun reconnectOnNetworkRestored() { android.util.Log.d("SignalRNtfObserver", "Network restored, starting reconnection sequence...") scope.launch { // 1. Сначала делаем HTTP ping запрос чтобы "разбудить" сетевой стек android.util.Log.d("SignalRNtfObserver", "Sending HTTP ping to wake up network...") val pingSuccess = sendHttpPing() android.util.Log.d("SignalRNtfObserver", "HTTP ping result: $pingSuccess") // 2. Небольшая пауза для стабилизации kotlinx.coroutines.delay(1000) // 3. Принудительное переподключение SignalR android.util.Log.d("SignalRNtfObserver", "Forcing SignalR reconnect...") signalrClient.reconnect() // 4. Ждем пока SignalR подключится (максимум 10 секунд) var waitCount = 0 while (signalrClient.status.value != ConnectionStatus.CONNECTED && waitCount < 20) { kotlinx.coroutines.delay(500) waitCount++ } if (signalrClient.status.value == ConnectionStatus.CONNECTED) { android.util.Log.d("SignalRNtfObserver", "SignalR reconnected, syncing missed messages...") // 5. Синхронизируем пропущенные сообщения syncMissedMessages() } else { android.util.Log.w("SignalRNtfObserver", "SignalR failed to reconnect within timeout") } // 6. Запускаем синхронизацию отложенных сообщений MessageSyncWorker.scheduleSync(context) android.util.Log.d("SignalRNtfObserver", "Outgoing sync worker scheduled") } } /** * Отправляет HTTP ping запрос для активации сетевого соединения */ private suspend fun sendHttpPing(): Boolean { return try { val baseUrl = serverConfig.getBaseUrl() val token = tokenManager.getToken() if (baseUrl.isBlank() || token == null) { android.util.Log.w("SignalRNtfObserver", "Cannot ping: missing baseUrl or token") return false } val url = "${baseUrl.removeSuffix("/api/")}/api/auth/refresh" val request = Request.Builder() .url(url) .post(okhttp3.RequestBody.create(null, "{}")) .addHeader("Authorization", "Bearer $token") .build() val response = pingClient.newCall(request).execute() val success = response.isSuccessful || response.code == 401 // 401 OK для refresh android.util.Log.d("SignalRNtfObserver", "HTTP ping to $url: ${response.code}") response.close() success } catch (e: Exception) { android.util.Log.e("SignalRNtfObserver", "HTTP ping failed", e) false } } private fun connectSignalR() { val token = tokenManager.getToken() val baseUrl = serverConfig.getBaseUrl() if (token == null || baseUrl.isBlank()) { android.util.Log.w("SignalRNtfObserver", "Cannot connect SignalR: token=${token != null}, baseUrl=$baseUrl") return } val isOnline = networkManager.isOnline.value if (!isOnline) { android.util.Log.w("SignalRNtfObserver", "Cannot connect SignalR: network is offline") return } // Проверяем текущее состояние SignalR val currentStatus = signalrClient.status.value if (currentStatus == ConnectionStatus.CONNECTED) { android.util.Log.d("SignalRNtfObserver", "SignalR already connected, skipping") return } android.util.Log.d("SignalRNtfObserver", "Connecting SignalR with token: ${token.take(10)}..., baseUrl: $baseUrl") signalrClient.connect(baseUrl.removeSuffix("/api/"), token) // Если подключение не удалось в течение 5 секунд - пробуем снова scope.launch { kotlinx.coroutines.delay(5000) if (signalrClient.status.value == ConnectionStatus.DISCONNECTED) { android.util.Log.w("SignalRNtfObserver", "SignalR connection timeout, retrying with reconnect()...") signalrClient.reconnect() } } } private fun joinAllChats() { scope.launch { try { val chats = chatRepository.getChats() android.util.Log.d("SignalRNtfObserver", "Joining ${chats.size} chat rooms") chats.forEach { chat -> signalrClient.joinChat(chat.id) } } catch (e: Exception) { android.util.Log.e("SignalRNtfObserver", "Failed to join chats", e) } } } fun refresh() { scope.launch { try { val chats = chatRepository.getChats() val total = chats.sumOf { it.unreadCount } activeChatTracker.setTotalUnreadCount(total) android.util.Log.d("SignalRNtfObserver", "Refreshed chats: ${chats.size}, total unread: $total") } catch (e: Exception) { android.util.Log.e("SignalRNtfObserver", "Failed to refresh", e) } } } }