mirror of
https://github.com/permissionlesstech/bitchat-android.git
synced 2026-07-25 11:05:20 +00:00
Centralized relay manager (#182)
* noise: send nonce * broadcast manager * add missing file
This commit is contained in:
@@ -316,6 +316,15 @@ class BluetoothMeshService(private val context: Context) {
|
|||||||
peerManager.updatePeerLastSeen(peerID)
|
peerManager.updatePeerLastSeen(peerID)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// Network information for relay manager
|
||||||
|
override fun getNetworkSize(): Int {
|
||||||
|
return peerManager.getActivePeerCount()
|
||||||
|
}
|
||||||
|
|
||||||
|
override fun getBroadcastRecipient(): ByteArray {
|
||||||
|
return SpecialRecipients.BROADCAST
|
||||||
|
}
|
||||||
|
|
||||||
override fun handleNoiseHandshake(routed: RoutedPacket, step: Int): Boolean {
|
override fun handleNoiseHandshake(routed: RoutedPacket, step: Int): Boolean {
|
||||||
return runBlocking { securityManager.handleNoiseHandshake(routed, step) }
|
return runBlocking { securityManager.handleNoiseHandshake(routed, step) }
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -177,12 +177,7 @@ class MessageHandler(private val myPeerID: String) {
|
|||||||
// Notify delegate to handle peer management
|
// Notify delegate to handle peer management
|
||||||
val isFirstAnnounce = delegate?.addOrUpdatePeer(peerID, nickname) ?: false
|
val isFirstAnnounce = delegate?.addOrUpdatePeer(peerID, nickname) ?: false
|
||||||
|
|
||||||
// Relay announce if TTL > 0
|
// Announce relay is now handled by centralized PacketRelayManager
|
||||||
if (packet.ttl > 1u) {
|
|
||||||
val relayPacket = packet.copy(ttl = (packet.ttl - 1u).toUByte())
|
|
||||||
delay(Random.nextLong(100, 300))
|
|
||||||
delegate?.relayPacket(RoutedPacket(relayPacket, peerID, routed.relayAddress))
|
|
||||||
}
|
|
||||||
|
|
||||||
return isFirstAnnounce
|
return isFirstAnnounce
|
||||||
}
|
}
|
||||||
@@ -196,20 +191,15 @@ class MessageHandler(private val myPeerID: String) {
|
|||||||
if (peerID == myPeerID) return
|
if (peerID == myPeerID) return
|
||||||
|
|
||||||
val recipientID = packet.recipientID?.takeIf { !it.contentEquals(delegate?.getBroadcastRecipient()) }
|
val recipientID = packet.recipientID?.takeIf { !it.contentEquals(delegate?.getBroadcastRecipient()) }
|
||||||
var recipientIDString = ""
|
|
||||||
if (recipientID != null) {
|
|
||||||
recipientIDString = recipientID.toHexString()
|
|
||||||
}
|
|
||||||
if (recipientID == null) {
|
if (recipientID == null) {
|
||||||
// BROADCAST MESSAGE
|
// BROADCAST MESSAGE
|
||||||
handleBroadcastMessage(routed)
|
handleBroadcastMessage(routed)
|
||||||
} else if (recipientID.toHexString() == myPeerID) {
|
} else if (recipientID.toHexString() == myPeerID) {
|
||||||
// PRIVATE MESSAGE FOR US
|
// PRIVATE MESSAGE FOR US
|
||||||
handlePrivateMessage(packet, peerID)
|
handlePrivateMessage(packet, peerID)
|
||||||
} else if (packet.ttl > 0u) {
|
|
||||||
// RELAY MESSAGE
|
|
||||||
relayMessage(routed)
|
|
||||||
}
|
}
|
||||||
|
// Message relay is now handled by centralized PacketRelayManager
|
||||||
}
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
@@ -248,8 +238,7 @@ class MessageHandler(private val myPeerID: String) {
|
|||||||
delegate?.onMessageReceived(messageWithCurrentTime)
|
delegate?.onMessageReceived(messageWithCurrentTime)
|
||||||
}
|
}
|
||||||
|
|
||||||
// Relay broadcast messages
|
// Broadcast message relay is now handled by centralized PacketRelayManager
|
||||||
relayMessage(routed)
|
|
||||||
|
|
||||||
} catch (e: Exception) {
|
} catch (e: Exception) {
|
||||||
Log.e(TAG, "Failed to process broadcast message: ${e.message}")
|
Log.e(TAG, "Failed to process broadcast message: ${e.message}")
|
||||||
@@ -308,11 +297,7 @@ class MessageHandler(private val myPeerID: String) {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
// Relay if TTL > 0
|
// Leave message relay is now handled by centralized PacketRelayManager
|
||||||
if (packet.ttl > 1u) {
|
|
||||||
val relayPacket = packet.copy(ttl = (packet.ttl - 1u).toUByte())
|
|
||||||
delegate?.relayPacket(routed.copy(packet = relayPacket))
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
@@ -333,11 +318,8 @@ class MessageHandler(private val myPeerID: String) {
|
|||||||
} catch (e: Exception) {
|
} catch (e: Exception) {
|
||||||
Log.e(TAG, "Failed to decrypt delivery ACK: ${e.message}")
|
Log.e(TAG, "Failed to decrypt delivery ACK: ${e.message}")
|
||||||
}
|
}
|
||||||
} else if (packet.ttl > 0u && String(packet.senderID).replace("\u0000", "") != myPeerID) {
|
|
||||||
// Relay
|
|
||||||
val relayPacket = packet.copy(ttl = (packet.ttl - 1u).toUByte())
|
|
||||||
delegate?.relayPacket(routed.copy(packet = relayPacket))
|
|
||||||
}
|
}
|
||||||
|
// Delivery ACK relay is now handled by centralized PacketRelayManager
|
||||||
}
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
@@ -358,39 +340,8 @@ class MessageHandler(private val myPeerID: String) {
|
|||||||
} catch (e: Exception) {
|
} catch (e: Exception) {
|
||||||
Log.e(TAG, "Failed to decrypt read receipt: ${e.message}")
|
Log.e(TAG, "Failed to decrypt read receipt: ${e.message}")
|
||||||
}
|
}
|
||||||
} else if (packet.ttl > 0u && String(packet.senderID).replace("\u0000", "") != myPeerID) {
|
|
||||||
// Relay
|
|
||||||
val relayPacket = packet.copy(ttl = (packet.ttl - 1u).toUByte())
|
|
||||||
delegate?.relayPacket(routed.copy(packet = relayPacket))
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
/**
|
|
||||||
* Relay message with adaptive probability (same as iOS)
|
|
||||||
*/
|
|
||||||
private suspend fun relayMessage(routed: RoutedPacket) {
|
|
||||||
val packet = routed.packet
|
|
||||||
if (packet.ttl == 0u.toUByte()) return
|
|
||||||
|
|
||||||
val relayPacket = packet.copy(ttl = (packet.ttl - 1u).toUByte())
|
|
||||||
|
|
||||||
// Check network size and apply adaptive relay probability
|
|
||||||
val networkSize = delegate?.getNetworkSize() ?: 1
|
|
||||||
val relayProb = when {
|
|
||||||
networkSize <= 10 -> 1.0
|
|
||||||
networkSize <= 30 -> 0.85
|
|
||||||
networkSize <= 50 -> 0.7
|
|
||||||
networkSize <= 100 -> 0.55
|
|
||||||
else -> 0.4
|
|
||||||
}
|
|
||||||
|
|
||||||
val shouldRelay = relayPacket.ttl >= 4u || networkSize <= 3 || Random.nextDouble() < relayProb
|
|
||||||
|
|
||||||
if (shouldRelay) {
|
|
||||||
val delay = Random.nextLong(50, 500) // Random delay like iOS
|
|
||||||
delay(delay)
|
|
||||||
delegate?.relayPacket(routed.copy(packet = relayPacket))
|
|
||||||
}
|
}
|
||||||
|
// Read receipt relay is now handled by centralized PacketRelayManager
|
||||||
}
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
|
|||||||
@@ -24,6 +24,9 @@ class PacketProcessor(private val myPeerID: String) {
|
|||||||
// Delegate for callbacks
|
// Delegate for callbacks
|
||||||
var delegate: PacketProcessorDelegate? = null
|
var delegate: PacketProcessorDelegate? = null
|
||||||
|
|
||||||
|
// Packet relay manager for centralized relay decisions
|
||||||
|
private val packetRelayManager = PacketRelayManager(myPeerID)
|
||||||
|
|
||||||
// Coroutines
|
// Coroutines
|
||||||
private val processorScope = CoroutineScope(Dispatchers.IO + SupervisorJob())
|
private val processorScope = CoroutineScope(Dispatchers.IO + SupervisorJob())
|
||||||
|
|
||||||
@@ -51,6 +54,11 @@ class PacketProcessor(private val myPeerID: String) {
|
|||||||
// Cache actors to reuse them
|
// Cache actors to reuse them
|
||||||
private val actors = mutableMapOf<String, kotlinx.coroutines.channels.SendChannel<RoutedPacket>>()
|
private val actors = mutableMapOf<String, kotlinx.coroutines.channels.SendChannel<RoutedPacket>>()
|
||||||
|
|
||||||
|
init {
|
||||||
|
// Set up the packet relay manager delegate immediately
|
||||||
|
setupRelayManager()
|
||||||
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Process received packet - main entry point for all incoming packets
|
* Process received packet - main entry point for all incoming packets
|
||||||
* SURGICAL FIX: Route to per-peer actor for serialized processing
|
* SURGICAL FIX: Route to per-peer actor for serialized processing
|
||||||
@@ -73,6 +81,25 @@ class PacketProcessor(private val myPeerID: String) {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Set up the packet relay manager with its delegate
|
||||||
|
*/
|
||||||
|
fun setupRelayManager() {
|
||||||
|
packetRelayManager.delegate = object : PacketRelayManagerDelegate {
|
||||||
|
override fun getNetworkSize(): Int {
|
||||||
|
return delegate?.getNetworkSize() ?: 1
|
||||||
|
}
|
||||||
|
|
||||||
|
override fun getBroadcastRecipient(): ByteArray {
|
||||||
|
return delegate?.getBroadcastRecipient() ?: ByteArray(0)
|
||||||
|
}
|
||||||
|
|
||||||
|
override fun broadcastPacket(routed: RoutedPacket) {
|
||||||
|
delegate?.relayPacket(routed)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Handle received packet - core protocol logic (exact same as iOS)
|
* Handle received packet - core protocol logic (exact same as iOS)
|
||||||
*/
|
*/
|
||||||
@@ -107,9 +134,14 @@ class PacketProcessor(private val myPeerID: String) {
|
|||||||
Log.w(TAG, "Unknown message type: ${packet.type}")
|
Log.w(TAG, "Unknown message type: ${packet.type}")
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
// Update last seen timestamp
|
// Update last seen timestamp
|
||||||
if (validPacket)
|
if (validPacket) {
|
||||||
delegate?.updatePeerLastSeen(peerID)
|
delegate?.updatePeerLastSeen(peerID)
|
||||||
|
|
||||||
|
// CENTRALIZED RELAY LOGIC: Handle relay decisions for all packets not addressed to us
|
||||||
|
packetRelayManager.handlePacketRelay(routed)
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
@@ -175,11 +207,7 @@ class PacketProcessor(private val myPeerID: String) {
|
|||||||
handleReceivedPacket(RoutedPacket(reassembledPacket, routed.peerID, routed.relayAddress))
|
handleReceivedPacket(RoutedPacket(reassembledPacket, routed.peerID, routed.relayAddress))
|
||||||
}
|
}
|
||||||
|
|
||||||
// Relay fragment regardless of reassembly
|
// Fragment relay is now handled by centralized PacketRelayManager
|
||||||
if (routed.packet.ttl > 0u) {
|
|
||||||
val relayPacket = routed.packet.copy(ttl = (routed.packet.ttl - 1u).toUByte())
|
|
||||||
delegate?.relayPacket(RoutedPacket(relayPacket, routed.peerID, routed.relayAddress))
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
@@ -229,6 +257,9 @@ class PacketProcessor(private val myPeerID: String) {
|
|||||||
}
|
}
|
||||||
actors.clear()
|
actors.clear()
|
||||||
|
|
||||||
|
// Shutdown the relay manager
|
||||||
|
packetRelayManager.shutdown()
|
||||||
|
|
||||||
// Cancel the main scope
|
// Cancel the main scope
|
||||||
processorScope.cancel()
|
processorScope.cancel()
|
||||||
|
|
||||||
@@ -246,6 +277,10 @@ interface PacketProcessorDelegate {
|
|||||||
// Peer management
|
// Peer management
|
||||||
fun updatePeerLastSeen(peerID: String)
|
fun updatePeerLastSeen(peerID: String)
|
||||||
|
|
||||||
|
// Network information
|
||||||
|
fun getNetworkSize(): Int
|
||||||
|
fun getBroadcastRecipient(): ByteArray
|
||||||
|
|
||||||
// Message type handlers
|
// Message type handlers
|
||||||
fun handleNoiseHandshake(routed: RoutedPacket, step: Int): Boolean
|
fun handleNoiseHandshake(routed: RoutedPacket, step: Int): Boolean
|
||||||
fun handleNoiseEncrypted(routed: RoutedPacket)
|
fun handleNoiseEncrypted(routed: RoutedPacket)
|
||||||
|
|||||||
@@ -0,0 +1,201 @@
|
|||||||
|
package com.bitchat.android.mesh
|
||||||
|
|
||||||
|
import android.util.Log
|
||||||
|
import com.bitchat.android.model.RoutedPacket
|
||||||
|
import com.bitchat.android.protocol.BitchatPacket
|
||||||
|
import com.bitchat.android.util.toHexString
|
||||||
|
import kotlinx.coroutines.*
|
||||||
|
import kotlin.random.Random
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Centralized packet relay management
|
||||||
|
*
|
||||||
|
* This class handles all relay decisions and logic for bitchat packets.
|
||||||
|
* All packets that aren't specifically addressed to us get processed here.
|
||||||
|
*/
|
||||||
|
class PacketRelayManager(private val myPeerID: String) {
|
||||||
|
|
||||||
|
companion object {
|
||||||
|
private const val TAG = "PacketRelayManager"
|
||||||
|
}
|
||||||
|
|
||||||
|
// Delegate for callbacks
|
||||||
|
var delegate: PacketRelayManagerDelegate? = null
|
||||||
|
|
||||||
|
// Coroutines
|
||||||
|
private val relayScope = CoroutineScope(Dispatchers.IO + SupervisorJob())
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Main entry point for relay decisions
|
||||||
|
* Only packets that aren't specifically addressed to us should be passed here
|
||||||
|
*/
|
||||||
|
suspend fun handlePacketRelay(routed: RoutedPacket) {
|
||||||
|
val packet = routed.packet
|
||||||
|
val peerID = routed.peerID ?: "unknown"
|
||||||
|
|
||||||
|
Log.d(TAG, "Evaluating relay for packet type ${packet.type} from $peerID (TTL: ${packet.ttl})")
|
||||||
|
|
||||||
|
// Double-check this packet isn't addressed to us
|
||||||
|
if (isPacketAddressedToMe(packet)) {
|
||||||
|
Log.d(TAG, "Packet addressed to us, skipping relay")
|
||||||
|
return
|
||||||
|
}
|
||||||
|
|
||||||
|
// Skip our own packets
|
||||||
|
if (peerID == myPeerID) {
|
||||||
|
Log.d(TAG, "Packet from ourselves, skipping relay")
|
||||||
|
return
|
||||||
|
}
|
||||||
|
|
||||||
|
// Check TTL and decrement
|
||||||
|
if (packet.ttl == 0u.toUByte()) {
|
||||||
|
Log.d(TAG, "TTL expired, not relaying packet")
|
||||||
|
return
|
||||||
|
}
|
||||||
|
|
||||||
|
// Decrement TTL by 1
|
||||||
|
val relayPacket = packet.copy(ttl = (packet.ttl - 1u).toUByte())
|
||||||
|
Log.d(TAG, "Decremented TTL from ${packet.ttl} to ${relayPacket.ttl}")
|
||||||
|
|
||||||
|
// Apply relay logic based on packet type
|
||||||
|
val shouldRelay = shouldRelayPacket(relayPacket, peerID)
|
||||||
|
|
||||||
|
if (shouldRelay) {
|
||||||
|
relayPacket(RoutedPacket(relayPacket, peerID, routed.relayAddress))
|
||||||
|
} else {
|
||||||
|
Log.d(TAG, "Relay decision: NOT relaying packet type ${packet.type}")
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Check if a packet is specifically addressed to us
|
||||||
|
*/
|
||||||
|
private fun isPacketAddressedToMe(packet: BitchatPacket): Boolean {
|
||||||
|
val recipientID = packet.recipientID
|
||||||
|
|
||||||
|
// No recipient means broadcast (not addressed to us specifically)
|
||||||
|
if (recipientID == null) {
|
||||||
|
return false
|
||||||
|
}
|
||||||
|
|
||||||
|
// Check if it's a broadcast recipient
|
||||||
|
val broadcastRecipient = delegate?.getBroadcastRecipient()
|
||||||
|
if (broadcastRecipient != null && recipientID.contentEquals(broadcastRecipient)) {
|
||||||
|
return false
|
||||||
|
}
|
||||||
|
|
||||||
|
// Check if recipient matches our peer ID
|
||||||
|
val recipientIDString = recipientID.toHexString()
|
||||||
|
return recipientIDString == myPeerID
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Determine if we should relay this packet based on type and network conditions
|
||||||
|
*/
|
||||||
|
private fun shouldRelayPacket(packet: BitchatPacket, fromPeerID: String): Boolean {
|
||||||
|
// Always relay if TTL is high enough (indicates important message)
|
||||||
|
if (packet.ttl >= 4u) {
|
||||||
|
Log.d(TAG, "High TTL (${packet.ttl}), relaying")
|
||||||
|
return true
|
||||||
|
}
|
||||||
|
|
||||||
|
// Get network size for adaptive relay probability
|
||||||
|
val networkSize = delegate?.getNetworkSize() ?: 1
|
||||||
|
|
||||||
|
// Small networks always relay to ensure connectivity
|
||||||
|
if (networkSize <= 3) {
|
||||||
|
Log.d(TAG, "Small network ($networkSize peers), relaying")
|
||||||
|
return true
|
||||||
|
}
|
||||||
|
|
||||||
|
// Apply adaptive relay probability based on network size
|
||||||
|
val relayProb = when {
|
||||||
|
networkSize <= 10 -> 1.0 // Always relay in small networks
|
||||||
|
networkSize <= 30 -> 0.85 // High probability for medium networks
|
||||||
|
networkSize <= 50 -> 0.7 // Moderate probability
|
||||||
|
networkSize <= 100 -> 0.55 // Lower probability for large networks
|
||||||
|
else -> 0.4 // Lowest probability for very large networks
|
||||||
|
}
|
||||||
|
|
||||||
|
val shouldRelay = Random.nextDouble() < relayProb
|
||||||
|
Log.d(TAG, "Network size: $networkSize, Relay probability: $relayProb, Decision: $shouldRelay")
|
||||||
|
|
||||||
|
return shouldRelay
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Relay message with adaptive probability and timing (same as iOS)
|
||||||
|
* Moved from MessageHandler.kt
|
||||||
|
*/
|
||||||
|
suspend fun relayMessage(routed: RoutedPacket) {
|
||||||
|
val packet = routed.packet
|
||||||
|
|
||||||
|
if (packet.ttl == 0u.toUByte()) {
|
||||||
|
Log.d(TAG, "TTL expired, not relaying message")
|
||||||
|
return
|
||||||
|
}
|
||||||
|
|
||||||
|
val relayPacket = packet.copy(ttl = (packet.ttl - 1u).toUByte())
|
||||||
|
|
||||||
|
// Check network size and apply adaptive relay probability
|
||||||
|
val networkSize = delegate?.getNetworkSize() ?: 1
|
||||||
|
val relayProb = when {
|
||||||
|
networkSize <= 10 -> 1.0
|
||||||
|
networkSize <= 30 -> 0.85
|
||||||
|
networkSize <= 50 -> 0.7
|
||||||
|
networkSize <= 100 -> 0.55
|
||||||
|
else -> 0.4
|
||||||
|
}
|
||||||
|
|
||||||
|
val shouldRelay = relayPacket.ttl >= 4u || networkSize <= 3 || Random.nextDouble() < relayProb
|
||||||
|
|
||||||
|
if (shouldRelay) {
|
||||||
|
val delay = Random.nextLong(50, 500) // Random delay like iOS
|
||||||
|
Log.d(TAG, "Relaying message after ${delay}ms delay")
|
||||||
|
delay(delay)
|
||||||
|
relayPacket(routed.copy(packet = relayPacket))
|
||||||
|
} else {
|
||||||
|
Log.d(TAG, "Relay decision: NOT relaying message (network size: $networkSize, prob: $relayProb)")
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Actually broadcast the packet for relay
|
||||||
|
*/
|
||||||
|
private fun relayPacket(routed: RoutedPacket) {
|
||||||
|
Log.d(TAG, "🔄 Relaying packet type ${routed.packet.type} with TTL ${routed.packet.ttl}")
|
||||||
|
delegate?.broadcastPacket(routed)
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Get debug information
|
||||||
|
*/
|
||||||
|
fun getDebugInfo(): String {
|
||||||
|
return buildString {
|
||||||
|
appendLine("=== Packet Relay Manager Debug Info ===")
|
||||||
|
appendLine("Relay Scope Active: ${relayScope.isActive}")
|
||||||
|
appendLine("My Peer ID: $myPeerID")
|
||||||
|
appendLine("Network Size: ${delegate?.getNetworkSize() ?: "unknown"}")
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Shutdown the relay manager
|
||||||
|
*/
|
||||||
|
fun shutdown() {
|
||||||
|
Log.d(TAG, "Shutting down PacketRelayManager")
|
||||||
|
relayScope.cancel()
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Delegate interface for packet relay manager callbacks
|
||||||
|
*/
|
||||||
|
interface PacketRelayManagerDelegate {
|
||||||
|
// Network information
|
||||||
|
fun getNetworkSize(): Int
|
||||||
|
fun getBroadcastRecipient(): ByteArray
|
||||||
|
|
||||||
|
// Packet operations
|
||||||
|
fun broadcastPacket(routed: RoutedPacket)
|
||||||
|
}
|
||||||
Reference in New Issue
Block a user