Ask every transport's registry before the archive guard drops a sender

The archive guard in the first commit checked only the Bluetooth peer
registry. The Wi-Fi Aware transport feeds the same archive from a
registry of its own. With Wi-Fi Aware on, a sender known only there was
never archived for sync, and with Bluetooth switched off nothing was.
Codex flagged this on the PR.

The holder now keeps one liveness probe per transport, registered the
way it already tracks gossip owners, and the shared manager's delegate
asks all of them. The probes live in the holder because the Bluetooth
service has no reference to the Wi-Fi service.

Three tests pin the holder, including a Wi-Fi-only sender's broadcast
being archived and served. The Wi-Fi Aware register and unregister
lines have no unit test, since that service has none.
This commit is contained in:
heyaim 2026-09-03 00:23:39 -05:00
parent b9bf303d9a
commit e1da7f8960
4 changed files with 153 additions and 9 deletions

View File

@ -171,10 +171,10 @@ class BluetoothMeshService(private val context: Context) : TransportBridgeServic
}
)
com.bitchat.android.service.MeshServiceHolder.setGossipManager(
gossipSyncManager,
hasLivePeer = { peerID -> peerManager.getPeerInfo(peerID) != null }
) { packet ->
com.bitchat.android.service.MeshServiceHolder.registerLivenessProbe("BLE") { peerID ->
peerManager.getPeerInfo(peerID) != null
}
com.bitchat.android.service.MeshServiceHolder.setGossipManager(gossipSyncManager) { packet ->
signPacketBeforeBroadcast(packet)
}
if (isBleTransportEnabled()) {

View File

@ -6,6 +6,7 @@ import com.bitchat.android.mesh.UnifiedMeshService
import com.bitchat.android.model.RoutedPacket
import com.bitchat.android.protocol.BitchatPacket
import com.bitchat.android.sync.GossipSyncManager
import java.util.concurrent.ConcurrentHashMap
/**
* Process-wide holder to share a single BluetoothMeshService instance
@ -19,10 +20,17 @@ object MeshServiceHolder {
private val activeGossipOwners = mutableSetOf<String>()
// One liveness probe per transport, keyed like the gossip owners above. The shared gossip
// manager archives a broadcast only when its sender is in a live peer registry, and every
// transport keeps its own registry: Bluetooth registers its lookup, Wi-Fi Aware registers
// its own while it runs. A sender known to any transport is live. The Bluetooth service
// registers before any delegate exists, so the empty map is never consulted in practice;
// if it were, the answer is true, which archives as the manager did before the guard.
private val livenessProbes = ConcurrentHashMap<String, (String) -> Boolean>()
@Synchronized
fun setGossipManager(
mgr: GossipSyncManager,
hasLivePeer: (String) -> Boolean = { true },
signer: (BitchatPacket) -> BitchatPacket
) {
val previous = sharedGossipSyncManager
@ -30,12 +38,28 @@ object MeshServiceHolder {
try { previous?.stop() } catch (_: Exception) { }
}
sharedGossipSyncManager = mgr
mgr.delegate = TransportGossipDelegate(signer, hasLivePeer)
mgr.delegate = TransportGossipDelegate(signer)
if (activeGossipOwners.isNotEmpty()) {
mgr.start()
}
}
fun registerLivenessProbe(owner: String, probe: (String) -> Boolean) {
livenessProbes[owner] = probe
}
fun unregisterLivenessProbe(owner: String) {
livenessProbes.remove(owner)
}
/** True when any transport's live peer registry holds this peer. */
fun hasLivePeer(peerID: String): Boolean {
if (livenessProbes.isEmpty()) return true
return livenessProbes.values.any { probe ->
try { probe(peerID) } catch (_: Exception) { false }
}
}
@Synchronized
fun startSharedGossip(owner: String) {
val wasIdle = activeGossipOwners.isEmpty()
@ -54,10 +78,9 @@ object MeshServiceHolder {
}
private class TransportGossipDelegate(
private val signer: (BitchatPacket) -> BitchatPacket,
private val livePeer: (String) -> Boolean
private val signer: (BitchatPacket) -> BitchatPacket
) : GossipSyncManager.Delegate {
override fun hasLivePeer(peerID: String): Boolean = livePeer(peerID)
override fun hasLivePeer(peerID: String): Boolean = MeshServiceHolder.hasLivePeer(peerID)
override fun sendPacket(packet: BitchatPacket) {
TransportBridgeService.broadcastFromLocal(RoutedPacket(packet))
@ -144,6 +167,7 @@ object MeshServiceHolder {
try { sharedGossipSyncManager?.stop() } catch (_: Exception) { }
sharedGossipSyncManager = null
activeGossipOwners.clear()
livenessProbes.clear()
meshService = null
unifiedMeshService = null
}

View File

@ -515,6 +515,9 @@ class WifiAwareMeshService(private val context: Context) : MeshService, Transpor
TransportBridgeService.register("WIFI", this)
meshCore.startCore()
com.bitchat.android.service.MeshServiceHolder.registerLivenessProbe("WIFI") { peerID ->
meshCore.getPeerInfo(peerID) != null
}
com.bitchat.android.service.MeshServiceHolder.startSharedGossip("WIFI")
startPeriodicConnectionMaintenance()
connectionTracker.start()
@ -532,6 +535,7 @@ class WifiAwareMeshService(private val context: Context) : MeshService, Transpor
// Unregister from bridge
TransportBridgeService.unregister("WIFI")
com.bitchat.android.service.MeshServiceHolder.stopSharedGossip("WIFI")
com.bitchat.android.service.MeshServiceHolder.unregisterLivenessProbe("WIFI")
try { com.bitchat.android.services.AppStateStore.clearTransportPeers("WIFI") } catch (_: Exception) { }
try { com.bitchat.android.services.AppStateStore.clearTransportDirectPeers("WIFI") } catch (_: Exception) { }
@ -578,6 +582,7 @@ class WifiAwareMeshService(private val context: Context) : MeshService, Transpor
isActive = false
TransportBridgeService.unregister("WIFI")
com.bitchat.android.service.MeshServiceHolder.stopSharedGossip("WIFI")
com.bitchat.android.service.MeshServiceHolder.unregisterLivenessProbe("WIFI")
try { com.bitchat.android.services.AppStateStore.clearTransportPeers("WIFI") } catch (_: Exception) { }
try { com.bitchat.android.services.AppStateStore.clearTransportDirectPeers("WIFI") } catch (_: Exception) { }
val oldPublishSession = publishSession

View File

@ -0,0 +1,115 @@
package com.bitchat.android.service
import com.bitchat.android.model.RequestSyncPacket
import com.bitchat.android.protocol.BitchatPacket
import com.bitchat.android.protocol.MessageType
import com.bitchat.android.protocol.SpecialRecipients
import com.bitchat.android.sync.GossipSyncManager
import kotlinx.coroutines.CoroutineScope
import kotlinx.coroutines.Dispatchers
import kotlinx.coroutines.SupervisorJob
import kotlinx.coroutines.cancel
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
/**
* The shared gossip manager archives a broadcast only when its sender is in a live peer
* registry, and every transport keeps its own registry. A sender known only to the Wi-Fi
* Aware registry must count as live, or its broadcasts would never be archived for sync.
*/
class MeshServiceHolderLivenessTest {
private lateinit var scope: CoroutineScope
private lateinit var manager: GossipSyncManager
private val blePeer = "aaaaaaaaaaaaaaaa"
private val wifiSender = ByteArray(8) { 0x22 }
private val wifiPeer = wifiSender.joinToString("") { "%02x".format(it) }
private val stranger = "cccccccccccccccc"
private val requester = "aabbccddeeff0011"
private val config = object : GossipSyncManager.ConfigProvider {
override fun seenCapacity(): Int = 100
override fun gcsMaxBytes(): Int = 400
override fun gcsTargetFpr(): Double = 0.01
}
@Before
fun setUp() {
scope = CoroutineScope(SupervisorJob() + Dispatchers.Unconfined)
manager = GossipSyncManager(myPeerID = "1122334455667788", scope = scope, configProvider = config)
MeshServiceHolder.unregisterLivenessProbe("WIFI")
MeshServiceHolder.registerLivenessProbe("BLE") { it == blePeer }
MeshServiceHolder.setGossipManager(manager) { it }
}
@After
fun tearDown() {
MeshServiceHolder.unregisterLivenessProbe("WIFI")
MeshServiceHolder.unregisterLivenessProbe("BLE")
scope.cancel()
}
private fun broadcastFromWifiSender(): BitchatPacket = BitchatPacket(
version = 1u,
type = MessageType.MESSAGE.value,
senderID = wifiSender,
recipientID = SpecialRecipients.BROADCAST,
timestamp = (System.currentTimeMillis() - 1000L).toULong(),
payload = "over wifi".toByteArray(),
signature = ByteArray(64) { 0x33 },
ttl = 7u
)
/** A filter the requester builds when it holds nothing: everything we have is missing. */
private fun requestForNothingHeld() = RequestSyncPacket(p = 7, m = 1, data = ByteArray(0))
@Test
fun `a peer known only to the Wi-Fi Aware registry is live once its probe is registered`() {
val delegate = manager.delegate!!
assertFalse("before the Wi-Fi probe exists only Bluetooth answers", delegate.hasLivePeer(wifiPeer))
MeshServiceHolder.registerLivenessProbe("WIFI") { it == wifiPeer }
assertTrue(delegate.hasLivePeer(wifiPeer))
assertTrue("the Bluetooth registry still counts", delegate.hasLivePeer(blePeer))
assertFalse("a peer in no registry is absent", delegate.hasLivePeer(stranger))
}
@Test
fun `removing a transport's probe makes its peers absent again`() {
val delegate = manager.delegate!!
MeshServiceHolder.registerLivenessProbe("WIFI") { it == wifiPeer }
assertTrue(delegate.hasLivePeer(wifiPeer))
MeshServiceHolder.unregisterLivenessProbe("WIFI")
assertFalse(delegate.hasLivePeer(wifiPeer))
assertTrue(delegate.hasLivePeer(blePeer))
}
@Test
fun `a broadcast from a Wi-Fi-only sender is archived and served`() {
MeshServiceHolder.registerLivenessProbe("WIFI") { it == wifiPeer }
val original = broadcastFromWifiSender()
manager.onPublicPacketSeen(original)
val sent = mutableListOf<BitchatPacket>()
manager.delegate = object : GossipSyncManager.Delegate {
override fun sendPacket(packet: BitchatPacket) = Unit
override fun sendPacketToPeer(peerID: String, packet: BitchatPacket) {
sent += packet
}
override fun signPacketForBroadcast(packet: BitchatPacket): BitchatPacket = packet
}
manager.handleRequestSync(requester, requestForNothingHeld())
assertEquals("the Wi-Fi-only sender's message must be served", 1, sent.size)
assertTrue(sent[0].payload.contentEquals(original.payload))
}
}