Browse Source

client(android): 通话翻译三处修复——同页第二场聋、call.stop 偶发丢失、译文积压追赶

全部在 Android 原生(azure_speech / bluetooth_manager 的 Kotlin),iOS 一行未动。
2026-09-22 ALN-AL80 + Echo-one 真机,配合协议栈 RFCOMM 日志逐条定位与复验。

同页第二次开始没反应(azure_speech)
- AST 就绪状态记在插件上(astSessionReady,与 iOS 同义),建桥与重新 enable 时都继承;
  BesCallPcmBridge.disable() 不再把 ready 打回 false。
- 原因:markReady 只在 initialize 末尾发一次,同页面再开始只走 startContinuousTranslation,
  disable() 清掉的 ready 再也没人置真,上行帧全扣在环形缓冲里。实测进页面第一场好、
  之后每场 ready:false / downlinkFrames:0,退出重进又好一场。

结束指令偶发没到耳机(bluetooth_manager)
- 控制命令走单线程队列:距上一条命令 ≥30ms、距上一个实时音频包 ≥500ms 才写;
  下行音频改走 writeRealtime 直发。
- 原因:BluetoothSocket 写入经本地管道,协议栈按信用一次把管道里攒着的字节合成一帧
  (PORT_WriteData p_len=4 / 172 / 256),耳机一帧只解开头一条命令。8 次收尾丢 3 次,
  耳机不回 BB D2、继续推音频,下一次 call.start 应答 BB D1 07 01(已在通话模式)。
  修后 8 轮 0 丢,边播译文边点结束也在 0.5s 内收到 D2。

译文积压越说越慢(bluetooth_manager CallTranslationDownlink 重写)
- 积压 >1s 裁掉超过 150ms 的静音;>2s 起 WSOLA 变速不变调,4s 达 1.3x;不丢内容
  (兜底丢弃常量默认关)。新增 getCallDownlinkStats 与收尾总账日志。
- 原因:译文音频突发产出、只能 1 倍速播,说话人不停顿时源时长≈原话时长,队列只涨不缩,
  实测 A 路积压 6.2s。修后同量压测最大 3.2s,绝大部分时间 <1.5s,光裁静音就省了 10%,
  变速最高只用到 1.11x。

其它
- translation_controller 收尾那几行日志提到 warn(release 包能看到 call.stop 有没有发)。

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
main
Rodger-Wang 3 weeks ago
parent
commit
d1db3e5951
  1. 16
      apps/client/lib/modules/translation/controllers/translation_controller.dart
  2. 44
      apps/client/local_plugins/azure_speech/android/src/main/kotlin/com/yunqiinnovation/azure_speech/AzureSpeechPlugin.kt
  3. 6
      apps/client/local_plugins/azure_speech/android/src/main/kotlin/com/yunqiinnovation/azure_speech/BesCallPcmBridge.kt
  4. 2
      apps/client/local_plugins/bluetooth_manager/android/src/main/kotlin/com/example/bluetooth_manager/AiAudioDownlink.kt
  5. 63
      apps/client/local_plugins/bluetooth_manager/android/src/main/kotlin/com/example/bluetooth_manager/BluetoothManager.kt
  6. 2
      apps/client/local_plugins/bluetooth_manager/android/src/main/kotlin/com/example/bluetooth_manager/BluetoothManagerPlugin.kt
  7. 386
      apps/client/local_plugins/bluetooth_manager/android/src/main/kotlin/com/example/bluetooth_manager/CallTranslationDownlink.kt

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")

386
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
// legB = 己方听(对端译文),legA = 对端听(己方译文)
private val legBQueue = LinkedBlockingQueue<ByteArray>(QUEUE_CAP)
private val legAQueue = LinkedBlockingQueue<ByteArray>(QUEUE_CAP)
/** 一条腿:源 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
// 两条腿各自独立的编码器状态,不能共用(与解码侧 spk/mic 双 handle 同理)
private var encHandleA: Long = 0L
private var encHandleB: Long = 0L
// 统计
var pushedFrames = 0L // 收到的源帧数(20ms)
var sentFrames = 0L // 实际发出的帧数
var trimmedSamples = 0L // 裁掉的静音样本
var droppedSamples = 0L // 兜底丢弃的样本
var maxBacklogSamples = 0L
var lastRate = 1.0
// 攒够 640 字节才能编一帧,余量留到下次
private var pendingA = ByteArray(0)
private var pendingB = ByteArray(0)
private val pendingLock = Any()
/** 还没播出去的源音频(队列里的 + 变速器里没读完的),单位样本 */
fun backlogSamples(): Long = srcSamples + stretcher.unreadSamples()
fun backlogSec(): Double = backlogSamples() / SAMPLE_RATE.toDouble()
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
}
}
val rest = merged.copyOfRange(bulkLen, merged.size)
if (isA) pendingA = rest else pendingB = rest
frames = out
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
}
val queue = if (isA) legAQueue else legBQueue
for (f in frames) {
// 队列满说明下行速度跟不上产出(通常是耳机没在收),丢最旧的保实时
if (!queue.offer(f)) {
queue.poll()
queue.offer(f)
Log.w(TAG, "${l.name} 路积压超过 ${HARD_CAP_SEC}s,丢弃最旧音频至 ${target / SAMPLE_RATE.toDouble()}s")
}
}
if (isA) pushedA += frames.size else pushedB += frames.size
val bl = l.backlogSamples()
if (bl > l.maxBacklogSamples) l.maxBacklogSamples = bl
}
}
/** 积压时把长静音裁到只剩 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++
}
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))
}
}
/**
* 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))
}
private fun releaseEncoders() {
if (encHandleA != 0L) {
runCatching { G722Codec.release(encHandleA) }
encHandleA = 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
}
if (encHandleB != 0L) {
runCatching { G722Codec.release(encHandleB) }
encHandleB = 0L
}
}
}

Loading…
Cancel
Save