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 2884b84d..abc91994 100644 --- a/app/src/main/java/com/bitchat/android/mesh/BluetoothMeshService.kt +++ b/app/src/main/java/com/bitchat/android/mesh/BluetoothMeshService.kt @@ -316,6 +316,15 @@ class BluetoothMeshService(private val context: Context) { 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 { return runBlocking { securityManager.handleNoiseHandshake(routed, step) } } diff --git a/app/src/main/java/com/bitchat/android/mesh/MessageHandler.kt b/app/src/main/java/com/bitchat/android/mesh/MessageHandler.kt index e552e641..de4a885b 100644 --- a/app/src/main/java/com/bitchat/android/mesh/MessageHandler.kt +++ b/app/src/main/java/com/bitchat/android/mesh/MessageHandler.kt @@ -177,12 +177,7 @@ class MessageHandler(private val myPeerID: String) { // Notify delegate to handle peer management val isFirstAnnounce = delegate?.addOrUpdatePeer(peerID, nickname) ?: false - // Relay announce if TTL > 0 - 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)) - } + // Announce relay is now handled by centralized PacketRelayManager return isFirstAnnounce } @@ -196,20 +191,15 @@ class MessageHandler(private val myPeerID: String) { if (peerID == myPeerID) return val recipientID = packet.recipientID?.takeIf { !it.contentEquals(delegate?.getBroadcastRecipient()) } - var recipientIDString = "" - if (recipientID != null) { - recipientIDString = recipientID.toHexString() - } + if (recipientID == null) { // BROADCAST MESSAGE handleBroadcastMessage(routed) } else if (recipientID.toHexString() == myPeerID) { // PRIVATE MESSAGE FOR US 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) } - // Relay broadcast messages - relayMessage(routed) + // Broadcast message relay is now handled by centralized PacketRelayManager } catch (e: Exception) { Log.e(TAG, "Failed to process broadcast message: ${e.message}") @@ -308,11 +297,7 @@ class MessageHandler(private val myPeerID: String) { } } - // Relay if TTL > 0 - if (packet.ttl > 1u) { - val relayPacket = packet.copy(ttl = (packet.ttl - 1u).toUByte()) - delegate?.relayPacket(routed.copy(packet = relayPacket)) - } + // Leave message relay is now handled by centralized PacketRelayManager } /** @@ -333,11 +318,8 @@ class MessageHandler(private val myPeerID: String) { } catch (e: Exception) { 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) { 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 } /** 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 fde5ea86..cc2331dc 100644 --- a/app/src/main/java/com/bitchat/android/mesh/PacketProcessor.kt +++ b/app/src/main/java/com/bitchat/android/mesh/PacketProcessor.kt @@ -24,6 +24,9 @@ class PacketProcessor(private val myPeerID: String) { // Delegate for callbacks var delegate: PacketProcessorDelegate? = null + // Packet relay manager for centralized relay decisions + private val packetRelayManager = PacketRelayManager(myPeerID) + // Coroutines private val processorScope = CoroutineScope(Dispatchers.IO + SupervisorJob()) @@ -51,6 +54,11 @@ class PacketProcessor(private val myPeerID: String) { // Cache actors to reuse them private val actors = mutableMapOf>() + init { + // Set up the packet relay manager delegate immediately + setupRelayManager() + } + /** * Process received packet - main entry point for all incoming packets * 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) */ @@ -107,9 +134,14 @@ class PacketProcessor(private val myPeerID: String) { Log.w(TAG, "Unknown message type: ${packet.type}") } } + // Update last seen timestamp - if (validPacket) + if (validPacket) { 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)) } - // Relay fragment regardless of reassembly - if (routed.packet.ttl > 0u) { - val relayPacket = routed.packet.copy(ttl = (routed.packet.ttl - 1u).toUByte()) - delegate?.relayPacket(RoutedPacket(relayPacket, routed.peerID, routed.relayAddress)) - } + // Fragment relay is now handled by centralized PacketRelayManager } /** @@ -229,6 +257,9 @@ class PacketProcessor(private val myPeerID: String) { } actors.clear() + // Shutdown the relay manager + packetRelayManager.shutdown() + // Cancel the main scope processorScope.cancel() @@ -246,6 +277,10 @@ interface PacketProcessorDelegate { // Peer management fun updatePeerLastSeen(peerID: String) + // Network information + fun getNetworkSize(): Int + fun getBroadcastRecipient(): ByteArray + // Message type handlers fun handleNoiseHandshake(routed: RoutedPacket, step: Int): Boolean fun handleNoiseEncrypted(routed: RoutedPacket) diff --git a/app/src/main/java/com/bitchat/android/mesh/PacketRelayManager.kt b/app/src/main/java/com/bitchat/android/mesh/PacketRelayManager.kt new file mode 100644 index 00000000..24ea0119 --- /dev/null +++ b/app/src/main/java/com/bitchat/android/mesh/PacketRelayManager.kt @@ -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) +}