From 7f2e684572d7a8f4be2b5b90d897cebb61ef6f91 Mon Sep 17 00:00:00 2001 From: callebtc <93376500+callebtc@users.noreply.github.com> Date: Mon, 27 Jul 2026 22:33:28 +0200 Subject: [PATCH 1/2] feat: retry scheduler for queued private messages and session re-establishment --- .../android/mesh/BluetoothMeshService.kt | 11 + .../bitchat/android/services/MessageRouter.kt | 218 ++++++++++++++++-- .../com/bitchat/android/ui/ChatViewModel.kt | 9 + .../com/bitchat/android/util/AppConstants.kt | 7 + .../android/services/MessageRouterTest.kt | 200 ++++++++++++++++ 5 files changed, 424 insertions(+), 21 deletions(-) create mode 100644 app/src/test/kotlin/com/bitchat/android/services/MessageRouterTest.kt diff --git a/app/src/main/java/com/bitchat/android/mesh/BluetoothMeshService.kt b/app/src/main/java/com/bitchat/android/mesh/BluetoothMeshService.kt index b17294f1..61924738 100644 --- a/app/src/main/java/com/bitchat/android/mesh/BluetoothMeshService.kt +++ b/app/src/main/java/com/bitchat/android/mesh/BluetoothMeshService.kt @@ -137,6 +137,17 @@ class BluetoothMeshService(private val context: Context) : TransportBridgeServic messageHandler.packetProcessor = packetProcessor //startPeriodicDebugLogging() + // Flush queued private messages as soon as a BLE Noise session authenticates, + // instead of relying on the foreground-only UI poll. + encryptionService.onSessionEstablished = { peerID -> + Log.d(TAG, "BLE Noise session established with ${peerID.take(8)}") + try { + com.bitchat.android.services.MessageRouter + .tryGetInstance() + ?.onSessionEstablished(peerID) + } catch (_: Exception) { } + } + // Initialize sync manager (needs serviceScope) gossipSyncManager = GossipSyncManager( myPeerID = myPeerID, diff --git a/app/src/main/java/com/bitchat/android/services/MessageRouter.kt b/app/src/main/java/com/bitchat/android/services/MessageRouter.kt index 652e0c93..7905aee8 100644 --- a/app/src/main/java/com/bitchat/android/services/MessageRouter.kt +++ b/app/src/main/java/com/bitchat/android/services/MessageRouter.kt @@ -6,6 +6,15 @@ import com.bitchat.android.favorites.FavoriteControlMessage import com.bitchat.android.mesh.MeshService import com.bitchat.android.model.ReadReceipt import com.bitchat.android.nostr.NostrTransport +import com.bitchat.android.util.AppConstants +import kotlinx.coroutines.CoroutineScope +import kotlinx.coroutines.Dispatchers +import kotlinx.coroutines.SupervisorJob +import kotlinx.coroutines.cancel +import kotlinx.coroutines.delay +import kotlinx.coroutines.isActive +import kotlinx.coroutines.launch +import java.util.concurrent.ConcurrentHashMap /** * Routes messages between local mesh transports and Nostr, matching iOS behavior. @@ -22,9 +31,27 @@ class MessageRouter private constructor( DROPPED } + private data class QueuedMessage( + val content: String, + val nickname: String, + val messageID: String, + val enqueuedAtMs: Long + ) + + private data class ConversationRetry( + val handshakeAttempts: Int, + val nextHandshakeAttemptAtMs: Long + ) + companion object { private const val TAG = "MessageRouter" + private const val OUTBOX_TICK_MS = AppConstants.Router.OUTBOX_TICK_MS + private const val OUTBOX_MESSAGE_TTL_MS = AppConstants.Router.OUTBOX_MESSAGE_TTL_MS + private const val OUTBOX_MAX_PER_PEER = AppConstants.Router.OUTBOX_MAX_PER_PEER + private val HANDSHAKE_RETRY_BACKOFF_MS = AppConstants.Router.HANDSHAKE_RETRY_BACKOFF_MS + @Volatile private var INSTANCE: MessageRouter? = null + internal var disableSchedulerForTesting = false fun tryGetInstance(): MessageRouter? = INSTANCE fun getInstance(context: Context, mesh: MeshService): MessageRouter { val instance = INSTANCE ?: synchronized(this) { @@ -44,10 +71,30 @@ class MessageRouter private constructor( instance.nostr.senderPeerID = mesh.myPeerID return instance } + + internal fun resetForTesting() { + INSTANCE?.schedulerScope?.cancel() + INSTANCE = null + } } - // Outbox: peerID -> queued (content, nickname, messageID) - private val outbox = mutableMapOf>>() + // Outbox: conversationID -> queued messages, oldest first + private val outbox = ConcurrentHashMap>() + + // Per-conversation handshake retry state for queued messages + private val retryState = ConcurrentHashMap() + + private val schedulerScope = CoroutineScope(Dispatchers.Default + SupervisorJob()) + + // Injectable clock for tests + internal var clock: () -> Long = { System.currentTimeMillis() } + + // Called with the messageID of queued messages that expired or were evicted + var onMessageExpired: ((String) -> Unit)? = null + + init { + if (!disableSchedulerForTesting) startOutboxScheduler() + } // Listener for favorites changes to flush outbox when npub mapping appears/changes private val favoriteListener = object: com.bitchat.android.favorites.FavoritesChangeListener { @@ -88,10 +135,9 @@ class MessageRouter private constructor( return RouteResult.NOSTR } else { Log.d(TAG, "Queued PM for ${conversationID} (no mesh, no Nostr mapping) msg_id=${messageID.take(8)}…") - val q = outbox.getOrPut(conversationID) { mutableListOf() } - q.add(Triple(content, recipientNickname, messageID)) + enqueue(conversationID, QueuedMessage(content, recipientNickname, messageID, clock())) Log.d(TAG, "Initiating noise handshake after queueing PM for ${conversationID.take(16)}…") - if (hasMesh) meshTarget?.let { mesh.initiateNoiseHandshake(it) } + if (hasMesh) meshTarget?.let { kickHandshake(conversationID, it, immediate = true) } return RouteResult.QUEUED } } @@ -145,23 +191,27 @@ class MessageRouter private constructor( val queued = outbox[conversationID] ?: outbox[peerID] ?: return if (queued.isEmpty()) return Log.d(TAG, "Flushing outbox for ${conversationID.take(16)}… count=${queued.size}") - val iterator = queued.iterator() - while (iterator.hasNext()) { - val (content, nickname, messageID) = iterator.next() - val resolution = ContactDirectory.resolve(conversationID) - val meshTarget = resolution.meshPeerID - val nostrTarget = resolution.noiseKeyHex ?: conversationID - if (meshTarget != null && isReady(mesh, meshTarget)) { - mesh.sendPrivateMessage(content, meshTarget, nickname, messageID) - iterator.remove() - } else if (canSendViaNostr(nostrTarget)) { - nostr.sendPrivateMessage(content, nostrTarget, nickname, messageID) - iterator.remove() + synchronized(queued) { + val iterator = queued.iterator() + while (iterator.hasNext()) { + val entry = iterator.next() + val resolution = ContactDirectory.resolve(conversationID) + val meshTarget = resolution.meshPeerID + val nostrTarget = resolution.noiseKeyHex ?: conversationID + if (meshTarget != null && isReady(mesh, meshTarget)) { + mesh.sendPrivateMessage(entry.content, meshTarget, entry.nickname, entry.messageID) + iterator.remove() + } else if (canSendViaNostr(nostrTarget)) { + nostr.sendPrivateMessage(entry.content, nostrTarget, entry.nickname, entry.messageID) + iterator.remove() + } } } if (queued.isEmpty()) { - outbox.remove(conversationID) - outbox.remove(peerID) + outbox.remove(conversationID, queued) + outbox.remove(peerID, queued) + retryState.remove(conversationID) + retryState.remove(peerID) } } @@ -170,6 +220,102 @@ class MessageRouter private constructor( outbox.keys.toList().forEach { flushOutboxFor(it) } } + @Synchronized + private fun enqueue(conversationID: String, entry: QueuedMessage) { + val queue = outbox.getOrPut(conversationID) { mutableListOf() } + queue.add(entry) + while (queue.size > OUTBOX_MAX_PER_PEER) { + val evicted = queue.removeAt(0) + Log.w(TAG, "Outbox full for ${conversationID.take(16)}…; evicting oldest msg_id=${evicted.messageID.take(8)}…") + notifyExpired(evicted.messageID) + } + } + + private fun notifyExpired(messageID: String) { + try { onMessageExpired?.invoke(messageID) } catch (_: Exception) { } + } + + /** + * Initiate a Noise handshake for a conversation with queued messages, applying + * exponential backoff between attempts. [immediate] resets the backoff (peer just + * appeared or a new message was queued). Kicks are suppressed while a previous + * attempt is still inside its backoff window, so alias duplicates and frequent + * peer-list updates cannot spam handshakes. + */ + @Synchronized + private fun kickHandshake(conversationID: String, meshTarget: String, immediate: Boolean) { + val now = clock() + val current = retryState[conversationID] + if (current != null && now < current.nextHandshakeAttemptAtMs) return + val attempts = if (immediate) 0 else (current?.handshakeAttempts ?: 0) + try { mesh.initiateNoiseHandshake(meshTarget) } catch (_: Exception) { } + val backoff = HANDSHAKE_RETRY_BACKOFF_MS[attempts.coerceAtMost(HANDSHAKE_RETRY_BACKOFF_MS.size - 1)] + retryState[conversationID] = ConversationRetry( + handshakeAttempts = attempts + 1, + nextHandshakeAttemptAtMs = now + backoff + ) + Log.d(TAG, "Handshake attempt ${attempts + 1} for ${conversationID.take(16)}…, next retry in ${backoff}ms") + } + + private fun startOutboxScheduler() { + schedulerScope.launch { + while (isActive) { + delay(OUTBOX_TICK_MS) + try { tickOutbox() } catch (e: Exception) { + Log.w(TAG, "Outbox scheduler tick failed: ${e.message}") + } + } + } + } + + /** + * One scheduler pass over the outbox: expire old entries, flush what can be sent, + * and re-initiate handshakes (with backoff) for peers that are connected but have + * no established session yet. + */ + internal fun tickOutbox(nowMs: Long = clock()) { + outbox.keys.toList().forEach { conversationID -> + expireOldEntries(conversationID, nowMs) + val queued = outbox[conversationID] ?: return@forEach + if (queued.isEmpty()) return@forEach + + val resolution = ContactDirectory.resolve(conversationID) + val meshTarget = resolution.meshPeerID + + if (meshTarget != null && isReady(mesh, meshTarget)) { + flushOutboxFor(conversationID) + return@forEach + } + if (canSendViaNostr(resolution.noiseKeyHex ?: conversationID)) { + flushOutboxFor(conversationID) + return@forEach + } + // Peer visible but no session: retry the handshake with backoff. + if (meshTarget != null && isConnected(mesh, meshTarget)) { + kickHandshake(conversationID, meshTarget, immediate = false) + } + } + } + + private fun expireOldEntries(conversationID: String, nowMs: Long) { + val queued = outbox[conversationID] ?: return + synchronized(queued) { + val iterator = queued.iterator() + while (iterator.hasNext()) { + val entry = iterator.next() + if (nowMs - entry.enqueuedAtMs > OUTBOX_MESSAGE_TTL_MS) { + Log.w(TAG, "Expiring queued PM for ${conversationID.take(16)}… msg_id=${entry.messageID.take(8)}…") + iterator.remove() + notifyExpired(entry.messageID) + } + } + } + if (queued.isEmpty()) { + outbox.remove(conversationID, queued) + retryState.remove(conversationID) + } + } + private fun canSendViaNostr(peerID: String): Boolean { return try { val resolution = ContactDirectory.resolve(peerID) @@ -208,20 +354,50 @@ class MessageRouter private constructor( // Called when mesh peer list changes; attempt to flush any matching outbox entries fun onPeersUpdated(peers: List) { peers.forEach { pid -> + kickHandshakeIfPending(pid) flushOutboxFor(pid) val noiseHex = try { mesh.getPeerInfo(pid)?.noisePublicKey?.let { ContactIdentityResolver.noiseKeyHex(it) } } catch (_: Exception) { null } - noiseHex?.let { flushOutboxFor(it) } + noiseHex?.let { + kickHandshakeIfPending(it) + flushOutboxFor(it) + } } } // Called when a Noise session becomes established; flush both the mesh peerID and its noiseHex alias fun onSessionEstablished(peerID: String) { + resetRetry(peerID) flushOutboxFor(peerID) val noiseHex = try { mesh.getPeerInfo(peerID)?.noisePublicKey?.let { ContactIdentityResolver.noiseKeyHex(it) } } catch (_: Exception) { null } - noiseHex?.let { flushOutboxFor(it) } + noiseHex?.let { + resetRetry(it) + flushOutboxFor(it) + } + } + + /** Reset handshake backoff for a conversation whose session just came up. */ + private fun resetRetry(peerID: String) { + retryState.remove(ContactDirectory.canonicalConversationId(peerID)) + retryState.remove(peerID) + } + + /** + * A peer (re)appeared: if we still owe them queued messages and there is no working + * session yet, restart the handshake immediately instead of waiting for the backoff. + */ + private fun kickHandshakeIfPending(peerID: String) { + val conversationID = ContactDirectory.canonicalConversationId(peerID) + val queued = outbox[conversationID] ?: outbox[peerID] ?: return + if (queued.isEmpty()) return + val resolution = ContactDirectory.resolve(conversationID) + val meshTarget = resolution.meshPeerID ?: return + if (isReady(mesh, meshTarget)) return + if (!isConnected(mesh, meshTarget)) return + Log.d(TAG, "Peer ${meshTarget.take(8)}… reappeared with ${queued.size} queued PM(s); re-initiating handshake") + kickHandshake(conversationID, meshTarget, immediate = true) } } diff --git a/app/src/main/java/com/bitchat/android/ui/ChatViewModel.kt b/app/src/main/java/com/bitchat/android/ui/ChatViewModel.kt index 1f2d7980..6b9fb236 100644 --- a/app/src/main/java/com/bitchat/android/ui/ChatViewModel.kt +++ b/app/src/main/java/com/bitchat/android/ui/ChatViewModel.kt @@ -232,6 +232,15 @@ class ChatViewModel( loadAndInitialize() ContactDirectory.initialize(getApplication()) { mesh } com.bitchat.android.services.AppStateStore.canonicalizePrivateChats() + // Mark queued private messages as failed when the router gives up on them + try { + com.bitchat.android.services.MessageRouter.getInstance(getApplication(), mesh).onMessageExpired = { messageID -> + messageManager.updateMessageDeliveryStatus( + messageID, + com.bitchat.android.model.DeliveryStatus.Failed("Message expired before delivery") + ) + } + } catch (_: Exception) { } // Hydrate UI state from process-wide AppStateStore to survive Activity recreation viewModelScope.launch { try { com.bitchat.android.services.AppStateStore.peers.collect { peers -> diff --git a/app/src/main/java/com/bitchat/android/util/AppConstants.kt b/app/src/main/java/com/bitchat/android/util/AppConstants.kt index 4dec775f..7cc4c923 100644 --- a/app/src/main/java/com/bitchat/android/util/AppConstants.kt +++ b/app/src/main/java/com/bitchat/android/util/AppConstants.kt @@ -134,6 +134,13 @@ object AppConstants { const val MAX_FILE_SIZE_BYTES: Long = 50L * 1024 * 1024 } + object Router { + const val OUTBOX_TICK_MS: Long = 2_000L + const val OUTBOX_MESSAGE_TTL_MS: Long = 86_400_000L // 24 hours + const val OUTBOX_MAX_PER_PEER: Int = 100 + val HANDSHAKE_RETRY_BACKOFF_MS: LongArray = longArrayOf(5_000L, 15_000L, 30_000L, 60_000L) + } + object Services { const val SEEN_MESSAGE_MAX_IDS: Int = 10_000 } diff --git a/app/src/test/kotlin/com/bitchat/android/services/MessageRouterTest.kt b/app/src/test/kotlin/com/bitchat/android/services/MessageRouterTest.kt new file mode 100644 index 00000000..4ec89fc4 --- /dev/null +++ b/app/src/test/kotlin/com/bitchat/android/services/MessageRouterTest.kt @@ -0,0 +1,200 @@ +package com.bitchat.android.services + +import android.content.Context +import android.os.Build +import com.bitchat.android.identity.SecureIdentityStateManager +import com.bitchat.android.mesh.MeshService +import com.bitchat.android.mesh.PeerInfo +import org.junit.After +import org.junit.Assert.assertEquals +import org.junit.Before +import org.junit.Test +import org.junit.runner.RunWith +import org.mockito.kotlin.any +import org.mockito.kotlin.anyOrNull +import org.mockito.kotlin.clearInvocations +import org.mockito.kotlin.eq +import org.mockito.kotlin.mock +import org.mockito.kotlin.never +import org.mockito.kotlin.times +import org.mockito.kotlin.verify +import org.mockito.kotlin.whenever +import org.robolectric.RobolectricTestRunner +import org.robolectric.RuntimeEnvironment +import org.robolectric.annotation.Config +import java.util.UUID + +@RunWith(RobolectricTestRunner::class) +@Config(sdk = [Build.VERSION_CODES.P], manifest = Config.NONE) +class MessageRouterTest { + + private val myPeerID = "1111222233334444" + private val peerID = "aaaabbbbccccdddd" + private val noiseKey = ByteArray(32) { 0x0B } + + private lateinit var mesh: MeshService + private lateinit var router: MessageRouter + private var fakeTime = 1_000_000L + private val expired = mutableListOf() + + @Before + fun setup() { + val context = RuntimeEnvironment.getApplication() + val prefs = context.getSharedPreferences( + "message-router-test-${UUID.randomUUID()}", + Context.MODE_PRIVATE + ) + val identityManager = SecureIdentityStateManager(prefs, testOnly = true) + ContactDirectory.identityManagerProvider = { identityManager } + + mesh = mock() + whenever(mesh.myPeerID).thenReturn(myPeerID) + whenever(mesh.getPeerNicknames()).thenReturn(mapOf(peerID to "peer")) + + ContactDirectory.initialize(context) { mesh } + + MessageRouter.disableSchedulerForTesting = true + MessageRouter.resetForTesting() + fakeTime = 1_000_000L + expired.clear() + + router = MessageRouter.getInstance(context, mesh) + router.clock = { fakeTime } + router.onMessageExpired = { expired.add(it) } + } + + @After + fun tearDown() { + MessageRouter.resetForTesting() + MessageRouter.disableSchedulerForTesting = false + ContactDirectory.identityManagerProvider = { SecureIdentityStateManager(it) } + } + + @Test + fun `queued message flushes after peer returns and session establishes`() { + peerOffline() + val result = router.sendPrivate("hello", peerID, "peer", "msg-1") + + assertEquals(MessageRouter.RouteResult.QUEUED, result) + verify(mesh, never()).sendPrivateMessage(any(), any(), any(), anyOrNull()) + verify(mesh, never()).initiateNoiseHandshake(any()) + + // Peer reappears without a session: handshake kicked immediately + peerConnectedNoSession() + router.onPeersUpdated(listOf(peerID)) + verify(mesh, times(1)).initiateNoiseHandshake(peerID) + verify(mesh, never()).sendPrivateMessage(any(), any(), any(), anyOrNull()) + + // Session established: queued message is sent + peerReady() + router.onSessionEstablished(peerID) + verify(mesh, times(1)).sendPrivateMessage("hello", peerID, "peer", "msg-1") + } + + @Test + fun `scheduler retries handshake with capped backoff`() { + peerConnectedNoSession() + val result = router.sendPrivate("hello", peerID, "peer", "msg-1") + assertEquals(MessageRouter.RouteResult.QUEUED, result) + verify(mesh, times(1)).initiateNoiseHandshake(peerID) // immediate kick at enqueue + clearInvocations(mesh) + + router.tickOutbox() // backoff (5s) not yet elapsed + verify(mesh, never()).initiateNoiseHandshake(any()) + + fakeTime += 6_000 + router.tickOutbox() // attempt 2, next in 15s + verify(mesh, times(1)).initiateNoiseHandshake(peerID) + + fakeTime += 7_000 + router.tickOutbox() // too early + verify(mesh, times(1)).initiateNoiseHandshake(peerID) + + fakeTime += 9_000 + router.tickOutbox() // attempt 3, next in 30s + verify(mesh, times(2)).initiateNoiseHandshake(peerID) + + fakeTime += 31_000 + router.tickOutbox() // attempt 4, next in 60s + verify(mesh, times(3)).initiateNoiseHandshake(peerID) + + fakeTime += 61_000 + router.tickOutbox() // attempt 5, capped at 60s + verify(mesh, times(4)).initiateNoiseHandshake(peerID) + } + + @Test + fun `expired entries are dropped and reported`() { + peerOffline() + router.sendPrivate("old message", peerID, "peer", "msg-old") + + fakeTime += 86_400_001L + router.tickOutbox() + + assertEquals(listOf("msg-old"), expired) + + // Nothing left to flush even when the peer becomes reachable + peerReady() + router.tickOutbox() + verify(mesh, never()).sendPrivateMessage(any(), any(), any(), anyOrNull()) + } + + @Test + fun `outbox cap evicts oldest and preserves order`() { + peerOffline() + repeat(101) { i -> + router.sendPrivate("content-$i", peerID, "peer", "msg-$i") + } + + assertEquals(listOf("msg-0"), expired) + + peerReady() + router.onSessionEstablished(peerID) + verify(mesh, times(100)).sendPrivateMessage(any(), eq(peerID), any(), any()) + verify(mesh, times(1)).sendPrivateMessage("content-1", peerID, "peer", "msg-1") + verify(mesh, times(1)).sendPrivateMessage("content-100", peerID, "peer", "msg-100") + verify(mesh, never()).sendPrivateMessage(eq("content-0"), any(), any(), anyOrNull()) + } + + @Test + fun `peer reappearance without pending messages does not kick handshake`() { + peerConnectedNoSession() + router.onPeersUpdated(listOf(peerID)) + verify(mesh, never()).initiateNoiseHandshake(any()) + } + + @Test + fun `established session flushes directly without handshake retry state`() { + peerReady() + val result = router.sendPrivate("direct", peerID, "peer", "msg-direct") + assertEquals(MessageRouter.RouteResult.MESH, result) + verify(mesh, times(1)).sendPrivateMessage("direct", peerID, "peer", "msg-direct") + verify(mesh, never()).initiateNoiseHandshake(any()) + } + + private fun peerOffline() { + whenever(mesh.getPeerInfo(peerID)).thenReturn(peerInfo(isConnected = false)) + whenever(mesh.hasEstablishedSession(peerID)).thenReturn(false) + } + + private fun peerConnectedNoSession() { + whenever(mesh.getPeerInfo(peerID)).thenReturn(peerInfo(isConnected = true)) + whenever(mesh.hasEstablishedSession(peerID)).thenReturn(false) + } + + private fun peerReady() { + whenever(mesh.getPeerInfo(peerID)).thenReturn(peerInfo(isConnected = true)) + whenever(mesh.hasEstablishedSession(peerID)).thenReturn(true) + } + + private fun peerInfo(isConnected: Boolean) = PeerInfo( + id = peerID, + nickname = "peer", + isConnected = isConnected, + isDirectConnection = true, + noisePublicKey = noiseKey, + signingPublicKey = ByteArray(32) { 0x0A }, + isVerifiedNickname = false, + lastSeen = System.currentTimeMillis() + ) +} From 3c150d91e0a455cc416b2c1f52924ce85400c656 Mon Sep 17 00:00:00 2001 From: callebtc <93376500+callebtc@users.noreply.github.com> Date: Mon, 27 Jul 2026 22:48:10 +0200 Subject: [PATCH 2/2] fix: unify outbox locking and tie retry scheduler to service lifecycle --- .../android/service/MeshForegroundService.kt | 1 + .../bitchat/android/services/MessageRouter.kt | 73 ++++++++++++------- .../android/services/MessageRouterTest.kt | 17 +++++ 3 files changed, 64 insertions(+), 27 deletions(-) diff --git a/app/src/main/java/com/bitchat/android/service/MeshForegroundService.kt b/app/src/main/java/com/bitchat/android/service/MeshForegroundService.kt index d9ad2d7f..2d84b16f 100644 --- a/app/src/main/java/com/bitchat/android/service/MeshForegroundService.kt +++ b/app/src/main/java/com/bitchat/android/service/MeshForegroundService.kt @@ -144,6 +144,7 @@ class MeshForegroundService : Service() { when (intent?.action) { ACTION_STOP -> { // Stop FGS and mesh cleanly + try { com.bitchat.android.services.MessageRouter.tryGetInstance()?.stopOutboxScheduler() } catch (_: Exception) { } try { unifiedMeshService?.stopServices() ?: meshService?.stopServices() } catch (_: Exception) { } try { MeshServiceHolder.clear() } catch (_: Exception) { } try { stopForeground(true) } catch (_: Exception) { } diff --git a/app/src/main/java/com/bitchat/android/services/MessageRouter.kt b/app/src/main/java/com/bitchat/android/services/MessageRouter.kt index 7905aee8..a09fa31e 100644 --- a/app/src/main/java/com/bitchat/android/services/MessageRouter.kt +++ b/app/src/main/java/com/bitchat/android/services/MessageRouter.kt @@ -66,9 +66,11 @@ class MessageRouter private constructor( } } } - // Always update mesh reference and sync peer ID + // Always update mesh reference and sync peer ID, and make sure the retry + // scheduler is running (it is stopped together with MeshForegroundService). instance.mesh = mesh instance.nostr.senderPeerID = mesh.myPeerID + instance.startOutboxScheduler() return instance } @@ -85,6 +87,7 @@ class MessageRouter private constructor( private val retryState = ConcurrentHashMap() private val schedulerScope = CoroutineScope(Dispatchers.Default + SupervisorJob()) + private var schedulerJob: kotlinx.coroutines.Job? = null // Injectable clock for tests internal var clock: () -> Long = { System.currentTimeMillis() } @@ -93,7 +96,7 @@ class MessageRouter private constructor( var onMessageExpired: ((String) -> Unit)? = null init { - if (!disableSchedulerForTesting) startOutboxScheduler() + startOutboxScheduler() } // Listener for favorites changes to flush outbox when npub mapping appears/changes @@ -185,26 +188,27 @@ class MessageRouter private constructor( } } - // Flush any queued messages for a specific peerID + // Flush any queued messages for a specific peerID. + // All outbox mutations happen under the router monitor so a concurrent enqueue cannot + // be lost between the empty check and the map removal. + @Synchronized fun flushOutboxFor(peerID: String) { val conversationID = ContactDirectory.canonicalConversationId(peerID) val queued = outbox[conversationID] ?: outbox[peerID] ?: return if (queued.isEmpty()) return Log.d(TAG, "Flushing outbox for ${conversationID.take(16)}… count=${queued.size}") - synchronized(queued) { - val iterator = queued.iterator() - while (iterator.hasNext()) { - val entry = iterator.next() - val resolution = ContactDirectory.resolve(conversationID) - val meshTarget = resolution.meshPeerID - val nostrTarget = resolution.noiseKeyHex ?: conversationID - if (meshTarget != null && isReady(mesh, meshTarget)) { - mesh.sendPrivateMessage(entry.content, meshTarget, entry.nickname, entry.messageID) - iterator.remove() - } else if (canSendViaNostr(nostrTarget)) { - nostr.sendPrivateMessage(entry.content, nostrTarget, entry.nickname, entry.messageID) - iterator.remove() - } + val iterator = queued.iterator() + while (iterator.hasNext()) { + val entry = iterator.next() + val resolution = ContactDirectory.resolve(conversationID) + val meshTarget = resolution.meshPeerID + val nostrTarget = resolution.noiseKeyHex ?: conversationID + if (meshTarget != null && isReady(mesh, meshTarget)) { + mesh.sendPrivateMessage(entry.content, meshTarget, entry.nickname, entry.messageID) + iterator.remove() + } else if (canSendViaNostr(nostrTarget)) { + nostr.sendPrivateMessage(entry.content, nostrTarget, entry.nickname, entry.messageID) + iterator.remove() } } if (queued.isEmpty()) { @@ -257,8 +261,11 @@ class MessageRouter private constructor( Log.d(TAG, "Handshake attempt ${attempts + 1} for ${conversationID.take(16)}…, next retry in ${backoff}ms") } + @Synchronized private fun startOutboxScheduler() { - schedulerScope.launch { + if (disableSchedulerForTesting) return + if (schedulerJob?.isActive == true) return + schedulerJob = schedulerScope.launch { while (isActive) { delay(OUTBOX_TICK_MS) try { tickOutbox() } catch (e: Exception) { @@ -268,11 +275,24 @@ class MessageRouter private constructor( } } + /** + * Stop retrying while the mesh transports are down. Persistent network work must + * follow the MeshForegroundService lifecycle; getInstance restarts the scheduler + * and rebinds the mesh reference when the service comes back. + */ + fun stopOutboxScheduler() { + schedulerJob?.cancel() + schedulerJob = null + } + + internal val isSchedulerRunning: Boolean get() = schedulerJob?.isActive == true + /** * One scheduler pass over the outbox: expire old entries, flush what can be sent, * and re-initiate handshakes (with backoff) for peers that are connected but have * no established session yet. */ + @Synchronized internal fun tickOutbox(nowMs: Long = clock()) { outbox.keys.toList().forEach { conversationID -> expireOldEntries(conversationID, nowMs) @@ -299,15 +319,13 @@ class MessageRouter private constructor( private fun expireOldEntries(conversationID: String, nowMs: Long) { val queued = outbox[conversationID] ?: return - synchronized(queued) { - val iterator = queued.iterator() - while (iterator.hasNext()) { - val entry = iterator.next() - if (nowMs - entry.enqueuedAtMs > OUTBOX_MESSAGE_TTL_MS) { - Log.w(TAG, "Expiring queued PM for ${conversationID.take(16)}… msg_id=${entry.messageID.take(8)}…") - iterator.remove() - notifyExpired(entry.messageID) - } + val iterator = queued.iterator() + while (iterator.hasNext()) { + val entry = iterator.next() + if (nowMs - entry.enqueuedAtMs > OUTBOX_MESSAGE_TTL_MS) { + Log.w(TAG, "Expiring queued PM for ${conversationID.take(16)}… msg_id=${entry.messageID.take(8)}…") + iterator.remove() + notifyExpired(entry.messageID) } } if (queued.isEmpty()) { @@ -389,6 +407,7 @@ class MessageRouter private constructor( * A peer (re)appeared: if we still owe them queued messages and there is no working * session yet, restart the handshake immediately instead of waiting for the backoff. */ + @Synchronized private fun kickHandshakeIfPending(peerID: String) { val conversationID = ContactDirectory.canonicalConversationId(peerID) val queued = outbox[conversationID] ?: outbox[peerID] ?: return diff --git a/app/src/test/kotlin/com/bitchat/android/services/MessageRouterTest.kt b/app/src/test/kotlin/com/bitchat/android/services/MessageRouterTest.kt index 4ec89fc4..48cbe2da 100644 --- a/app/src/test/kotlin/com/bitchat/android/services/MessageRouterTest.kt +++ b/app/src/test/kotlin/com/bitchat/android/services/MessageRouterTest.kt @@ -7,6 +7,8 @@ import com.bitchat.android.mesh.MeshService import com.bitchat.android.mesh.PeerInfo import org.junit.After import org.junit.Assert.assertEquals +import org.junit.Assert.assertFalse +import org.junit.Assert.assertTrue import org.junit.Before import org.junit.Test import org.junit.runner.RunWith @@ -172,6 +174,21 @@ class MessageRouterTest { verify(mesh, never()).initiateNoiseHandshake(any()) } + @Test + fun `scheduler stops with the mesh service and restarts on rebind`() { + MessageRouter.disableSchedulerForTesting = false + MessageRouter.resetForTesting() + val context = RuntimeEnvironment.getApplication() + val running = MessageRouter.getInstance(context, mesh) + assertTrue(running.isSchedulerRunning) + + running.stopOutboxScheduler() + assertFalse(running.isSchedulerRunning) + + val rebound = MessageRouter.getInstance(context, mesh) + assertTrue(rebound.isSchedulerRunning) + } + private fun peerOffline() { whenever(mesh.getPeerInfo(peerID)).thenReturn(peerInfo(isConnected = false)) whenever(mesh.hasEstablishedSession(peerID)).thenReturn(false)