diff --git a/local_plugins/ble_service/android/src/main/kotlin/com/yunqiinnovation/ble_service/BleService.kt b/local_plugins/ble_service/android/src/main/kotlin/com/yunqiinnovation/ble_service/BleService.kt index 072318720..e79364bae 100644 --- a/local_plugins/ble_service/android/src/main/kotlin/com/yunqiinnovation/ble_service/BleService.kt +++ b/local_plugins/ble_service/android/src/main/kotlin/com/yunqiinnovation/ble_service/BleService.kt @@ -16,7 +16,7 @@ import java.io.File import java.io.FileOutputStream import java.text.SimpleDateFormat import java.util.Date -import java.util.concurrent.LinkedBlockingQueue +import java.util.concurrent.LinkedBlockingDeque import android.os.Handler import java.util.concurrent.atomic.AtomicBoolean import android.util.Log @@ -133,9 +133,16 @@ object BleService { // 音频数据缓存 //private val audioDataBuffer = mutableListOf() - // 音频数据发送相关 - private val audioSendQueue = LinkedBlockingQueue() - private var audioSendThread: Thread? = null + // 音频数据发送相关:左右声道各一条独立缓存队列 + 独立发送线程,两线程错开 20ms 启动,各自按声道节拍下发。 + // 缓存兼作重连暂停期间的暂存与 F4/F5 降速积压;LinkedBlockingDeque 线程安全(生产者=编码回调线程,消费者=对应发送线程)。 + private val leftSendBuffer = LinkedBlockingDeque() + private val rightSendBuffer = LinkedBlockingDeque() + private var leftSendThread: Thread? = null + private var rightSendThread: Thread? = null + // 两线程写的是同一个 callWriteChar(共享可变 .value)+同一 gatt,真正的写必须串行化,否则并发写会互相覆盖发错数据。 + private val bleWriteLock = Any() + // 右声道发送线程相对左声道错开启动的毫秒数,避免两路同一时刻抢 BLE 写、平摊链路占用。 + private const val RIGHT_SEND_THREAD_STAGGER_MS = 40L private val audioSendHandler = Handler(Looper.getMainLooper()) private val isAudioSending = AtomicBoolean(false) // 上一包音频实际发送的时间戳(ms),仅用于发送节拍日志统计间隔;-1 表示尚未发送过 @@ -196,24 +203,28 @@ object BleService { // 音频下行发送间隔(ms),调试界面可在运行时调整;通过 setCallTranslationDebugParams() 修改 @Volatile private var audioSendIntervalNormal = DEFAULT_AUDIO_SEND_INTERVAL_NORMAL - // 左/右声道下发节拍倍数(以一个 audioSendIntervalNormal=40ms 拍为单位): - // 1 拍=40ms/包(全速)、4 拍=160ms/包、8 拍=320ms/包(三档)。 - // 由设备 F4(左)/F5(右) 上报的解码缓存空余字节数动态切换:空余越少越降速,缓解设备解码溢出丢包。 + // 左/右声道下发节拍倍数(以一个 audioSendIntervalNormal 拍为单位): + // 1 拍=1×拍长/包(全速)、4 拍、8 拍(三档)。由设备 F4(左)/F5(右) 上报的解码空余动态切换:空余越少越降速。 + // 各声道发送线程按 tick - lastSendTick(线程内局部) >= paceTicks 控制该声道下发频率。 @Volatile private var leftPaceTicks: Int = 1 @Volatile private var rightPaceTicks: Int = 1 - // 各声道上次成功发送所处的拍号;-1 表示尚未发送过(立即可发)。用于按 paceTicks 控制每声道下发频率。 - private var lastLeftSendTick = -1L - private var lastRightSendTick = -1L - // 左右声道各自的下行缓存:主队列取出的合包先按声道分流到这里,再由发送线程按各自档位取出下发。 - // 降速档位会在此积压;仅在音频发送线程内访问,无需额外加锁。 - private val leftHoldBuffer = ArrayDeque() - private val rightHoldBuffer = ArrayDeque() - // 缓存安全上限(纯防 OOM 兜底):超限时丢弃最旧包并告警;单包约 165B,10000 包≈1.6MB。 + // 声道缓存安全上限(纯防 OOM 兜底):现在缓存元素是单帧 opus,超限丢最旧帧并告警。 private const val CHANNEL_HOLD_BUFFER_MAX = 10000 - // 左右同一拍都满足发送条件时的轮转标记,保证两声道公平、避免一方饿死。 - private var holdDrainPreferLeft = true + // ---- 下行合包(已从 OpusAudioManager 上移到 BleService) ---- + // 每包合并的 opus 帧数(默认5,调试界面可调 1..20):发送线程每拍从声道队列取最多 bundleFrameCount 帧, + // 不足则用静音帧补齐、队列空则不发;组包布局 [4B 序号(小端)] + [1B 声道] + [N × opus 帧]。 + @Volatile + private var bundleFrameCount = 5 + // 合包头部字节数:bytes0-3 序号(小端 uint32) + byte4 声道 + private const val BUNDLE_HEADER_SIZE = 5 + // 左右声道各自的下行包序号(每成功发一包自增),供设备侧核对连续性;仅各自声道线程访问 + private var leftPacketSeq = 0 + private var rightPacketSeq = 0 + // 补齐用静音 opus 帧模板(OpusAudioManager 运行时编码零PCM生成,经 onSilenceFrameReady 送来);未就绪时为 null + @Volatile + private var silenceFrame: ByteArray? = null // ---- F4/F5 解码空余字节数 → 下发档位:简单三档(可在通话调试界面运行时调整) ---- // 每档 = (剩余字节数下限 门限, 发送间隔 sendMs)。设备上报某声道解码空余 freeBytes 后, // 从第一档往下取第一个满足 freeBytes>=门限 的档,用该档发送间隔下发(内部按 audioSendIntervalNormal 换算为拍数)。 @@ -354,11 +365,16 @@ object BleService { notifyAudioDataReceived(data, channel) } - override fun onAudioDataEncoded(data: ByteArray) { - // 将编码后的音频数据加入发送队列 - addAudioDataToSendQueue(data) + override fun onAudioFrameEncoded(channelPrefix: Byte, frame: ByteArray) { + // 单帧入对应声道队列,由发送线程按拍取 bundleFrameCount 帧合包下发 + enqueueOpusFrame(channelPrefix, frame) } - + + override fun onSilenceFrameReady(frame: ByteArray) { + // 缓存静音帧模板,用于合包时不足帧补齐 + silenceFrame = frame + } + override fun onDecodeStreamStateChanged(isStarted: Boolean) { Log.d(TAG, "Opus解码流状态变化: $isStarted") } @@ -1500,136 +1516,131 @@ object BleService { } /** - * 将一包已成型的音频合帧加入下行发送队列 - * 包格式:[4B 序号(大端 uint32)] + [1B 声道(0=左/1=右)] + [5 × opus(40B)] - * 双声道独立编码模式下,每次回调送来的就是一整包,直接入队,不再二次切片 - * @param data 待发送的合帧包 + * 将一帧 opus 按声道加入对应下行队列。合包(取 N 帧/补静音/组头)在发送线程按拍完成。 + * 重连窗口内不丢弃:继续入队暂存,回连成功后续发;仅超 OOM 兜底上限才丢最旧帧。 + * @param channelPrefix 声道标识(CHANNEL_LEFT=0 / CHANNEL_RIGHT=1) + * @param frame 一帧 opus 编码数据 */ - private fun addAudioDataToSendQueue(data: ByteArray) { - if (data.isEmpty()) return - // 重连窗口内不丢弃:继续入队,由发送线程在暂停期间抽干到左右声道缓存暂存,回连成功后按节拍继续下发。 - if (!audioSendQueue.offer(data)) { - Log.w(TAG, "音频发送队列已满,丢弃当前包 size=${data.size}") + private fun enqueueOpusFrame(channelPrefix: Byte, frame: ByteArray) { + if (frame.isEmpty()) return + val buffer = when (channelPrefix) { + OpusAudioManager.CHANNEL_RIGHT -> rightSendBuffer + // 左声道及未知声道统一归入左队列 + else -> leftSendBuffer } + while (buffer.size >= CHANNEL_HOLD_BUFFER_MAX) { + buffer.pollFirst() + CallLog.w(TAG, "声道积压缓存超上限($CHANNEL_HOLD_BUFFER_MAX),丢弃最旧帧 ch=$channelPrefix") + } + buffer.offerLast(frame) } /** - * 启动音频数据发送线程(优化版本) + * 启动左右两条独立发送线程(各发一路声道),右声道相对左声道错开 RIGHT_SEND_THREAD_STAGGER_MS 启动, + * 避免两路同一时刻抢 BLE 写。每条线程按自身声道的 paceTicks 控制下发频率,互不影响。 */ private fun startAudioSendThread() { - if (audioSendThread?.isAlive == true) { + if (leftSendThread?.isAlive == true || rightSendThread?.isAlive == true) { Log.d(TAG, "音频发送线程已在运行") return } - isAudioSending.set(true) - audioSendThread = Thread { - CallLog.i(TAG, "音频发送线程已启动") + leftSendThread = createChannelSendThread(isLeft = true, startOffsetMs = 0L).apply { start() } + rightSendThread = createChannelSendThread(isLeft = false, startOffsetMs = RIGHT_SEND_THREAD_STAGGER_MS).apply { start() } + CallLog.i(TAG, "音频发送线程已启动(左/右双线程,右错开 ${RIGHT_SEND_THREAD_STAGGER_MS}ms)") + } - // 固定 40ms 一拍的发送节拍。每一拍做一次"发不发/发哪路"的决策;左右各按自身档位(paceTicks) - // 控制频率(40/80/160ms),定时器始终稳定在 audioSendIntervalNormal 一拍,下发间隔不随 poll/发送耗时抖动。 - var nextTickAt = System.currentTimeMillis() + /** + * 构建单声道发送线程:固定 audioSendIntervalNormal 一拍(睡到绝对时刻,不随耗时抖动), + * 每拍按该声道 paceTicks 判断是否到点,到点则从该声道缓存取队首一包下发。 + * 重连暂停期间只跳过下发、不清缓存(译音继续积压待回连续发)。 + * @param isLeft true=左声道 / false=右声道 + * @param startOffsetMs 相对启动时刻的错开毫秒(右声道用来与左声道错峰) + */ + private fun createChannelSendThread(isLeft: Boolean, startOffsetMs: Long): Thread { + val chName = if (isLeft) "L" else "R" + return Thread { + CallLog.i(TAG, "音频发送线程[$chName]已启动 offset=${startOffsetMs}ms") + var nextTickAt = System.currentTimeMillis() + startOffsetMs var tick = 0L - + var lastSendTick = -1L + // 首拍错开:右线程先睡 startOffsetMs 再进入循环,使两路发送时刻相互错峰 + val initDelay = nextTickAt - System.currentTimeMillis() + if (initDelay > 0) { + try { Thread.sleep(initDelay) } catch (e: InterruptedException) { return@Thread } + } while (isAudioSending.get() && !Thread.currentThread().isInterrupted) { try { - // 0) 重连窗口:暂停下发但**不丢译音**——仍抽干主队列到左右声道缓存暂存 - // (缓存超大兜底 CHANNEL_HOLD_BUFFER_MAX 才丢最旧防 OOM),只是本拍不写 BLE; - // 待回连成功恢复后按节拍把暂存的译音继续发出。 - if (isDownlinkPaused) { - while (true) { - val more = audioSendQueue.poll() ?: break - routeToHoldBuffer(more) + // 重连暂停期间只跳过下发,不动缓存(译音继续积压,回连后续发) + if (!isDownlinkPaused) { + val buffer = if (isLeft) leftSendBuffer else rightSendBuffer + val paceTicks = if (isLeft) leftPaceTicks else rightPaceTicks + if (buffer.isNotEmpty() && (lastSendTick < 0 || tick - lastSendTick >= paceTicks)) { + if (trySendChannelChunk(isLeft)) lastSendTick = tick } - nextTickAt += audioSendIntervalNormal - tick++ - val sleepMs = nextTickAt - System.currentTimeMillis() - if (sleepMs > 0) Thread.sleep(sleepMs) else nextTickAt = System.currentTimeMillis() - continue - } - - // 1) 非阻塞抽干主队列,按声道分流到左右缓存(不在热路径阻塞,保证节拍稳定)。 - // 合包布局:[4B 序号(小端)] + [1B 声道(0=左/1=右)] + [N × opus] - while (true) { - val more = audioSendQueue.poll() ?: break - routeToHoldBuffer(more) - } - - // 2) 本拍最多发一包:左右各按自身档位(paceTicks)判断是否到点;两路都到点时按 holdDrainPreferLeft - // 轮转二选一,另一路顺延到下一拍(即"一拍一包、左右交替"的下发方式)。 - val leftDue = leftHoldBuffer.isNotEmpty() && - (lastLeftSendTick < 0 || tick - lastLeftSendTick >= leftPaceTicks) - val rightDue = rightHoldBuffer.isNotEmpty() && - (lastRightSendTick < 0 || tick - lastRightSendTick >= rightPaceTicks) - val pickLeft: Boolean? = when { - leftDue && rightDue -> { holdDrainPreferLeft = !holdDrainPreferLeft; holdDrainPreferLeft } - leftDue -> true - rightDue -> false - else -> null - } - when (pickLeft) { - true -> if (trySendChannelChunk(true)) lastLeftSendTick = tick - false -> if (trySendChannelChunk(false)) lastRightSendTick = tick - null -> {} } - - // 3) 固定 40ms 节拍:无论本拍是否发送都睡到下一拍绝对时刻,保证下发节拍稳定、不随耗时抖动。 + // 固定一拍:睡到下一拍绝对时刻,节拍稳定不随发送耗时抖动 nextTickAt += audioSendIntervalNormal tick++ val sleepMs = nextTickAt - System.currentTimeMillis() if (sleepMs > 0) { Thread.sleep(sleepMs) } else { - // 落后于节拍(发送/分流超时),重新对齐,避免之后疯狂追发 + // 落后于节拍(发送超时),重新对齐,避免之后疯狂追发 nextTickAt = System.currentTimeMillis() } } catch (e: InterruptedException) { - Log.d(TAG, "音频发送线程被中断") + Log.d(TAG, "音频发送线程[$chName]被中断") break } catch (e: Exception) { - Log.e(TAG, "音频发送线程异常: ${e.message}", e) + Log.e(TAG, "音频发送线程[$chName]异常: ${e.message}", e) } } - Log.i(TAG, "音频发送线程已停止") - }.apply { - name = "AudioSendThread" - start() - } + CallLog.i(TAG, "音频发送线程[$chName]已停止") + }.apply { name = "AudioSendThread-$chName" } } /** - * 从指定声道缓存取队首一包下发,返回是否成功写入。仅在音频发送线程内调用。 - * - 成功:已塞入底层发送缓冲,写调试录音并打印发送日志(Δ/包序/档位),返回 true; - * - 失败(拥塞):原包放回队头、记一次拥塞告警,返回 false,调用方据此不更新该声道发送拍号、下一拍重试,绝不丢包。 + * 从指定声道队列取最多 bundleFrameCount 帧,不足用静音帧补齐,组包后下发。仅在该声道发送线程内调用。 + * - 队列空:不发,返回 false; + * - 成功:写调试录音、打印发送日志、推进该声道包序,返回 true; + * - 失败(拥塞):取出的真实帧按原序放回队头(静音补齐帧丢弃、下拍重组),返回 false,下一拍重试,绝不丢帧。 */ private fun trySendChannelChunk(isLeft: Boolean): Boolean { - val buffer = if (isLeft) leftHoldBuffer else rightHoldBuffer - if (buffer.isEmpty()) return false - val audioChunk = buffer.removeFirst() - // 解析包头:声道(byte[4]) 与 序号(byte[0..3] 小端),供日志与重试归位使用 - val channelByte = if (audioChunk.size >= 5) audioChunk[4] else null - val seq = if (audioChunk.size >= 5) - ((audioChunk[3].toInt() and 0xFF) shl 24) or - ((audioChunk[2].toInt() and 0xFF) shl 16) or - ((audioChunk[1].toInt() and 0xFF) shl 8) or - (audioChunk[0].toInt() and 0xFF) - else -1 - val chTag = when (channelByte) { - OpusAudioManager.CHANNEL_LEFT -> "L" - OpusAudioManager.CHANNEL_RIGHT -> "R" - else -> "?($channelByte)" - } + val buffer = if (isLeft) leftSendBuffer else rightSendBuffer + val n = bundleFrameCount.coerceIn(1, 20) + // 取最多 n 帧真实数据 + val realFrames = ArrayList(n) + while (realFrames.size < n) { + val f = buffer.pollFirst() ?: break + realFrames.add(f) + } + if (realFrames.isEmpty()) return false // 队列空则不发 + + val channel = if (isLeft) OpusAudioManager.CHANNEL_LEFT else OpusAudioManager.CHANNEL_RIGHT + val chTag = if (isLeft) "L" else "R" + val seq = if (isLeft) leftPacketSeq else rightPacketSeq + + // 组包帧 = 真实帧 + 静音帧补齐到 n(静音帧未就绪时只发真实帧数,罕见的启动窗口) + val frames = ArrayList(n) + frames.addAll(realFrames) + val silence = silenceFrame + var padCount = 0 + if (silence != null) { + while (frames.size < n) { frames.add(silence); padCount++ } + } + val packet = buildBundlePacket(seq, channel, frames) // 同步写入,拿到真实的成功/失败:false=底层发送缓冲已满(拥塞),作为背压信号 - val sent = sendAudioChunkBlocking(audioChunk) + val sent = sendAudioChunkBlocking(packet) if (!sent) { - // 写失败:放回原声道队头,下一拍重试,绝不丢弃、声道内顺序不变 - buffer.addFirst(audioChunk) + // 写失败:真实帧按原顺序放回队头(倒序 offerFirst),下一拍重试;静音补齐帧丢弃,下次按新帧重组 + for (i in realFrames.indices.reversed()) buffer.offerFirst(realFrames[i]) writeFailRetryCount++ if (writeFailRetryCount == 1 || writeFailRetryCount % 50 == 0) { CallLog.w(TAG, "下行写入拥塞,重试中 ch=$chTag seq=$seq 连续失败=$writeFailRetryCount " + - "hold(L=${leftHoldBuffer.size},R=${rightHoldBuffer.size})") + "buf(L=${leftSendBuffer.size},R=${rightSendBuffer.size})") } return false } @@ -1637,80 +1648,77 @@ object BleService { CallLog.i(TAG, "下行写入已恢复,之前连续失败=$writeFailRetryCount") writeFailRetryCount = 0 } - // 调试录音文件按声道分别保存(左右各自独立编码,混写无法解码听辨): - // 剥离 5 字节包头后,按声道标记写入对应声道文件(仅成功时写,避免重试重复) - if (audioChunk.size > 5) { - val opusPayload = audioChunk.copyOfRange(5, audioChunk.size) - when (channelByte) { + // 成功后推进该声道包序 + if (isLeft) leftPacketSeq++ else rightPacketSeq++ + + // 调试录音:剥离 5 字节包头后按声道写入(含静音补齐部分,听辨即静音) + if (packet.size > BUNDLE_HEADER_SIZE) { + val opusPayload = packet.copyOfRange(BUNDLE_HEADER_SIZE, packet.size) + when (channel) { OpusAudioManager.CHANNEL_LEFT -> recordfile1Left?.saveAudioDataToWav(opusPayload) OpusAudioManager.CHANNEL_RIGHT -> recordfile1Right?.saveAudioDataToWav(opusPayload) - else -> {} - } - } - // 发送成功后才推进发送节拍统计与声道包序 - if (audioChunk.size >= 5) { - val now = System.currentTimeMillis() - val deltaMs = if (lastAudioSendTime < 0) 0 else now - lastAudioSendTime - lastAudioSendTime = now - // 按声道统计各自的包序增量:正常应为 +1,出现非 +1 即说明该声道丢包或乱序 - val seqGap = when (channelByte) { - OpusAudioManager.CHANNEL_LEFT -> { - val gap = if (lastLeftSeq < 0) 1L else seq.toLong() - lastLeftSeq - lastLeftSeq = seq.toLong() - gap - } - OpusAudioManager.CHANNEL_RIGHT -> { - val gap = if (lastRightSeq < 0) 1L else seq.toLong() - lastRightSeq - lastRightSeq = seq.toLong() - gap - } - else -> 0L - } - val gapTag = if (seqGap == 1L) "" else " !gap=$seqGap" - val sendLogLine = "音频下行发送 ch=$chTag seq=$seq Δ=${deltaMs}ms size=${audioChunk.size}B " + - "lastSeq(L=$lastLeftSeq,R=$lastRightSeq)$gapTag pace(L=${leftPaceTicks}x,R=${rightPaceTicks}x) " + - "hold(L=${leftHoldBuffer.size},R=${rightHoldBuffer.size}) queue=${audioSendQueue.size}" - // 每包都落文件:推送节拍(Δ)、档位(pace)、背压(hold/queue)、丢包(gap) 是控流排查的核心信号,需完整记录 - CallLog.i(TAG, sendLogLine) - - // [LAT-TRACE] 点3:右声道(对方译音)首包写到 BLE 耳机(每句首包)。 - // 与 azure_speech 的 [LAT-TRACE] 点2(app收到译音) 同时钟相减 = app 内部下行耗时。 - if (channelByte == OpusAudioManager.CHANNEL_RIGHT) { - if (now - latLastBleRightTs > 700) { - CallLog.i(TAG, "[LAT-TRACE] 3.译音写到BLE(右/对方) ts=$now seq=$seq") - } - latLastBleRightTs = now } } + + // 发送节拍统计与包序日志 + val now = System.currentTimeMillis() + val deltaMs = if (lastAudioSendTime < 0) 0 else now - lastAudioSendTime + lastAudioSendTime = now + val seqGap = when (channel) { + OpusAudioManager.CHANNEL_LEFT -> { + val gap = if (lastLeftSeq < 0) 1L else seq.toLong() - lastLeftSeq + lastLeftSeq = seq.toLong(); gap + } + OpusAudioManager.CHANNEL_RIGHT -> { + val gap = if (lastRightSeq < 0) 1L else seq.toLong() - lastRightSeq + lastRightSeq = seq.toLong(); gap + } + else -> 0L + } + val gapTag = if (seqGap == 1L) "" else " !gap=$seqGap" + CallLog.i(TAG, "音频下行发送 ch=$chTag seq=$seq Δ=${deltaMs}ms 帧=${realFrames.size}+静音$padCount size=${packet.size}B " + + "lastSeq(L=$lastLeftSeq,R=$lastRightSeq)$gapTag pace(L=${leftPaceTicks}x,R=${rightPaceTicks}x) " + + "buf(L=${leftSendBuffer.size},R=${rightSendBuffer.size})") + + // [LAT-TRACE] 点3:右声道(对方译音)首包写到 BLE(每句首包) + if (channel == OpusAudioManager.CHANNEL_RIGHT) { + if (now - latLastBleRightTs > 700) { + CallLog.i(TAG, "[LAT-TRACE] 3.译音写到BLE(右/对方) ts=$now seq=$seq") + } + latLastBleRightTs = now + } return true } /** - * 将一包数据按其声道标记分流到对应的下行缓存,由发送线程按各自档位取出下发。 - * 仅作为安全阀:当某声道缓存超过 CHANNEL_HOLD_BUFFER_MAX 时丢弃最旧包并告警(正常流控不会触发)。 + * 组下行包:[4B 序号(小端 uint32)] + [1B 声道(0=左/1=右)] + [N 帧 opus 顺次拼接] */ - private fun routeToHoldBuffer(chunk: ByteArray) { - val channelTag = if (chunk.size >= 5) chunk[4] else null - val buffer = when (channelTag) { - OpusAudioManager.CHANNEL_RIGHT -> rightHoldBuffer - // 左声道及无有效声道标记(异常包)统一归入左缓存处理 - else -> leftHoldBuffer + private fun buildBundlePacket(seq: Int, channel: Byte, frames: List): ByteArray { + var payloadSize = 0 + for (f in frames) payloadSize += f.size + val packet = ByteArray(BUNDLE_HEADER_SIZE + payloadSize) + packet[0] = (seq and 0xFF).toByte() + packet[1] = (seq ushr 8 and 0xFF).toByte() + packet[2] = (seq ushr 16 and 0xFF).toByte() + packet[3] = (seq ushr 24 and 0xFF).toByte() + packet[4] = channel + var offset = BUNDLE_HEADER_SIZE + for (f in frames) { + System.arraycopy(f, 0, packet, offset, f.size) + offset += f.size } - if (buffer.size >= CHANNEL_HOLD_BUFFER_MAX) { - buffer.removeFirst() - CallLog.w(TAG, "声道积压缓存超上限($CHANNEL_HOLD_BUFFER_MAX),丢弃最旧包 ch=${channelTag}") - } - buffer.addLast(chunk) + return packet } /** - * 停止音频数据发送线程 + * 停止左右两条音频发送线程并清空缓存 */ private fun stopAudioSendThread() { isAudioSending.set(false) - audioSendThread?.interrupt() + leftSendThread?.interrupt() + rightSendThread?.interrupt() resetDownlinkBuffers() - CallLog.i(TAG, "音频发送线程已停止,队列和缓冲区已清空") + CallLog.i(TAG, "音频发送线程(左/右)已停止,队列和缓冲区已清空") } /** @@ -1718,15 +1726,13 @@ object BleService { * 供停止发送线程(会话结束/清理)时调用;重连暂停期间**不**调用它——那时要保留缓存的译音待回连后续发。 */ private fun resetDownlinkBuffers() { - audioSendQueue.clear() - audioBuffer.clear() // 清空音频缓冲区 - leftHoldBuffer.clear() // 清空左声道下行缓存 - rightHoldBuffer.clear() // 清空右声道下行缓存 - leftPaceTicks = 1 // 档位复位为全速 40ms + leftSendBuffer.clear() // 清空左声道下行缓存队列 + rightSendBuffer.clear() // 清空右声道下行缓存队列 + audioBuffer.clear() // 清空音频缓冲区 + leftPaceTicks = 1 // 档位复位为全速 rightPaceTicks = 1 - lastLeftSendTick = -1L // 重置左右声道发送拍号 - lastRightSendTick = -1L - holdDrainPreferLeft = true + leftPacketSeq = 0 // 复位左右声道下行包序号 + rightPacketSeq = 0 lastAudioSendTime = -1L // 重置发送节拍计时,下次首包 Δ 从 0 开始 lastLeftSeq = -1L // 重置左右声道包序追踪 lastRightSeq = -1L @@ -1746,7 +1752,7 @@ object BleService { if (isDownlinkPaused) return isDownlinkPaused = true CallLog.i(TAG, "[CALL_RECONNECT] 下行音频已暂停(重连中),译音继续缓存待回连后续发 " + - "hold(L=${leftHoldBuffer.size},R=${rightHoldBuffer.size}) queue=${audioSendQueue.size}") + "buf(L=${leftSendBuffer.size},R=${rightSendBuffer.size})") } /** @@ -1755,16 +1761,14 @@ object BleService { */ private fun exitDownlinkPause() { if (!isDownlinkPaused) return - // 复位节拍/档位(不动缓存):从全速开始排空积压,F4/F5 会随即接管调速 + // 复位档位(不动缓存):从全速开始排空积压,F4/F5 会随即接管调速。 + // 两声道线程内的 lastSendTick 会因 tick 已推进而立即到点,无需外部复位。 leftPaceTicks = 1 rightPaceTicks = 1 - lastLeftSendTick = -1L - lastRightSendTick = -1L - holdDrainPreferLeft = true lastAudioSendTime = -1L isDownlinkPaused = false CallLog.i(TAG, "[CALL_RECONNECT] 下行音频已恢复,档位复位全速续发缓存译音 " + - "hold(L=${leftHoldBuffer.size},R=${rightHoldBuffer.size}) queue=${audioSendQueue.size}") + "buf(L=${leftSendBuffer.size},R=${rightSendBuffer.size})") } /** @@ -1851,10 +1855,15 @@ object BleService { Log.e(TAG, "蓝牙连接或特征值未准备就绪") false } else { - ch.value = chunk - val ok = gatt.writeCharacteristic(ch) == true - if (ok) bytesSentInCurrentSecond += chunk.size - ok + // 左右两条发送线程共用同一 callWriteChar(共享可变 .value)与同一 gatt, + // 必须串行化「设值+写入」这一段,否则并发写会互相覆盖 value,发出错误/重复数据。 + // 锁只在提交这一瞬持有(WRITE_TYPE_NO_RESPONSE 不等 ACK),配合右线程 20ms 错开,几乎不互相阻塞。 + synchronized(bleWriteLock) { + ch.value = chunk + val ok = gatt.writeCharacteristic(ch) == true + if (ok) bytesSentInCurrentSecond += chunk.size + ok + } } } catch (e: Exception) { Log.e(TAG, "发送音频数据块异常: ${e.message}", e) @@ -2014,8 +2023,8 @@ object BleService { Log.i(TAG, "[CALL_TRANS_DEBUG] 设置音频下行发送间隔: ${audioSendIntervalNormal}ms") } if (bundleFrameCount != null && bundleFrameCount > 0) { - opusAudioManager.setBundleFrameCount(bundleFrameCount) - Log.i(TAG, "[CALL_TRANS_DEBUG] 设置下行合包帧数: ${opusAudioManager.getBundleFrameCount()}") + this.bundleFrameCount = bundleFrameCount.coerceIn(1, 20) + Log.i(TAG, "[CALL_TRANS_DEBUG] 设置下行合包帧数: ${this.bundleFrameCount}") } // 三档控流:门限允许 0(第三档兜底),故用 >=0 判定;发送间隔要求 >0 if (t1FreeBytes != null && t1FreeBytes >= 0) tier1FreeBytes = t1FreeBytes @@ -2040,7 +2049,7 @@ object BleService { fun getCallTranslationDebugParams(): Map { return mapOf( "sendIntervalMs" to audioSendIntervalNormal, - "bundleFrameCount" to opusAudioManager.getBundleFrameCount(), + "bundleFrameCount" to bundleFrameCount, "tier1FreeBytes" to tier1FreeBytes, "tier1SendMs" to tier1SendMs, "tier2FreeBytes" to tier2FreeBytes, diff --git a/local_plugins/ble_service/android/src/main/kotlin/com/yunqiinnovation/ble_service/OpusAudioManager.kt b/local_plugins/ble_service/android/src/main/kotlin/com/yunqiinnovation/ble_service/OpusAudioManager.kt index bd1c1e30c..2ab962568 100644 --- a/local_plugins/ble_service/android/src/main/kotlin/com/yunqiinnovation/ble_service/OpusAudioManager.kt +++ b/local_plugins/ble_service/android/src/main/kotlin/com/yunqiinnovation/ble_service/OpusAudioManager.kt @@ -41,14 +41,8 @@ class OpusAudioManager { // 双声道下行包声道标识:0=左(本端翻译) / 1=右(对端翻译) const val CHANNEL_LEFT: Byte = 0 const val CHANNEL_RIGHT: Byte = 1 - - // 下行合包默认帧数:左右各自累积 bundleFrameCount 帧 opus 合并为 1 包后才下发,降低 BLE 发送频率 - // 调试界面可在运行时调整;通过 setBundleFrameCount() 修改 - const val DEFAULT_BUNDLE_FRAME_COUNT = 5 - @Volatile - private var bundleFrameCount = DEFAULT_BUNDLE_FRAME_COUNT - // 合包头部:bytes 0-3 序号(小端 uint32),byte 4 声道标识(0=左/1=右) - private const val BUNDLE_HEADER_SIZE = 5 + // 生成静音帧模板时喂入的零 PCM 长度(100ms@16k@16bit),足够编码器产出至少一帧 + private const val SILENCE_PCM_BYTES = 3200 } // ====================================================================================================== @@ -67,11 +61,18 @@ class OpusAudioManager { fun onAudioDataDecoded(data: ByteArray, channel: Int) /** - * 音频数据编码完成回调 - * @param data 编码后的音频数据 + * 单帧 opus 编码完成回调(合包已上移到 BleService:这里只回单帧 + 声道,由 BleService 入对应队列) + * @param channelPrefix 声道标识(CHANNEL_LEFT=0 / CHANNEL_RIGHT=1) + * @param frame 一帧 opus 编码数据 */ - fun onAudioDataEncoded(data: ByteArray) - + fun onAudioFrameEncoded(channelPrefix: Byte, frame: ByteArray) + + /** + * 静音帧模板生成完成回调(运行时用独立编码器编码零 PCM 得到,供 BleService 合包不足帧时补齐) + * @param frame 一帧静音 opus 数据(长度与真实帧一致) + */ + fun onSilenceFrameReady(frame: ByteArray) + /** * 解码流状态变化回调 * @param isStarted 是否已启动 @@ -111,12 +112,10 @@ class OpusAudioManager { private var rightEncodeManager: OpusManager? = null private var isDualEncoding = false - // 左右各自的 opus 帧缓存与发送序号;累满 bundleFrameCount 帧后打包下发 - private val bundleLock = Any() - private val leftBundleFrames = ArrayList(DEFAULT_BUNDLE_FRAME_COUNT) - private val rightBundleFrames = ArrayList(DEFAULT_BUNDLE_FRAME_COUNT) - private var leftBundleSeq = 0 - private var rightBundleSeq = 0 + // 静音帧模板:启动双声道编码时用独立编码器编码零 PCM,捕获首帧作为补齐用静音帧(经 onSilenceFrameReady 交给 BleService) + private var silenceEncodeManager: OpusManager? = null + @Volatile + private var capturedSilence: ByteArray? = null // 上下文:用于双声道编码模式下,分别为左右声道写入 mono PCM WAV 调试文件 private var appContext: Context? = null @@ -172,6 +171,8 @@ class OpusAudioManager { // 双声道独立编码器:左右各自维持独立的 OpusManager leftEncodeManager = OpusManager() rightEncodeManager = OpusManager() + // 静音帧模板专用编码器(仅编码一段零 PCM 取首帧) + silenceEncodeManager = OpusManager() // 启动队列处理 startAudioQueueProcessing() @@ -332,12 +333,6 @@ class OpusAudioManager { Log.d(TAG, "已清理音频数据缓存") } - // 重置合包发送序号 - synchronized(bundleLock) { - leftBundleSeq = 0 - rightBundleSeq = 0 - } - Log.i(TAG, "已停止Opus数据流解码") callback?.onDecodeStreamStateChanged(false) return true @@ -389,8 +384,9 @@ class OpusAudioManager { Log.d(TAG, "准备开始Opus数据流编码") encodeOpusManager?.startEncodeStream(object : OnEncodeStreamCallback { override fun onEncodeStream(data: ByteArray?) { + // 旧单声道编码路径(非通话翻译):单帧当左声道回给 BleService if (data != null) { - callback?.onAudioDataEncoded(data) + callback?.onAudioFrameEncoded(CHANNEL_LEFT, data) } } @@ -438,13 +434,13 @@ class OpusAudioManager { // ====================================================================================================== // 双声道独立编码方法(左右声道各用独立 OpusManager,互不污染预测状态) // 输出格式:每包 = [4B 序号(小端 uint32)] + [1B 声道(0=左/1=右)] + [N × opus 帧] - // 左右各自累积 bundleFrameCount 帧后打包一次,通过 onAudioDataEncoded 下发 + // 每编出一帧即通过 onAudioFrameEncoded(声道, 帧) 回给 BleService;合包(取N帧/补静音/组头)由 BleService 完成。 // ====================================================================================================== /** * 启动双声道独立编码流 * 左右各创建独立 OpusManager;左右各自累积 5 帧 opus 后合包, - * 通过 onAudioDataEncoded 输出一包(含 5 字节包头 + 5 帧 opus),由下游队列按节拍发送 + * 每编出一帧即通过 onAudioFrameEncoded 单帧回给 BleService;同时用独立编码器编码零 PCM 生成静音帧模板 */ fun startDualEncodeStream( hasHeader: Boolean = false, @@ -460,12 +456,7 @@ class OpusAudioManager { return true } - synchronized(bundleLock) { - leftBundleFrames.clear() - rightBundleFrames.clear() - leftBundleSeq = 0 - rightBundleSeq = 0 - } + capturedSilence = null openPcmWavWriters(sampleRate) @@ -479,7 +470,8 @@ class OpusAudioManager { return try { leftEncodeManager?.startEncodeStream(makeOption(), object : OnEncodeStreamCallback { override fun onEncodeStream(data: ByteArray?) { - if (data != null) emitChannelFrame(CHANNEL_LEFT, data) + // 单帧直接回给 BleService 入左队列,合包由 BleService 完成 + if (data != null) callback?.onAudioFrameEncoded(CHANNEL_LEFT, data) } override fun onStart() { Log.i(TAG, "左声道独立编码已开始") } override fun onComplete(outPath: String?) {} @@ -490,7 +482,7 @@ class OpusAudioManager { }) rightEncodeManager?.startEncodeStream(makeOption(), object : OnEncodeStreamCallback { override fun onEncodeStream(data: ByteArray?) { - if (data != null) emitChannelFrame(CHANNEL_RIGHT, data) + if (data != null) callback?.onAudioFrameEncoded(CHANNEL_RIGHT, data) } override fun onStart() { Log.i(TAG, "右声道独立编码已开始") } override fun onComplete(outPath: String?) {} @@ -502,6 +494,9 @@ class OpusAudioManager { isDualEncoding = true callback?.onEncodeStreamStateChanged(true) + // 生成静音帧模板(异步):独立编码器编码一段零 PCM,捕获首帧 → onSilenceFrameReady 交给 BleService 补齐用 + startSilenceFrameCapture(makeOption()) + CallLog.i(TAG, "双声道独立编码流已启动 sampleRate=$sampleRate packetSize=$packetSize") true } catch (e: Exception) { @@ -519,18 +514,41 @@ class OpusAudioManager { isDualEncoding = false try { leftEncodeManager?.stopEncodeStream() } catch (e: Exception) { Log.e(TAG, "停止左声道编码异常", e) } try { rightEncodeManager?.stopEncodeStream() } catch (e: Exception) { Log.e(TAG, "停止右声道编码异常", e) } - synchronized(bundleLock) { - leftBundleFrames.clear() - rightBundleFrames.clear() - leftBundleSeq = 0 - rightBundleSeq = 0 - } + try { silenceEncodeManager?.stopEncodeStream() } catch (e: Exception) { Log.e(TAG, "停止静音编码异常", e) } + capturedSilence = null closePcmWavWriters() callback?.onEncodeStreamStateChanged(false) CallLog.i(TAG, "双声道独立编码流已停止") return true } + /** + * 启动静音帧模板捕获:用独立编码器编码一段零 PCM,捕获输出的首帧作为补齐用静音帧。 + * 只喂一次零 PCM,取首帧即够用;捕获后经 onSilenceFrameReady 交给 BleService。异步,通常数十 ms 内就绪。 + */ + private fun startSilenceFrameCapture(option: OpusOption) { + try { + silenceEncodeManager?.startEncodeStream(option, object : OnEncodeStreamCallback { + override fun onEncodeStream(data: ByteArray?) { + if (data != null && capturedSilence == null) { + capturedSilence = data + callback?.onSilenceFrameReady(data) + CallLog.i(TAG, "静音帧模板已生成 size=${data.size}B") + } + } + override fun onStart() {} + override fun onComplete(outPath: String?) {} + override fun onError(code: Int, message: String?) { + Log.w(TAG, "静音编码器错误 [$code] $message") + } + }) + // 喂一段零 PCM 触发编码;产出首帧即被捕获,后续帧忽略 + silenceEncodeManager?.writeEncodeStream(ByteArray(SILENCE_PCM_BYTES)) + } catch (e: Exception) { + Log.w(TAG, "生成静音帧失败: ${e.message}") + } + } + /** * 是否正在进行双声道独立编码 */ @@ -538,7 +556,7 @@ class OpusAudioManager { /** * 写入左声道 mono PCM(本端翻译音频) - * 累满 bundleFrameCount 帧 opus 后通过 onAudioDataEncoded 下发:[4B 序号] + [1B 声道=0] + [N × opus] + * 每编出一帧即通过 onAudioFrameEncoded(左, 帧) 回给 BleService 入左队列 */ fun writeLeftPcm(pcm: ByteArray) { if (!isDualEncoding) { @@ -551,7 +569,7 @@ class OpusAudioManager { /** * 写入右声道 mono PCM(对端翻译音频) - * 累满 bundleFrameCount 帧 opus 后通过 onAudioDataEncoded 下发:[4B 序号] + [1B 声道=1] + [N × opus] + * 每编出一帧即通过 onAudioFrameEncoded(右, 帧) 回给 BleService 入右队列 */ fun writeRightPcm(pcm: ByteArray) { if (!isDualEncoding) { @@ -562,74 +580,6 @@ class OpusAudioManager { rightEncodeManager?.writeEncodeStream(pcm) } - /** - * 设置下行合包帧数(调试用)。取值范围 1..20,超出自动钳制。 - * 立即生效:后续累满新帧数即打包下发,已缓存的不足帧待补满后随新值打包。 - * @param count 每包累积的 opus 帧数 - */ - fun setBundleFrameCount(count: Int) { - val clamped = count.coerceIn(1, 20) - bundleFrameCount = clamped - Log.i(TAG, "设置下行合包帧数: $clamped") - } - - /** - * 获取当前下行合包帧数 - */ - fun getBundleFrameCount(): Int = bundleFrameCount - - /** - * 单帧入合包缓冲;左右各自累满 bundleFrameCount 帧后整包通过回调下发 - * 输出包布局:[seq(4B 小端 uint32)] + [channel(1B)] + [N 帧 opus 顺次拼接] - */ - private fun emitChannelFrame(channelPrefix: Byte, opusFrame: ByteArray) { - val frameCount = bundleFrameCount - val packet: ByteArray? = synchronized(bundleLock) { - when (channelPrefix) { - CHANNEL_LEFT -> { - leftBundleFrames.add(opusFrame) - if (leftBundleFrames.size >= frameCount) { - val p = buildBundlePacket(leftBundleSeq, CHANNEL_LEFT, leftBundleFrames) - leftBundleSeq++ - leftBundleFrames.clear() - p - } else null - } - CHANNEL_RIGHT -> { - rightBundleFrames.add(opusFrame) - if (rightBundleFrames.size >= frameCount) { - val p = buildBundlePacket(rightBundleSeq, CHANNEL_RIGHT, rightBundleFrames) - rightBundleSeq++ - rightBundleFrames.clear() - p - } else null - } - else -> { - Log.w(TAG, "未知声道前缀: $channelPrefix") - null - } - } - } - if (packet != null) callback?.onAudioDataEncoded(packet) - } - - private fun buildBundlePacket(seq: Int, channel: Byte, frames: List): ByteArray { - var payloadSize = 0 - for (f in frames) payloadSize += f.size - val packet = ByteArray(BUNDLE_HEADER_SIZE + payloadSize) - // 序号小端序:低字节在前 - packet[0] = (seq and 0xFF).toByte() - packet[1] = (seq ushr 8 and 0xFF).toByte() - packet[2] = (seq ushr 16 and 0xFF).toByte() - packet[3] = (seq ushr 24 and 0xFF).toByte() - packet[4] = channel - var offset = BUNDLE_HEADER_SIZE - for (f in frames) { - System.arraycopy(f, 0, packet, offset, f.size) - offset += f.size - } - return packet - } private fun openPcmWavWriters(sampleRate: Int) { val ctx = appContext ?: return