Compare commits

...

2 Commits

Author SHA1 Message Date
Rodger-Wang d85f0fbaed client(ios),docs: 译文积压追赶同步到 iOS 下行发送器;CLAUDE.md 记录 9-22 三条结论 3 weeks ago
Rodger-Wang d1db3e5951 client(android): 通话翻译三处修复——同页第二场聋、call.stop 偶发丢失、译文积压追赶 3 weeks ago
  1. 51
      CLAUDE.md
  2. 16
      apps/client/lib/modules/translation/controllers/translation_controller.dart
  3. 44
      apps/client/local_plugins/azure_speech/android/src/main/kotlin/com/yunqiinnovation/azure_speech/AzureSpeechPlugin.kt
  4. 6
      apps/client/local_plugins/azure_speech/android/src/main/kotlin/com/yunqiinnovation/azure_speech/BesCallPcmBridge.kt
  5. 2
      apps/client/local_plugins/bluetooth_manager/android/src/main/kotlin/com/example/bluetooth_manager/AiAudioDownlink.kt
  6. 63
      apps/client/local_plugins/bluetooth_manager/android/src/main/kotlin/com/example/bluetooth_manager/BluetoothManager.kt
  7. 2
      apps/client/local_plugins/bluetooth_manager/android/src/main/kotlin/com/example/bluetooth_manager/BluetoothManagerPlugin.kt
  8. 388
      apps/client/local_plugins/bluetooth_manager/android/src/main/kotlin/com/example/bluetooth_manager/CallTranslationDownlink.kt
  9. 399
      apps/client/local_plugins/bluetooth_manager/ios/bluetooth_manager/Sources/bluetooth_manager/CallTranslationDownlink.swift

51
CLAUDE.md

@ -850,6 +850,57 @@ xcrun devicectl device copy from --device <UDID> \
`180A`(设备信息)、`180F`(空壳电池)、`00001100-D102-11E1-9B23-00025B00AFAE`(私有,在用)、
`66666666-…`、`86868686-…`。
### Android 的 SPP 控制命令会被协议栈合帧(2026-09-22 ALN-AL80 真机定性,已修)
**这是 iOS 上 GATT 那三条铁律在 Android 上的对应物,坑的形状完全不同。** Android 的
`BluetoothSocket` 写入走的是到蓝牙进程的本地管道,协议栈只在链路有信用(耳机回 credit)时
才从管道取数据,而且一次把**管道里攒着的全部字节合成一个 RFCOMM 帧**发出去。恒玄固件
一帧只解开头那一条命令,排在后面的**静默丢失**。协议栈日志(`adb logcat` 里 `bt_rfcomm:
PORT_WriteData p_len=N`)能直接看到:
- 通话翻译期间下行音频包(84B)经常被合成 168/252/…/1512 字节的帧,即管道里最多攒过 18 包 ≈ 360ms;
- 8 次收尾有 3 次 `call.stop` 丢失:`p_len=4`(asr.stopped + call.stop 合成一帧)、`p_len=256`
(3 包音频 + 两条命令)、`p_len=172`(2 包音频 + 两条命令)。耳机既不回 `BB D2` 也不退通话模式、
**继续推音频**,下一次 `call.start` 的应答变成 `BB D1 07 01 02 03 02`(已在通话模式)——
这就是「点结束没反应 / 耳机没收到结束指令」,偶现是因为取决于点结束那一刻管道里有没有东西。
- 握手里 `AA 06` 紧跟 `AA 09` 的「`AA 09` 无回包」多半也是它(见上文命令号那节的误判记录)。
修法在 [BluetoothManager.kt](apps/client/local_plugins/bluetooth_manager/android/src/main/kotlin/com/example/bluetooth_manager/BluetoothManager.kt):
控制命令(`writeBytes`)走单线程队列,距上一条命令 ≥30ms、**距上一个实时音频包 ≥500ms** 才写
(后者是为了让管道里攒着的音频先发完,命令才能落在一帧的开头),最多等 5s 兜底;下行音频改走
`writeRealtime` 直发。修后 12 轮 0 丢,边播译文边点结束也在 0.5s 内收到 D2。
⚠️ 判据是 `Logger.w` 的 `call.stop 已发出` 之后 ~0.5s 内有没有 `BB D2`——再出现「已在通话模式」
的 `D1 07 01` 应答,先怀疑这里的间隔不够,别去查 Dart。
⚠️ `sendData` 在 Android 上**本来就不等真正写出**(`result.success` 先返回),排队没改变任何调用方语义;
但 Dart 若在发完命令后立刻 `disconnect`,排着队的命令会随链路一起没了。
### 通话翻译 AST 就绪状态:同页第二场必聋(2026-09-22 Android 定性,两端同一个坑)
原生直连桥(`BesCallPcmBridge`)的上行帧只有 `ready==true` 才推给 AST,否则扣在 3 秒环形缓冲里丢掉;
`ready` 只由 `initialize` 末尾的 `markReady()` 置真。**同一页面第二次点开始不会再 `initialize`**
(只走 `startContinuousTranslation`),而 `disable()` 又把 `ready` 清了 → 之后每一场都
`ready:false / downlinkFrames:0`,上行帧照常到(framesA/B 几百上千、丢帧 0),就是不进 AST。
真机六场:进页面后第一场好、同页之后每一场都聋,退出重进又好一场。iOS 9-21 那次只修了
「桥比 initialize 晚建」的竞态(`astSessionReady` 建桥时继承),**没覆盖「桥已存在、重新 enable」**。
现在 Android 两处都修了:就绪状态记在插件上,建桥 / 重新 enable 时都继承,`disable()` 不再清 `ready`。
⚠️ iOS 的 `BesCallPcmBridge.swift` `disable()` 仍会清 `ready`,同页第二场理论上一样聋——
iPhone 上每次都退出重进测就碰不到,要复现就同页连开两场看 `ready`。
### 译文积压追赶:下行裁静音 + WSOLA 变速(2026-09-22,Android 实测,iOS 已同步代码未验)
译文音频是模型**突发**产出的(每 3~6 秒一段文本,音频一次性到),播放只能 1 倍速;说话人不停顿时
译文时长 ≈ 原话时长(实测 A 路 133 秒会话产出 129 秒音频),下行队列只涨不缩——A 路积压到 **6.2 秒**,
用户说的"越说越慢"就是它。**上行没有积压**(桥峰值 8~12 帧),瓶颈只在 `CallTranslationDownlink`。
处理两层,都不丢内容:积压 >1s 裁掉超过 150ms 的静音;>2s 起 WSOLA 变速不变调,4s 达 1.3x 上限。
同量压测:最大积压 6.2s → 3.2s,绝大部分时间 <1.5s,**光裁静音就省了 10%**(TTS 句间停顿),
变速最高只用到 1.11x。兜底整段丢弃(`HARD_CAP_SEC`)默认关。
- 参数与算法两端一份(`CallTranslationDownlink.kt` / `.swift`),改一端要同改另一端;
WSOLA 的数学用 Python 复刻验证过(rate=1.0 逐样本恒等,1.3x 时长 0.773≈1/1.3,无跳变)。
- 收尾一条 `下行积压总账 A{源=..s 播出=..s 最大积压=..s 裁静音=..s 变速追回=..s}`,每 100 包一条
`积压 A=..s B=..s 速率 A=.. B=..`——排查"慢"先看这两行,再决定是模型切句还是播放积压。
- ⚠️ 同一房间用两台手机测试时,B 路(对端通话音)吃到的就是机主自己的声音,`en->zh` 那条腿会把
中文原话原样吐回来并合成中文 TTS 灌回耳朵,听感像"慢了 3 秒又念了一遍",那是测试条件不是 bug。
### 通话翻译的 PCM 已下沉到原生:bluetooth_manager ⇄ azure_speech 直连(2026-09-20)
苹果测试反馈「WiFi 下断断续续、长句基本不行」正是此前记录的技术债触发条件(Dart 链路每秒约 200 次跨界,

16
apps/client/lib/modules/translation/controllers/translation_controller.dart

@ -3025,17 +3025,27 @@ class TranslationController extends GetxController with WidgetsBindingObserver {
final session = _callSession;
_callSession = null;
final alive = session != null && session.state == DeviceConnectionState.ready;
// ⚠️ 这一段全用 warn 级:release 包只落 warning 以上,2026-09-22 真机有一场
// 耳机没回 call.stop 的 ACK(下一场 call.start 应答成了 `BB D1 07 01 …`),
// 而当时这里全是 info,分不清是没发、发了被拒、还是 session 已经不是那个了。
if (!alive) {
Logger.w('Translation',
'[CALL-BRIDGE] 收尾时设备会话不可用(session=${session == null ? "null" : session.state}),'
'asr.stopped / call.stop 都不发');
}
try {
if (alive) await session.invokeFeature(DeviceFeatures.asrStopped);
} catch (_) {}
} catch (e) {
Logger.w('Translation', '[CALL-BRIDGE] asr.stopped 发送失败: $e');
}
unawaited(BackgroundKeepalive.release('callTranslation'));
await _detachBesCallTranslationBridge(); // 关掉两路 sink = 停译文下行
if (!alive) return;
try {
await session.invokeFeature(DeviceFeatures.callStop);
Logger.info('[CALL-BRIDGE] 已退出设备通话模式');
Logger.w('Translation', '[CALL-BRIDGE] call.stop 已发出,等耳机回 BB D2');
} catch (e) {
Logger.error('[CALL-BRIDGE] 退出设备通话模式失败: $e');
Logger.w('Translation', '[CALL-BRIDGE] 退出设备通话模式失败: $e');
}
}

44
apps/client/local_plugins/azure_speech/android/src/main/kotlin/com/yunqiinnovation/azure_speech/AzureSpeechPlugin.kt

@ -162,6 +162,28 @@ class AzureSpeechPlugin : BleService.Callback, FlutterPlugin, ActivityAware,
private var currentAstProvider: String = "azure"
/** 通话翻译 PCM 原生直连桥,按会话开关(见 setNativeCallBridge) */
private var callPcmBridge: BesCallPcmBridge? = null
/**
* AST 两路 helper 是否已就绪(与 iOS 的 astSessionReady 同义)。
*
* ⚠️ **不能只靠桥上那一句 markReady**。就绪只在 initialize 末尾发一次,而桥的
* 生命周期跟它是两条线:
* - 桥比 initialize 晚建(进程里第一场会话):那句打在 null 上静默丢失;
* - 同一页面第二场会话:不再 initialize,只 `startContinuousTranslation` + 重新 enable 桥,
* 再也没人 markReady。
* 两种都表现为「点了开始没反应」:上行帧照常到桥(framesA/B 几百上千、丢帧 0),全扣在
* 环形缓冲里,`ready:false`、`downlinkFrames:0`。2026-09-22 ALN-AL80 真机六场会话:
* 进页面后第一场好、同页之后每一场都聋,退出重进又好一场。
* 所以就绪状态记在插件上,建桥 / 重新 enable 时都从这里继承。
*/
@Volatile private var astSessionReady = false
/** AST 就绪:置标志 + 通知桥。桥还没建也没关系,建/enable 的时候会自己补上。 */
private fun markAstBridgeReady() {
astSessionReady = true
callPcmBridge?.markReady()
}
private fun ensureCallPcmBridge(): BesCallPcmBridge {
callPcmBridge?.let { return it }
val b = BesCallPcmBridge(
@ -173,6 +195,8 @@ class AzureSpeechPlugin : BleService.Callback, FlutterPlugin, ActivityAware,
onOverflow = { dropped -> sendAstEvent(mapOf("type" to "nativeBridgeOverflow", "dropped" to dropped)) },
)
callPcmBridge = b
// 桥比 AST 初始化晚建:把错过的那次 markReady 补回来,否则这一整场会话都不会 ready
if (astSessionReady) b.markReady()
return b
}
// ---------- EMAI 会话(百炼多模态)原生客户端 ----------
@ -1473,6 +1497,7 @@ class AzureSpeechPlugin : BleService.Callback, FlutterPlugin, ActivityAware,
}
"dispose" -> {
astSessionReady = false
callPcmBridge?.onSessionDisposed()
try {
// 让所有等待 delay(1200L) 的旧 init 协程在恢复时直接 return,避免落到已 dispose 的 helper 上
@ -1589,7 +1614,15 @@ class AzureSpeechPlugin : BleService.Callback, FlutterPlugin, ActivityAware,
// Dart 只收 TTS 的 final 标记与门限/丢帧的诊断事件。
"setNativeCallBridge" -> {
val enabled = call.argument<Boolean>("enabled") ?: false
if (enabled) ensureCallPcmBridge().enable() else callPcmBridge?.disable()
if (enabled) {
val b = ensureCallPcmBridge()
b.enable()
// 同页面第二场:helper 早就 initialize 过、不会再发 markReady,这里按插件上的
// 就绪状态补一次(幂等)。见 astSessionReady 的说明。
if (astSessionReady) b.markReady()
} else {
callPcmBridge?.disable()
}
result.success(true)
}
@ -1628,6 +1661,7 @@ class AzureSpeechPlugin : BleService.Callback, FlutterPlugin, ActivityAware,
Log.d(tag, "initializeIntegrated: supportedLanguages=$supportedLanguages")
currentAstProvider = provider
// 原生直连桥:新会话重新刻音色 → 门限复位;helper 起来之前的帧先攒着(见 BesCallPcmBridge)
astSessionReady = false
callPcmBridge?.onSessionInit()
val subscriptionKey = call.argument<String>("subscriptionKey") ?: ""
@ -1731,7 +1765,7 @@ class AzureSpeechPlugin : BleService.Callback, FlutterPlugin, ActivityAware,
callback = callbackB
) ?: false
callPcmBridge?.markReady()
markAstBridgeReady()
result.success(initA && initB)
}
} else if (currentAstProvider == "azure" ) {
@ -1792,7 +1826,7 @@ class AzureSpeechPlugin : BleService.Callback, FlutterPlugin, ActivityAware,
serviceConfig = serviceConfigB,
callback = callbackB
)
callPcmBridge?.markReady()
markAstBridgeReady()
result.success(true)
}
} else if (currentAstProvider == "volcano" ) {
@ -1838,7 +1872,7 @@ class AzureSpeechPlugin : BleService.Callback, FlutterPlugin, ActivityAware,
result.success(false); return@launch
}
doubaoAstHelperB?.initialize(cfgB, callbackB)
callPcmBridge?.markReady()
markAstBridgeReady()
result.success(true)
}
} else if (currentAstProvider == "alibaba" ) {
@ -1887,7 +1921,7 @@ class AzureSpeechPlugin : BleService.Callback, FlutterPlugin, ActivityAware,
result.success(false); return@launch
}
val initB = bailianAstHelperB?.initialize(bailianConfigB, callbackB) ?: false
callPcmBridge?.markReady()
markAstBridgeReady()
result.success(initA && initB)
}
}

6
apps/client/local_plugins/azure_speech/android/src/main/kotlin/com/yunqiinnovation/azure_speech/BesCallPcmBridge.kt

@ -74,7 +74,11 @@ class BesCallPcmBridge(
ingress.execute {
Log.d(TAG, "已关闭 A=$framesA B=$framesB 丢帧=$dropped B门限前=$gatedB 下行=$downlinkFrames 峰值积压=$maxPendingSeen")
holdA.clear(); holdB.clear()
ready = false
// ⚠️ 刻意**不**把 ready 打回 false:它表达的是「AST helper 起来了没」,只由
// onSessionInit / onSessionDisposed 翻假。原来在这里清掉,而同页面第二场会话不会再
// initialize、也就没人再 markReady,于是之后每一场都 ready:false、上行全扣在缓冲里
//(2026-09-22 ALN-AL80 真机:进页面第一场好、之后每场聋)。
// AzureSpeechPlugin 在重新 enable 时也会按插件上的就绪状态补一次 markReady,双保险。
}
}

2
apps/client/local_plugins/bluetooth_manager/android/src/main/kotlin/com/example/bluetooth_manager/AiAudioDownlink.kt

@ -196,7 +196,7 @@ object AiAudioDownlink {
// 整包一帧真实数据都没有就不发,避免空闲时白占 SPP 带宽
if (!hasData) return
BluetoothManager.writeBytes(packetBuf)
BluetoothManager.writeRealtime(packetBuf)
sentPackets++
if (sentPackets == 1L || sentPackets % 100 == 0L) {
Log.d(TAG, "AI 已下发 $sentPackets 包 队列=${outQueue.size} 间隔=${intervalMs}ms")

63
apps/client/local_plugins/bluetooth_manager/android/src/main/kotlin/com/example/bluetooth_manager/BluetoothManager.kt

@ -22,6 +22,7 @@ import io.flutter.plugin.common.MethodChannel.Result
import io.flutter.plugin.common.PluginRegistry
import java.io.IOException
import java.util.UUID
import java.util.concurrent.Executors
import java.util.concurrent.LinkedBlockingQueue
import java.util.concurrent.ThreadPoolExecutor
import java.util.concurrent.TimeUnit
@ -710,10 +711,72 @@ object BluetoothManager : PluginRegistry.RequestPermissionsResultListener {
aiCallback?.invoke(data)
}
// ---------- 控制命令的发送时机 ----------
//
// Android 的 BluetoothSocket 写入走的是到蓝牙进程的本地管道;协议栈只在链路有信用
// (耳机回 credit)时才从管道取数据,而且一次把**攒着的全部字节合成一个 RFCOMM 帧**发出去
// (btsock_rfc 一次读到 MTU)。恒玄固件一帧只解开头那一条命令,排在后面的静默丢失。
// 2026-09-22 ALN-AL80 真机(协议栈 `PORT_WriteData p_len=N`):
// - 通话翻译期间下行音频包(84B)经常被合成 168/252/…/1512 字节的帧,即管道里最多攒过
// 18 包 ≈ 360ms;
// - 8 次收尾有 3 次 call.stop 丢失:`p_len=4`(asr.stopped+call.stop 合成一帧)、
// `p_len=256`(3 包音频 + 两条命令)、`p_len=172`(2 包音频 + 两条命令)。耳机既不回
// BB D2 也不退通话模式、继续推音频,下一次 call.start 应答成 `BB D1 07 01 …`
// (已在通话模式)——就是「关闭没反应 / 耳机没收到结束指令」。
// 握手里 AA 06 紧跟 AA 09 的「AA 09 无回包」多半也是它。
//
// iOS 没这个问题:每次 GATT 写各自是一个 ATT 请求,不会合并。所以只改 Android:
// - 控制命令走单线程队列(保序),发送前等到「距上一条命令 ≥ CTRL_WRITE_GAP_MS 且
// 距上一个实时音频包 ≥ CTRL_AFTER_REALTIME_MS」——后者是为了让管道里攒着的音频先发完,
// 命令才能落在一帧的开头;最多等 CTRL_WAIT_CAP_MS 兜底(音频一直不停时也别永远卡着);
// - 实时音频包(通话/AI 下行,本来就 20ms/60ms 一包)走 [writeRealtime] 直发,不排队不等。
// ⚠️ Dart 那边 sendData 本来就不等真正写出(原来也是 result.success 后就返回),
// 所以排队不改变任何调用方的语义;收尾顺序是先停下行再发 call.stop,代价只是耳机
// 晚约半秒退出通话模式。
private const val CTRL_WRITE_GAP_MS = 30L
private const val CTRL_AFTER_REALTIME_MS = 500L
private const val CTRL_WAIT_CAP_MS = 5000L
private val ctrlWriteExecutor = Executors.newSingleThreadExecutor { r -> Thread(r, "bt.ctrl.write") }
@Volatile private var lastCtrlWriteAtNs = 0L
@Volatile private var lastRealtimeWriteAtNs = 0L
/** 控制命令(AA xx…):串行 + 等管道里的音频发完。异步,调用方不等结果(原来也不等)。 */
fun writeBytes(message: ByteArray) {
if (thread == null) {
publishBluetoothStatus(0)
return
}
ctrlWriteExecutor.execute {
val startedNs = System.nanoTime()
while (true) {
val now = System.nanoTime()
val sinceCtrl = if (lastCtrlWriteAtNs == 0L) Long.MAX_VALUE else (now - lastCtrlWriteAtNs) / 1_000_000L
val sinceRt = if (lastRealtimeWriteAtNs == 0L) Long.MAX_VALUE else (now - lastRealtimeWriteAtNs) / 1_000_000L
val need = maxOf(CTRL_WRITE_GAP_MS - sinceCtrl, CTRL_AFTER_REALTIME_MS - sinceRt)
if (need <= 0L || (now - startedNs) / 1_000_000L >= CTRL_WAIT_CAP_MS) break
try { Thread.sleep(minOf(need, 50L)) } catch (_: InterruptedException) { Thread.currentThread().interrupt(); break }
}
val waitedMs = (System.nanoTime() - startedNs) / 1_000_000L
val t = thread
if (t == null) {
Log.w("BesCtrlWrite", "控制命令未发出(链路已断)cmd=0x${"%02X".format(message.getOrElse(1) { 0 })}")
publishBluetoothStatus(0)
return@execute
}
t.write(message)
lastCtrlWriteAtNs = System.nanoTime()
if (waitedMs >= 100L) {
Log.d("BesCtrlWrite", "控制命令 0x${"%02X".format(message.getOrElse(1) { 0 })} 等了 ${waitedMs}ms 才发(让管道里的音频先发完)")
}
}
}
/** 实时音频包(通话/AI 下行发送器专用):直发,不排队。见 writeBytes 的说明。 */
fun writeRealtime(message: ByteArray) {
val t = thread
if (t != null) {
t.write(message)
lastRealtimeWriteAtNs = System.nanoTime()
} else {
publishBluetoothStatus(0)
}

2
apps/client/local_plugins/bluetooth_manager/android/src/main/kotlin/com/example/bluetooth_manager/BluetoothManagerPlugin.kt

@ -150,6 +150,8 @@ class BluetoothManagerPlugin : FlutterPlugin, MethodCallHandler, ActivityAware {
CallTranslationDownlink.stop()
result.success(true)
}
// 与 iOS 同名:Dart 的 DeviceSession.diagnostics() 会取它写进日志([DOWNLINK] 取样)
"getCallDownlinkStats" -> result.success(CallTranslationDownlink.stats())
"pushTtsPcm" -> {
val leg = call.argument<String>("leg")
val pcm = call.argument<ByteArray>("pcm")

388
apps/client/local_plugins/bluetooth_manager/android/src/main/kotlin/com/example/bluetooth_manager/CallTranslationDownlink.kt

@ -1,15 +1,18 @@
package com.example.bluetooth_manager
import android.util.Log
import java.util.concurrent.LinkedBlockingQueue
import java.util.concurrent.ScheduledFuture
import java.util.concurrent.ScheduledThreadPoolExecutor
import java.util.concurrent.TimeUnit
import kotlin.math.max
import kotlin.math.min
import kotlin.math.roundToInt
import kotlin.math.sqrt
/**
* 通话翻译的**下行 TTS 音频**发送器(恒玄/BES 专用)。
*
* 上层(Dart 的 AST 流水线)只负责丢 16k/16bit/mono 的 PCM,
* 上层(AST 流水线 / 原生直连桥)只负责丢 16k/16bit/mono 的 PCM,
* 编码、节流、组包全在这里做——这样 Dart 不需要碰 G.722,也不需要自己掐 20ms 时钟。
*
* 包格式(与 deepvoice CallTranslationHelper.sendAudioData 一致):
@ -22,37 +25,85 @@ import java.util.concurrent.TimeUnit
* 一帧固定 640 字节 PCM(320 样本 = 20ms)→ 40 字节码流,
* 所以默认 20ms 的发送间隔正好是实时速率。耳机通过 `0xD6/0xE9`
* 要求改节奏时(16/20/25ms),调 [setIntervalMs] 跟上,否则耳机侧会缓冲溢出或断续。
*
* ## 积压追赶(2026-09-22)
*
* 译文音频是模型**突发**产出的(每 3~6 秒一段文本,音频一次性到),播放却只能 1 倍速;
* 说话人不停顿时,译文时长 ≈ 原话时长(这场实测 A 路 133 秒会话产出 129 秒音频),
* 队列只涨不缩——真机压测 A 路积压到 **6.2 秒**,用户听到的就是「越说越慢」。
* 上行那边没有积压(桥峰值 8 帧),瓶颈就在这里。
*
* 处理分两层,都不丢内容:
* 1. **裁静音**:积压超过 [SILENCE_TRIM_START_SEC] 时,把源音频里超过 [SILENCE_KEEP_MS]
* 的静音段裁到只剩 [SILENCE_KEEP_MS](TTS 句间常有 200~400ms 停顿),无损。
* 2. **变速不变调追赶**(WSOLA,见 [Wsola]):积压 ≤ [CATCHUP_START_SEC] 时 1.0x 原样播;
* 到 [CATCHUP_FULL_SEC] 线性提到 [MAX_RATE](1.3x 是语音可懂度的常用上限)。
* 1.3x 时每秒追回 0.3 秒,配合说话人的自然停顿收敛到 2~3 秒以内。
* 兜底 [HARD_CAP_SEC](默认 0 = 关):积压超过它才整段丢最旧的音频——这会丢内容,
* 产品上默认不启用。
*
* ⚠️ 这里改的只是 Android;iOS 的 CallTranslationDownlink.swift 仍是原样 1 倍速播放。
*/
object CallTranslationDownlink {
private const val TAG = "BesCallDownlink"
private const val SAMPLE_RATE = 16000
/** 一帧 PCM:320 样本 × 2 字节 = 20ms @16kHz */
private const val FRAME_PCM_BYTES = 640
private const val FRAME_SAMPLES = 320
private const val FRAME_PCM_BYTES = FRAME_SAMPLES * 2
/** 一帧编码后码流长度 */
private const val FRAME_ENCODED_BYTES = 40
private const val PACKET_SIZE = 84
/** 队列上限:约 4096 × 20ms,远超实际需要,仅作失控保护 */
private const val QUEUE_CAP = 4096
// ---- 追赶策略参数 ----
private const val CATCHUP_START_SEC = 2.0
private const val CATCHUP_FULL_SEC = 4.0
private const val MAX_RATE = 1.3
private const val SILENCE_TRIM_START_SEC = 1.0
private const val SILENCE_KEEP_MS = 150
private const val SILENCE_RMS = 300.0
private const val HARD_CAP_SEC = 0.0
/** 一条腿:源 PCM 队列 + 变速器 + 编码器 + 统计。所有字段只在 synchronized(this) 里动。 */
private class Leg(val name: String) {
val src = ArrayDeque<ShortArray>() // 未播的源 PCM,20ms 一块
var srcSamples = 0L // src 里的样本数
var partial = ShortArray(FRAME_SAMPLES) // 不满 20ms 的尾巴
var partialLen = 0
var silentRun = 0 // 队尾连续静音样本数(裁静音用)
val stretcher = Wsola()
var encHandle = 0L
// legB = 己方听(对端译文),legA = 对端听(己方译文)
private val legBQueue = LinkedBlockingQueue<ByteArray>(QUEUE_CAP)
private val legAQueue = LinkedBlockingQueue<ByteArray>(QUEUE_CAP)
// 统计
var pushedFrames = 0L // 收到的源帧数(20ms)
var sentFrames = 0L // 实际发出的帧数
var trimmedSamples = 0L // 裁掉的静音样本
var droppedSamples = 0L // 兜底丢弃的样本
var maxBacklogSamples = 0L
var lastRate = 1.0
// 两条腿各自独立的编码器状态,不能共用(与解码侧 spk/mic 双 handle 同理)
private var encHandleA: Long = 0L
private var encHandleB: Long = 0L
/** 还没播出去的源音频(队列里的 + 变速器里没读完的),单位样本 */
fun backlogSamples(): Long = srcSamples + stretcher.unreadSamples()
fun backlogSec(): Double = backlogSamples() / SAMPLE_RATE.toDouble()
// 攒够 640 字节才能编一帧,余量留到下次
private var pendingA = ByteArray(0)
private var pendingB = ByteArray(0)
private val pendingLock = Any()
fun reset() {
src.clear(); srcSamples = 0; partialLen = 0; silentRun = 0
stretcher.reset()
pushedFrames = 0; sentFrames = 0; trimmedSamples = 0; droppedSamples = 0
maxBacklogSamples = 0; lastRate = 1.0
}
}
private val legA = Leg("A") // 对端听
private val legB = Leg("B") // 己方听
private val scheduler = ScheduledThreadPoolExecutor(1)
private var sendTask: ScheduledFuture<*>? = null
private val packetBuf = ByteArray(PACKET_SIZE)
private val pcmBytes = ByteArray(FRAME_PCM_BYTES)
@Volatile
private var running = false
@ -61,8 +112,6 @@ object CallTranslationDownlink {
private var intervalMs: Long = 20L
private var sentPackets = 0L
private var pushedA = 0L
private var pushedB = 0L
/** 静音帧(deepvoice 原样照搬,耳机侧认这个码流) */
private val silenceFrame = byteArrayOf(
@ -80,19 +129,13 @@ object CallTranslationDownlink {
fun start() {
if (running) return
running = true
legAQueue.clear()
legBQueue.clear()
synchronized(pendingLock) {
pendingA = ByteArray(0)
pendingB = ByteArray(0)
for (leg in arrayOf(legA, legB)) synchronized(leg) {
leg.reset()
if (leg.encHandle == 0L) leg.encHandle = G722Codec.init(SAMPLE_RATE)
}
sentPackets = 0L
pushedA = 0L
pushedB = 0L
if (encHandleA == 0L) encHandleA = G722Codec.init(16000)
if (encHandleB == 0L) encHandleB = G722Codec.init(16000)
scheduleSendTask()
Log.d(TAG, "下行发送启动, interval=${intervalMs}ms")
Log.d(TAG, "下行发送启动, interval=${intervalMs}ms 追赶策略: >${CATCHUP_START_SEC}s 起变速, ${CATCHUP_FULL_SEC}s 达 ${MAX_RATE}x, 裁静音 >${SILENCE_TRIM_START_SEC}s")
}
@Synchronized
@ -101,16 +144,46 @@ object CallTranslationDownlink {
running = false
sendTask?.cancel(false)
sendTask = null
legAQueue.clear()
legBQueue.clear()
synchronized(pendingLock) {
pendingA = ByteArray(0)
pendingB = ByteArray(0)
val a = summary(legA)
val b = summary(legB)
for (leg in arrayOf(legA, legB)) synchronized(leg) {
leg.src.clear(); leg.srcSamples = 0; leg.partialLen = 0
leg.stretcher.reset()
if (leg.encHandle != 0L) {
runCatching { G722Codec.release(leg.encHandle) }
leg.encHandle = 0L
}
}
releaseEncoders()
Log.d(TAG, "下行发送停止, 共发出 $sentPackets 包 (pushA=$pushedA pushB=$pushedB)")
Log.d(TAG, "下行发送停止, 共发出 $sentPackets 包 (pushA=${legA.pushedFrames} pushB=${legB.pushedFrames})")
Log.w(TAG, "下行积压总账 A{$a} B{$b}")
}
private fun summary(leg: Leg): String = synchronized(leg) {
val pushedSec = leg.pushedFrames * 0.02
val sentSec = leg.sentFrames * 0.02
"源=%.1fs 播出=%.1fs 最大积压=%.1fs 裁静音=%.1fs 变速追回=%.1fs 兜底丢=%.1fs".format(
pushedSec, sentSec, leg.maxBacklogSamples / SAMPLE_RATE.toDouble(),
leg.trimmedSamples / SAMPLE_RATE.toDouble(),
max(0.0, pushedSec - sentSec - leg.trimmedSamples / SAMPLE_RATE.toDouble() - leg.droppedSamples / SAMPLE_RATE.toDouble() - leg.backlogSec()),
leg.droppedSamples / SAMPLE_RATE.toDouble()
)
}
/** 给 Dart 的 diagnostics() 用(与 iOS 的 getCallDownlinkStats 同名同用途,键不要求一致)。 */
fun stats(): Map<String, Any> = mapOf(
"running" to running,
"sentPackets" to sentPackets,
"intervalMs" to intervalMs,
"backlogA" to synchronized(legA) { legA.backlogSec() },
"backlogB" to synchronized(legB) { legB.backlogSec() },
"maxBacklogA" to synchronized(legA) { legA.maxBacklogSamples / SAMPLE_RATE.toDouble() },
"maxBacklogB" to synchronized(legB) { legB.maxBacklogSamples / SAMPLE_RATE.toDouble() },
"rateA" to synchronized(legA) { legA.lastRate },
"rateB" to synchronized(legB) { legB.lastRate },
"trimmedA" to synchronized(legA) { legA.trimmedSamples / SAMPLE_RATE.toDouble() },
"trimmedB" to synchronized(legB) { legB.trimmedSamples / SAMPLE_RATE.toDouble() },
)
/**
* 耳机要求调整下行节奏(`0xD6`/`0xE9` 的 data[6])。
* 只在运行中才重排任务,停止状态下仅记住新值。
@ -141,74 +214,227 @@ object CallTranslationDownlink {
fun pushPcm(leg: String, pcm: ByteArray) {
if (!running || pcm.isEmpty()) return
val isA = leg.equals("A", ignoreCase = true) || leg.equals("uplink", ignoreCase = true)
val frames: List<ByteArray>
synchronized(pendingLock) {
val merged = if (isA) pendingA + pcm else pendingB + pcm
val frameCount = merged.size / FRAME_PCM_BYTES
if (frameCount == 0) {
if (isA) pendingA = merged else pendingB = merged
return
}
// JNI 的 encode 本身支持多帧,一次调用编完整段,省掉逐帧过 JNI 的开销
val bulkLen = frameCount * FRAME_PCM_BYTES
val handle = if (isA) encHandleA else encHandleB
val encoded = try {
G722Codec.encodeSync(handle, merged.copyOfRange(0, bulkLen))
} catch (e: Exception) {
Log.w(TAG, "G.722 编码失败: ${e.message}")
null
val l = if (isA) legA else legB
synchronized(l) {
val n = pcm.size / 2
var i = 0
while (i < n) {
val take = min(FRAME_SAMPLES - l.partialLen, n - i)
var k = 0
while (k < take) {
val p = (i + k) * 2
l.partial[l.partialLen + k] = ((pcm[p].toInt() and 0xFF) or (pcm[p + 1].toInt() shl 8)).toShort()
k++
}
l.partialLen += take
i += take
if (l.partialLen == FRAME_SAMPLES) {
val block = l.partial.copyOf()
l.partialLen = 0
l.pushedFrames++
if (!trimSilence(l, block)) {
l.src.addLast(block)
l.srcSamples += FRAME_SAMPLES
}
}
}
val out = ArrayList<ByteArray>(frameCount)
if (encoded != null) {
var o = 0
while (o + FRAME_ENCODED_BYTES <= encoded.size) {
out.add(encoded.copyOfRange(o, o + FRAME_ENCODED_BYTES))
o += FRAME_ENCODED_BYTES
if (HARD_CAP_SEC > 0.0) {
val cap = (HARD_CAP_SEC * SAMPLE_RATE).toLong()
if (l.srcSamples > cap) {
val target = ((HARD_CAP_SEC - 2.0).coerceAtLeast(1.0) * SAMPLE_RATE).toLong()
while (l.srcSamples > target && l.src.isNotEmpty()) {
l.srcSamples -= l.src.removeFirst().size
l.droppedSamples += FRAME_SAMPLES
}
Log.w(TAG, "${l.name} 路积压超过 ${HARD_CAP_SEC}s,丢弃最旧音频至 ${target / SAMPLE_RATE.toDouble()}s")
}
}
val rest = merged.copyOfRange(bulkLen, merged.size)
if (isA) pendingA = rest else pendingB = rest
frames = out
val bl = l.backlogSamples()
if (bl > l.maxBacklogSamples) l.maxBacklogSamples = bl
}
val queue = if (isA) legAQueue else legBQueue
for (f in frames) {
// 队列满说明下行速度跟不上产出(通常是耳机没在收),丢最旧的保实时
if (!queue.offer(f)) {
queue.poll()
queue.offer(f)
}
}
/** 积压时把长静音裁到只剩 SILENCE_KEEP_MS;返回 true 表示这一块被裁掉了。 */
private fun trimSilence(l: Leg, block: ShortArray): Boolean {
var acc = 0.0
for (s in block) acc += s.toDouble() * s.toDouble()
val rms = sqrt(acc / block.size)
if (rms >= SILENCE_RMS) {
l.silentRun = 0
return false
}
l.silentRun += block.size
val keep = SILENCE_KEEP_MS * SAMPLE_RATE / 1000
if (l.silentRun > keep && l.backlogSec() > SILENCE_TRIM_START_SEC) {
l.trimmedSamples += block.size
return true
}
return false
}
/** 按积压深度算这一帧的播放速率。 */
private fun rateFor(backlogSec: Double): Double = when {
backlogSec <= CATCHUP_START_SEC -> 1.0
backlogSec >= CATCHUP_FULL_SEC -> MAX_RATE
else -> 1.0 + (MAX_RATE - 1.0) * (backlogSec - CATCHUP_START_SEC) / (CATCHUP_FULL_SEC - CATCHUP_START_SEC)
}
/** 取这条腿的下一帧编码码流;没东西可播返回 null。 */
private fun nextEncodedFrame(l: Leg): ByteArray? = synchronized(l) {
if (l.encHandle == 0L) return null
val backlog = l.backlogSec()
val rate = rateFor(backlog)
l.lastRate = rate
val out = l.stretcher.nextFrame(rate, l) ?: return null
var i = 0
while (i < FRAME_SAMPLES) {
val v = out[i].toInt()
pcmBytes[i * 2] = (v and 0xFF).toByte()
pcmBytes[i * 2 + 1] = ((v shr 8) and 0xFF).toByte()
i++
}
if (isA) pushedA += frames.size else pushedB += frames.size
val enc = try {
G722Codec.encodeSync(l.encHandle, pcmBytes)
} catch (e: Exception) {
Log.w(TAG, "G.722 编码失败: ${e.message}"); null
}
l.sentFrames++
if (enc == null || enc.size < FRAME_ENCODED_BYTES) silenceFrame else enc
}
private fun sendOnePacket() {
if (!running) return
val legB = legBQueue.poll()
val legA = legAQueue.poll()
val b = nextEncodedFrame(legB)
val a = nextEncodedFrame(legA)
// 两路都没内容就不发,避免通话里灌满无谓的静音包
if (legB == null && legA == null) return
if (b == null && a == null) return
packetBuf[0] = 0xAA.toByte()
packetBuf[1] = 0x56.toByte()
packetBuf[2] = PACKET_SIZE.toByte()
packetBuf[3] = 0x03.toByte()
System.arraycopy(legB ?: silenceFrame, 0, packetBuf, 4, FRAME_ENCODED_BYTES)
System.arraycopy(legA ?: silenceFrame, 0, packetBuf, 44, FRAME_ENCODED_BYTES)
BluetoothManager.writeBytes(packetBuf)
System.arraycopy(b ?: silenceFrame, 0, packetBuf, 4, FRAME_ENCODED_BYTES)
System.arraycopy(a ?: silenceFrame, 0, packetBuf, 44, FRAME_ENCODED_BYTES)
BluetoothManager.writeRealtime(packetBuf)
sentPackets++
if (sentPackets == 1L || sentPackets % 100 == 0L) {
Log.d(TAG, "已下发 $sentPackets 包 (队列 A=${legAQueue.size} B=${legBQueue.size})")
Log.d(TAG, "已下发 $sentPackets 包 (积压 A=%.1fs B=%.1fs 速率 A=%.2f B=%.2f)".format(
synchronized(legA) { legA.backlogSec() }, synchronized(legB) { legB.backlogSec() },
legA.lastRate, legB.lastRate))
}
}
private fun releaseEncoders() {
if (encHandleA != 0L) {
runCatching { G722Codec.release(encHandleA) }
encHandleA = 0L
/**
* WSOLA 变速不变调(Verhelst & Roelands)。16k 单声道,输出恒为 20ms 一帧。
*
* 合成侧固定步长 [S](10ms)、帧长 [N](20ms)、Hann 50% 重叠相加;分析侧步长 = S × rate,
* 每帧在名义位置 ±[T] 内搜索与「上一帧自然延续」最相似的起点(互相关),保证拼接处波形连续。
* rate=1.0 且偏移=0 时 Hann 50% 重叠相加恒等于原信号,所以不用在直通/变速间切换。
* 附加时延 ≈ N + T 样本(≈26ms)。
*
* 输入不够时:源队列已空 → 补零把尾巴冲出来(TTS 末尾本来就是静音);源队列还有 → 先取。
*/
private class Wsola(
private val N: Int = 320,
private val S: Int = 160,
private val T: Int = 96,
) {
private var inBuf = ShortArray(N * 8)
private var inLen = 0
private var anaPos = 0.0
private var prevEnd = -1
private val win = FloatArray(N) { i -> (0.5 - 0.5 * kotlin.math.cos(2.0 * Math.PI * i / N)).toFloat() }
private val acc = FloatArray(N)
private val out = ShortArray(FRAME_SAMPLES)
private var padded = 0
fun reset() {
inLen = 0; anaPos = 0.0; prevEnd = -1; padded = 0
java.util.Arrays.fill(acc, 0f)
}
/** 变速器里还没播掉的样本(不含补的零) */
fun unreadSamples(): Long = max(0L, (inLen - padded - anaPos.roundToInt()).toLong())
/** 取 20ms 输出;返回 null 表示这条腿当前没有东西可播。 */
fun nextFrame(rate: Double, l: Leg): ShortArray? {
var produced = 0
while (produced < FRAME_SAMPLES) {
val nominal = anaPos.roundToInt()
val need = nominal + T + N
// 先从源队列取
while (inLen < need && l.src.isNotEmpty()) {
append(l.src.removeFirst()); l.srcSamples -= N
}
if (inLen < need) {
// 源队列空了:还有实际音频没播完就补零冲出来,否则这条腿此刻没东西
val realLeft = inLen - padded - nominal
if (produced == 0 && realLeft <= 0) return null
val pad = need - inLen
ensure(need); java.util.Arrays.fill(inBuf, inLen, need, 0.toShort()); inLen = need; padded += pad
}
val start = pickStart(nominal)
var i = 0
while (i < N) { acc[i] += win[i] * inBuf[start + i]; i++ }
// 输出前 S 个
var j = 0
while (j < S) {
val v = acc[j].roundToInt().coerceIn(-32768, 32767)
out[produced + j] = v.toShort(); j++
}
System.arraycopy(acc, S, acc, 0, N - S)
java.util.Arrays.fill(acc, N - S, N, 0f)
produced += S
prevEnd = start + S
anaPos += S * rate
compact()
}
return out
}
private fun pickStart(nominal: Int): Int {
if (prevEnd < 0 || prevEnd + S > inLen) return nominal.coerceIn(0, inLen - N)
var best = 0
var bestScore = Double.NEGATIVE_INFINITY
val lo = max(-T, -nominal)
val hi = min(T, inLen - N - nominal)
var d = lo
while (d <= hi) {
val c = nominal + d
var xy = 0.0; var yy = 1e-6
var k = 0
while (k < S) {
val x = inBuf[prevEnd + k].toDouble(); val y = inBuf[c + k].toDouble()
xy += x * y; yy += y * y; k++
}
val score = xy / sqrt(yy)
if (score > bestScore) { bestScore = score; best = d }
d++
}
return nominal + best
}
private fun append(block: ShortArray) {
// 之前补的零已经被新到的音频"夹"在中间,当成内容算(最多多算 26ms)
padded = 0
ensure(inLen + block.size)
System.arraycopy(block, 0, inBuf, inLen, block.size)
inLen += block.size
}
private fun ensure(cap: Int) {
if (inBuf.size < cap) inBuf = inBuf.copyOf(max(cap, inBuf.size * 2))
}
if (encHandleB != 0L) {
runCatching { G722Codec.release(encHandleB) }
encHandleB = 0L
/** 丢掉已经用不到的输入前缀,避免无限增长。 */
private fun compact() {
val keepFrom = min(anaPos.roundToInt(), prevEnd) - T - N
if (keepFrom > N * 4) {
System.arraycopy(inBuf, keepFrom, inBuf, 0, inLen - keepFrom)
inLen -= keepFrom
anaPos -= keepFrom
prevEnd -= keepFrom
}
}
}
}

399
apps/client/local_plugins/bluetooth_manager/ios/bluetooth_manager/Sources/bluetooth_manager/CallTranslationDownlink.swift

@ -4,7 +4,7 @@ import os.log
/// 通话翻译的**下行 TTS 音频**发送器(恒玄/BES 专用),Android
/// `CallTranslationDownlink.kt` 的 iOS 对应实现,行为逐条对齐。
///
/// 上层(Dart 的 AST 流水线)只丢 16k/16bit/mono 的 PCM,
/// 上层(Dart 的 AST 流水线 / 原生直连桥)只丢 16k/16bit/mono 的 PCM,
/// 编码、节流、组包都在这里做。
///
/// 包格式:`AA 56 54 03 [40B legB][40B legA]` 共 84 字节,mode=0x03 双声道
@ -15,29 +15,85 @@ import os.log
/// 一帧固定 640 字节 PCM(320 样本 = 20ms)→ 40 字节码流,
/// 所以默认 20ms 的发送间隔正好是实时速率;耳机用 `0xD6/0xE9`
/// 要求改节奏时(16/20/25ms)调 [setIntervalMs] 跟上。
///
/// ## 积压追赶(2026-09-22,与 Android 同一套参数)
///
/// 译文音频是模型**突发**产出的(每 3~6 秒一段文本,音频一次性到),播放却只能 1 倍速;
/// 说话人不停顿时译文时长 ≈ 原话时长,队列只涨不缩——Android 真机压测 A 路积压到 6.2 秒,
/// 用户听到的就是「越说越慢」。处理分两层,都不丢内容:
/// 1. **裁静音**:积压超过 [silenceTrimStartSec] 时,把源音频里超过 [silenceKeepMs] 的静音段
/// 裁到只剩 [silenceKeepMs](TTS 句间常有 200~400ms 停顿),无损。同量压测光这一步就省 10%。
/// 2. **变速不变调追赶**(WSOLA,见 [Wsola]):积压 ≤ [catchupStartSec] 时 1.0x 原样播;
/// 到 [catchupFullSec] 线性提到 [maxRate](1.3x 是语音可懂度的常用上限)。
/// 兜底 [hardCapSec](默认 0 = 关):积压超过它才整段丢最旧的音频——会丢内容,产品上默认不启用。
/// Android 同量压测:最大积压 6.2s → 3.2s,变速最高只用到 1.11x。
public final class CallTranslationDownlink {
public static let shared = CallTranslationDownlink()
private static let sampleRate = 16000
/// 一帧 PCM:320 样本 × 2 字节 = 20ms @16kHz
private static let frameSamples = 320
private static let framePcmBytes = 640
private static let frameEncodedBytes = 40
private static let packetSize = 84
/// 队列上限,仅作失控保护
private static let queueCap = 4096
// legB = 己方听(对端译文),legA = 对端听(己方译文)
private var legAQueue: [Data] = []
private var legBQueue: [Data] = []
// ---- 追赶策略参数(与 Android 一致,改一端要同改另一端)----
private static let catchupStartSec = 2.0
private static let catchupFullSec = 4.0
private static let maxRate = 1.3
private static let silenceTrimStartSec = 1.0
private static let silenceKeepMs = 150
private static let silenceRms = 300.0
private static let hardCapSec = 0.0
/// 一条腿:源 PCM 队列 + 变速器 + 编码器 + 统计。只在 [queue] 上访问。
private final class Leg {
let name: String
var src: [[Int16]] = [] // 未播的源 PCM,20ms 一块
var srcHead = 0 // src 的读游标(避免 removeFirst 的 O(n))
var srcSamples = 0 // 未播样本数
var partial = [Int16](repeating: 0, count: CallTranslationDownlink.frameSamples)
var partialLen = 0
var silentRun = 0
let stretcher = Wsola()
var encoder: G722Codec?
var pushedFrames = 0
var sentFrames = 0
var trimmedSamples = 0
var droppedSamples = 0
var maxBacklogSamples = 0
var lastRate = 1.0
init(_ name: String) { self.name = name }
func backlogSamples() -> Int { srcSamples + stretcher.unreadSamples() }
func backlogSec() -> Double { Double(backlogSamples()) / Double(CallTranslationDownlink.sampleRate) }
func srcIsEmpty() -> Bool { srcHead >= src.count }
func srcPop() -> [Int16] {
let b = src[srcHead]
srcHead += 1
srcSamples -= b.count
if srcHead >= 64 && srcHead * 2 >= src.count {
src.removeFirst(srcHead); srcHead = 0
}
return b
}
func srcPush(_ b: [Int16]) { src.append(b); srcSamples += b.count }
func srcClear() { src.removeAll(); srcHead = 0; srcSamples = 0 }
// 两条腿各自独立的编码器状态,不能共用
private var encoderA: G722Codec?
private var encoderB: G722Codec?
func reset() {
srcClear(); partialLen = 0; silentRun = 0
stretcher.reset()
pushedFrames = 0; sentFrames = 0; trimmedSamples = 0; droppedSamples = 0
maxBacklogSamples = 0; lastRate = 1.0
}
}
// 攒够 640 字节才能编一帧,余量留到下次
private var pendingA = Data()
private var pendingB = Data()
private let legA = Leg("A") // 对端听
private let legB = Leg("B") // 己方听
/// 串行队列保护队列与 pending,同时承载定时发包
/// 串行队列保护所有状态,同时承载定时发包
private let queue = DispatchQueue(label: "bes.call.downlink")
private var timer: DispatchSourceTimer?
@ -56,11 +112,9 @@ public final class CallTranslationDownlink {
private var tickGapSumMs: Double = 0
private var tickGapCount: Int = 0
private var tickGapMaxMs: Double = 0
/// 两路各自的入队/出队总数,用来判断是产出多了还是发得慢了
private var pushedFramesA = 0
private var pushedFramesB = 0
private var maxQueueA = 0
private var maxQueueB = 0
/// 取样窗口内的积压峰值(帧),resetPeaks 会清
private var peakQueueA = 0
private var peakQueueB = 0
/// 静音帧(与 deepvoice 一致,耳机侧认这个码流)
private let silenceFrame = Data([
@ -76,19 +130,18 @@ public final class CallTranslationDownlink {
queue.async {
guard !self.running else { return }
self.running = true
self.legAQueue.removeAll()
self.legBQueue.removeAll()
self.pendingA = Data()
self.pendingB = Data()
for leg in [self.legA, self.legB] {
leg.reset()
if leg.encoder == nil { leg.encoder = G722Codec(sampleRate: 16000) }
}
self.sentPackets = 0
if self.encoderA == nil { self.encoderA = G722Codec(sampleRate: 16000) }
if self.encoderB == nil { self.encoderB = G722Codec(sampleRate: 16000) }
self.lastTickAt = 0
self.tickGapSumMs = 0; self.tickGapCount = 0; self.tickGapMaxMs = 0
self.pushedFramesA = 0; self.pushedFramesB = 0
self.maxQueueA = 0; self.maxQueueB = 0
self.peakQueueA = 0; self.peakQueueB = 0
self.scheduleTimer()
os_log("[BesCallDownlink] 下行发送启动, interval=%dms", log: Self.dlLog, type: .info, self.intervalMs)
os_log("[BesCallDownlink] 下行发送启动, interval=%dms 追赶策略: >%.1fs 起变速, %.1fs 达 %.1fx, 裁静音 >%.1fs",
log: Self.dlLog, type: .info, self.intervalMs,
Self.catchupStartSec, Self.catchupFullSec, Self.maxRate, Self.silenceTrimStartSec)
}
}
@ -98,29 +151,43 @@ public final class CallTranslationDownlink {
self.running = false
self.timer?.cancel()
self.timer = nil
self.legAQueue.removeAll()
self.legBQueue.removeAll()
self.pendingA = Data()
self.pendingB = Data()
self.encoderA = nil
self.encoderB = nil
let a = self.summary(self.legA)
let b = self.summary(self.legB)
let pushedA = self.legA.pushedFrames, pushedB = self.legB.pushedFrames
for leg in [self.legA, self.legB] {
leg.srcClear(); leg.partialLen = 0
leg.stretcher.reset()
leg.encoder = nil
}
// 清自己的队列还不够:已经递给 GATT 层的那些包也得丢掉,
// 否则它们会一直堵在控制命令(查电量/查版本)前面
BluetoothManager.sharedInstance?.flushRealtimeQueue()
let avgGap = self.tickGapCount > 0 ? self.tickGapSumMs / Double(self.tickGapCount) : 0
let bleStats = BluetoothManager.sharedInstance?.bleWriteStats() ?? [:]
os_log("""
[BesCallDownlink] 收尾: 发出=%d 包 | 入队 A=%d B=%d | 队列峰值 A=%d B=%d \
[BesCallDownlink] 收尾: 发出=%d 包 | 入队 A=%d B=%d \
| 定时器实际间隔 平均=%.1fms 最大=%.1fms (期望 %dms) | GATT %{public}@
""",
log: Self.dlLog, type: .info,
self.sentPackets, self.pushedFramesA, self.pushedFramesB,
self.maxQueueA, self.maxQueueB,
self.sentPackets, pushedA, pushedB,
avgGap, self.tickGapMaxMs, self.intervalMs,
String(describing: bleStats))
os_log("[BesCallDownlink] 下行积压总账 A{%{public}@} B{%{public}@}", log: Self.dlLog, type: .default, a, b)
}
}
/// 只在 [queue] 上调用
private func summary(_ leg: Leg) -> String {
let sr = Double(Self.sampleRate)
let pushedSec = Double(leg.pushedFrames) * 0.02
let sentSec = Double(leg.sentFrames) * 0.02
let trimmed = Double(leg.trimmedSamples) / sr
let dropped = Double(leg.droppedSamples) / sr
let catchup = max(0, pushedSec - sentSec - trimmed - dropped - leg.backlogSec())
return String(format: "源=%.1fs 播出=%.1fs 最大积压=%.1fs 裁静音=%.1fs 变速追回=%.1fs 兜底丢=%.1fs",
pushedSec, sentSec, Double(leg.maxBacklogSamples) / sr, trimmed, catchup, dropped)
}
/// 下行发送器的实时统计,供 Dart 在切前后台等时刻取样。
///
/// 这些数字原来只在 os_log 里,而排查后台卡顿时 `idevicesyslog` 极不稳定
@ -130,19 +197,27 @@ public final class CallTranslationDownlink {
///
/// **gapAvgMs 是判断 iOS 后台节流的关键**:期望等于 intervalMs(16/20/25),
/// 明显变大就说明 DispatchSourceTimer 被系统合并触发了,下行必然断续。
/// queueA/B 现在是「未播的源音频」折算的帧数(积压),不再是编码队列长度。
public func stats() -> [String: Any] {
var out: [String: Any] = [:]
queue.sync {
let avg = self.tickGapCount > 0
? self.tickGapSumMs / Double(self.tickGapCount) : 0
let sr = Double(Self.sampleRate)
out = [
"running": self.running,
"intervalMs": self.intervalMs,
"sentPackets": self.sentPackets,
"queueA": self.legAQueue.count,
"queueB": self.legBQueue.count,
"maxQueueA": self.maxQueueA,
"maxQueueB": self.maxQueueB,
"queueA": self.legA.backlogSamples() / Self.frameSamples,
"queueB": self.legB.backlogSamples() / Self.frameSamples,
"maxQueueA": self.peakQueueA,
"maxQueueB": self.peakQueueB,
"backlogA": (self.legA.backlogSec() * 10).rounded() / 10,
"backlogB": (self.legB.backlogSec() * 10).rounded() / 10,
"rateA": (self.legA.lastRate * 100).rounded() / 100,
"rateB": (self.legB.lastRate * 100).rounded() / 100,
"trimmedA": (Double(self.legA.trimmedSamples) / sr * 10).rounded() / 10,
"trimmedB": (Double(self.legB.trimmedSamples) / sr * 10).rounded() / 10,
"gapAvgMs": (avg * 10).rounded() / 10,
"gapMaxMs": (self.tickGapMaxMs * 10).rounded() / 10,
]
@ -155,7 +230,7 @@ public final class CallTranslationDownlink {
public func resetPeaks() {
queue.async {
self.tickGapSumMs = 0; self.tickGapCount = 0; self.tickGapMaxMs = 0
self.maxQueueA = 0; self.maxQueueB = 0
self.peakQueueA = 0; self.peakQueueB = 0
}
}
@ -180,42 +255,91 @@ public final class CallTranslationDownlink {
/// 推一段待下发的 TTS PCM。
/// [leg] "A" = 己方译文给对端听;"B" = 对端译文给己方听。
/// [pcm] 16kHz / 16bit / mono / little-endian,长度任意。
public func pushPcm(leg: String, pcm: Data) {
guard !pcm.isEmpty else { return }
let isA = leg.uppercased() == "A" || leg.lowercased() == "uplink"
queue.async {
guard self.running else { return }
var merged = isA ? self.pendingA : self.pendingB
merged.append(pcm)
let frameCount = merged.count / Self.framePcmBytes
guard frameCount > 0 else {
if isA { self.pendingA = merged } else { self.pendingB = merged }
return
let l = isA ? self.legA : self.legB
let n = pcm.count / 2
var i = 0
pcm.withUnsafeBytes { (raw: UnsafeRawBufferPointer) in
while i < n {
let take = min(Self.frameSamples - l.partialLen, n - i)
var k = 0
while k < take {
let p = (i + k) * 2
let v = UInt16(raw[p]) | (UInt16(raw[p + 1]) << 8)
l.partial[l.partialLen + k] = Int16(bitPattern: v)
k += 1
}
l.partialLen += take
i += take
if l.partialLen == Self.frameSamples {
let block = l.partial
l.partialLen = 0
l.pushedFrames += 1
if !self.trimSilence(l, block) { l.srcPush(block) }
}
}
}
let encoder = isA ? self.encoderA : self.encoderB
var offset = 0
for _ in 0..<frameCount {
let frame = merged.subdata(in: offset..<(offset + Self.framePcmBytes))
offset += Self.framePcmBytes
if let encoded = encoder?.encode(frame), encoded.count >= Self.frameEncodedBytes {
let chunk = encoded.prefix(Self.frameEncodedBytes)
if isA {
// 队列满说明下行跟不上产出(通常是耳机没在收),丢最旧的保实时
if self.legAQueue.count >= Self.queueCap { self.legAQueue.removeFirst() }
self.legAQueue.append(Data(chunk))
self.pushedFramesA += 1
if self.legAQueue.count > self.maxQueueA { self.maxQueueA = self.legAQueue.count }
} else {
if self.legBQueue.count >= Self.queueCap { self.legBQueue.removeFirst() }
self.legBQueue.append(Data(chunk))
self.pushedFramesB += 1
if self.legBQueue.count > self.maxQueueB { self.maxQueueB = self.legBQueue.count }
if Self.hardCapSec > 0 {
let cap = Int(Self.hardCapSec * Double(Self.sampleRate))
if l.srcSamples > cap {
let target = Int(max(1.0, Self.hardCapSec - 2.0) * Double(Self.sampleRate))
while l.srcSamples > target && !l.srcIsEmpty() {
_ = l.srcPop(); l.droppedSamples += Self.frameSamples
}
os_log("[BesCallDownlink] %{public}@ 路积压超过 %.1fs,丢弃最旧音频至 %.1fs",
log: Self.dlLog, type: .default, l.name, Self.hardCapSec, Double(target) / Double(Self.sampleRate))
}
}
let rest = merged.subdata(in: offset..<merged.count)
if isA { self.pendingA = rest } else { self.pendingB = rest }
let bl = l.backlogSamples()
if bl > l.maxBacklogSamples { l.maxBacklogSamples = bl }
let frames = bl / Self.frameSamples
if isA { if frames > self.peakQueueA { self.peakQueueA = frames } }
else { if frames > self.peakQueueB { self.peakQueueB = frames } }
}
}
/// 积压时把长静音裁到只剩 silenceKeepMs;返回 true 表示这一块被裁掉了。只在 [queue] 上调用。
private func trimSilence(_ l: Leg, _ block: [Int16]) -> Bool {
var acc = 0.0
for s in block { let d = Double(s); acc += d * d }
let rms = (acc / Double(block.count)).squareRoot()
if rms >= Self.silenceRms {
l.silentRun = 0
return false
}
l.silentRun += block.count
let keep = Self.silenceKeepMs * Self.sampleRate / 1000
if l.silentRun > keep && l.backlogSec() > Self.silenceTrimStartSec {
l.trimmedSamples += block.count
return true
}
return false
}
/// 按积压深度算这一帧的播放速率。
private func rateFor(_ backlogSec: Double) -> Double {
if backlogSec <= Self.catchupStartSec { return 1.0 }
if backlogSec >= Self.catchupFullSec { return Self.maxRate }
return 1.0 + (Self.maxRate - 1.0) * (backlogSec - Self.catchupStartSec) / (Self.catchupFullSec - Self.catchupStartSec)
}
/// 取这条腿的下一帧编码码流;没东西可播返回 nil。只在 [queue] 上调用。
private func nextEncodedFrame(_ l: Leg) -> Data? {
guard let encoder = l.encoder else { return nil }
let rate = rateFor(l.backlogSec())
l.lastRate = rate
guard let out = l.stretcher.nextFrame(rate: rate, leg: l) else { return nil }
let pcm = out.withUnsafeBufferPointer { Data(buffer: $0) } // Int16 小端 = 本机字节序
l.sentFrames += 1
if let enc = encoder.encode(pcm), enc.count >= Self.frameEncodedBytes {
return Data(enc.prefix(Self.frameEncodedBytes))
}
return silenceFrame
}
/// 只在 [queue] 上调用
@ -233,27 +357,152 @@ public final class CallTranslationDownlink {
}
lastTickAt = now
let legB: Data? = legBQueue.isEmpty ? nil : legBQueue.removeFirst()
let legA: Data? = legAQueue.isEmpty ? nil : legAQueue.removeFirst()
let b = nextEncodedFrame(legB)
let a = nextEncodedFrame(legA)
// 两路都没内容就不发,避免通话里灌满无谓的静音包
if legB == nil && legA == nil { return }
if b == nil && a == nil { return }
var packet = Data(capacity: Self.packetSize)
packet.append(contentsOf: [0xAA, 0x56, UInt8(Self.packetSize), 0x03])
packet.append(legB ?? silenceFrame)
packet.append(legA ?? silenceFrame)
packet.append(b ?? silenceFrame)
packet.append(a ?? silenceFrame)
BluetoothManager.writeRealtimeAudio(packet)
sentPackets += 1
// 每 100 包(约 2 秒)一条。队列长度单调上涨 = 发送跟不上产出。
// 每 100 包(约 2 秒)一条。积压持续 >2s 说明产出快于播放、追赶已介入。
if sentPackets == 1 || sentPackets % 100 == 0 {
let avgGap = tickGapCount > 0 ? tickGapSumMs / Double(tickGapCount) : 0
let ble = BluetoothManager.sharedInstance?.bleWriteStats() ?? [:]
os_log("[BesCallDownlink] 已下发 %d 包 (队列 A=%d B=%d) 间隔 平均=%.1fms 最大=%.1fms | GATT 丢=%{public}@ 阻塞=%{public}@ 待发=%{public}@",
os_log("[BesCallDownlink] 已下发 %d 包 (积压 A=%.1fs B=%.1fs 速率 A=%.2f B=%.2f) 间隔 平均=%.1fms 最大=%.1fms | GATT 丢=%{public}@ 阻塞=%{public}@ 待发=%{public}@",
log: Self.dlLog, type: .info,
sentPackets, legAQueue.count, legBQueue.count, avgGap, tickGapMaxMs,
sentPackets, legA.backlogSec(), legB.backlogSec(), legA.lastRate, legB.lastRate,
avgGap, tickGapMaxMs,
String(describing: ble["dropped"] ?? 0),
String(describing: ble["blocked"] ?? 0),
String(describing: ble["queued"] ?? 0))
}
}
/// WSOLA 变速不变调(Verhelst & Roelands),与 Android `Wsola` 同一套参数与步骤。
///
/// 合成侧固定步长 S(10ms)、帧长 N(20ms)、Hann 50% 重叠相加;分析侧步长 = S × rate,
/// 每帧在名义位置 ±T 内搜索与「上一帧自然延续」最相似的起点(互相关),保证拼接处波形连续。
/// rate=1.0 且偏移=0 时 Hann 50% 重叠相加恒等于原信号(Python 复刻验证:逐样本误差 0),
/// 所以不用在直通/变速间切换。附加时延 ≈ N + T 样本(≈26ms)。
///
/// 输入不够时:源队列已空 → 补零把尾巴冲出来(TTS 末尾本来就是静音);源队列还有 → 先取。
private final class Wsola {
private let N = 320
private let S = 160
private let T = 96
private var inBuf = [Int16](repeating: 0, count: 320 * 8)
private var inLen = 0
private var anaPos = 0.0
private var prevEnd = -1
private let win: [Float]
private var acc: [Float]
private var out: [Int16]
private var padded = 0
init() {
let n = 320
win = (0..<n).map { i in Float(0.5 - 0.5 * cos(2.0 * Double.pi * Double(i) / Double(n))) }
acc = [Float](repeating: 0, count: n)
out = [Int16](repeating: 0, count: CallTranslationDownlink.frameSamples)
}
func reset() {
inLen = 0; anaPos = 0; prevEnd = -1; padded = 0
for i in 0..<acc.count { acc[i] = 0 }
}
/// 变速器里还没播掉的样本(不含补的零)
func unreadSamples() -> Int { max(0, inLen - padded - Int(anaPos.rounded())) }
/// 取 20ms 输出;返回 nil 表示这条腿当前没有东西可播。
func nextFrame(rate: Double, leg: Leg) -> [Int16]? {
var produced = 0
let frame = CallTranslationDownlink.frameSamples
while produced < frame {
let nominal = Int(anaPos.rounded())
let need = nominal + T + N
while inLen < need && !leg.srcIsEmpty() {
append(leg.srcPop())
}
if inLen < need {
let realLeft = inLen - padded - nominal
if produced == 0 && realLeft <= 0 { return nil }
let pad = need - inLen
ensure(need)
for i in inLen..<need { inBuf[i] = 0 }
inLen = need; padded += pad
}
let start = pickStart(nominal)
for i in 0..<N { acc[i] += win[i] * Float(inBuf[start + i]) }
for j in 0..<S {
let v = Int(acc[j].rounded())
out[produced + j] = Int16(clamping: v)
}
for i in 0..<(N - S) { acc[i] = acc[i + S] }
for i in (N - S)..<N { acc[i] = 0 }
produced += S
prevEnd = start + S
anaPos += Double(S) * rate
compact()
}
return out
}
private func pickStart(_ nominal: Int) -> Int {
if prevEnd < 0 || prevEnd + S > inLen { return min(max(nominal, 0), inLen - N) }
var best = 0
var bestScore = -Double.infinity
let lo = max(-T, -nominal)
let hi = min(T, inLen - N - nominal)
if lo > hi { return min(max(nominal, 0), inLen - N) }
inBuf.withUnsafeBufferPointer { p in
var d = lo
while d <= hi {
let c = nominal + d
var xy = 0.0, yy = 1e-6
var k = 0
while k < S {
let x = Double(p[prevEnd + k]), y = Double(p[c + k])
xy += x * y; yy += y * y; k += 1
}
let score = xy / yy.squareRoot()
if score > bestScore { bestScore = score; best = d }
d += 1
}
}
return nominal + best
}
private func append(_ block: [Int16]) {
// 之前补的零已经被新到的音频"夹"在中间,当成内容算(最多多算 26ms)
padded = 0
ensure(inLen + block.count)
for i in 0..<block.count { inBuf[inLen + i] = block[i] }
inLen += block.count
}
private func ensure(_ cap: Int) {
if inBuf.count < cap {
inBuf.append(contentsOf: [Int16](repeating: 0, count: max(cap, inBuf.count * 2) - inBuf.count))
}
}
/// 丢掉已经用不到的输入前缀,避免无限增长。
private func compact() {
let keepFrom = min(Int(anaPos.rounded()), prevEnd) - T - N
if keepFrom > N * 4 {
let remain = inLen - keepFrom
if remain > 0 {
for i in 0..<remain { inBuf[i] = inBuf[keepFrom + i] }
}
inLen = max(0, remain)
anaPos -= Double(keepFrom)
prevEnd -= keepFrom
}
}
}
}

Loading…
Cancel
Save