diff --git a/apps/client/lib/modules/translation/controllers/translation_controller.dart b/apps/client/lib/modules/translation/controllers/translation_controller.dart index aac0a3da..2a8e3b5c 100644 --- a/apps/client/lib/modules/translation/controllers/translation_controller.dart +++ b/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'); } } diff --git a/apps/client/local_plugins/azure_speech/android/src/main/kotlin/com/yunqiinnovation/azure_speech/AzureSpeechPlugin.kt b/apps/client/local_plugins/azure_speech/android/src/main/kotlin/com/yunqiinnovation/azure_speech/AzureSpeechPlugin.kt index 8d820ff5..6af1c40b 100644 --- a/apps/client/local_plugins/azure_speech/android/src/main/kotlin/com/yunqiinnovation/azure_speech/AzureSpeechPlugin.kt +++ b/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("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("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) } } diff --git a/apps/client/local_plugins/azure_speech/android/src/main/kotlin/com/yunqiinnovation/azure_speech/BesCallPcmBridge.kt b/apps/client/local_plugins/azure_speech/android/src/main/kotlin/com/yunqiinnovation/azure_speech/BesCallPcmBridge.kt index 40265bea..0e111fed 100644 --- a/apps/client/local_plugins/azure_speech/android/src/main/kotlin/com/yunqiinnovation/azure_speech/BesCallPcmBridge.kt +++ b/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,双保险。 } } diff --git a/apps/client/local_plugins/bluetooth_manager/android/src/main/kotlin/com/example/bluetooth_manager/AiAudioDownlink.kt b/apps/client/local_plugins/bluetooth_manager/android/src/main/kotlin/com/example/bluetooth_manager/AiAudioDownlink.kt index 7e69d94f..18d71e91 100644 --- a/apps/client/local_plugins/bluetooth_manager/android/src/main/kotlin/com/example/bluetooth_manager/AiAudioDownlink.kt +++ b/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") diff --git a/apps/client/local_plugins/bluetooth_manager/android/src/main/kotlin/com/example/bluetooth_manager/BluetoothManager.kt b/apps/client/local_plugins/bluetooth_manager/android/src/main/kotlin/com/example/bluetooth_manager/BluetoothManager.kt index 05280055..eecfa5a0 100644 --- a/apps/client/local_plugins/bluetooth_manager/android/src/main/kotlin/com/example/bluetooth_manager/BluetoothManager.kt +++ b/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) } diff --git a/apps/client/local_plugins/bluetooth_manager/android/src/main/kotlin/com/example/bluetooth_manager/BluetoothManagerPlugin.kt b/apps/client/local_plugins/bluetooth_manager/android/src/main/kotlin/com/example/bluetooth_manager/BluetoothManagerPlugin.kt index 562b1392..f4825bcd 100644 --- a/apps/client/local_plugins/bluetooth_manager/android/src/main/kotlin/com/example/bluetooth_manager/BluetoothManagerPlugin.kt +++ b/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("leg") val pcm = call.argument("pcm") diff --git a/apps/client/local_plugins/bluetooth_manager/android/src/main/kotlin/com/example/bluetooth_manager/CallTranslationDownlink.kt b/apps/client/local_plugins/bluetooth_manager/android/src/main/kotlin/com/example/bluetooth_manager/CallTranslationDownlink.kt index 021db893..e92e03a4 100644 --- a/apps/client/local_plugins/bluetooth_manager/android/src/main/kotlin/com/example/bluetooth_manager/CallTranslationDownlink.kt +++ b/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() // 未播的源 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(QUEUE_CAP) - private val legAQueue = LinkedBlockingQueue(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 = 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 - 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(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 + } } } }