Browse Source

上传ble协议优化

newdev_chengguofeng
liwei1dao 3 months ago
parent
commit
d15b59a8d7
  1. 383
      local_plugins/ble_service/android/src/main/kotlin/com/yunqiinnovation/ble_service/BleService.kt
  2. 172
      local_plugins/ble_service/android/src/main/kotlin/com/yunqiinnovation/ble_service/OpusAudioManager.kt

383
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<Byte>()
// 音频数据发送相关
private val audioSendQueue = LinkedBlockingQueue<ByteArray>()
private var audioSendThread: Thread? = null
// 音频数据发送相关:左右声道各一条独立缓存队列 + 独立发送线程,两线程错开 20ms 启动,各自按声道节拍下发。
// 缓存兼作重连暂停期间的暂存与 F4/F5 降速积压;LinkedBlockingDeque 线程安全(生产者=编码回调线程,消费者=对应发送线程)。
private val leftSendBuffer = LinkedBlockingDeque<ByteArray>()
private val rightSendBuffer = LinkedBlockingDeque<ByteArray>()
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<ByteArray>()
private val rightHoldBuffer = ArrayDeque<ByteArray>()
// 缓存安全上限(纯防 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<ByteArray>(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<ByteArray>(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>): 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<String, Any> {
return mapOf(
"sendIntervalMs" to audioSendIntervalNormal,
"bundleFrameCount" to opusAudioManager.getBundleFrameCount(),
"bundleFrameCount" to bundleFrameCount,
"tier1FreeBytes" to tier1FreeBytes,
"tier1SendMs" to tier1SendMs,
"tier2FreeBytes" to tier2FreeBytes,

172
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<ByteArray>(DEFAULT_BUNDLE_FRAME_COUNT)
private val rightBundleFrames = ArrayList<ByteArray>(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>): 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

Loading…
Cancel
Save