From ea5c5bb0da4ef2b1ce52b23d6e1f44ec8b607c6d Mon Sep 17 00:00:00 2001 From: Taksh Date: Sat, 15 Aug 2026 12:43:21 +0530 Subject: [PATCH 1/3] Serialize access to the per-peer actor map MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit This class exists to give each peer one actor so that peer's packets are processed sequentially — the comment says it prevents a session-management race. The map holding those actors was a plain mutableMapOf, mutated from whichever BLE callback thread delivered the packet. So the guarantee did not hold. Two threads could both miss on getOrPut for the same peer and each create an actor; the second put replaces the first, which keeps running with whatever it was already handed. That is the concurrent session handling the actors were introduced to stop, and it is also a plain HashMap being structurally modified from several threads. Also drops peerActors, declared next to it and referenced nowhere. --- .../bitchat/android/mesh/PacketProcessor.kt | 47 ++++++++++++------- 1 file changed, 31 insertions(+), 16 deletions(-) diff --git a/app/src/main/java/com/bitchat/android/mesh/PacketProcessor.kt b/app/src/main/java/com/bitchat/android/mesh/PacketProcessor.kt index 37a295e8..037aa248 100644 --- a/app/src/main/java/com/bitchat/android/mesh/PacketProcessor.kt +++ b/app/src/main/java/com/bitchat/android/mesh/PacketProcessor.kt @@ -20,6 +20,7 @@ class PacketProcessor(private val myPeerID: String) { companion object { private const val TAG = "PacketProcessor" + } // Delegate for callbacks @@ -40,8 +41,6 @@ class PacketProcessor(private val myPeerID: String) { // Per-peer actors to serialize packet processing // Each peer gets its own actor that processes packets sequentially // This prevents race conditions in session management - private val peerActors = mutableMapOf>() - @OptIn(ObsoleteCoroutinesApi::class) private fun getOrCreateActorForPeer(peerID: String) = processorScope.actor( capacity = Channel.UNLIMITED @@ -51,9 +50,20 @@ class PacketProcessor(private val myPeerID: String) { } } - // Cache actors to reuse them - private val actors = mutableMapOf>() - + // Cache actors to reuse them. + // + // Access is serialized on [actorsLock]. The whole point of this class is + // one actor per peer; a plain map mutated from the BLE callback threads + // let two threads each create one for the same peer, which quietly + // reinstated the concurrent session handling the actors exist to prevent. + // + // Access-ordered so the eldest entry is the least recently *used* rather + // than the least recently created — evicting by insertion order would + // drop a long-lived active peer in favour of a burst of new IDs. + private val actorsLock = Any() + private val actors = + mutableMapOf>() + init { // Set up the packet relay manager delegate immediately setupRelayManager() @@ -72,7 +82,9 @@ class PacketProcessor(private val myPeerID: String) { } // Get or create actor for this peer - val actor = actors.getOrPut(peerID) { getOrCreateActorForPeer(peerID) } + val actor = synchronized(actorsLock) { + actors.getOrPut(peerID) { getOrCreateActorForPeer(peerID) } + } // Send packet to peer's dedicated actor for serialized processing processorScope.launch { @@ -245,12 +257,13 @@ class PacketProcessor(private val myPeerID: String) { return buildString { appendLine("=== Packet Processor Debug Info ===") appendLine("Processor Scope Active: ${processorScope.isActive}") - appendLine("Active Peer Actors: ${actors.size}") + val peerIDs = synchronized(actorsLock) { actors.keys.toList() } + appendLine("Active Peer Actors: ${peerIDs.size}") appendLine("My Peer ID: $myPeerID") - - if (actors.isNotEmpty()) { + + if (peerIDs.isNotEmpty()) { appendLine("Peer Actors:") - actors.keys.forEach { peerID -> + peerIDs.forEach { peerID -> appendLine(" - $peerID") } } @@ -261,13 +274,15 @@ class PacketProcessor(private val myPeerID: String) { * Shutdown the processor and all peer actors */ fun shutdown() { - Log.d(TAG, "Shutting down PacketProcessor and ${actors.size} peer actors") - - // Close all peer actors gracefully - actors.values.forEach { actor -> - actor.close() + synchronized(actorsLock) { + Log.d(TAG, "Shutting down PacketProcessor and ${actors.size} peer actors") + + // Close all peer actors gracefully + actors.values.forEach { actor -> + actor.close() + } + actors.clear() } - actors.clear() // Shutdown the relay manager packetRelayManager.shutdown() From 7e037ee3f3f592a265328506bf86d6bbbb2b9585 Mon Sep 17 00:00:00 2001 From: Taksh Date: Sat, 15 Aug 2026 12:44:19 +0530 Subject: [PATCH 2/3] Bound the per-peer actor map and evict least-recently-used MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Each distinct sender ID allocated a coroutine and an unbounded channel that were released only at shutdown. Peer IDs are ephemeral and rotate, and the ID is taken straight off the wire, so the map grows without limit in ordinary use and can be grown deliberately by anyone in radio range. Cap it at 128 live actors. The map is access-ordered, so the entry dropped is the peer least recently *used*: evicting by insertion order would discard a long-lived active peer in favour of a burst of new IDs, which is the case worth protecting. Eviction closes the channel rather than cancelling it, so the actor drains what it already holds before finishing. Two of the four tests fail on main. The concurrency one does not — that race is timing-dependent — so it is there to hold the fix rather than to demonstrate the bug. --- .../bitchat/android/mesh/PacketProcessor.kt | 39 ++++- .../mesh/PacketProcessorActorBoundsTest.kt | 148 ++++++++++++++++++ 2 files changed, 184 insertions(+), 3 deletions(-) create mode 100644 app/src/test/kotlin/com/bitchat/android/mesh/PacketProcessorActorBoundsTest.kt diff --git a/app/src/main/java/com/bitchat/android/mesh/PacketProcessor.kt b/app/src/main/java/com/bitchat/android/mesh/PacketProcessor.kt index 037aa248..f7740d9e 100644 --- a/app/src/main/java/com/bitchat/android/mesh/PacketProcessor.kt +++ b/app/src/main/java/com/bitchat/android/mesh/PacketProcessor.kt @@ -21,6 +21,17 @@ class PacketProcessor(private val myPeerID: String) { companion object { private const val TAG = "PacketProcessor" + /** + * Live per-peer actors to retain. + * + * Every distinct sender ID used to allocate a coroutine and an + * unbounded channel that were only released at shutdown. Peer IDs are + * ephemeral and rotate, and the ID is read straight off the wire, so + * the set of keys grows without limit during ordinary use and can be + * grown deliberately by anyone in radio range. Bounded here, evicting + * least-recently-used, which is the peer least likely to be mid-session. + */ + internal const val MAX_PEER_ACTORS = 128 } // Delegate for callbacks @@ -62,8 +73,23 @@ class PacketProcessor(private val myPeerID: String) { // drop a long-lived active peer in favour of a burst of new IDs. private val actorsLock = Any() private val actors = - mutableMapOf>() - + object : LinkedHashMap>( + 16, 0.75f, true + ) { + override fun removeEldestEntry( + eldest: MutableMap.MutableEntry> + ): Boolean { + if (size <= MAX_PEER_ACTORS) return false + // Closing lets the actor drain what it already holds and then + // finish, rather than cancelling mid-packet. + eldest.value.close() + // Deliberately not formatPeerForLog: that reaches into the + // delegate, and this runs under actorsLock. + Log.d(TAG, "Evicting least-recently-used peer actor for ${eldest.key}") + return true + } + } + init { // Set up the packet relay manager delegate immediately setupRelayManager() @@ -250,6 +276,13 @@ class PacketProcessor(private val myPeerID: String) { // delegate?.handleDeliveryAck(routed) // } + /** + * Number of live per-peer actors. Test seam; `getDebugInfo` reports the + * same number alongside the cap. + */ + internal val activePeerActorCount: Int + get() = synchronized(actorsLock) { actors.size } + /** * Get debug information */ @@ -258,7 +291,7 @@ class PacketProcessor(private val myPeerID: String) { appendLine("=== Packet Processor Debug Info ===") appendLine("Processor Scope Active: ${processorScope.isActive}") val peerIDs = synchronized(actorsLock) { actors.keys.toList() } - appendLine("Active Peer Actors: ${peerIDs.size}") + appendLine("Active Peer Actors: ${peerIDs.size} (cap $MAX_PEER_ACTORS)") appendLine("My Peer ID: $myPeerID") if (peerIDs.isNotEmpty()) { diff --git a/app/src/test/kotlin/com/bitchat/android/mesh/PacketProcessorActorBoundsTest.kt b/app/src/test/kotlin/com/bitchat/android/mesh/PacketProcessorActorBoundsTest.kt new file mode 100644 index 00000000..bf0d8a82 --- /dev/null +++ b/app/src/test/kotlin/com/bitchat/android/mesh/PacketProcessorActorBoundsTest.kt @@ -0,0 +1,148 @@ +package com.bitchat.android.mesh + +import com.bitchat.android.model.RoutedPacket +import com.bitchat.android.protocol.BitchatPacket +import com.bitchat.android.protocol.MessageType +import com.bitchat.android.protocol.SpecialRecipients +import org.junit.After +import org.junit.Assert.assertEquals +import org.junit.Assert.assertTrue +import org.junit.Test +import java.util.concurrent.CountDownLatch +import java.util.concurrent.Executors +import java.util.concurrent.TimeUnit + +/** + * The per-peer actor map is keyed by an ID read straight off the wire and was + * only ever emptied at shutdown, so it grew for the life of the process. It + * was also a plain map mutated from the BLE callback threads, which meant the + * one-actor-per-peer property this class exists to provide was not actually + * guaranteed. + */ +class PacketProcessorActorBoundsTest { + private val processors = mutableListOf() + + @After + fun tearDown() { + processors.forEach(PacketProcessor::shutdown) + } + + @Test + fun `actor count stays at the cap however many peers are seen`() { + val processor = processor() + + // Peer IDs are ephemeral and rotate, and nothing stops a peer in radio + // range from minting new ones, so "distinct senders" is unbounded. + repeat(PacketProcessor.MAX_PEER_ACTORS + 200) { index -> + processor.processPacket(packetFrom(peerID(index))) + } + + assertEquals(PacketProcessor.MAX_PEER_ACTORS, processor.activePeerActorCount) + } + + @Test + fun `a peer still being used is not evicted by a burst of new ones`() { + val processor = processor() + val busy = peerID(0) + + processor.processPacket(packetFrom(busy)) + repeat(PacketProcessor.MAX_PEER_ACTORS * 2) { index -> + processor.processPacket(packetFrom(peerID(index + 1))) + // Keep touching the long-lived peer. Evicting by insertion order + // would drop it anyway; access order is what saves it. + processor.processPacket(packetFrom(busy)) + } + + assertEquals(PacketProcessor.MAX_PEER_ACTORS, processor.activePeerActorCount) + assertTrue( + "The peer in active use must survive eviction", + processor.getDebugInfo().contains(busy) + ) + } + + @Test + fun `concurrent packets from one peer create exactly one actor`() { + val processor = processor() + val peer = peerID(7) + val threads = 16 + val pool = Executors.newFixedThreadPool(threads) + val start = CountDownLatch(1) + val done = CountDownLatch(threads) + + repeat(threads) { + pool.execute { + start.await() + repeat(50) { processor.processPacket(packetFrom(peer)) } + done.countDown() + } + } + start.countDown() + assertTrue(done.await(10, TimeUnit.SECONDS)) + pool.shutdown() + + // Two actors for one peer would process that peer's packets in + // parallel — exactly the session race the actors were introduced to + // stop. Timing-dependent, so this locks the fix in rather than + // reproducing the failure on demand. + assertEquals(1, processor.activePeerActorCount) + } + + @Test + fun `shutdown releases every actor`() { + val processor = PacketProcessor(MY_PEER_ID).also { it.delegate = NoopDelegate() } + + repeat(10) { processor.processPacket(packetFrom(peerID(it))) } + assertEquals(10, processor.activePeerActorCount) + + processor.shutdown() + + assertEquals(0, processor.activePeerActorCount) + } + + private fun processor(): PacketProcessor = + PacketProcessor(MY_PEER_ID).also { + it.delegate = NoopDelegate() + processors += it + } + + private fun peerID(index: Int): String = String.format("%016x", index + 1) + + private fun packetFrom(peerID: String): RoutedPacket { + val packet = BitchatPacket( + version = 1u, + type = MessageType.MESSAGE.value, + senderID = peerID.hexToBytes(), + recipientID = SpecialRecipients.BROADCAST, + timestamp = 1u, + payload = byteArrayOf(0x01), + ttl = 7u + ) + return RoutedPacket(packet, peerID, "direct-link") + } + + private class NoopDelegate : PacketProcessorDelegate { + override fun validatePacketSecurity(packet: BitchatPacket, peerID: String) = true + override fun updatePeerLastSeen(peerID: String) = Unit + override fun getPeerNickname(peerID: String): String? = null + override fun getNetworkSize() = 1 + override fun getBroadcastRecipient(): ByteArray = SpecialRecipients.BROADCAST + override fun handleNoiseHandshake(routed: RoutedPacket) = true + override fun handleNoiseEncrypted(routed: RoutedPacket) = true + override suspend fun handleAnnounce(routed: RoutedPacket) = true + override fun handleMessage(routed: RoutedPacket) = Unit + override fun handleLeave(routed: RoutedPacket) = Unit + override fun handleFragment(packet: BitchatPacket): BitchatPacket? = null + override fun handleRequestSync(routed: RoutedPacket) = Unit + override fun sendAnnouncementToPeer(peerID: String) = Unit + override fun sendCachedMessages(peerID: String) = Unit + override fun relayPacket(routed: RoutedPacket) = Unit + override fun sendToPeer(peerID: String, routed: RoutedPacket) = false + } + + private fun String.hexToBytes(): ByteArray = + chunked(2).map { it.toInt(16).toByte() }.toByteArray() + + private companion object { + const val MY_PEER_ID = "1111222233334444" + } +}