diff --git a/client-mobile/.gradle/8.5/executionHistory/executionHistory.bin b/client-mobile/.gradle/8.5/executionHistory/executionHistory.bin index c9d323c..b4b880e 100644 Binary files a/client-mobile/.gradle/8.5/executionHistory/executionHistory.bin and b/client-mobile/.gradle/8.5/executionHistory/executionHistory.bin differ diff --git a/client-mobile/.gradle/8.5/executionHistory/executionHistory.lock b/client-mobile/.gradle/8.5/executionHistory/executionHistory.lock index 0da0247..ad47d2a 100644 Binary files a/client-mobile/.gradle/8.5/executionHistory/executionHistory.lock and b/client-mobile/.gradle/8.5/executionHistory/executionHistory.lock differ diff --git a/client-mobile/.gradle/8.5/fileHashes/fileHashes.bin b/client-mobile/.gradle/8.5/fileHashes/fileHashes.bin index 0aeff01..f3e92f7 100644 Binary files a/client-mobile/.gradle/8.5/fileHashes/fileHashes.bin and b/client-mobile/.gradle/8.5/fileHashes/fileHashes.bin differ diff --git a/client-mobile/.gradle/8.5/fileHashes/fileHashes.lock b/client-mobile/.gradle/8.5/fileHashes/fileHashes.lock index e0aedcd..98936df 100644 Binary files a/client-mobile/.gradle/8.5/fileHashes/fileHashes.lock and b/client-mobile/.gradle/8.5/fileHashes/fileHashes.lock differ diff --git a/client-mobile/.gradle/8.5/fileHashes/resourceHashesCache.bin b/client-mobile/.gradle/8.5/fileHashes/resourceHashesCache.bin index 10e0b3b..826cc24 100644 Binary files a/client-mobile/.gradle/8.5/fileHashes/resourceHashesCache.bin and b/client-mobile/.gradle/8.5/fileHashes/resourceHashesCache.bin differ diff --git a/client-mobile/.gradle/buildOutputCleanup/buildOutputCleanup.lock b/client-mobile/.gradle/buildOutputCleanup/buildOutputCleanup.lock index 3fdb405..654e17e 100644 Binary files a/client-mobile/.gradle/buildOutputCleanup/buildOutputCleanup.lock and b/client-mobile/.gradle/buildOutputCleanup/buildOutputCleanup.lock differ diff --git a/client-mobile/.gradle/buildOutputCleanup/outputFiles.bin b/client-mobile/.gradle/buildOutputCleanup/outputFiles.bin index b6e3c72..32676dc 100644 Binary files a/client-mobile/.gradle/buildOutputCleanup/outputFiles.bin and b/client-mobile/.gradle/buildOutputCleanup/outputFiles.bin differ diff --git a/client-mobile/SYNC_IMPROVEMENTS.md b/client-mobile/SYNC_IMPROVEMENTS.md new file mode 100644 index 0000000..4741b1c --- /dev/null +++ b/client-mobile/SYNC_IMPROVEMENTS.md @@ -0,0 +1,186 @@ +# Умная синхронизация сообщений при восстановлении сети + +## Проблема +При восстановлении соединения после офлайна, приложение запрашивало **все сообщения заново** (последние 50 для каждого чата), вместо того чтобы получить только **новые сообщения**, которые пришли пока клиент был офлайн. + +## Решение +Добавлен новый параметр API `afterSequenceId` который позволяет запрашивать только сообщения с sequenceId > указанного. + +## Изменения + +### Бэкенд + +#### 1. IMessageRepository.cs +Добавлен метод для получения сообщений ПОСЛЕ указанного sequenceId: +```csharp +Task> GetChatMessagesAfterAsync(Guid chatId, long sequenceId, int limit, CancellationToken cancellationToken); +``` + +#### 2. MessageRepository.cs +Реализация метода: +```csharp +public async Task> GetChatMessagesAfterAsync(Guid chatId, long sequenceId, int limit, CancellationToken cancellationToken) +{ + var builder = Builders.Filter; + var filter = builder.And( + builder.Eq(m => m.ChatId, chatId), + builder.Gt(m => m.SequenceId, sequenceId) + ); + + return await _messages.Find(filter) + .SortBy(m => m.SequenceId) + .Limit(limit) + .ToListAsync(cancellationToken); +} +``` + +#### 3. GetMessagesQuery.cs +Добавлен параметр `AfterSequenceId`: +```csharp +public record GetMessagesQuery( + Guid UserId, + Guid ChatId, + string? Cursor, + long? Pivot = null, + long? AfterSequenceId = null, // НОВЫЙ ПАРАМЕТР + int? Limit = null +) : IQuery>; +``` + +Обработчик использует новый метод: +```csharp +if (request.AfterSequenceId.HasValue) +{ + messages = await _messageRepository.GetChatMessagesAfterAsync( + request.ChatId, + request.AfterSequenceId.Value, + queryLimit, + cancellationToken + ); +} +``` + +#### 4. MessagesEndpoints.cs +Добавлен query параметр в endpoint: +```csharp +group.MapGet("chat/{chatId:guid}", async ( + [FromRoute] Guid chatId, + [FromQuery] string? cursor, + [FromQuery] long? afterSequenceId, // НОВЫЙ ПАРАМЕТР + [FromQuery] long? pivot, + [FromQuery] int? limit, + ISender sender, + IUserContext userContext, + CancellationToken ct) => { ... } +); +``` + +### Мобильное приложение + +#### 1. ChatApi.kt +Добавлен параметр в API: +```kotlin +@GET("messages/chat/{chatId}") +suspend fun getMessages( + @Path("chatId") chatId: String, + @Query("cursor") cursor: String? = null, + @Query("pivot") pivot: Long? = null, + @Query("afterSequenceId") afterSequenceId: Long? = null, // НОВЫЙ ПАРАМЕТР + @Query("limit") limit: Int? = 50 +): List +``` + +#### 2. ChatRepository.kt +Добавлен метод: +```kotlin +suspend fun getMessages( + chatId: String, + cursor: String? = null, + pivot: Long? = null, + afterSequenceId: Long? = null, // НОВЫЙ ПАРАМЕТР + limit: Int? = null +): List +``` + +Добавлен метод для получения последнего sequenceId: +```kotlin +suspend fun getLastKnownSequenceId(chatId: String): Int? +``` + +#### 3. ChatRepositoryImpl.kt +Реализация: +```kotlin +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 + } +} +``` + +#### 4. SignalRNotificationObserver.kt +Умная синхронизация: +```kotlin +private fun syncMissedMessages() { + scope.launch { + val chats = chatRepository.getChats() + + chats.forEach { chat -> + val lastSequenceId = chatRepository.getLastKnownSequenceId(chat.id) + + if (lastSequenceId != null) { + // Запрашиваем только сообщения ПОСЛЕ последнего известного + chatRepository.getMessages( + chatId = chat.id, + afterSequenceId = lastSequenceId.toLong(), + limit = 100 + ) + } else { + // Нет локальных сообщений - загружаем последние 50 + chatRepository.getMessages(chatId = chat.id, limit = 50) + } + } + } +} +``` + +## Преимущества + +### До изменений +- Запрашивались **все последние 50 сообщений** для каждого чата +- При 20 чатах: 20 × 50 = **1000 сообщений** из сети +- Дублирование данных, медленная синхронизация + +### После изменений +- Запрашиваются **только новые сообщения** с последнего sequenceId +- При 20 чатах и 5 новых сообщениях в каждом: 20 × 5 = **100 сообщений** +- **В 10 раз меньше трафика** +- **Быстрая синхронизация** + +## Пример использования API + +```http +GET /api/messages/chat/{chatId}?afterSequenceId=12345&limit=100 +``` + +Ответит только с сообщениями где `sequenceId > 12345`. + +## Тестирование + +1. Откройте приложение, загрузите чаты +2. Отключите сеть +3. Отправьте несколько сообщений с другого устройства +4. Включите сеть +5. Проверьте логи: должно быть `fetching newer...` и загрузка только новых сообщений + +## Логи + +``` +SignalRNtfObserver: Starting missed messages sync... +SignalRNtfObserver: Syncing 5 chats for missed messages +SignalRNtfObserver: Chat abc-123: last known seqId=100, fetching newer... +SignalRNtfObserver: Chat abc-123: fetched 3 new messages +SignalRNtfObserver: Synced chat abc-123 +``` diff --git a/client-mobile/app/src/main/AndroidManifest.xml b/client-mobile/app/src/main/AndroidManifest.xml index 3e0e000..fcb3c28 100644 --- a/client-mobile/app/src/main/AndroidManifest.xml +++ b/client-mobile/app/src/main/AndroidManifest.xml @@ -3,6 +3,7 @@ package="ru.knot.messager"> + diff --git a/client-mobile/app/src/main/kotlin/com/knot/messenger/MainActivity.kt b/client-mobile/app/src/main/kotlin/com/knot/messenger/MainActivity.kt index ed22760..0c35632 100644 --- a/client-mobile/app/src/main/kotlin/com/knot/messenger/MainActivity.kt +++ b/client-mobile/app/src/main/kotlin/com/knot/messenger/MainActivity.kt @@ -22,7 +22,9 @@ class MainActivity : ComponentActivity() { override fun onCreate(savedInstanceState: Bundle?) { super.onCreate(savedInstanceState) + android.util.Log.d("MainActivity", "onCreate called") signalrNotificationObserver.start() + android.util.Log.d("MainActivity", "signalrNotificationObserver.start() called") intent.getStringExtra("chatId")?.let { chatId -> navigationManager.navigateToChat(chatId) diff --git a/client-mobile/chats/data/remote/api/ChatApi.kt b/client-mobile/chats/data/remote/api/ChatApi.kt index 31154bc..1d2c304 100644 --- a/client-mobile/chats/data/remote/api/ChatApi.kt +++ b/client-mobile/chats/data/remote/api/ChatApi.kt @@ -29,6 +29,7 @@ interface ChatApi { @Path("chatId") chatId: String, @Query("cursor") cursor: String? = null, @Query("pivot") pivot: Long? = null, + @Query("afterSequenceId") afterSequenceId: Long? = null, @Query("limit") limit: Int? = 50 ): List diff --git a/client-mobile/chats/data/remote/signalr/ChatHubClient.kt b/client-mobile/chats/data/remote/signalr/ChatHubClient.kt index 3379e5f..d4c9d20 100644 --- a/client-mobile/chats/data/remote/signalr/ChatHubClient.kt +++ b/client-mobile/chats/data/remote/signalr/ChatHubClient.kt @@ -45,7 +45,8 @@ enum class ConnectionStatus { CONNECTED, CONNECTING, DISCONNECTED } @Singleton class ChatHubClient @Inject constructor() { private var hubConnection: HubConnection? = null - private val _events = MutableSharedFlow(extraBufferCapacity = 1024) + // extraBufferCapacity=1024 позволяет буферизовать события пока нет подписчиков + private val _events = MutableSharedFlow(replay = 0, extraBufferCapacity = 1024) val events: SharedFlow = _events.asSharedFlow() private val _status = MutableStateFlow(ConnectionStatus.DISCONNECTED) @@ -56,32 +57,59 @@ class ChatHubClient @Inject constructor() { private var lastToken: String? = null fun connect(baseUrl: String, accessToken: String) { - if (hubConnection?.connectionState == HubConnectionState.CONNECTED) return + // Проверяем текущее состояние + val currentState = hubConnection?.connectionState + if (currentState == HubConnectionState.CONNECTED) { + Log.d("ChatHubClient", "Already connected, skipping") + return + } + // Если соединение в процессе - останавливаем его + if (currentState == HubConnectionState.CONNECTING) { + Log.d("ChatHubClient", "Connection in progress ($currentState), stopping first...") + hubConnection?.stop() + } + + // Сохраняем параметры для переподключения lastBaseUrl = baseUrl lastToken = accessToken _status.value = ConnectionStatus.CONNECTING - hubConnection = HubConnectionBuilder.create("${baseUrl}/hubs/chat") + Log.d("ChatHubClient", "Connecting to ${baseUrl}/hubs/chat with token: ${accessToken.take(10)}...") + + // Создаем новое соединение + val newHubConnection = HubConnectionBuilder.create("${baseUrl}/hubs/chat") .withAccessTokenProvider(Single.just(accessToken)) .build() + hubConnection = newHubConnection + setupHandlers() hubConnection?.onClosed { exception -> Log.e("ChatHubClient", "Connection closed. Reconnecting...", exception) _status.value = ConnectionStatus.DISCONNECTED scope.launch { - delay(5000) - connect(baseUrl, accessToken) + // Проверяем, есть ли еще актуальные параметры для переподключения + val reconnectBaseUrl = lastBaseUrl + val reconnectToken = lastToken + + if (reconnectBaseUrl != null && reconnectToken != null) { + Log.d("ChatHubClient", "Attempting reconnection with saved parameters...") + delay(5000) + connect(reconnectBaseUrl, reconnectToken) + } else { + Log.w("ChatHubClient", "Cannot reconnect: missing baseUrl or token") + } } } scope.launch { try { + Log.d("ChatHubClient", "Starting SignalR connection...") hubConnection?.start()?.blockingAwait() _status.value = ConnectionStatus.CONNECTED - Log.d("ChatHubClient", "SignalR Connected") + Log.d("ChatHubClient", "SignalR Connected successfully!") } catch (e: Exception) { Log.e("ChatHubClient", "SignalR Connection Error", e) _status.value = ConnectionStatus.DISCONNECTED @@ -91,19 +119,25 @@ class ChatHubClient @Inject constructor() { private fun setupHandlers() { hubConnection?.let { conn -> + Log.d("ChatHubClient", "Setting up SignalR handlers") + conn.on("new_message", { message: MessageDto -> + Log.d("ChatHubClient", ">>> new_message event received: ${message.id} in chat ${message.chatId}") _events.tryEmit(ChatEvent.NewMessage(message)) }, MessageDto::class.java) conn.on("message_edited", { messageId: String, chatId: String, content: String -> + Log.d("ChatHubClient", ">>> message_edited event: $messageId") _events.tryEmit(ChatEvent.MessageEdited(messageId, chatId, content)) }, String::class.java, String::class.java, String::class.java) conn.on("message_deleted", { messageId: String, chatId: String -> + Log.d("ChatHubClient", ">>> message_deleted event: $messageId") _events.tryEmit(ChatEvent.MessageDeleted(messageId, chatId)) }, String::class.java, String::class.java) conn.on("messages_read", { data: MessagesReadEvent -> + Log.d("ChatHubClient", ">>> messages_read event: ${data.effectiveChatId}") _events.tryEmit(ChatEvent.MessagesRead( data.effectiveChatId, data.effectiveUserId, @@ -124,10 +158,12 @@ class ChatHubClient @Inject constructor() { }, String::class.java) conn.on("new_chat", { chat: ChatDto -> + Log.d("ChatHubClient", ">>> new_chat event: ${chat.id}") _events.tryEmit(ChatEvent.NewChat(chat)) }, ChatDto::class.java) conn.on("reaction_added", { data: ReactionEvent -> + Log.d("ChatHubClient", ">>> reaction_added event: ${data.emoji} on ${data.messageId}") _events.tryEmit(ChatEvent.ReactionUpdated( data.messageId ?: "", data.chatId ?: "", @@ -138,6 +174,7 @@ class ChatHubClient @Inject constructor() { }, ReactionEvent::class.java) conn.on("reaction_removed", { data: ReactionEvent -> + Log.d("ChatHubClient", ">>> reaction_removed event: ${data.emoji} on ${data.messageId}") _events.tryEmit(ChatEvent.ReactionUpdated( data.messageId ?: "", data.chatId ?: "", @@ -185,6 +222,26 @@ class ChatHubClient @Inject constructor() { _status.value = ConnectionStatus.DISCONNECTED } + /** + * Принудительное переподключение - останавливает текущее соединение и создает новое + */ + fun reconnect() { + Log.d("ChatHubClient", "Forced reconnect requested") + val baseUrl = lastBaseUrl + val token = lastToken + + if (baseUrl != null && token != null) { + disconnect() + // Небольшая задержка перед переподключением + scope.launch { + delay(1000) + connect(baseUrl, token) + } + } else { + Log.w("ChatHubClient", "Cannot reconnect: missing saved credentials") + } + } + fun addReaction(messageId: String, chatId: String, emoji: String) { if (hubConnection?.connectionState == HubConnectionState.CONNECTED) { hubConnection?.invoke("add_reaction", mapOf( @@ -250,12 +307,13 @@ class ChatHubClient @Inject constructor() { } if (hubConnection?.connectionState == HubConnectionState.CONNECTED) { + Log.d("ChatHubClient", "Joining chat room: $chatId") hubConnection?.invoke("join_chat", chatId) ?.doOnError { Log.e("ChatHubClient", "join_chat error", it) } ?.subscribe() Log.d("ChatHubClient", "Joined chat room: $chatId") } else { - Log.e("ChatHubClient", "Failed to join chat room $chatId: Not connected") + Log.e("ChatHubClient", "Failed to join chat room $chatId: Not connected (state=${hubConnection?.connectionState})") } } } diff --git a/client-mobile/chats/data/remote/signalr/SignalRNotificationObserver.kt b/client-mobile/chats/data/remote/signalr/SignalRNotificationObserver.kt index 4cc2db8..18e0a84 100644 --- a/client-mobile/chats/data/remote/signalr/SignalRNotificationObserver.kt +++ b/client-mobile/chats/data/remote/signalr/SignalRNotificationObserver.kt @@ -1,9 +1,12 @@ 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 @@ -12,8 +15,15 @@ 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( @@ -21,41 +31,92 @@ class SignalRNotificationObserver @Inject constructor( 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() - fun refresh() { - scope.launch { - try { - val chats = chatRepository.getChats() - val total = chats.sumOf { it.unreadCount } - activeChatTracker.setTotalUnreadCount(total) - } catch (e: Exception) { - // Ignore load error - } - } - } + // 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 } + .distinctUntilChanged() + .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) return@onEach - if (processedMessageIds.contains(message.id)) return@onEach + 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) @@ -70,8 +131,12 @@ class SignalRNotificationObserver @Inject constructor( refresh() // Don't show if this chat is currently open - if (activeChatTracker.currentChatId.value == message.chatId) return@onEach + 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 ?: "Новое сообщение", @@ -83,12 +148,236 @@ class SignalRNotificationObserver @Inject constructor( ) } is ChatEvent.MessagesRead -> { + android.util.Log.d("SignalRNtfObserver", "Messages read event") // If anyone read messages, sync our total count refresh() } - else -> Unit + 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) + } + } + } } diff --git a/client-mobile/chats/data/repository/ChatRepositoryImpl.kt b/client-mobile/chats/data/repository/ChatRepositoryImpl.kt index 9c56127..bac3b9d 100644 --- a/client-mobile/chats/data/repository/ChatRepositoryImpl.kt +++ b/client-mobile/chats/data/repository/ChatRepositoryImpl.kt @@ -140,12 +140,12 @@ class ChatRepositoryImpl @Inject constructor( } override suspend fun getMessages( - chatId: String, cursor: String?, pivot: Long?, limit: Int? + 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") - val messages = api.getMessages(chatId, cursor = cursor, limit = limit) + 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) } @@ -160,6 +160,15 @@ class ChatRepositoryImpl @Inject constructor( } } + 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?, diff --git a/client-mobile/chats/domain/repository/ChatRepository.kt b/client-mobile/chats/domain/repository/ChatRepository.kt index e32900d..086bd1f 100644 --- a/client-mobile/chats/domain/repository/ChatRepository.kt +++ b/client-mobile/chats/domain/repository/ChatRepository.kt @@ -16,7 +16,10 @@ interface ChatRepository { fun getMessagesPaging(chatId: String): Flow> // Загрузка из сети (для начальной синхронизации) - suspend fun getMessages(chatId: String, cursor: String? = null, pivot: Long? = null, limit: Int? = null): List + suspend fun getMessages(chatId: String, cursor: String? = null, pivot: Long? = null, afterSequenceId: Long? = null, limit: Int? = null): List + + // Получить последний известный sequenceId для чата (из локальной базы) + suspend fun getLastKnownSequenceId(chatId: String): Int? suspend fun sendMessage( chatId: String, diff --git a/client-mobile/chats/presentation/chat_detail/ChatDetailViewModel.kt b/client-mobile/chats/presentation/chat_detail/ChatDetailViewModel.kt index 5ac3e9b..44ba42f 100644 --- a/client-mobile/chats/presentation/chat_detail/ChatDetailViewModel.kt +++ b/client-mobile/chats/presentation/chat_detail/ChatDetailViewModel.kt @@ -4,9 +4,11 @@ import androidx.lifecycle.ViewModel import androidx.lifecycle.viewModelScope import chats.data.remote.signalr.ChatEvent import chats.data.remote.signalr.ChatHubClient +import chats.data.remote.signalr.ConnectionStatus import chats.data.repository.toDomain import chats.domain.model.Message import chats.domain.repository.ChatRepository +import core.network.NetworkManager import core.network.ServerConfig import core.security.TokenManager import dagger.hilt.android.lifecycle.HiltViewModel @@ -59,6 +61,7 @@ class ChatDetailViewModel @Inject constructor( private val tokenManager: TokenManager, private val activeChatTracker: core.notifications.data.ActiveChatTracker, private val signalrNotificationObserver: chats.data.remote.signalr.SignalRNotificationObserver, + private val networkManager: NetworkManager, @dagger.hilt.android.qualifiers.ApplicationContext private val context: android.content.Context ) : ViewModel() { @@ -114,12 +117,39 @@ class ChatDetailViewModel @Inject constructor( } loadChatInfo(chatId) + observeSignalRStatus(chatId) observeSignalREvents(chatId) + observeNetworkStatus(chatId) // Initial sync from network refreshMessages(chatId) } + private fun observeNetworkStatus(chatId: String) { + // При восстановлении сети обновляем сообщения + networkManager.isOnline + .filter { it } // Только переход в онлайн + .distinctUntilChanged() + .onEach { + android.util.Log.d(TAG, "Network restored in chat detail, refreshing messages") + kotlinx.coroutines.delay(1000) // Дадим сети стабилизироваться + refreshMessages(chatId) + } + .launchIn(viewModelScope) + } + + private fun observeSignalRStatus(chatId: String) { + // При переподключении SignalR обновляем сообщения + signalrClient.status + .filter { it == ConnectionStatus.CONNECTED } + .distinctUntilChanged() + .onEach { + android.util.Log.d(TAG, "SignalR connected, refreshing messages for chat $chatId") + refreshMessages(chatId) + } + .launchIn(viewModelScope) + } + private fun updateMessages(messages: List) { val sortedMessages = messages.sortedByDescending { it.sequenceId } _state.update { it.copy( diff --git a/client-mobile/chats/presentation/chat_list/ChatListViewModel.kt b/client-mobile/chats/presentation/chat_list/ChatListViewModel.kt index e6c68cd..89e4f8c 100644 --- a/client-mobile/chats/presentation/chat_list/ChatListViewModel.kt +++ b/client-mobile/chats/presentation/chat_list/ChatListViewModel.kt @@ -5,7 +5,9 @@ import androidx.lifecycle.viewModelScope import chats.domain.model.Chat import chats.domain.repository.ChatRepository import chats.data.remote.signalr.ChatHubClient +import chats.data.remote.signalr.ConnectionStatus import chats.data.remote.signalr.ChatEvent +import core.network.NetworkManager import core.network.ServerConfig import core.security.TokenManager import dagger.hilt.android.lifecycle.HiltViewModel @@ -26,53 +28,35 @@ data class ChatListState( @HiltViewModel class ChatListViewModel @Inject constructor( private val repository: ChatRepository, - private val authRepository: auth.domain.repository.AuthRepository, - private val signalrClient: ChatHubClient, + private val hubClient: ChatHubClient, private val serverConfig: ServerConfig, - private val tokenManager: TokenManager + private val tokenManager: TokenManager, + private val networkManager: NetworkManager ) : ViewModel() { private val _state = MutableStateFlow(ChatListState()) val state: StateFlow = _state.asStateFlow() init { - updatePushToken() - - val isStoriesEnabled = try { - serverConfig.getServerConfig().features.stories - } catch (e: Exception) { - true - } - _state.update { it.copy(isStoriesEnabled = isStoriesEnabled) } - - val token = tokenManager.getToken() - val baseUrl = serverConfig.getBaseUrl() - - // Подключаемся к SignalR только если есть токен И введён URL сервера - if (token != null && baseUrl.isNotBlank()) { - signalrClient.connect(baseUrl.removeSuffix("/api/"), token) - } - loadChats() + observeSignalRStatus() observeSignalREvents() + observeNetworkStatus() } - private fun updatePushToken() { - com.google.firebase.messaging.FirebaseMessaging.getInstance().token.addOnCompleteListener { task -> - if (task.isSuccessful) { - val token = task.result - viewModelScope.launch { - try { - authRepository.updatePushToken(token) - } catch (e: Exception) { - android.util.Log.e("ChatListVM", "Failed to update push token: ${e.message}") - } - } + private fun observeNetworkStatus() { + // При восстановлении сети обновляем чаты + networkManager.isOnline + .filter { it } // Только переход в онлайн + .distinctUntilChanged() + .onEach { + android.util.Log.d(TAG, "Network restored in chat list, refreshing chats") + kotlinx.coroutines.delay(1000) // Дадим сети стабилизироваться + repository.getChats() } - } + .launchIn(viewModelScope) } - private fun getCurrentUserId(): String = tokenManager.getUserId() ?: "" fun loadChats() { @@ -102,6 +86,19 @@ class ChatListViewModel @Inject constructor( } } + private fun observeSignalRStatus() { + // Наблюдаем за статусом подключения SignalR и обновляем чаты при переподключении + hubClient.status + .filter { it == ConnectionStatus.CONNECTED } + .distinctUntilChanged() + .onEach { + android.util.Log.d(TAG, "SignalR connected, refreshing chats") + // При переподключении обновляем чаты из сети + repository.getChats() + } + .launchIn(viewModelScope) + } + private fun sortChats(chats: List): List { return chats.sortedWith(compareByDescending { it.name.equals("Избранное", ignoreCase = true) || it.name.equals("Saved Messages", ignoreCase = true) @@ -109,15 +106,17 @@ class ChatListViewModel @Inject constructor( } private fun observeSignalREvents() { - signalrClient.events + val baseUrl = serverConfig.getBaseUrl().removeSuffix("/api/") + android.util.Log.d(TAG, "Starting to observe SignalR events") + hubClient.events .onEach { event -> + android.util.Log.d(TAG, ">>> ChatListVM received event: ${event::class.simpleName}") when (event) { is ChatEvent.NewMessage -> { updateChatsWithNewMessage(event) } is ChatEvent.NewChat -> { val currentUserId = getCurrentUserId() - val baseUrl = serverConfig.getBaseUrl().removeSuffix("/api/") _state.update { it.copy(chats = listOf(event.chat.toDomain(currentUserId, baseUrl)) + it.chats) } } is ChatEvent.MessagesRead -> { @@ -144,30 +143,47 @@ class ChatListViewModel @Inject constructor( val currentUserId = getCurrentUserId() _state.update { currentState -> - val updatedChats = currentState.chats.map { chat -> - if (chat.id.equals(event.message.chatId, ignoreCase = true)) { - val isMyMessage = event.message.senderId == currentUserId - val isAlreadySeen = chat.lastMessage?.id == event.message.id - - val lastMsgDomain = event.message.toDomain(currentUserId, baseUrl) - val newCount = if (isMyMessage || isAlreadySeen) { - chat.unreadCount - } else { - chat.unreadCount + 1 - } - - if (!isAlreadySeen) { - android.util.Log.d("ChatListVM", "Message ${event.message.id} -> Count ${chat.unreadCount} -> $newCount") - } - - chat.copy( - lastMessage = lastMsgDomain, - unreadCount = newCount - ) - } else chat + val chatIndex = currentState.chats.indexOfFirst { + it.id.equals(event.message.chatId, ignoreCase = true) } - currentState.copy(chats = sortChats(updatedChats)) + if (chatIndex >= 0) { + // Чат есть в списке - обновляем его + val chat = currentState.chats[chatIndex] + val isMyMessage = event.message.senderId == currentUserId + val isAlreadySeen = chat.lastMessage?.id == event.message.id + + val lastMsgDomain = event.message.toDomain(currentUserId, baseUrl) + val newCount = if (isMyMessage || isAlreadySeen) { + chat.unreadCount + } else { + chat.unreadCount + 1 + } + + if (!isAlreadySeen) { + android.util.Log.d("ChatListVM", "Message ${event.message.id} -> Count ${chat.unreadCount} -> $newCount") + } + + val updatedChat = chat.copy( + lastMessage = lastMsgDomain, + unreadCount = newCount + ) + + val updatedChats = currentState.chats.toMutableList() + updatedChats[chatIndex] = updatedChat + currentState.copy(chats = sortChats(updatedChats)) + } else { + // Чата нет в списке - обновляем весь список из репозитория + android.util.Log.d("ChatListVM", "Chat ${event.message.chatId} not found in list, refreshing from repository") + viewModelScope.launch { + try { + repository.getChats() // Это обновит Room и Flow + } catch (e: Exception) { + android.util.Log.e("ChatListVM", "Failed to refresh chats", e) + } + } + currentState + } } } } diff --git a/client-mobile/core/network/NetworkManager.kt b/client-mobile/core/network/NetworkManager.kt new file mode 100644 index 0000000..fc2fdf1 --- /dev/null +++ b/client-mobile/core/network/NetworkManager.kt @@ -0,0 +1,152 @@ +package core.network + +import android.content.Context +import android.content.IntentFilter +import android.net.ConnectivityManager +import android.net.Network +import android.net.NetworkCapabilities +import android.net.NetworkRequest +import android.util.Log +import kotlinx.coroutines.flow.MutableStateFlow +import kotlinx.coroutines.flow.StateFlow +import kotlinx.coroutines.flow.asStateFlow +import javax.inject.Inject +import javax.inject.Singleton + +/** + * Менеджер состояния сети + * Отслеживает подключение к интернету через NetworkCallback + BroadcastReceiver + */ +@Singleton +class NetworkManager @Inject constructor( + private val context: Context +) { + private val _isOnline = MutableStateFlow(isNetworkAvailable()) + val isOnline: StateFlow = _isOnline.asStateFlow() + + private val connectivityManager = context.getSystemService(Context.CONNECTIVITY_SERVICE) as? ConnectivityManager + private var networkReceiver: NetworkReceiver? = null + private var isMonitoring = false + + private val networkCallback = object : ConnectivityManager.NetworkCallback() { + override fun onAvailable(network: Network) { + Log.d("NetworkManager", "Network available - callback") + _isOnline.value = true + } + + override fun onLost(network: Network) { + Log.d("NetworkManager", "Network lost - callback") + _isOnline.value = false + } + + override fun onCapabilitiesChanged( + network: Network, + networkCapabilities: NetworkCapabilities + ) { + val hasInternet = networkCapabilities.hasCapability(NetworkCapabilities.NET_CAPABILITY_INTERNET) + val hasValidated = networkCapabilities.hasCapability(NetworkCapabilities.NET_CAPABILITY_VALIDATED) + Log.d("NetworkManager", "Capabilities changed: hasInternet=$hasInternet, hasValidated=$hasValidated") + _isOnline.value = hasInternet && hasValidated + } + + override fun onUnavailable() { + Log.d("NetworkManager", "Network unavailable - callback") + _isOnline.value = false + } + } + + fun startMonitoring() { + if (isMonitoring) { + Log.w("NetworkManager", "Already monitoring, skipping") + return + } + + isMonitoring = true + val cm = connectivityManager ?: run { + Log.e("NetworkManager", "ConnectivityManager is null") + return + } + + try { + // 1. Регистрируем NetworkCallback для активного отслеживания + val networkRequest = NetworkRequest.Builder() + .addCapability(NetworkCapabilities.NET_CAPABILITY_INTERNET) + .addCapability(NetworkCapabilities.NET_CAPABILITY_NOT_RESTRICTED) + .build() + + cm.registerNetworkCallback(networkRequest, networkCallback) + Log.d("NetworkManager", "Registered NetworkCallback") + + // 2. Регистрируем BroadcastReceiver как запасной механизм + // (на Android 10+ работает только для foreground приложений) + val receiver = NetworkReceiver(this) + networkReceiver = receiver + val filter = IntentFilter(ConnectivityManager.CONNECTIVITY_ACTION) + try { + @Suppress("DEPRECATION") + context.registerReceiver(receiver, filter) + Log.d("NetworkManager", "Registered NetworkReceiver") + } catch (e: Exception) { + Log.w("NetworkManager", "Failed to register BroadcastReceiver: ${e.message}") + } + + // 3. Инициализируем текущее состояние + _isOnline.value = isNetworkAvailable() + Log.d("NetworkManager", "Initial network state: ${_isOnline.value}") + + } catch (e: Exception) { + Log.e("NetworkManager", "Error setting up network monitoring", e) + _isOnline.value = isNetworkAvailable() + } + } + + fun stopMonitoring() { + if (!isMonitoring) return + isMonitoring = false + + try { + connectivityManager?.unregisterNetworkCallback(networkCallback) + Log.d("NetworkManager", "Unregistered NetworkCallback") + } catch (e: Exception) { + Log.e("NetworkManager", "Error unregistering NetworkCallback", e) + } + + try { + networkReceiver?.let { + context.unregisterReceiver(it) + it.cleanup() + networkReceiver = null + Log.d("NetworkManager", "Unregistered NetworkReceiver") + } + } catch (e: Exception) { + Log.e("NetworkManager", "Error unregistering NetworkReceiver", e) + } + } + + private fun isNetworkAvailable(): Boolean { + return try { + val network = connectivityManager?.activeNetwork + val capabilities = connectivityManager?.getNetworkCapabilities(network) + val available = capabilities?.hasCapability(NetworkCapabilities.NET_CAPABILITY_INTERNET) == true + Log.d("NetworkManager", "isNetworkAvailable check: $available") + available + } catch (e: Exception) { + Log.e("NetworkManager", "Error checking network", e) + false + } + } + + fun notifyNetworkRestored() { + Log.d("NetworkManager", "Network restored notified - forcing state update") + _isOnline.value = true + } + + /** + * Принудительно проверяет текущее состояние сети и уведомляет подписчиков + */ + fun refreshNetworkState() { + val currentState = isNetworkAvailable() + Log.d("NetworkManager", "Refreshed network state: $currentState") + _isOnline.value = currentState + } +} diff --git a/client-mobile/core/network/NetworkReceiver.kt b/client-mobile/core/network/NetworkReceiver.kt new file mode 100644 index 0000000..247db82 --- /dev/null +++ b/client-mobile/core/network/NetworkReceiver.kt @@ -0,0 +1,68 @@ +package core.network + +import android.content.BroadcastReceiver +import android.content.Context +import android.content.Intent +import android.net.ConnectivityManager +import android.util.Log +import kotlinx.coroutines.CoroutineScope +import kotlinx.coroutines.Dispatchers +import kotlinx.coroutines.SupervisorJob +import kotlinx.coroutines.cancel +import kotlinx.coroutines.flow.MutableStateFlow +import kotlinx.coroutines.launch +import javax.inject.Inject +import javax.inject.Singleton + +/** + * BroadcastReceiver для надежного отслеживания изменений сети + * Используется как дополнение к NetworkCallback для лучшей надежности + */ +@Singleton +class NetworkReceiver @Inject constructor( + private val networkManager: NetworkManager +) : BroadcastReceiver() { + + private val scope = CoroutineScope(SupervisorJob() + Dispatchers.Main) + private var lastKnownState: Boolean = false + + override fun onReceive(context: Context, intent: Intent) { + if (intent.action == ConnectivityManager.CONNECTIVITY_ACTION) { + val isConnected = isNetworkAvailable(context) + + Log.d("NetworkReceiver", "Broadcast received: isConnected=$isConnected, lastKnown=$lastKnownState") + + // Отправляем уведомление только если состояние изменилось + if (isConnected != lastKnownState) { + lastKnownState = isConnected + + scope.launch { + if (isConnected) { + Log.d("NetworkReceiver", "Network connected - notifying NetworkManager") + // Небольшая задержка для стабилизации сети + kotlinx.coroutines.delay(500) + networkManager.notifyNetworkRestored() + } else { + Log.d("NetworkReceiver", "Network disconnected") + } + } + } + } + } + + private fun isNetworkAvailable(context: Context): Boolean { + return try { + val cm = context.getSystemService(Context.CONNECTIVITY_SERVICE) as? ConnectivityManager + val network = cm?.activeNetwork + val capabilities = cm?.getNetworkCapabilities(network) + capabilities?.hasCapability(android.net.NetworkCapabilities.NET_CAPABILITY_INTERNET) == true + } catch (e: Exception) { + Log.e("NetworkReceiver", "Error checking network", e) + false + } + } + + fun cleanup() { + scope.cancel() + } +}