This commit is contained in:
callebtc
2025-10-21 13:25:39 +02:00
parent 6fb45f730a
commit 27c4d363b4
@@ -114,85 +114,60 @@ class BluetoothPacketBroadcaster(
private val transferJobs = ConcurrentHashMap<String, Job>() private val transferJobs = ConcurrentHashMap<String, Job>()
// Replace simple actor FIFO with a priority-aware channel + processor // Replace simple actor FIFO with a priority-aware channel + processor
@Volatile private val broadcasterChannel = Channel<BroadcastRequest>(capacity = Channel.UNLIMITED)
private var broadcasterChannel = Channel<BroadcastRequest>(capacity = Channel.UNLIMITED)
private val sequenceCounter = AtomicLong(0) private val sequenceCounter = AtomicLong(0)
// Track processor state for recovery
@Volatile
private var processorJob: Job? = null
private val processorLock = Any()
init { init {
startPriorityProcessor() startPriorityProcessor()
} }
private fun startPriorityProcessor() { private fun startPriorityProcessor() {
synchronized(processorLock) { broadcasterScope.launch {
// Cancel existing processor if running Log.d(TAG, "🎭 Priority packet broadcaster processor started")
processorJob?.cancel()
// Create new channel for fresh start // Min-heap by (primaryPriority, secondaryPriority, sequence)
broadcasterChannel = Channel(capacity = Channel.UNLIMITED) val priorityQueue = 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()
}
}
processorJob = broadcasterScope.launch { var processorException: Exception? = null
var processorException: Exception? = null try {
Log.d(TAG, "🎭 Priority packet broadcaster processor started") while (isActive) {
// If priority queue is empty, suspend to receive at least one item
if (priorityQueue.isEmpty()) {
val first = broadcasterChannel.receiveCatching().getOrNull() ?: break
priorityQueue.offer(computeQueuedBroadcast(first))
}
// Min-heap by (primaryPriority, secondaryPriority, sequence) // Drain any immediately available items without suspending
val priorityQueue = PriorityQueue<QueuedBroadcast>(11) { a, b -> while (true) {
when { val received = broadcasterChannel.tryReceive().getOrNull() ?: break
a.primaryPriority != b.primaryPriority -> a.primaryPriority - b.primaryPriority priorityQueue.offer(computeQueuedBroadcast(received))
a.secondaryPriority != b.secondaryPriority -> a.secondaryPriority - b.secondaryPriority
else -> (a.sequence - b.sequence).coerceIn(Int.MIN_VALUE.toLong(), Int.MAX_VALUE.toLong()).toInt()
} }
}
try { // Process one highest-priority item
while (isActive) { val next = priorityQueue.poll() ?: continue
// If priority queue is empty, suspend to receive at least one item try {
if (priorityQueue.isEmpty()) { broadcastSinglePacketInternal(next.request.routed, next.request.gattServer, next.request.characteristic)
val first = broadcasterChannel.receiveCatching().getOrNull() ?: break } catch (e: Exception) {
priorityQueue.offer(computeQueuedBroadcast(first)) // BLE error during broadcast - log and save for channel closure
} Log.e(TAG, "❌ Broadcast failed: ${e.message}", e)
processorException = e
// Drain any immediately available items without suspending throw e // Re-throw to exit loop and close channel
while (true) {
val received = broadcasterChannel.tryReceive().getOrNull() ?: break
priorityQueue.offer(computeQueuedBroadcast(received))
}
// Process one highest-priority item
val next = priorityQueue.poll() ?: continue
try {
broadcastSinglePacketInternal(next.request.routed, next.request.gattServer, next.request.characteristic)
} catch (e: Exception) {
// Log the broadcast failure - this is likely a transient BLE error
Log.e(TAG, "❌ Broadcast failed, will attempt recovery: ${e.message}", e)
processorException = e
// Don't throw - instead break and trigger recovery
break
}
}
} catch (e: Exception) {
Log.e(TAG, "❌ Priority processor loop terminated due to exception: ${e.message}", e)
processorException = e
} finally {
// Close the current channel
broadcasterChannel.close(processorException)
Log.w(TAG, "🎭 Priority packet broadcaster processor terminated, attempting recovery...")
// Schedule automatic recovery after a short delay
if (processorException != null && broadcasterScope.isActive) {
broadcasterScope.launch {
delay(1000) // Wait 1 second before recovery
Log.d(TAG, "🔄 Attempting to restart priority processor...")
startPriorityProcessor()
}
} else {
Log.d(TAG, "🎭 Priority packet broadcaster processor shut down (no recovery scheduled)")
} }
} }
} catch (e: Exception) {
Log.e(TAG, "❌ Priority processor loop terminated: ${e.message}", e)
processorException = e
} finally {
// CRITICAL: Close channel so producers fail fast instead of accumulating packets
// This triggers the fallback path in broadcastSinglePacket()
broadcasterChannel.close(processorException)
Log.d(TAG, "🎭 Priority packet broadcaster processor terminated, channel closed")
} }
} }
} }