This commit is contained in:
callebtc
2025-10-20 15:26:43 +02:00
parent f633509848
commit 35b197bba4
@@ -20,7 +20,8 @@ import kotlinx.coroutines.launch
import kotlinx.coroutines.channels.Channel import kotlinx.coroutines.channels.Channel
import kotlinx.coroutines.Job import kotlinx.coroutines.Job
import java.util.concurrent.ConcurrentHashMap import java.util.concurrent.ConcurrentHashMap
import kotlinx.coroutines.channels.actor import java.util.PriorityQueue
import java.util.concurrent.atomic.AtomicLong
/** /**
* Handles packet broadcasting to connected devices using actor pattern for serialization * Handles packet broadcasting to connected devices using actor pattern for serialization
@@ -99,25 +100,90 @@ class BluetoothPacketBroadcaster(
val gattServer: BluetoothGattServer?, val gattServer: BluetoothGattServer?,
val characteristic: BluetoothGattCharacteristic? val characteristic: BluetoothGattCharacteristic?
) )
// Internal queued item with priority metadata for fair ordering
private data class QueuedBroadcast(
val request: BroadcastRequest,
val primaryPriority: Int,
val secondaryPriority: Int,
val sequence: Long
)
// Actor scope for the broadcaster // Actor scope for the broadcaster
private val broadcasterScope = CoroutineScope(Dispatchers.IO + SupervisorJob()) private val broadcasterScope = CoroutineScope(Dispatchers.IO + SupervisorJob())
private val transferJobs = ConcurrentHashMap<String, Job>() private val transferJobs = ConcurrentHashMap<String, Job>()
// SERIALIZATION: Actor to serialize all broadcast operations // Replace simple actor FIFO with a priority-aware channel + processor
@OptIn(kotlinx.coroutines.ObsoleteCoroutinesApi::class) private val broadcasterChannel = Channel<BroadcastRequest>(capacity = Channel.UNLIMITED)
private val broadcasterActor = broadcasterScope.actor<BroadcastRequest>( private val sequenceCounter = AtomicLong(0)
capacity = Channel.UNLIMITED
) { init {
Log.d(TAG, "🎭 Created packet broadcaster actor") startPriorityProcessor()
try { }
for (request in channel) {
broadcastSinglePacketInternal(request.routed, request.gattServer, request.characteristic) private fun startPriorityProcessor() {
broadcasterScope.launch {
Log.d(TAG, "🎭 Created priority packet broadcaster processor")
// Min-heap by (primaryPriority, secondaryPriority, sequence)
val queue = PriorityQueue<QueuedBroadcast>(11) { a, b ->
when {
a.primaryPriority != b.primaryPriority -> a.primaryPriority - b.primaryPriority
a.secondaryPriority != b.secondaryPriority -> a.secondaryPriority - b.secondaryPriority
else -> (a.sequence - b.sequence).coerceIn(Int.MIN_VALUE.toLong(), Int.MAX_VALUE.toLong()).toInt()
}
}
try {
while (isActive) {
// If queue is empty, suspend to receive at least one item
if (queue.isEmpty()) {
val first = broadcasterChannel.receiveCatching().getOrNull() ?: break
queue.offer(computeQueuedBroadcast(first))
}
// Drain any immediately available items without suspending
while (true) {
val received = broadcasterChannel.tryReceive().getOrNull() ?: break
queue.offer(computeQueuedBroadcast(received))
}
// Process one highest-priority item
val next = queue.poll() ?: continue
broadcastSinglePacketInternal(next.request.routed, next.request.gattServer, next.request.characteristic)
}
} catch (e: Exception) {
Log.w(TAG, "Priority processor loop ended: ${'$'}{e.message}")
} finally {
Log.d(TAG, "🎭 Priority packet broadcaster processor terminated")
} }
} finally {
Log.d(TAG, "🎭 Packet broadcaster actor terminated")
} }
} }
private fun computeQueuedBroadcast(req: BroadcastRequest): QueuedBroadcast {
val pkt = req.routed.packet
val type = MessageType.fromValue(pkt.type)
// Primary priority: 0 = normal, 1 = fragments, 2 = file transfer (non-fragment)
val primary = when (type) {
MessageType.FRAGMENT -> 1
MessageType.FILE_TRANSFER -> 2
else -> 0
}
// Secondary priority: for fragments, use total fragment count (smaller = higher priority)
val secondary = if (type == MessageType.FRAGMENT) {
try {
val payload = com.bitchat.android.model.FragmentPayload.decode(pkt.payload)
// If decode fails, deprioritize heavily
payload?.total ?: Int.MAX_VALUE
} catch (_: Exception) {
Int.MAX_VALUE
}
} else 0
val seq = sequenceCounter.getAndIncrement()
return QueuedBroadcast(req, primary, secondary, seq)
}
fun broadcastPacket( fun broadcastPacket(
routed: RoutedPacket, routed: RoutedPacket,
@@ -267,7 +333,7 @@ class BluetoothPacketBroadcaster(
// Submit broadcast request to actor for serialized processing // Submit broadcast request to actor for serialized processing
broadcasterScope.launch { broadcasterScope.launch {
try { try {
broadcasterActor.send(BroadcastRequest(routed, gattServer, characteristic)) broadcasterChannel.send(BroadcastRequest(routed, gattServer, characteristic))
} catch (e: Exception) { } catch (e: Exception) {
Log.w(TAG, "Failed to send broadcast request to actor: ${e.message}") Log.w(TAG, "Failed to send broadcast request to actor: ${e.message}")
// Fallback to direct processing if actor fails // Fallback to direct processing if actor fails
@@ -427,7 +493,7 @@ class BluetoothPacketBroadcaster(
return buildString { return buildString {
appendLine("=== Packet Broadcaster Debug Info ===") appendLine("=== Packet Broadcaster Debug Info ===")
appendLine("Broadcaster Scope Active: ${broadcasterScope.isActive}") appendLine("Broadcaster Scope Active: ${broadcasterScope.isActive}")
appendLine("Actor Channel Closed: ${broadcasterActor.isClosedForSend}") appendLine("Actor Channel Closed: ${broadcasterChannel.isClosedForSend}")
appendLine("Connection Scope Active: ${connectionScope.isActive}") appendLine("Connection Scope Active: ${connectionScope.isActive}")
} }
} }
@@ -439,11 +505,11 @@ class BluetoothPacketBroadcaster(
Log.d(TAG, "Shutting down BluetoothPacketBroadcaster actor") Log.d(TAG, "Shutting down BluetoothPacketBroadcaster actor")
// Close the actor gracefully // Close the actor gracefully
broadcasterActor.close() broadcasterChannel.close()
// Cancel the broadcaster scope // Cancel the broadcaster scope
broadcasterScope.cancel() broadcasterScope.cancel()
Log.d(TAG, "BluetoothPacketBroadcaster shutdown complete") Log.d(TAG, "BluetoothPacketBroadcaster shutdown complete")
} }
} }