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 42d40a77..d63ea952 100644 --- a/app/src/main/java/com/bitchat/android/mesh/BluetoothMeshService.kt +++ b/app/src/main/java/com/bitchat/android/mesh/BluetoothMeshService.kt @@ -886,6 +886,7 @@ class BluetoothMeshService(private val context: Context) : TransportBridgeServic */ fun sendMessage(content: String, mentions: List = emptyList(), channel: String? = null) { if (content.isEmpty()) return + val bridgePolicyAtSend = BridgeMeshPort.outboundPolicy() serviceScope.launch { val packet = BitchatPacket( @@ -912,7 +913,8 @@ class BluetoothMeshService(private val context: Context) : TransportBridgeServic content, myPeerID, packet.timestamp.toLong(), - nickname + nickname, + bridgePolicyAtSend ) } } diff --git a/app/src/main/java/com/bitchat/android/mesh/BridgeMeshPort.kt b/app/src/main/java/com/bitchat/android/mesh/BridgeMeshPort.kt index 6793b54a..6f4eada1 100644 --- a/app/src/main/java/com/bitchat/android/mesh/BridgeMeshPort.kt +++ b/app/src/main/java/com/bitchat/android/mesh/BridgeMeshPort.kt @@ -3,6 +3,31 @@ package com.bitchat.android.mesh import com.bitchat.android.model.IdentityAnnouncement import com.bitchat.android.protocol.BitchatPacket +/** + * Immutable privacy decision captured when a public message is accepted. + * + * A later opt-in must not authorize a message that was composed while + * bridging was disabled or nearby-only. Publication therefore requires both + * this send-time decision and the current policy to allow bridging. + */ +class BridgeOutboundPolicy internal constructor( + internal val enabled: Boolean, + internal val nearbyOnly: Boolean +) { + val allowsBridging: Boolean + get() = enabled && !nearbyOnly + + fun permitsPublication(current: BridgeOutboundPolicy): Boolean = + allowsBridging && current.allowsBridging + + companion object { + val Denied = BridgeOutboundPolicy(enabled = false, nearbyOnly = false) + + internal fun capture(enabled: Boolean, nearbyOnly: Boolean): BridgeOutboundPolicy = + BridgeOutboundPolicy(enabled = enabled, nearbyOnly = nearbyOnly) + } +} + /** * Transport-facing bridge boundary. * @@ -12,12 +37,14 @@ import com.bitchat.android.protocol.BitchatPacket */ interface BridgeMeshDelegate { fun advertisedCell(): String? + fun outboundPolicy(): BridgeOutboundPolicy fun bridgeOutgoing( content: String, senderPeerId: String, timestampMs: Long, - nickname: String? + nickname: String?, + policyAtSend: BridgeOutboundPolicy ) fun handleAuthenticatedRadioMessage(messageId: String) @@ -37,13 +64,17 @@ object BridgeMeshPort : BridgeMeshDelegate { override fun advertisedCell(): String? = delegate?.advertisedCell() + override fun outboundPolicy(): BridgeOutboundPolicy = + delegate?.outboundPolicy() ?: BridgeOutboundPolicy.Denied + override fun bridgeOutgoing( content: String, senderPeerId: String, timestampMs: Long, - nickname: String? + nickname: String?, + policyAtSend: BridgeOutboundPolicy ) { - delegate?.bridgeOutgoing(content, senderPeerId, timestampMs, nickname) + delegate?.bridgeOutgoing(content, senderPeerId, timestampMs, nickname, policyAtSend) } override fun handleAuthenticatedRadioMessage(messageId: String) { diff --git a/app/src/main/java/com/bitchat/android/mesh/MeshCore.kt b/app/src/main/java/com/bitchat/android/mesh/MeshCore.kt index c6380e7d..0d9806f8 100644 --- a/app/src/main/java/com/bitchat/android/mesh/MeshCore.kt +++ b/app/src/main/java/com/bitchat/android/mesh/MeshCore.kt @@ -511,6 +511,7 @@ class MeshCore( fun sendMessage(content: String, mentions: List = emptyList(), channel: String? = null) { if (content.isEmpty()) return + val bridgePolicyAtSend = BridgeMeshPort.outboundPolicy() scope.launch { val packet = BitchatPacket( version = 1u, @@ -532,7 +533,8 @@ class MeshCore( content, myPeerID, packet.timestamp.toLong(), - nickname + nickname, + bridgePolicyAtSend ) } } diff --git a/app/src/main/java/com/bitchat/android/services/bridge/BridgeCollections.kt b/app/src/main/java/com/bitchat/android/services/bridge/BridgeCollections.kt index 50a482c7..c0a0610f 100644 --- a/app/src/main/java/com/bitchat/android/services/bridge/BridgeCollections.kt +++ b/app/src/main/java/com/bitchat/android/services/bridge/BridgeCollections.kt @@ -19,6 +19,37 @@ internal class BoundedIdSet(private val capacity: Int) { fun clear() = values.clear() } +/** + * Owns one logical relay subscription while assigning a distinct wire ID to + * every replacement. Relay CLOSE/REQ writes may execute out of order, but a + * delayed close can only affect the retired generation. + * + * Callers keep this object confined to their coordinator dispatcher. + */ +internal class RelaySubscriptionSlot(private val idPrefix: String) { + private var generation = 0L + private var activeId: String? = null + + fun replace(close: (String) -> Unit, open: (String) -> Unit): String { + activeId?.let(close) + val replacementId = "$idPrefix-${++generation}" + activeId = replacementId + try { + open(replacementId) + } catch (error: Throwable) { + activeId = null + throw error + } + return replacementId + } + + fun close(close: (String) -> Unit) { + val id = activeId ?: return + activeId = null + close(id) + } +} + /** * A small insertion-ordered expiring set. Callers own synchronization; bridge * coordinators keep each instance confined to their serial dispatcher. diff --git a/app/src/main/java/com/bitchat/android/services/bridge/CourierCoordinator.kt b/app/src/main/java/com/bitchat/android/services/bridge/CourierCoordinator.kt index f1dbd5e6..153ebaea 100644 --- a/app/src/main/java/com/bitchat/android/services/bridge/CourierCoordinator.kt +++ b/app/src/main/java/com/bitchat/android/services/bridge/CourierCoordinator.kt @@ -57,6 +57,7 @@ internal class CourierCoordinator( private val pendingDrops = mutableListOf() private val signatureAttemptTimes = mutableListOf() private var subscribedTags: Set = emptySet() + private val courierSubscription = RelaySubscriptionSlot("mesh-bridge-courier") @Volatile private var enabled = false private val publishedDropKeys = @@ -209,18 +210,26 @@ internal class CourierCoordinator( .toSet() val allTags = myTags + peerTags if (allTags == subscribedTags) return - relayManager.unsubscribe(COURIER_SUBSCRIPTION) - subscribedTags = allTags - if (allTags.isEmpty()) return - relayManager.subscribe( - filter = com.bitchat.android.nostr.NostrFilter.courierDrops( - allTags, - since = now - CourierEnvelope.MAX_LIFETIME_MS - ), - id = COURIER_SUBSCRIPTION, - handler = { event -> scope.launch { handleDropEvent(event) } }, - targetRelayUrls = NostrRelayManager.defaultRelays() - ) + if (allTags.isEmpty()) { + courierSubscription.close(relayManager::unsubscribe) + subscribedTags = emptySet() + } else { + courierSubscription.replace( + close = relayManager::unsubscribe, + open = { subscriptionId -> + relayManager.subscribe( + filter = com.bitchat.android.nostr.NostrFilter.courierDrops( + allTags, + since = now - CourierEnvelope.MAX_LIFETIME_MS + ), + id = subscriptionId, + handler = { event -> scope.launch { handleDropEvent(event) } }, + targetRelayUrls = NostrRelayManager.defaultRelays() + ) + } + ) + subscribedTags = allTags + } } private fun handleDropEvent(event: NostrEvent) { @@ -365,7 +374,7 @@ internal class CourierCoordinator( } private fun closeSubscription() { - relayManager.unsubscribe(COURIER_SUBSCRIPTION) + courierSubscription.close(relayManager::unsubscribe) subscribedTags = emptySet() } @@ -411,7 +420,6 @@ internal class CourierCoordinator( private fun ByteArray.toHex(): String = joinToString("") { "%02x".format(it) } private companion object { - const val COURIER_SUBSCRIPTION = "mesh-bridge-courier" const val MAX_TRACKED_IDS = 512 const val MAX_WATCHED_PEERS = 16 const val MAX_PENDING_DROPS = 20 diff --git a/app/src/main/java/com/bitchat/android/services/bridge/MeshBridgeService.kt b/app/src/main/java/com/bitchat/android/services/bridge/MeshBridgeService.kt index 6354e7cf..283c0d42 100644 --- a/app/src/main/java/com/bitchat/android/services/bridge/MeshBridgeService.kt +++ b/app/src/main/java/com/bitchat/android/services/bridge/MeshBridgeService.kt @@ -8,6 +8,7 @@ import com.bitchat.android.geohash.GeohashChannelLevel import com.bitchat.android.geohash.LocationChannelManager import com.bitchat.android.mesh.MeshService import com.bitchat.android.mesh.BridgeMeshDelegate +import com.bitchat.android.mesh.BridgeOutboundPolicy import com.bitchat.android.model.BitchatMessage import com.bitchat.android.model.IdentityAnnouncement import com.bitchat.android.model.NostrCarrierPacket @@ -58,7 +59,6 @@ object MeshBridgeService : BridgeMeshDelegate { private const val TAG = "MeshBridgeService" private const val PREFS = "bitchat_bridge" private const val KEY_ENABLED = "bridge_enabled_v1" - private const val BRIDGE_SUBSCRIPTION = "mesh-bridge-rendezvous" private const val CELL_PRECISION = 6 private const val MAX_EVENT_AGE_MS = 15L * 60 * 1000 private const val MAX_CONTENT_BYTES = 16_000 @@ -85,6 +85,8 @@ object MeshBridgeService : BridgeMeshDelegate { val activeCell: StateFlow = _activeCell.asStateFlow() private val _bridgedParticipants = MutableStateFlow>(emptyList()) val bridgedParticipants: StateFlow> = _bridgedParticipants.asStateFlow() + private val outboundPolicy = AtomicReference(BridgeOutboundPolicy.Denied) + private val rendezvousSubscription = RelaySubscriptionSlot("mesh-bridge-rendezvous") @Volatile private var appContext: Context? = null @@ -137,6 +139,12 @@ object MeshBridgeService : BridgeMeshDelegate { val prekeyManager = PrekeyManager.getInstance(application) prefs = application.getSharedPreferences(PREFS, Context.MODE_PRIVATE) _isEnabled.value = loadEnabledWithMigration(prefs!!) + outboundPolicy.set( + BridgeOutboundPolicy.capture( + enabled = _isEnabled.value, + nearbyOnly = _nearbyOnly.value + ) + ) PeerCapabilities.setBridgeEnabled(_isEnabled.value) prekeyCoordinator = PrekeyCoordinator( manager = prekeyManager, @@ -206,7 +214,8 @@ object MeshBridgeService : BridgeMeshDelegate { } fun setEnabled(enabled: Boolean) { - if (_isEnabled.value == enabled) return + if (outboundPolicy.get().enabled == enabled) return + outboundPolicy.set(BridgeOutboundPolicy.capture(enabled, nearbyOnly = false)) _isEnabled.value = enabled prefs?.edit { putBoolean(KEY_ENABLED, enabled) } PeerCapabilities.setBridgeEnabled(enabled) @@ -232,20 +241,26 @@ object MeshBridgeService : BridgeMeshDelegate { } fun setNearbyOnly(enabled: Boolean) { + outboundPolicy.updateAndGet { current -> + BridgeOutboundPolicy.capture(current.enabled, nearbyOnly = enabled) + } _nearbyOnly.value = enabled } /** Cell included in announce TLV 0x06 while the bridge switch is on. */ override fun advertisedCell(): String? = _activeCell.value.takeIf { _isEnabled.value } + override fun outboundPolicy(): BridgeOutboundPolicy = outboundPolicy.get() + override fun bridgeOutgoing( content: String, senderPeerId: String, timestampMs: Long, - nickname: String? + nickname: String?, + policyAtSend: BridgeOutboundPolicy ) { scope.launch { - if (!_isEnabled.value || _nearbyOnly.value) return@launch + if (!policyAtSend.permitsPublication(outboundPolicy.get())) return@launch val cell = _activeCell.value ?: currentCell() ?: return@launch if (content.toByteArray(Charsets.UTF_8).size > MAX_CONTENT_BYTES) return@launch val identity = NostrIdentityBridge.deriveBridgeIdentity(cell, requireContext()) @@ -257,6 +272,7 @@ object MeshBridgeService : BridgeMeshDelegate { meshSenderId = senderPeerId, meshTimestampMs = timestampMs ) + if (!policyAtSend.permitsPublication(outboundPolicy.get())) return@launch publishedEventIds.add(event.id) injectedEventIds.add(event.id) if (relayManager?.isConnected?.value == true) { @@ -381,22 +397,28 @@ object MeshBridgeService : BridgeMeshDelegate { if (cell == null) return val cells = linkedSetOf(cell).apply { addAll(Geohash.neighborsSamePrecision(cell)) } if (changed || forceSubscriptions || cells != subscribedCells) { - relayManager?.unsubscribe(BRIDGE_SUBSCRIPTION) - subscribedCells = cells val targets = linkedSetOf() cells.forEach { subscribedCell -> relayManager?.ensureGeohashRelaysConnected(subscribedCell) targets += relayManager?.getRelaysForGeohash(subscribedCell).orEmpty() } - relayManager?.subscribe( - filter = NostrFilter.bridgeRendezvous( - cells, - since = clock() - MAX_EVENT_AGE_MS - ), - id = BRIDGE_SUBSCRIPTION, - handler = { event -> scope.launch { handleRendezvousEvent(event) } }, - targetRelayUrls = targets.toList() - ) + relayManager?.let { manager -> + rendezvousSubscription.replace( + close = manager::unsubscribe, + open = { subscriptionId -> + manager.subscribe( + filter = NostrFilter.bridgeRendezvous( + cells, + since = clock() - MAX_EVENT_AGE_MS + ), + id = subscriptionId, + handler = { event -> scope.launch { handleRendezvousEvent(event) } }, + targetRelayUrls = targets.toList() + ) + } + ) + subscribedCells = cells + } publishPresence() } } @@ -664,7 +686,9 @@ object MeshBridgeService : BridgeMeshDelegate { } private fun closeSubscriptions() { - relayManager?.unsubscribe(BRIDGE_SUBSCRIPTION) + relayManager?.let { manager -> + rendezvousSubscription.close(manager::unsubscribe) + } subscribedCells = emptySet() } diff --git a/app/src/test/kotlin/com/bitchat/android/services/bridge/BridgeRaceRegressionTest.kt b/app/src/test/kotlin/com/bitchat/android/services/bridge/BridgeRaceRegressionTest.kt new file mode 100644 index 00000000..338437f4 --- /dev/null +++ b/app/src/test/kotlin/com/bitchat/android/services/bridge/BridgeRaceRegressionTest.kt @@ -0,0 +1,57 @@ +package com.bitchat.android.services.bridge + +import com.bitchat.android.mesh.BridgeOutboundPolicy +import org.junit.Assert.assertFalse +import org.junit.Assert.assertNotEquals +import org.junit.Assert.assertTrue +import org.junit.Test + +class BridgeRaceRegressionTest { + @Test + fun `later opt in cannot authorize a send that was previously denied`() { + val policyAtSend = BridgeOutboundPolicy.capture(enabled = false, nearbyOnly = false) + val policyAtPublish = BridgeOutboundPolicy.capture(enabled = true, nearbyOnly = false) + + assertFalse(policyAtSend.permitsPublication(policyAtPublish)) + } + + @Test + fun `later opt out still blocks a send that was previously allowed`() { + val policyAtSend = BridgeOutboundPolicy.capture(enabled = true, nearbyOnly = false) + val policyAtPublish = BridgeOutboundPolicy.capture(enabled = false, nearbyOnly = false) + + assertFalse(policyAtSend.permitsPublication(policyAtPublish)) + } + + @Test + fun `publication requires bridge permission at send and publish time`() { + val allowed = BridgeOutboundPolicy.capture(enabled = true, nearbyOnly = false) + val nearbyOnly = BridgeOutboundPolicy.capture(enabled = true, nearbyOnly = true) + + assertTrue(allowed.permitsPublication(allowed)) + assertFalse(nearbyOnly.permitsPublication(allowed)) + assertFalse(allowed.permitsPublication(nearbyOnly)) + } + + @Test + fun `delayed close targets the retired subscription generation`() { + val slot = RelaySubscriptionSlot("bridge") + val activeOnRelay = mutableSetOf() + val delayedCloses = mutableListOf() + + val first = slot.replace( + close = delayedCloses::add, + open = activeOnRelay::add + ) + val second = slot.replace( + close = delayedCloses::add, + open = activeOnRelay::add + ) + + delayedCloses.forEach(activeOnRelay::remove) + + assertNotEquals(first, second) + assertFalse(first in activeOnRelay) + assertTrue(second in activeOnRelay) + } +}