Browse Source

优化控流逻辑代码

newdev_chengguofeng
liwei1dao 4 months ago
parent
commit
58d86f98b4
  1. 3
      .gitignore
  2. 7
      local_plugins/agent_service/android/src/main/kotlin/com/yunqiinnovation/agent_service/BleAgent.kt
  3. 6
      local_plugins/azure_speech/android/src/main/kotlin/com/yunqiinnovation/azure_speech/AzureSpeechPlugin.kt
  4. 10
      local_plugins/azure_speech/android/src/main/protos/products/understanding/ast/ast_service.proto
  5. 58
      local_plugins/ble_service/android/src/main/kotlin/com/yunqiinnovation/ble_service/BleCommandSender.kt
  6. 19
      local_plugins/ble_service/android/src/main/kotlin/com/yunqiinnovation/ble_service/BleConst.kt
  7. 279
      local_plugins/ble_service/android/src/main/kotlin/com/yunqiinnovation/ble_service/BleService.kt
  8. 28
      local_plugins/ble_service/android/src/main/kotlin/com/yunqiinnovation/ble_service/BleServicePlugin.kt
  9. 6
      local_plugins/ota/android/src/main/kotlin/com/example/ota/OtaPlugin.kt

3
.gitignore

@ -235,3 +235,6 @@ fastlane/readme.md
.qoder/
ios/Runner.xcodeproj/project.pbxproj
.superpowers/
# Java heap dump files
*.hprof

7
local_plugins/agent_service/android/src/main/kotlin/com/yunqiinnovation/agent_service/BleAgent.kt

@ -206,13 +206,6 @@ object BleAgent : BleService.Callback, AgentServiceListener {
override fun onWakeupVoiceUpgradeFailed() {
}
override fun onCallSleepReceived(channel: Int) {
}
override fun onCallResumeReceived(channel: Int) {
}
//============================================================================================
// AgentServiceListener 接口实现

6
local_plugins/azure_speech/android/src/main/kotlin/com/yunqiinnovation/azure_speech/AzureSpeechPlugin.kt

@ -2067,12 +2067,6 @@ class AzureSpeechPlugin : BleService.Callback, FlutterPlugin, ActivityAware,
override fun onWakeupVoiceUpgradeFailed() {
}
override fun onCallSleepReceived(channel: Int) {
}
override fun onCallResumeReceived(channel: Int) {
}
// TODO: 双端翻译时启用副芯片回调
// override fun onSecondaryChipConnectionStateChangedDetail(state: Int, deviceAddress: String?, reason: String?) { }
// override fun onSecondaryChipConnectionsCountChanged(count: Int, addresses: List<String>) { }

10
local_plugins/azure_speech/android/src/main/protos/products/understanding/ast/ast_service.proto

@ -15,6 +15,16 @@ message ReqParams {
string target_language = 3; // 目标语言
string speaker_id = 4;
// 断句/VAD 控制。字段号沿用 understanding.ReqParams(au_base.proto),端到端 AST 底层复用同一套 ASR VAD 参数。
// proto3 标量默认值(0/空串)不会序列化发送,仅在显式 set 时下发;不设则保持服务端默认行为。
// 注意:字段号需与火山服务端 ast.ReqParams 一致才生效,若服务端 tag 不同则会被忽略(不报错、不影响其它字段)。
int32 end_silence_time = 62; // 句尾静音判停时长(ms)
int32 vad_silence_time = 63; // VAD 静音时长(ms)
string vad_mode = 65; // VAD 模式
int32 vad_segment_duration = 66; // VAD 分段时长(ms)
int32 end_window_size = 67; // 句尾判停窗口(ms):越大越不易"没说完就断句"(核心调节项,火山默认约800)
int32 force_to_speech_time = 68; // 强制判为语音的时长(ms)
data.speech.understanding.Corpus corpus = 100;
}

58
local_plugins/ble_service/android/src/main/kotlin/com/yunqiinnovation/ble_service/BleCommandSender.kt

@ -86,16 +86,11 @@ class BleCommandSender {
fun notifyWakeupSignalReceived()
/**
* 通知声道通话休眠指令接收
* @param channel 休眠声道:BleConst.AUDIO_CHANNEL_LEFT(左) / BleConst.AUDIO_CHANNEL_RIGHT(右)
* 通知设备上报某声道解码缓存空余字节数(F4=左 / F5=右)。App 据此动态调整该声道下发速率。
* @param channel BleConst.AUDIO_CHANNEL_LEFT(左) / BleConst.AUDIO_CHANNEL_RIGHT(右)
* @param freeBytes 该声道解码缓存空余字节数
*/
fun notifyCallSleepReceived(channel: Int)
/**
* 通知声道通话恢复指令接收
* @param channel 恢复声道:BleConst.AUDIO_CHANNEL_LEFT(左) / BleConst.AUDIO_CHANNEL_RIGHT(右)
*/
fun notifyCallResumeReceived(channel: Int)
fun notifyCallDecodeFreeReported(channel: Int, freeBytes: Int)
/**
* 通知开始 AI 单次对话
@ -389,10 +384,8 @@ class BleCommandSender {
BleConst.CODEC_CONTROL_A2DP_PLAY -> "A2DP播放模式"
BleConst.CODEC_CONTROL_CALL_RECORD_PLAY -> "通话记录播放模式"
BleConst.CODEC_CONTROL_ENCODE_ON -> "已打开编码"
BleConst.CODEC_CONTROL_CALL_SLEEP_LEFT -> "左声道通话休眠"
BleConst.CODEC_CONTROL_CALL_SLEEP_RIGHT -> "右声道通话休眠"
BleConst.CODEC_CONTROL_CALL_RESUME_LEFT -> "左声道通话恢复"
BleConst.CODEC_CONTROL_CALL_RESUME_RIGHT -> "右声道通话恢复"
BleConst.CODEC_CONTROL_DECODE_FREE_LEFT -> "左声道解码空余上报"
BleConst.CODEC_CONTROL_DECODE_FREE_RIGHT -> "右声道解码空余上报"
else -> "未知状态($codecStatus)"
}
@ -599,26 +592,23 @@ class BleCommandSender {
}
BleConst.CMD_CONTROL_CODEC -> {
// 设备主动上报编解码控制状态,关心左/右声道通话休眠 (0xF0 / 0xF1) 与恢复 (0xF2 / 0xF3)
// 帧格式: 0xCC 0x05 [len] [codecStatus, channelMode, ...] [crc]
// 设备主动上报编解码控制状态,关心左/右声道解码缓存空余字节数上报 (0xF4 左 / 0xF5 右),用于下行流控
// 帧格式: 0xCC 0x05 [len] [codecStatus, ...] [crc];F4/F5 载荷为 [code][int32 大端]
val codecStatus =
if (notifyData.isNotEmpty()) notifyData[0].toInt() and 0xFF else -1
when (codecStatus) {
BleConst.CODEC_CONTROL_CALL_SLEEP_LEFT -> {
Log.i(TAG, "收到主动上报: 左声道通话休眠 (codecStatus=0xF0)")
callback?.notifyCallSleepReceived(BleConst.AUDIO_CHANNEL_LEFT)
}
BleConst.CODEC_CONTROL_CALL_SLEEP_RIGHT -> {
Log.i(TAG, "收到主动上报: 右声道通话休眠 (codecStatus=0xF1)")
callback?.notifyCallSleepReceived(BleConst.AUDIO_CHANNEL_RIGHT)
}
BleConst.CODEC_CONTROL_CALL_RESUME_LEFT -> {
Log.i(TAG, "收到主动上报: 左声道通话恢复 (codecStatus=0xF2)")
callback?.notifyCallResumeReceived(BleConst.AUDIO_CHANNEL_LEFT)
BleConst.CODEC_CONTROL_DECODE_FREE_LEFT,
BleConst.CODEC_CONTROL_DECODE_FREE_RIGHT -> {
// F4=左 / F5=右,紧跟一个大端 int32 空余字节数;notifyData[0]=code,[1..4]=int32,共需 >=5 字节
val channel = if (codecStatus == BleConst.CODEC_CONTROL_DECODE_FREE_LEFT)
BleConst.AUDIO_CHANNEL_LEFT else BleConst.AUDIO_CHANNEL_RIGHT
if (notifyData.size >= 5) {
val freeBytes = readInt32BE(notifyData, 1)
Log.i(TAG, "收到主动上报: ${if (channel == BleConst.AUDIO_CHANNEL_LEFT) "左" else "右"}声道解码空余=$freeBytes 字节 (code=0x${codecStatus.toString(16)})")
callback?.notifyCallDecodeFreeReported(channel, freeBytes)
} else {
Log.w(TAG, "解码空余上报数据长度不足: ${notifyData.size}(期望>=5),忽略 code=0x${codecStatus.toString(16)}")
}
BleConst.CODEC_CONTROL_CALL_RESUME_RIGHT -> {
Log.i(TAG, "收到主动上报: 右声道通话恢复 (codecStatus=0xF3)")
callback?.notifyCallResumeReceived(BleConst.AUDIO_CHANNEL_RIGHT)
}
else -> {
Log.d(TAG, "收到主动上报编解码控制: codecStatus=0x${codecStatus.toString(16)}")
@ -824,4 +814,14 @@ class BleCommandSender {
}
return crc.toByte()
}
/**
* 从 data[offset] 起按大端读取一个 int32(设备解码空余字节数 F4/F5 采用大端)。
* 若后续固件改为小端,只需把此处移位顺序反转即可,调用方无需改动。
*/
private fun readInt32BE(data: ByteArray, offset: Int): Int =
((data[offset].toInt() and 0xFF) shl 24) or
((data[offset + 1].toInt() and 0xFF) shl 16) or
((data[offset + 2].toInt() and 0xFF) shl 8) or
(data[offset + 3].toInt() and 0xFF)
}

19
local_plugins/ble_service/android/src/main/kotlin/com/yunqiinnovation/ble_service/BleConst.kt

@ -121,14 +121,17 @@ object BleConst {
const val CODEC_CONTROL_A2DP_PLAY = 0xA2
/** mic和dac(音乐或者通话远端)声音 */
const val CODEC_CONTROL_CALL_RECORD_PLAY = 0xA3
/** 左声道通话休眠(休眠期间暂停下发左声道,仅传输右声道,直到收到 F2 恢复指令才恢复) */
const val CODEC_CONTROL_CALL_SLEEP_LEFT = 0xF0
/** 右声道通话休眠(休眠期间暂停下发右声道,仅传输左声道,直到收到 F3 恢复指令才恢复) */
const val CODEC_CONTROL_CALL_SLEEP_RIGHT = 0xF1
/** 左声道通话恢复(收到后立即恢复左声道下发,按序补发休眠期间积压的左声道包,与 F0 配对) */
const val CODEC_CONTROL_CALL_RESUME_LEFT = 0xF2
/** 右声道通话恢复(收到后立即恢复右声道下发,按序补发休眠期间积压的右声道包,与 F1 配对) */
const val CODEC_CONTROL_CALL_RESUME_RIGHT = 0xF3
/**
* 设备主动上报"左声道"解码缓存空余字节数(下行流控)。
* 帧载荷: [0xF4][int32 大端],共 5 字节,单位字节。空余越小说明设备解码越来不及,App 据此把
* 左声道下发节拍在 40ms / 80ms / 160ms 三档间动态切换,缓解设备端解码缓存溢出丢包。
*/
const val CODEC_CONTROL_DECODE_FREE_LEFT = 0xF4
/**
* 设备主动上报"右声道"解码缓存空余字节数(下行流控)。
* 帧载荷: [0xF5][int32 大端],共 5 字节,单位字节。含义与 F4 一致、声道相反。
*/
const val CODEC_CONTROL_DECODE_FREE_RIGHT = 0xF5
/** 打开编码指令 */
const val CODEC_CONTROL_ENCODE_ON = 0xB1
/** 左声道 */

279
local_plugins/ble_service/android/src/main/kotlin/com/yunqiinnovation/ble_service/BleService.kt

@ -21,7 +21,6 @@ import android.os.Handler
import java.util.concurrent.atomic.AtomicBoolean
import android.util.Log
import androidx.annotation.RequiresPermission
import java.util.concurrent.TimeUnit
/**
* BLE服务类:提供蓝牙低功耗设备的扫描、连接和通信功能
@ -71,12 +70,6 @@ object BleService {
// 升级唤醒词失败相关回调
fun onWakeupVoiceUpgradeFailed()
// 声道通话休眠相关回调
fun onCallSleepReceived(channel: Int)
// 声道通话恢复相关回调
fun onCallResumeReceived(channel: Int)
}
// ======================================================================================================
@ -153,8 +146,6 @@ object BleService {
// 延迟追踪:右声道(对方译音)每句首包写到 BLE 的时刻,配合 azure_speech 侧 [LAT-TRACE] 点1/点2
// 相减即得 app 内部"收到译音→写到耳机"耗时;间隔>700ms 视为新一句,避免逐包刷屏
private var latLastBleRightTs = 0L
// 写失败(拥塞)后的退避时长(ms):给底层发送缓冲腾空的时间,再重试同一包,期间绝不丢包
private const val WRITE_RETRY_BACKOFF_MS = 5L
// 音频数据分块发送的常量
private val AUDIO_CHUNK_SIZE = 120 // 每次发送120字节
@ -162,23 +153,33 @@ object BleService {
// 音频下行发送间隔(ms),调试界面可在运行时调整;通过 setCallTranslationDebugParams() 修改
@Volatile
private var audioSendIntervalNormal = DEFAULT_AUDIO_SEND_INTERVAL_NORMAL
// 左/右声道休眠标记(true 表示该声道处于休眠,暂停下发)。
// 收到 F0(左)/F1(右) 时置 true,收到 F2(左)/F3(右) 恢复指令时置 false。
// 休眠期间该声道的包暂不发送(转入积压缓存),仅传输另一声道,恢复后按序补发,不丢包。
// 左/右声道下发节拍倍数(以一个 audioSendIntervalNormal=40ms 拍为单位):
// 1 拍=40ms/包(全速)、4 拍=160ms/包、8 拍=320ms/包(三档)。
// 由设备 F4(左)/F5(右) 上报的解码缓存空余字节数动态切换:空余越少越降速,缓解设备解码溢出丢包。
@Volatile
private var leftChannelSleeping: Boolean = false
private var leftPaceTicks: Int = 1
@Volatile
private var rightChannelSleeping: Boolean = false
// 声道休眠期间该声道的包不丢弃,转入对应积压缓存,待声道唤醒后按 FIFO 顺序补发。
// 仅在音频发送线程内访问,无需额外加锁。
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>()
// 积压缓存安全上限(防止设备持续上报休眠导致内存无限增长)。
// 1s 休眠约积压 16 包,留足余量;超限时丢弃最旧包并告警(仅在异常场景触发)。
// 要求绝对不丢包,故大幅抬高上限作为纯防 OOM 的兜底;单包约 165B,10000 包≈1.6MB。
// 缓存安全上限(纯防 OOM 兜底):超限时丢弃最旧包并告警;单包约 165B,10000 包≈1.6MB。
private const val CHANNEL_HOLD_BUFFER_MAX = 10000
// 左右积压同时可补发时的轮转标记,保证两声道公平补发、避免一方饿死。
// 左右同一拍都满足发送条件时的轮转标记,保证两声道公平、避免一方饿死。
private var holdDrainPreferLeft = true
// ---- F4/F5 解码空余字节数 → 下发档位阈值(带迟滞)----
// 设备解码缓存总量≈1840B、单包160B(≈11.5包),按包步进上报。边界取 160(单包)整数倍,
// 正好落在相邻上报值中间,避免临界抖动。降档敏感(及早减速防溢出)、升档保守(回升够多才提速),
// 中间留 2 包迟滞带,消除边界反复横跳。
// 阈值整体上调(补偿空余上报的网络延迟、留提前量):相对初始已上调 4 档(+640B≈4包),更早降速、更晚升速。
private const val DECODE_FREE_CRIT_DOWN = 1120 // 空余 ≤7包:紧急降到 320ms(8x)
private const val DECODE_FREE_CRIT_UP = 1440 // 处于 320ms 时,回升到 ≥9包 才升回 160ms(4x)
private const val DECODE_FREE_NORMAL_DOWN = 1920 // 处于 40ms 时,跌到 ≤12包 才降到 160ms(4x)
private const val DECODE_FREE_NORMAL_UP = 2880 // 升到 ≥18包 才回 40ms(1x) 全速
// 音频数据缓冲区,用于累积数据到80字节再发送
private val audioBuffer = mutableListOf<Byte>()
// 重发机制相关常量
@ -383,19 +384,11 @@ object BleService {
}
/**
* 通知声道通话休眠指令接收
*/
override fun notifyCallSleepReceived(channel: Int) {
// 使用外部类的方法来处理声道通话休眠指令,避免无限递归
this@BleService.notifyCallSleepReceived(channel)
}
/**
* 通知声道通话恢复指令接收
* 通知设备上报某声道解码缓存空余字节数(F4=左 / F5=右),用于下行流控调速
*/
override fun notifyCallResumeReceived(channel: Int) {
// 使用外部类的方法来处理声道通话恢复指令,避免无限递归
this@BleService.notifyCallResumeReceived(channel)
override fun notifyCallDecodeFreeReported(channel: Int, freeBytes: Int) {
// 使用外部类的方法处理,避免无限递归
this@BleService.onDecodeFreeReported(channel, freeBytes)
}
/**
@ -1038,9 +1031,9 @@ object BleService {
} else {
Log.w(TAG, "通话音频服务特征未完整找到")
}
// 在服务发现完成后,请求最大MTU
//requestMaxMtu(g)
// 服务发现完成后,请求高优先级连接参数(最短连接间隔),提升 BLE 双向吞吐、
// 缓解上下行互相挤占导致的收发失败。连接断开后系统自动复位,无需手动还原。
requestHighConnectionPriority(g)
}
override fun onCharacteristicChanged(g: BluetoothGatt, c: BluetoothGattCharacteristic) {
@ -1306,6 +1299,22 @@ object BleService {
return true
}
/**
* 请求高优先级连接参数(CONNECTION_PRIORITY_HIGH,最短连接间隔)。
* BLE 是半双工,上下行共享同一连接事件;下行音频写入过多会挤占设备上行 notify 的时隙,
* 表现为"设备端发数据失败"。提高连接优先级可缩短连接间隔、增大双向带宽,缓解互相挤占。
* 代价是更耗电,适合实时音频通话场景;连接断开后系统自动复位,无需手动还原。
* 注意:此调用不占用 GATT 操作队列(走 HCI 连接参数更新),不会干扰正在进行的 CCCD 串行注册。
*/
private fun requestHighConnectionPriority(gatt: BluetoothGatt) {
try {
val ok = gatt.requestConnectionPriority(BluetoothGatt.CONNECTION_PRIORITY_HIGH)
Log.i(TAG, "请求高优先级连接(CONNECTION_PRIORITY_HIGH): ${if (ok) "已发起" else "失败"}")
} catch (e: Exception) {
Log.e(TAG, "请求高优先级连接异常: ${e.message}", e)
}
}
/**
* 请求最大MTU值
* @param gatt BluetoothGatt实例
@ -1419,40 +1428,71 @@ object BleService {
audioSendThread = Thread {
Log.i(TAG, "音频发送线程已启动")
// 固定节拍的下一拍绝对时刻(ms)。每发完一包就 += 间隔,按这个时刻补偿性 sleep,
// 把 poll 等待和发送耗时一起抵消掉,保证实际发送周期稳定贴近 audioSendIntervalNormal。
// 0 表示尚未起拍,由首包发送时初始化。
var nextSendAt = 0L
// 固定 40ms 一拍的发送节拍。每一拍做一次"发不发/发哪路"的决策;左右各按自身档位(paceTicks)
// 控制频率(40/80/160ms),定时器始终稳定在 audioSendIntervalNormal 一拍,下发间隔不随 poll/发送耗时抖动。
var nextTickAt = System.currentTimeMillis()
var tick = 0L
while (isAudioSending.get() && !Thread.currentThread().isInterrupted) {
try {
// 1) 从主队列取新包,按声道分流到积压缓存。
// 合包布局:[4B 序号(大端)] + [1B 声道(0=左/1=右)] + [N × opus]
// 阻塞至多 20ms 等待新数据,避免空转;随后把队列里已到达的剩余包一次性分流。
audioSendQueue.poll(20, TimeUnit.MILLISECONDS)?.let { first ->
routeToHoldBuffer(first)
// 1) 非阻塞抽干主队列,按声道分流到左右缓存(不在热路径阻塞,保证节拍稳定)。
// 合包布局:[4B 序号(小端)] + [1B 声道(0=左/1=右)] + [N × opus]
while (true) {
val more = audioSendQueue.poll() ?: break
routeToHoldBuffer(more)
}
}
// 2) 选出一包可发送的数据:声道未休眠且缓存非空。
// 左右都可发时按 holdDrainPreferLeft 轮转,保证两声道公平补发、避免一方饿死。
val leftReady = leftHoldBuffer.isNotEmpty() && !leftChannelSleeping
val rightReady = rightHoldBuffer.isNotEmpty() && !rightChannelSleeping
val audioChunk: ByteArray? = when {
leftReady && rightReady -> {
holdDrainPreferLeft = !holdDrainPreferLeft
if (holdDrainPreferLeft) leftHoldBuffer.removeFirst()
else rightHoldBuffer.removeFirst()
// 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, "音频发送线程被中断")
break
} catch (e: Exception) {
Log.e(TAG, "音频发送线程异常: ${e.message}", e)
}
}
Log.i(TAG, "音频发送线程已停止")
}.apply {
name = "AudioSendThread"
start()
}
leftReady -> leftHoldBuffer.removeFirst()
rightReady -> rightHoldBuffer.removeFirst()
else -> null // 两声道都在休眠或都无积压,本轮不发送
}
if (audioChunk != null) {
/**
* 从指定声道缓存取队首一包下发,返回是否成功写入。仅在音频发送线程内调用。
* - 成功:已塞入底层发送缓冲,写调试录音并打印发送日志(Δ/包序/档位),返回 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)
@ -1469,22 +1509,16 @@ object BleService {
// 同步写入,拿到真实的成功/失败:false=底层发送缓冲已满(拥塞),作为背压信号
val sent = sendAudioChunkBlocking(audioChunk)
if (!sent) {
// 写失败:把这一包放回原声道队头,下轮重试,绝不丢弃(满足“绝对不丢包”)。
// 包仍在队头,声道内顺序不变。
if (channelByte == OpusAudioManager.CHANNEL_RIGHT) rightHoldBuffer.addFirst(audioChunk)
else leftHoldBuffer.addFirst(audioChunk)
// 写失败:放回原声道队头,下一拍重试,绝不丢弃、声道内顺序不变
buffer.addFirst(audioChunk)
writeFailRetryCount++
// 节流告警,避免拥塞期刷屏
if (writeFailRetryCount == 1 || writeFailRetryCount % 50 == 0) {
Log.w(TAG, "下行写入拥塞,重试中 ch=$chTag seq=$seq 连续失败=$writeFailRetryCount " +
"hold(L=${leftHoldBuffer.size},R=${rightHoldBuffer.size})")
}
// 退避,给底层缓冲腾空时间;重置节拍,待恢复后重新起拍
nextSendAt = 0L
Thread.sleep(WRITE_RETRY_BACKOFF_MS)
} else {
return false
}
if (writeFailRetryCount > 0) {
Log.i(TAG, "下行写入已恢复,之前连续失败=$writeFailRetryCount")
writeFailRetryCount = 0
@ -1514,7 +1548,7 @@ object BleService {
}
val gapTag = if (seqGap == 1L) "" else " !gap=$seqGap"
Log.i(TAG, "音频下行发送 ch=$chTag seq=$seq Δ=${deltaMs}ms size=${audioChunk.size}B " +
"lastSeq(L=$lastLeftSeq,R=$lastRightSeq)$gapTag " +
"lastSeq(L=$lastLeftSeq,R=$lastRightSeq)$gapTag pace(L=${leftPaceTicks}x,R=${rightPaceTicks}x) " +
"hold(L=${leftHoldBuffer.size},R=${rightHoldBuffer.size}) queue=${audioSendQueue.size}")
// [LAT-TRACE] 点3:右声道(对方译音)首包写到 BLE 耳机(每句首包)。
@ -1526,38 +1560,12 @@ object BleService {
latLastBleRightTs = now
}
}
// 固定节拍:睡到 nextSendAt 这个绝对时刻,把 poll/发送耗时算进同一周期,间隔不随它们抖动。
val nowMs = System.currentTimeMillis()
if (nextSendAt == 0L) nextSendAt = nowMs
nextSendAt += audioSendIntervalNormal
val sleepMs = nextSendAt - nowMs
if (sleepMs > 0) {
Thread.sleep(sleepMs)
} else {
// 落后于节拍(发送/积压处理超时),重新对齐,避免之后疯狂追发
nextSendAt = System.currentTimeMillis()
}
}
}
} catch (e: InterruptedException) {
Log.d(TAG, "音频发送线程被中断")
break
} catch (e: Exception) {
Log.e(TAG, "音频发送线程异常: ${e.message}", e)
}
}
Log.i(TAG, "音频发送线程已停止")
}.apply {
name = "AudioSendThread"
start()
}
return true
}
/**
* 将一包数据按其声道标记分流到对应的积压缓存。
* 休眠声道的包不会丢弃,先入缓存等唤醒后补发;非休眠声道的包同样入缓存,由发送线程统一按节拍取出。
* 仅作为安全阀:当某声道缓存超过 CHANNEL_HOLD_BUFFER_MAX 时丢弃最旧包并告警(正常 1s 休眠不会触发)。
* 将一包数据按其声道标记分流到对应的下行缓存,由发送线程按各自档位取出下发。
* 仅作为安全阀:当某声道缓存超过 CHANNEL_HOLD_BUFFER_MAX 时丢弃最旧包并告警(正常流控不会触发)。
*/
private fun routeToHoldBuffer(chunk: ByteArray) {
val channelTag = if (chunk.size >= 5) chunk[4] else null
@ -1581,10 +1589,13 @@ object BleService {
audioSendThread?.interrupt()
audioSendQueue.clear()
audioBuffer.clear() // 清空音频缓冲区
leftHoldBuffer.clear() // 清空左声道积压缓存
rightHoldBuffer.clear() // 清空右声道积压缓存
leftChannelSleeping = false
rightChannelSleeping = false
leftHoldBuffer.clear() // 清空左声道下行缓存
rightHoldBuffer.clear() // 清空右声道下行缓存
leftPaceTicks = 1 // 档位复位为全速 40ms
rightPaceTicks = 1
lastLeftSendTick = -1L // 重置左右声道发送拍号
lastRightSendTick = -1L
holdDrainPreferLeft = true
lastAudioSendTime = -1L // 重置发送节拍计时,下次首包 Δ 从 0 开始
lastLeftSeq = -1L // 重置左右声道包序追踪
lastRightSeq = -1L
@ -2184,58 +2195,44 @@ object BleService {
}
/**
* 收到设备主动上报的左/右声道通话休眠 (F0/F1) 时调用。
* 将对应声道置为休眠:该声道的包转入积压缓存暂不发送、仅传输另一声道,
* 直到收到对应的恢复指令 (F2/F3) 才恢复并按序补发,不丢包。
* 重复触发幂等,不会叠加。
* 收到设备主动上报的某声道解码缓存空余字节数 (F4=左 / F5=右) 时调用,用于下行流控。
* 空余越少说明设备解码越来不及,则把该声道下发节拍降速(40→80→160ms),缓解解码缓存溢出丢包;
* 空余恢复后再升回全速。仅更新对应声道档位,由发送线程在下一拍按新档位执行。
* @param channel BleConst.AUDIO_CHANNEL_LEFT(左) / BleConst.AUDIO_CHANNEL_RIGHT(右)
* @param freeBytes 该声道解码缓存空余字节数
*/
private fun notifyCallSleepReceived(channel: Int) {
private fun onDecodeFreeReported(channel: Int, freeBytes: Int) {
when (channel) {
BleConst.AUDIO_CHANNEL_LEFT -> {
leftChannelSleeping = true
Log.i(TAG, "收到左声道休眠指令(F0),暂停左声道下发,等待恢复指令(F2)")
val pace = nextPaceTicks(leftPaceTicks, freeBytes)
if (leftPaceTicks != pace) {
Log.i(TAG, "左声道解码空余=$freeBytes 字节,下发档位 ${leftPaceTicks}x→${pace}x (${pace * audioSendIntervalNormal}ms/包)")
}
leftPaceTicks = pace
}
BleConst.AUDIO_CHANNEL_RIGHT -> {
rightChannelSleeping = true
Log.i(TAG, "收到右声道休眠指令(F1),暂停右声道下发,等待恢复指令(F3)")
val pace = nextPaceTicks(rightPaceTicks, freeBytes)
if (rightPaceTicks != pace) {
Log.i(TAG, "右声道解码空余=$freeBytes 字节,下发档位 ${rightPaceTicks}x→${pace}x (${pace * audioSendIntervalNormal}ms/包)")
}
else -> Log.w(TAG, "收到未知声道休眠指令: channel=$channel")
}
for (callback in callbacks) {
try {
callback.onCallSleepReceived(channel)
} catch (e: Exception) {
Log.e(TAG, "分发声道通话休眠回调异常", e)
rightPaceTicks = pace
}
else -> Log.w(TAG, "收到未知声道解码空余上报: channel=$channel")
}
}
/**
* 收到设备主动上报的左/右声道通话恢复 (F2/F3) 时调用。
* 解除对应声道的休眠标记,发送线程随即按 FIFO 顺序补发休眠期间积压的该声道包,不丢包。
* 未处于休眠时收到恢复指令为幂等操作。
* @param channel BleConst.AUDIO_CHANNEL_LEFT(左) / BleConst.AUDIO_CHANNEL_RIGHT(右)
* 带迟滞的三档流控:由当前档位 currentPace 与最新空余 freeBytes 决定新档位(拍)。三档=拍数 1/4/8 → 40/160/320ms/包。
* 降档敏感、升档保守,中间区按当前档位保持,避免边界反复横跳:
* 空余 ≤CRIT_DOWN → 8(320ms);≥NORMAL_UP → 1(40ms);
* 中间区:40ms 档跌破 NORMAL_DOWN 才降 160ms;320ms 档升过 CRIT_UP 才回 160ms;160ms 档维持。
*/
private fun notifyCallResumeReceived(channel: Int) {
when (channel) {
BleConst.AUDIO_CHANNEL_LEFT -> {
leftChannelSleeping = false
Log.i(TAG, "收到左声道恢复指令(F2),恢复左声道下发,补发积压 ${leftHoldBuffer.size} 包")
}
BleConst.AUDIO_CHANNEL_RIGHT -> {
rightChannelSleeping = false
Log.i(TAG, "收到右声道恢复指令(F3),恢复右声道下发,补发积压 ${rightHoldBuffer.size} 包")
}
else -> Log.w(TAG, "收到未知声道恢复指令: channel=$channel")
}
for (callback in callbacks) {
try {
callback.onCallResumeReceived(channel)
} catch (e: Exception) {
Log.e(TAG, "分发声道通话恢复回调异常", e)
}
}
private fun nextPaceTicks(currentPace: Int, freeBytes: Int): Int = when {
freeBytes <= DECODE_FREE_CRIT_DOWN -> 8
freeBytes >= DECODE_FREE_NORMAL_UP -> 1
currentPace == 8 -> if (freeBytes >= DECODE_FREE_CRIT_UP) 4 else 8
currentPace == 1 -> if (freeBytes <= DECODE_FREE_NORMAL_DOWN) 4 else 1
else -> 4
}
/**

28
local_plugins/ble_service/android/src/main/kotlin/com/yunqiinnovation/ble_service/BleServicePlugin.kt

@ -460,31 +460,5 @@ class BleServicePlugin : FlutterPlugin, MethodCallHandler, ActivityAware,
sendEvent(statusEventSink, resultMap, "发送升级失败异常")
}
override fun onCallSleepReceived(channel: Int) {
val channelName = when (channel) {
BleConst.AUDIO_CHANNEL_LEFT -> "左声道"
BleConst.AUDIO_CHANNEL_RIGHT -> "右声道"
else -> "未知声道($channel)"
}
Log.i(TAG, "[Flutter推送] 声道通话休眠: $channelName (channel=$channel)")
sendEvent(
statusEventSink,
mapOf("type" to "callSleep", "channel" to channel),
"发送声道通话休眠事件异常"
)
}
override fun onCallResumeReceived(channel: Int) {
val channelName = when (channel) {
BleConst.AUDIO_CHANNEL_LEFT -> "左声道"
BleConst.AUDIO_CHANNEL_RIGHT -> "右声道"
else -> "未知声道($channel)"
}
Log.i(TAG, "[Flutter推送] 声道通话恢复: $channelName (channel=$channel)")
sendEvent(
statusEventSink,
mapOf("type" to "callResume", "channel" to channel),
"发送声道通话恢复事件异常"
)
}
// 注:F4/F5 解码空余字节数仅用于 native 内部下行流控(动态调整左右下发速率),不再向 Flutter 推送声道休眠/恢复事件。
}

6
local_plugins/ota/android/src/main/kotlin/com/example/ota/OtaPlugin.kt

@ -412,10 +412,4 @@ class OtaPlugin : BleService.Callback, FlutterPlugin, MethodCallHandler {
override fun onWakeupVoiceUpgradeFailed() {
}
override fun onCallSleepReceived(channel: Int) {
}
override fun onCallResumeReceived(channel: Int) {
}
}

Loading…
Cancel
Save