diff --git a/lib/data/services/ble_manager.dart b/lib/data/services/ble_manager.dart index 7c7477378..03e38014a 100644 --- a/lib/data/services/ble_manager.dart +++ b/lib/data/services/ble_manager.dart @@ -265,8 +265,12 @@ class BleManager extends GetxService { _eventSub?.cancel(); _eventSub = Jielihome.instance.events.listen( (event) { - // 总开关日志:每一个事件都打一次,便于诊断 native → dart 通道是否畅通 - print('[BLE] [EVT] ${event.runtimeType}'); + // 总开关日志:每一个事件都打一次,便于诊断 native → dart 通道是否畅通。 + // 通话翻译 PCM 上推 (TranslationAudioEvent) 是高频事件(20ms 一帧), + // 日志量大且无诊断价值,单独跳过;其它事件保留 runtimeType 打印。 + if (event is! TranslationAudioEvent) { + print('[BLE] [EVT] ${event.runtimeType}'); + } if (event is AdapterStatusEvent) { print('[BLE] [EVT] 适配器状态 enabled=${event.enabled} hasBle=${event.hasBle}'); } else if (event is ScanStatusEvent) { diff --git a/local_plugins/azure_speech/android/src/main/kotlin/com/yunqiinnovation/azure_speech/AstCallbacks.kt b/local_plugins/azure_speech/android/src/main/kotlin/com/yunqiinnovation/azure_speech/AstCallbacks.kt index 93c3d20c5..60ff8eb8d 100644 --- a/local_plugins/azure_speech/android/src/main/kotlin/com/yunqiinnovation/azure_speech/AstCallbacks.kt +++ b/local_plugins/azure_speech/android/src/main/kotlin/com/yunqiinnovation/azure_speech/AstCallbacks.kt @@ -209,6 +209,9 @@ class AzureAstCallback( "error" to error ) ) + // 与豆包 / 百炼对齐:合成失败时也要 markEnd,让下游 bridge 把已累积的半段 PCM 立即刷出, + // 否则会粘到下一段 utterance 的 isFinal 才一起下发。 + audioWriter?.markEnd() } override fun onSynthesisProgress( @@ -459,6 +462,9 @@ class IflytekAstCallback( "error" to error ) ) + // 与豆包 / 百炼对齐:合成失败时也要 markEnd,让下游 bridge 把已累积的半段 PCM 立即刷出, + // 否则会粘到下一段 utterance 的 isFinal 才一起下发。 + audioWriter?.markEnd() } override fun onSynthesisProgress( @@ -536,7 +542,35 @@ class DoubaoAstCallback( private val tag = "DoubaoAstCallback" private val doubaoFinalSourceTextCache: MutableMap = mutableMapOf() + // ─── 段尾信号诊断 ───────────────────────────────────────────── + // 怀疑豆包端到端没有给出段尾 → bridge 永远等不到 isFinal=true → 音频卡在缓冲。 + // 这里按 session 维护一份计数,把"豆包到底回调了哪几个方法、各自多少次/多少字节" + // 全部打印出来。一段正常 utterance 期望看到: + // 1) onSessionStarted — 1 次 + // 2) onPartialSourceText/onPartialText — 多次(流式) + // 3) onPartialAudio — 多次(流式 TTS) + // 4) onFinalSourceText / onFinalTranslatedText — 各 1 次 + // 5) onSessionFinished — 1 次(→ markEnd 触发下游 isFinal=true) + // 实测如果第 5 步缺失,就说明豆包没出段尾信号。 + private data class SessionStat( + @Volatile var partialAudioBytes: Long = 0, + @Volatile var partialAudioFrames: Long = 0, + @Volatile var partialTextCount: Long = 0, + @Volatile var partialSourceTextCount: Long = 0, + @Volatile var finalSourceTextHit: Boolean = false, + @Volatile var finalTranslatedTextHit: Boolean = false, + @Volatile var sessionFinishedHit: Boolean = false, + @Volatile var sessionErrorHit: Boolean = false, + @Volatile var firstAudioLogged: Boolean = false, + @Volatile var startedAtMs: Long = System.currentTimeMillis(), + ) + private val sessionStats: MutableMap = java.util.concurrent.ConcurrentHashMap() + private fun statOf(sessionId: String): SessionStat = + sessionStats.getOrPut(sessionId) { SessionStat() } + override fun onSessionStarted(sessionId: String) { + sessionStats[sessionId] = SessionStat() + FileLogger.i(tag, "[E2E-DIAG] [$serviceId/$direction] onSessionStarted sid=$sessionId") eventSender.send( mapOf( "type" to "serviceInitialized", @@ -548,6 +582,8 @@ class DoubaoAstCallback( } override fun onPartialSourceText(sessionId: String, text: String) { + val s = statOf(sessionId) + s.partialSourceTextCount++ eventSender.send( mapOf( "type" to "recognizing", @@ -561,6 +597,12 @@ class DoubaoAstCallback( } override fun onFinalSourceText(sessionId: String, finalText: String) { + val s = statOf(sessionId) + s.finalSourceTextHit = true + FileLogger.i( + tag, + "[E2E-DIAG] [$serviceId/$direction] onFinalSourceText sid=$sessionId text=\"$finalText\"" + ) val key = "$serviceId:$sessionId" doubaoFinalSourceTextCache[key] = finalText eventSender.send( @@ -576,6 +618,8 @@ class DoubaoAstCallback( } override fun onPartialText(sessionId: String, text: String) { + val s = statOf(sessionId) + s.partialTextCount++ eventSender.send( mapOf( "type" to "translatedInterim", @@ -590,6 +634,23 @@ class DoubaoAstCallback( } override fun onPartialAudio(sessionId: String, data: ByteArray) { + val s = statOf(sessionId) + s.partialAudioBytes += data.size + s.partialAudioFrames++ + if (!s.firstAudioLogged) { + s.firstAudioLogged = true + FileLogger.i( + tag, + "[E2E-DIAG] [$serviceId/$direction] onPartialAudio FIRST sid=$sessionId size=${data.size}" + ) + } + // 每 50 帧或累计每 ~16KB 再打一次进度,避免淹没 logcat + if (s.partialAudioFrames % 50L == 0L) { + FileLogger.d( + tag, + "[E2E-DIAG] [$serviceId/$direction] onPartialAudio progress sid=$sessionId frames=${s.partialAudioFrames} bytes=${s.partialAudioBytes}" + ) + } audioWriter?.write(data) } @@ -598,6 +659,17 @@ class DoubaoAstCallback( finalText: String, finalAudio: ByteArray ) { + val s = statOf(sessionId) + s.sessionFinishedHit = true + val durMs = System.currentTimeMillis() - s.startedAtMs + FileLogger.i( + tag, + "[E2E-DIAG] [$serviceId/$direction] onSessionFinished sid=$sessionId " + + "dur=${durMs}ms partialAudio=${s.partialAudioBytes}B/${s.partialAudioFrames}f " + + "partialText=${s.partialTextCount} partialSrc=${s.partialSourceTextCount} " + + "finalSrcHit=${s.finalSourceTextHit} finalTransHit=${s.finalTranslatedTextHit} " + + "finalText=\"$finalText\" finalAudio=${finalAudio.size}B → markEnd()" + ) eventSender.send( mapOf( "type" to "translated", @@ -615,9 +687,18 @@ class DoubaoAstCallback( // } // 火山豆包段尾:触发下游 RCSP runtime 立即整段下发。 audioWriter?.markEnd() + sessionStats.remove(sessionId) } override fun onSessionError(sessionId: String, code: Int, message: String) { + val s = statOf(sessionId) + s.sessionErrorHit = true + FileLogger.w( + tag, + "[E2E-DIAG] [$serviceId/$direction] onSessionError sid=$sessionId code=$code msg=$message " + + "partialAudio=${s.partialAudioBytes}B/${s.partialAudioFrames}f " + + "finalSrcHit=${s.finalSourceTextHit} finishedHit=${s.sessionFinishedHit} → markEnd()" + ) eventSender.send( mapOf( "type" to "error", @@ -630,9 +711,20 @@ class DoubaoAstCallback( ) // 错误也要把当前 buffer 释放掉,避免残留。 audioWriter?.markEnd() + sessionStats.remove(sessionId) } override fun onFinalTranslatedText(sessionId: String, finalText: String) { + val s = statOf(sessionId) + s.finalTranslatedTextHit = true + // 注意:onFinalTranslatedText 是"字幕段"结束(TranslationSubtitleEnd),早于真正 + // 的 TTS 段结束(TTSSentenceEnd)。所以这里**不能**调 markEnd —— 此刻 TTS 音频还 + // 没来;要等 onTtsSegmentEnd。 + FileLogger.i( + tag, + "[E2E-DIAG] [$serviceId/$direction] onFinalTranslatedText sid=$sessionId " + + "text=\"$finalText\" partialAudio=${s.partialAudioBytes}B/${s.partialAudioFrames}f" + ) val key = "$serviceId:$sessionId" val original = doubaoFinalSourceTextCache.remove(key) ?: "" eventSender.send( @@ -647,6 +739,25 @@ class DoubaoAstCallback( ) ) } + + /** + * TTS 单句段尾。连续翻译模式下豆包不会发 [onSessionFinished],只有这一个事件能告诉我们 + * "本段 utterance 的 TTS 音频全部到达"。**必须**在这里 markEnd,否则下游 bridge 永远 + * 等不到 isFinal=true,PCM 累积到 2MB 兜底限才下发。 + */ + override fun onTtsSegmentEnd(sessionId: String) { + val s = statOf(sessionId) + FileLogger.i( + tag, + "[E2E-DIAG] [$serviceId/$direction] onTtsSegmentEnd sid=$sessionId " + + "audio=${s.partialAudioBytes}B/${s.partialAudioFrames}f → markEnd()" + ) + audioWriter?.markEnd() + // 下一段重新累计;保留 session-level 文本统计(仍可能继续触发同 sessionId 的字幕事件)。 + s.partialAudioBytes = 0 + s.partialAudioFrames = 0 + s.firstAudioLogged = false + } } /** diff --git a/local_plugins/azure_speech/android/src/main/kotlin/com/yunqiinnovation/azure_speech/DoubaoE2ETranslateHelper.kt b/local_plugins/azure_speech/android/src/main/kotlin/com/yunqiinnovation/azure_speech/DoubaoE2ETranslateHelper.kt index 7f5c9a2bd..d2beac3a0 100644 --- a/local_plugins/azure_speech/android/src/main/kotlin/com/yunqiinnovation/azure_speech/DoubaoE2ETranslateHelper.kt +++ b/local_plugins/azure_speech/android/src/main/kotlin/com/yunqiinnovation/azure_speech/DoubaoE2ETranslateHelper.kt @@ -53,6 +53,24 @@ class DoubaoE2ETranslateHelper( /** 会话失败或取消 */ fun onSessionError(sessionId: String, code: Int, message: String) + + /** + * 单句 TTS 段尾(服务端 `Type.TTSSentenceEnd`)。 + * + * # 与 [onSessionFinished] 的区别 + * 两者层级完全不同: + * - [onSessionFinished] 是**整个会话生命周期结束**(对应服务端 `SessionFinished`)。 + * 连续翻译模式 (`startContinuousTranslation`) 下 WebSocket 长连接 + sessionId + * 全程保持,服务端永远不发 `SessionFinished`,所以这个回调实际上不会触发。 + * - [onTtsSegmentEnd] 是**单句 utterance 的 TTS 音频全部到达**。每段都触发一次, + * 是连续翻译里唯一能感知"这一句音频完结了"的事件。 + * + * # 谁需要它 + * 接耳机 RCSP 的实现必须 override 并在这里 `audioWriter.markEnd()`,把"段尾 + * isFinal=true"信号传给下游 bridge,否则 PCM 会一直累积到 2MB 兜底才下发。 + * default no-op 是为不破坏与 TTS 段尾无关的旧实现的兼容性。 + */ + fun onTtsSegmentEnd(sessionId: String) {} } private var conf: Config = Config() @@ -264,6 +282,15 @@ class DoubaoE2ETranslateHelper( return } + // 连续翻译模式下 WS 长连接 + sessionId 全程保持,服务端不发 SessionFinished; + // 每段 utterance 的真正音频段尾是 TTSSentenceEnd。下游(接耳机)需要这个信号 + // 触发 markEnd 把 bridge buffer 整段下发。 + if (event == Type.TTSSentenceEnd) { + Log.d(TAG, "onMessage: TTSSentenceEnd → onTtsSegmentEnd sid=$sessionId") + callback?.onTtsSegmentEnd(sessionId) + return + } + if (data.isNotEmpty()) { recvAudio.write(data) Log.d(TAG, "onMessage: partial audio appended size=${data.size}") diff --git a/local_plugins/device_jieli/android/libs/jl_bluetooth_rcsp_V4.2.0_beta2_40214_20251224.aar b/local_plugins/device_jieli/android/libs/jl_bluetooth_rcsp_V4.2.0_beta2_40214_20251224.aar deleted file mode 100644 index 15eb8cf5e..000000000 Binary files a/local_plugins/device_jieli/android/libs/jl_bluetooth_rcsp_V4.2.0_beta2_40214_20251224.aar and /dev/null differ diff --git a/local_plugins/device_jieli/android/src/main/kotlin/com/jielihome/jielihome/feature/record/JieliDeviceRecordPort.kt b/local_plugins/device_jieli/android/src/main/kotlin/com/jielihome/jielihome/feature/record/JieliDeviceRecordPort.kt index 1a26be54b..8a2a953b0 100644 --- a/local_plugins/device_jieli/android/src/main/kotlin/com/jielihome/jielihome/feature/record/JieliDeviceRecordPort.kt +++ b/local_plugins/device_jieli/android/src/main/kotlin/com/jielihome/jielihome/feature/record/JieliDeviceRecordPort.kt @@ -287,7 +287,7 @@ class JieliDeviceRecordPort( val addr = device?.address if (impl != null) { runCatching { - Log.i(TAG, "[APP->SDK] exitMode addr=$addr mode=${TranslationMode.MODE_CALL_TRANSLATION_WITH_STEREO}") + Log.i(TAG, "[APP->SDK] exitMode addr=$addr mode=${TranslationMode.MODE_CALL_RECORD}") impl.exitMode(object : OnRcspActionCallback { override fun onSuccess(d: BluetoothDevice?, t: Int?) { Log.i(TAG, "[APP->SDK] exitMode onSuccess addr=${d?.address} t=$t") diff --git a/local_plugins/device_jieli/android/src/main/kotlin/com/jielihome/jielihome/feature/translation/runtime/JieliAITranslationBridge.kt b/local_plugins/device_jieli/android/src/main/kotlin/com/jielihome/jielihome/feature/translation/runtime/JieliAITranslationBridge.kt new file mode 100644 index 000000000..cfb04d34f --- /dev/null +++ b/local_plugins/device_jieli/android/src/main/kotlin/com/jielihome/jielihome/feature/translation/runtime/JieliAITranslationBridge.kt @@ -0,0 +1,281 @@ +package com.jielihome.jielihome.feature.translation.runtime + +import android.util.Log +import com.jieli.bluetooth.bean.translation.AudioData +import com.jieli.bluetooth.bean.translation.TranslationMode +import com.jieli.bluetooth.bean.translation.TranslationResult +import com.jieli.bluetooth.constant.Constants +import com.jieli.bluetooth.interfaces.rcsp.translation.AITranslationCallback +import com.jieli.bluetooth.interfaces.rcsp.translation.IAITranslationApi +import com.jieli.jl_audio_decode.callback.OnStateCallback +import com.jieli.jl_audio_decode.opus.OpusManager +import com.jielihome.jielihome.audio.OpusStreamDecoder +import com.jielihome.jielihome.feature.translation.TranslationStreams +import java.io.ByteArrayOutputStream +import java.io.File +import java.util.concurrent.atomic.AtomicLong + +/** + * 杰理 SDK 的 [IAITranslationApi] 实现 —— call translation 路径专用的「上下行接力桥」。 + * + * # 为什么要有这个 bridge + * 早先的实现(旧 [RcspTranslationRuntime])走的是"自己手写 writeAudioData 队列"的方案: + * 1. 用 [com.jieli.bluetooth.interfaces.rcsp.translation.TranslationCallback.onReceiveAudioData] + * 拿上行 OPUS,自己解码; + * 2. 把外部翻译服务回送的 PCM 周期切片 → OPUS 编码 → 用自己的 WriteScheduler 串行调 + * [com.jieli.bluetooth.impl.rcsp.translation.TranslationImpl.writeAudioData] 推回耳机。 + * + * 实测踩到的两个坑: + * - 周期切片(每 1s)→ 多段独立 OPUS → 耳机端解码器频繁 reset → 杂音 / 断续; + * - 自己维护 in-flight=1 的写队列,溢出丢最老 → utterance 中段被截掉。 + * + * 官方 demo 的真正用法([com.jieli.bt.sdk.tool.translation.AITranslationImpl]):实现 + * [IAITranslationApi],把整段 utterance 编一次 OPUS,包成 [AudioData] 塞进 + * [TranslationResult.translationTTSData] 然后调 [AITranslationCallback.onTranslateResult], + * **SDK 内部** 接管 cmd=52 切包 / 发送时序 / 缓冲水位。 + * + * # 上下行职责 + * - **上行**(headset → app):SDK 主动调 [writeAudio],audioData.type = OPUS / PCM; + * 按 [AudioData.source] 分流到 up/down/stereo decoder,解码后通过 [onPcm] 抛出去。 + * - **下行**(app → headset):[feedTtsPcm] 累积外部翻译服务回送的 PCM,**纯 isFinal 驱动**: + * `isFinal=true` 时把累积 PCM 编成一段 OPUS,组装 [TranslationResult] 丢给 SDK。 + * 不再有周期 flush。如果上层一直不发 isFinal,缓冲到 [PCM_BUFFER_HARD_LIMIT_BYTES] + * 兜底强制 flush,避免内存爆。 + * + * # 与 [NoOpAITranslationApi] 的区别 + * - [NoOpAITranslationApi]:纯录音通路(JieliAssistantPort / JieliDeviceRecordPort 等) + * 只想拿原始 PCM,不希望 SDK 触发 AI 流程。 + * - 本类:通话翻译需要"上行解码 + 下行 SDK 接管 TTS 注入",要走 SDK 的标准 AI hook。 + * + * # 线程 + * - SDK 回调线程不固定;解码器 / 编码器调用都不假设线程。 + * - [feedTtsPcm] 可被任意线程调用,per-leg 缓冲用 [bufferLock] 互斥。 + * - OPUS 文件编码异步([OpusManager.encodeFile] 内部线程),完成后回到 SDK callback。 + */ +class JieliAITranslationBridge( + private val mode: TranslationMode, + private val tempDir: File, + /** 解码后的上行 PCM 出口;source 为 SDK 原值(SOURCE_E_SCO_UP_LINK / DOWN_LINK / MIX)。 */ + private val onPcm: (source: Int, pcm: ByteArray) -> Unit, + /** 解码 / 编码失败的统一出口。 */ + private val onError: (code: Int, msg: String?) -> Unit, + private val opusPacketSize: Int = if (mode.channel == 2) 80 else 200, +) : IAITranslationApi { + + companion object { + private const val TAG = "JieliAITranslationBridge" + /** 单段 utterance PCM 上限:≈ 60s 16k mono 16bit = 1.92MB。超限强制 flush 兜底。 */ + private const val PCM_BUFFER_HARD_LIMIT_BYTES = 2 * 1024 * 1024 + } + + private val isPcmMode = mode.audioType == Constants.AUDIO_TYPE_PCM + private val sampleRateHz = mode.sampleRate.takeIf { it > 0 } ?: 16000 + + @Volatile private var sdkCallback: AITranslationCallback? = null + @Volatile private var working = false + + /** 上行 OPUS 解码器:通话翻译有 up/down 两路单声道,stereo 模式只一路双声道。 */ + private val upDecoder: OpusStreamDecoder? = if (isPcmMode) null else OpusStreamDecoder( + channel = 1, + packetSize = opusPacketSize, + sampleRate = sampleRateHz, + onPcm = { pcm -> onPcm(AudioData.SOURCE_E_SCO_UP_LINK, pcm) }, + onError = { c, m -> onError(c, "upDecoder: $m") }, + ) + private val downDecoder: OpusStreamDecoder? = + if (!isPcmMode && mode.mode == TranslationMode.MODE_CALL_TRANSLATION) OpusStreamDecoder( + channel = 1, + packetSize = opusPacketSize, + sampleRate = sampleRateHz, + onPcm = { pcm -> onPcm(AudioData.SOURCE_E_SCO_DOWN_LINK, pcm) }, + onError = { c, m -> onError(c, "downDecoder: $m") }, + ) else null + private val stereoDecoder: OpusStreamDecoder? = + if (!isPcmMode && mode.mode == TranslationMode.MODE_CALL_TRANSLATION_WITH_STEREO) OpusStreamDecoder( + channel = 2, packetSize = 80, + sampleRate = sampleRateHz, + onPcm = { pcm -> onPcm(AudioData.SOURCE_E_SCO_MIX, pcm) }, + onError = { c, m -> onError(c, "stereoDecoder: $m") }, + ) else null + + /** Per-leg PCM 缓冲:key = outputStreamId(OUT_UPLINK / OUT_DOWNLINK / OUT_SPEAKER)。 */ + private val pcmBuffers = mutableMapOf() + private val bufferLock = Any() + private val encodeSeq = AtomicLong(0) + + /** 调试统计 */ + @Volatile private var rxFirstLogged = false + private val rxFirstPerSource = java.util.concurrent.ConcurrentHashMap() + + // ─── IAITranslationApi 实现 ─────────────────────────────────────────── + + override fun isWorking(): Boolean = working + + override fun startTranslating(mode: TranslationMode, callback: AITranslationCallback) { + Log.i(TAG, "[SDK->APP] startTranslating mode=${mode.mode} type=${mode.audioType} sr=${mode.sampleRate} ch=${mode.channel}") + sdkCallback = callback + if (!tempDir.exists()) tempDir.mkdirs() + upDecoder?.start() + downDecoder?.start() + stereoDecoder?.start() + working = true + callback.onStart() + } + + override fun stopTranslating() { + Log.i(TAG, "[SDK->APP] stopTranslating") + working = false + runCatching { upDecoder?.stop() } + runCatching { downDecoder?.stop() } + runCatching { stereoDecoder?.stop() } + // 兜底:剩余 buffer 直接丢,不补发 onTranslateResult(utterance 已被打断)。 + synchronized(bufferLock) { + pcmBuffers.values.forEach { it.reset() } + pcmBuffers.clear() + } + sdkCallback?.onStop(0, "stopTranslating") + sdkCallback = null + } + + override fun writeAudio(data: AudioData) { + if (!working) return + if (data.type != mode.audioType) { + // SDK 极少数情况下会推与 mode.audioType 不一致的帧(例如握手期残留);忽略不报错。 + return + } + val payload = data.audioData ?: return + if (payload.isEmpty()) return + + // 首帧(总) + if (!rxFirstLogged) { + rxFirstLogged = true + Log.i(TAG, "[SDK->APP] writeAudio FIRST source=${data.source} type=${data.type} size=${payload.size}") + } + // 首帧(按 source) + if (rxFirstPerSource.putIfAbsent(data.source, true) == null) { + Log.i(TAG, "[SDK->APP] writeAudio FIRST source=${data.source} size=${payload.size}") + } + + if (isPcmMode) { + onPcm(data.source, payload) + return + } + when (data.source) { + AudioData.SOURCE_E_SCO_UP_LINK -> upDecoder?.feedEncoded(payload) + AudioData.SOURCE_E_SCO_DOWN_LINK -> downDecoder?.feedEncoded(payload) + AudioData.SOURCE_E_SCO_MIX -> stereoDecoder?.feedEncoded(payload) + else -> Log.w(TAG, "[SDK->APP] writeAudio unknown source=${data.source} size=${payload.size}") + } + } + + // ─── 下行:外部翻译服务的 TTS PCM 接力 ───────────────────────────────── + + /** + * 把外部翻译服务回送的 TTS PCM 接力给 SDK。 + * + * 策略: + * - PCM 模式:每次调用作为一段 [AudioData] 直接交给 SDK。 + * - OPUS 模式:累积 per-leg buffer;`isFinal=true` 或缓冲到 [PCM_BUFFER_HARD_LIMIT_BYTES] + * 时整段编码 → 通过 [AITranslationCallback.onTranslateResult] 交给 SDK,SDK 内部 + * 完成 cmd=52 切包 / 发送 / 速率控制。 + * + * @return 是否接受本帧(false 表示丢弃,例如未启动 / 缓冲爆 / SDK callback 缺失) + */ + fun feedTtsPcm(outputStreamId: String, pcm: ByteArray, isFinal: Boolean): Boolean { + if (!working) return false + val cb = sdkCallback ?: return false + val source = sourceOf(outputStreamId) + + if (isPcmMode) { + postTranslationResult(cb, source, Constants.AUDIO_TYPE_PCM, pcm) + return true + } + // OPUS 模式:累积 + isFinal/超限触发整段编码 + val pendingFlush: ByteArray? = synchronized(bufferLock) { + val buf = pcmBuffers.getOrPut(outputStreamId) { ByteArrayOutputStream() } + if (pcm.isNotEmpty()) { + if (buf.size() + pcm.size > PCM_BUFFER_HARD_LIMIT_BYTES) { + Log.w(TAG, "feedTtsPcm leg=$outputStreamId buffer hit hard limit ${PCM_BUFFER_HARD_LIMIT_BYTES}, force flush") + val out = buf.toByteArray() + pcm + buf.reset() + return@synchronized out + } + buf.write(pcm) + } + if (isFinal) { + val out = buf.toByteArray() + buf.reset() + if (out.isEmpty()) null else out + } else null + } + if (pendingFlush != null) { + encodeAndDeliverAsync(cb, source, outputStreamId, pendingFlush) + } + return true + } + + private fun postTranslationResult( + cb: AITranslationCallback, + source: Int, + audioType: Int, + bytes: ByteArray, + ) { + val result = TranslationResult().also { + it.translationTTSData = AudioData(source, audioType, bytes) + } + runCatching { cb.onTranslateResult(result) }.onFailure { + Log.w(TAG, "onTranslateResult threw: ${it.message}") + } + } + + private fun encodeAndDeliverAsync( + cb: AITranslationCallback, + source: Int, + leg: String, + pcmBytes: ByteArray, + ) { + val seq = encodeSeq.incrementAndGet() + val pcmFile = File(tempDir, "tts_${source}_$seq.pcm") + val opusFile = File(tempDir, "tts_${source}_$seq.opus") + try { + pcmFile.writeBytes(pcmBytes) + } catch (e: Throwable) { + Log.w(TAG, "encodeAndDeliver leg=$leg seq=$seq write pcm failed: ${e.message}") + return + } + // OpusManager 默认 OpusOption(带 head),与 demo `MachineTranslation.tryToTTS` 一致; + // SDK 接收端按"完整带 head 的 OPUS 段"解析,所以一次 onTranslateResult = 一段 utterance。 + val encoder = OpusManager() + encoder.encodeFile(pcmFile.absolutePath, opusFile.absolutePath, object : OnStateCallback { + override fun onStart() {} + override fun onComplete(path: String?) { + try { + val opusBytes = if (opusFile.exists()) opusFile.readBytes() else null + if (opusBytes == null || opusBytes.isEmpty()) { + Log.w(TAG, "encodeAndDeliver leg=$leg seq=$seq: empty opus output") + } else { + Log.d(TAG, "encodeAndDeliver leg=$leg seq=$seq pcm=${pcmBytes.size}B opus=${opusBytes.size}B → onTranslateResult") + postTranslationResult(cb, source, Constants.AUDIO_TYPE_OPUS, opusBytes) + } + } finally { + runCatching { encoder.release() } + runCatching { pcmFile.delete() } + runCatching { opusFile.delete() } + } + } + override fun onError(code: Int, message: String?) { + Log.w(TAG, "encodeAndDeliver leg=$leg seq=$seq encodeFile error code=$code msg=$message") + runCatching { encoder.release() } + runCatching { pcmFile.delete() } + runCatching { opusFile.delete() } + onError(code, "encodeFile $leg: $message") + } + }) + } + + private fun sourceOf(outputStreamId: String): Int = when (outputStreamId) { + TranslationStreams.OUT_UPLINK -> AudioData.SOURCE_E_SCO_UP_LINK + TranslationStreams.OUT_DOWNLINK -> AudioData.SOURCE_E_SCO_DOWN_LINK + else -> AudioData.SOURCE_PHONE_MIC + } +} diff --git a/local_plugins/device_jieli/android/src/main/kotlin/com/jielihome/jielihome/feature/translation/runtime/NoOpAITranslationApi.kt b/local_plugins/device_jieli/android/src/main/kotlin/com/jielihome/jielihome/feature/translation/runtime/NoOpAITranslationApi.kt index 448c4deed..8cb7eee7d 100644 --- a/local_plugins/device_jieli/android/src/main/kotlin/com/jielihome/jielihome/feature/translation/runtime/NoOpAITranslationApi.kt +++ b/local_plugins/device_jieli/android/src/main/kotlin/com/jielihome/jielihome/feature/translation/runtime/NoOpAITranslationApi.kt @@ -6,8 +6,21 @@ import com.jieli.bluetooth.interfaces.rcsp.translation.AITranslationCallback import com.jieli.bluetooth.interfaces.rcsp.translation.IAITranslationApi /** - * SDK 的 [TranslationImpl] 构造函数要求一个 [IAITranslationApi]。 - * 我们走的是「外部翻译服务」路线,SDK 自带的 AI 流程不启用,所以这里给个空实现。 + * SDK 的 `TranslationImpl` 构造函数要求一个 [IAITranslationApi]。本类是给**纯录音通路** + * 占位用的空实现:上层只想拿原始 OPUS / PCM 流(自己挂 `TranslationCallback.onReceiveAudioData`), + * 不希望 SDK 触发任何 AI 翻译流程 / TTS 注入。 + * + * # 谁在用 + * - [com.jielihome.jielihome.feature.assistant.JieliAssistantPort](AI 助理上行通路) + * - [com.jielihome.jielihome.feature.record.JieliDeviceRecordPort](通话录音 MODE_CALL_RECORD) + * - [com.jielihome.jielihome.feature.translation.mode.RecordModeHandler](MODE_RECORD) + * - [com.jielihome.jielihome.feature.translation.mode.RecordingTranslationModeHandler](MODE_RECORDING_TRANSLATION) + * - [com.jielihome.jielihome.feature.translation.TranslationFeature](兜底 TranslationImpl 实例) + * + * # 谁**不在**用了 + * Call translation(mode 3 / 6)由 [JieliAITranslationBridge] 接管——那条路要的是 + * "SDK 接管 TTS 切包注入",必须给真正的 IAITranslationApi 实现,避免自己写 writeAudioData + * 队列触发的"耳机端解码器频繁 reset → 杂音 / 断续"。 */ internal class NoOpAITranslationApi : IAITranslationApi { override fun isWorking(): Boolean = false diff --git a/local_plugins/device_jieli/android/src/main/kotlin/com/jielihome/jielihome/feature/translation/runtime/RcspTranslationRuntime.kt b/local_plugins/device_jieli/android/src/main/kotlin/com/jielihome/jielihome/feature/translation/runtime/RcspTranslationRuntime.kt index 4a6c3545c..aedd7bdc5 100644 --- a/local_plugins/device_jieli/android/src/main/kotlin/com/jielihome/jielihome/feature/translation/runtime/RcspTranslationRuntime.kt +++ b/local_plugins/device_jieli/android/src/main/kotlin/com/jielihome/jielihome/feature/translation/runtime/RcspTranslationRuntime.kt @@ -1,6 +1,7 @@ package com.jielihome.jielihome.feature.translation.runtime import android.bluetooth.BluetoothDevice +import android.util.Log import com.jieli.bluetooth.bean.translation.AudioData import com.jieli.bluetooth.bean.translation.TranslationMode import com.jieli.bluetooth.constant.Constants @@ -8,42 +9,28 @@ import com.jieli.bluetooth.impl.JL_BluetoothManager import com.jieli.bluetooth.impl.rcsp.translation.TranslationImpl import com.jieli.bluetooth.interfaces.rcsp.callback.OnRcspActionCallback import com.jieli.bluetooth.interfaces.rcsp.translation.TranslationCallback -import android.util.Log -import com.jielihome.jielihome.audio.OpusStreamDecoder -import com.jielihome.jielihome.feature.translation.TranslationStreams -import com.jieli.jl_audio_decode.callback.OnStateCallback -import com.jieli.jl_audio_decode.opus.OpusManager -import java.io.ByteArrayOutputStream import java.io.File -import java.util.ArrayDeque -import java.util.concurrent.Executors -import java.util.concurrent.ScheduledExecutorService -import java.util.concurrent.TimeUnit -import java.util.concurrent.atomic.AtomicBoolean /** - * RCSP 翻译模式运行时。 - * - * audioType: - * - [Constants.AUDIO_TYPE_OPUS](默认):上行 OPUS → 解码 PCM;下行 PCM → **整段** 编码 OPUS 写回 - * - [Constants.AUDIO_TYPE_PCM]:上下行直接走 PCM,不经编解码 + * RCSP 翻译模式运行时(call translation 路径专用)。 * - * # TTS 回灌策略(与官方 demo `MachineTranslation.tryToTTS` / `OpusHelper.encodeFile` 对齐) + * # 角色 + * 把"进入/退出 SDK 翻译模式"的 RCSP 状态机生命周期,跟「上下行音频接力」绑在一起。 + * 真正的音频路由由 [JieliAITranslationBridge] 这个 [com.jieli.bluetooth.interfaces.rcsp.translation.IAITranslationApi] + * 实现负责(见它头注释)。 * - * 之前是流式 [com.jielihome.jielihome.audio.OpusStreamEncoder]:每出一帧 OPUS 立即包成 - * AudioData 调 [TranslationImpl.writeAudioData],导致耳机端 RCSP 把每帧当成"独立的一段 - * TTS 起点",解码器频繁 reset → 杂音 / 断续。 + * # 与旧实现的区别 + * 之前的 runtime 大约 500 行:自己用 [TranslationCallback.onReceiveAudioData] 拿上行 OPUS、 + * 自己 OPUS 编码 TTS、自己维护 `WriteScheduler` 串行队列调 [TranslationImpl.writeAudioData]。 + * 实测耳机端解码器频繁 reset 导致杂音 / 断续。 * - * 现在按 demo:[feedTtsPcm] 把 PCM 累积到 per-leg 缓冲;上层在每段 utterance 末尾把 - * `isFinal=true` 透传过来;runtime 这一刻: - * 1. 把累积 PCM 写到临时 .pcm 文件; - * 2. `OpusManager.encodeFile(pcmFile, opusFile, callback)` 离线整段编码(带 head 的默认 OpusOption); - * 3. 整段读出 → **一个** [AudioData] → [WriteScheduler] 单次入队下发。 + * 现在按官方 demo `AITranslationImpl` 的标准做法: + * 1. 构造 [TranslationImpl] 时传入真正的 [JieliAITranslationBridge],SDK 通过它把上行 + * OPUS 推给我们解码; + * 2. 下行 TTS PCM 累积成段后通过 [AITranslationCallback.onTranslateResult] 交给 SDK, + * 由 SDK 内部完成 cmd=52 切包 / 写时序 / 缓冲水位([TranslationImpl.PushDataWrapper])。 * - * # 文件命名 - * 临时文件位于 [tempDir],按 `tts__.` 命名;编码完成异步删除。 - * - * # source 字段写回方向 + * # 上下行 source 字段语义 * - 通话翻译给「对端听」 → AudioData.source = SOURCE_E_SCO_UP_LINK * - 通话翻译给「本机听」 → AudioData.source = SOURCE_E_SCO_DOWN_LINK * - 录音/音视频/面对面 → AudioData.source = SOURCE_PHONE_MIC(SDK 按 mode 自分发) @@ -56,90 +43,44 @@ class RcspTranslationRuntime( /** 解码(或直传 PCM)后的音频上行;source 为 SDK 原值 */ private val onPcm: (source: Int, pcm: ByteArray) -> Unit, private val onError: (code: Int, msg: String?) -> Unit, + /** OPUS 解码 packetSize;mono call=200,stereo=80 */ private val opusPacketSize: Int = if (mode.channel == 2) 80 else 200, ) { - private val translationImpl = TranslationImpl(btManager, NoOpAITranslationApi(), device) - private val isPcmMode = mode.audioType == Constants.AUDIO_TYPE_PCM - private val sampleRateHz = mode.sampleRate.takeIf { it > 0 } ?: 16000 + companion object { + private const val TAG = "RcspTranslationRuntime" + } - /** 入栈解码器:仅 OPUS 模式下创建 */ - private val upDecoder: OpusStreamDecoder? = if (isPcmMode) null else OpusStreamDecoder( - channel = 1, - packetSize = opusPacketSize, - sampleRate = sampleRateHz, - onPcm = { pcm -> onPcm(AudioData.SOURCE_E_SCO_UP_LINK, pcm) }, - onError = { c, m -> onError(c, "upDecoder: $m") }, + /** SDK 的 AI hook 实现:上行解码、下行 TTS 接力都在这里。 */ + private val bridge = JieliAITranslationBridge( + mode = mode, + tempDir = tempDir, + onPcm = onPcm, + onError = onError, + opusPacketSize = opusPacketSize, ) - private val downDecoder: OpusStreamDecoder? = - if (!isPcmMode && mode.mode == TranslationMode.MODE_CALL_TRANSLATION) OpusStreamDecoder( - channel = 1, - packetSize = opusPacketSize, - sampleRate = sampleRateHz, - onPcm = { pcm -> onPcm(AudioData.SOURCE_E_SCO_DOWN_LINK, pcm) }, - onError = { c, m -> onError(c, "downDecoder: $m") }, - ) else null - - private val stereoDecoder: OpusStreamDecoder? = - if (!isPcmMode && mode.mode == TranslationMode.MODE_CALL_TRANSLATION_WITH_STEREO) OpusStreamDecoder( - channel = 2, packetSize = 80, - sampleRate = sampleRateHz, - onPcm = { pcm -> onPcm(AudioData.SOURCE_E_SCO_MIX, pcm) }, - onError = { c, m -> onError(c, "stereoDecoder: $m") }, - ) else null - - /** 调试统计:SDK onReceiveAudioData 被回调多少次、按 (source,type) 分桶 */ - @Volatile private var rxCount = 0L - @Volatile private var rxTypeMismatch = 0L - @Volatile private var rxNullPayload = 0L - @Volatile private var rxLastReportMs = 0L - /** 首帧 flag:首次收到任意 AudioData 时立即打一条 INFO 级别日志(不受 1s 节流影响)。 */ - @Volatile private var rxFirstLogged = false - /** 首帧(按 source 维度):uplink / downlink / mix 分别各打一条,确认链路通畅。 */ - private val rxFirstPerSource = java.util.concurrent.ConcurrentHashMap() + private val translationImpl = TranslationImpl(btManager, bridge, device) + private val isPcmMode = mode.audioType == Constants.AUDIO_TYPE_PCM + /** + * 仅订阅 mode 变化事件用于日志 / 异常上报;不在这里消费音频 + * (音频走 [bridge] 的 `IAITranslationApi.writeAudio`,避免双路重复消费)。 + */ private val translationCallback = object : TranslationCallback { override fun onModeChange(d: BluetoothDevice, m: TranslationMode) { - Log.i(TAG, "[SDK<-DEV] onModeChange addr=${d.address} mode=${m.mode} type=${m.audioType} sr=${m.sampleRate} ch=${m.channel} strategy=${m.recordingStrategy}") + Log.i( + TAG, + "[SDK<-DEV] onModeChange addr=${d.address} mode=${m.mode} type=${m.audioType} " + + "sr=${m.sampleRate} ch=${m.channel} strategy=${m.recordingStrategy}" + ) + if (m.mode == TranslationMode.MODE_IDLE && mode.mode != TranslationMode.MODE_IDLE) { + onError(-1, "headset exited mode=${mode.mode} → MODE_IDLE") + } } override fun onReceiveAudioData(d: BluetoothDevice, data: AudioData) { - rxCount++ - val payloadSize = data.audioData?.size ?: 0 - // 首帧(总):立即打点,避免节流导致"完全静默"误判 - if (!rxFirstLogged) { - rxFirstLogged = true - Log.i(TAG, "[SDK<-DEV] onReceiveAudioData FIRST FRAME addr=${d.address} source=${data.source} type=${data.type} size=$payloadSize expect.type=${mode.audioType}") - } - // 首帧(按 source):分别打点,快速判断上行/下行是否都在推 - if (rxFirstPerSource.putIfAbsent(data.source, true) == null) { - Log.i(TAG, "[SDK<-DEV] onReceiveAudioData FIRST source=${data.source} type=${data.type} size=$payloadSize") - } - val now = System.currentTimeMillis() - if (now - rxLastReportMs >= 1000L) { - Log.d(TAG, "[SDK<-DEV] RX stats (last 1s): total=$rxCount typeMismatch=$rxTypeMismatch nullPayload=$rxNullPayload expect.type=${mode.audioType} last.source=${data.source} last.type=${data.type} payloadSize=$payloadSize") - rxCount = 0; rxTypeMismatch = 0; rxNullPayload = 0; rxLastReportMs = now - } - if (data.type != mode.audioType) { - rxTypeMismatch++ - return - } - val payload = data.audioData - if (payload == null) { - rxNullPayload++ - return - } - if (isPcmMode) { - onPcm(data.source, payload) - return - } - when (data.source) { - AudioData.SOURCE_E_SCO_UP_LINK -> upDecoder?.feedEncoded(payload) - AudioData.SOURCE_E_SCO_DOWN_LINK -> downDecoder?.feedEncoded(payload) - AudioData.SOURCE_E_SCO_MIX -> stereoDecoder?.feedEncoded(payload) - else -> Log.w(TAG, "[SDK<-DEV] onReceiveAudioData unknown source=${data.source} dropped") - } + // 音频走 bridge.writeAudio;这里不消费,避免与 bridge 双路重复解码。 } override fun onError(d: BluetoothDevice, code: Int, msg: String) { @@ -148,129 +89,7 @@ class RcspTranslationRuntime( } } - companion object { - private const val TAG = "RcspTranslationRuntime" - /** 每腿 writeAudioData 缓冲队列上限。整段下发后队列里几乎不会堆积,留 8 个槽足够。 */ - private const val WRITE_QUEUE_LIMIT = 8 - /** 单段 PCM 上限:≈ 60s 16k mono 16bit = 1.92MB,正常 utterance 远低于此值。 */ - private const val PCM_BUFFER_HARD_LIMIT_BYTES = 2 * 1024 * 1024 - /** - * 周期性 flush 间隔:每隔此毫秒数巡检一次 buffer,**只要非空就 flush**。 - * - * 历史实测:200ms 时耳机端 RCSP 解码器会被频繁 reset 直至死机;1s 长期稳定。 - * 当前实验:500ms —— 通话翻译"播报中间断一下"的优化尝试。需真机长时压测 - * (≥10 分钟连续讲话)确认无解码器 reset 风暴,否则回退到 1_000L。 - */ - private const val PERIODIC_FLUSH_MS = 1000L - /** - * 首段优先 flush 阈值(路线 B):每个 utterance 的**首段**达到此毫秒数对应字节量 - * 就立即 flush,不等周期。让用户开口到对方听到的延迟尽可能小(< 200ms)。 - * - * 触发后转入周期模式直到 isFinal=true 重置。火山不发 isFinal 的会话里, - * 只有第一句享受首段优先;后续句统一走 [PERIODIC_FLUSH_MS] 周期。这是为了 - * 避免连续讲话被反复识别成"新 utterance"导致每 100ms 切一段、reset 频率失控。 - */ - private const val FIRST_SEGMENT_FLUSH_MS = 100L - } - - /** - * 串行化的 writeAudioData 调度器(每个 source 一个)。 - * - * RCSP/SPP 写入要求:上一帧 [writeAudioData] 的回调返回(成功或失败)之前不要再发, - * 否则 SDK 会拒收新请求并报 "Operation in progress"。这里维护一个 in-flight ≤ 1 - * 的 FIFO 队列,溢出时丢最老的帧(保实时性,不堆积延迟)。 - */ - private class WriteScheduler( - private val tag: String, - private val send: (AudioData, OnRcspActionCallback) -> Unit, - ) { - private val queue = ArrayDeque(WRITE_QUEUE_LIMIT) - private val inFlight = AtomicBoolean(false) - private val lock = Any() - @Volatile private var droppedSinceReport = 0L - @Volatile private var lastReportMs = 0L - - fun enqueue(data: AudioData) { - val toSend: AudioData? = synchronized(lock) { - if (inFlight.compareAndSet(false, true)) { - data - } else { - if (queue.size >= WRITE_QUEUE_LIMIT) { - queue.pollFirst() // 丢最老 - droppedSinceReport++ - val now = System.currentTimeMillis() - if (now - lastReportMs >= 1000L) { - Log.w(TAG, "writeAudioData[$tag] queue full, dropped=$droppedSinceReport (last 1s)") - droppedSinceReport = 0 - lastReportMs = now - } - } - queue.offerLast(data) - null - } - } - if (toSend != null) dispatch(toSend) - } - - private fun dispatch(data: AudioData) { - send(data, object : OnRcspActionCallback { - override fun onSuccess(d: BluetoothDevice?, ok: Boolean?) = onComplete() - override fun onError(d: BluetoothDevice?, err: com.jieli.bluetooth.bean.base.BaseError?) { - val msg = err?.message.orEmpty() - if (!msg.contains("Operation in progress", ignoreCase = true)) { - Log.w(TAG, "writeAudioData[$tag] error: code=${err?.code} msg=$msg") - } - onComplete() - } - }) - } - - private fun onComplete() { - val next: AudioData? = synchronized(lock) { - val n = queue.pollFirst() - if (n == null) inFlight.set(false) - n - } - if (next != null) dispatch(next) - } - - fun clear() { - synchronized(lock) { - queue.clear() - inFlight.set(false) - } - } - } - - private val uplinkWriter = WriteScheduler("uplink") { data, cb -> - translationImpl.writeAudioData(data, cb) - } - private val downlinkWriter = WriteScheduler("downlink") { data, cb -> - translationImpl.writeAudioData(data, cb) - } - private val phoneMicWriter = WriteScheduler("phoneMic") { data, cb -> - translationImpl.writeAudioData(data, cb) - } - /** 兜底:PCM 直传 / 未识别 source 走这条;和 OPUS 编码器互不抢 in-flight 槽位。 */ - private val rawWriter = WriteScheduler("raw") { data, cb -> - translationImpl.writeAudioData(data, cb) - } - - /** - * Per-leg PCM 缓冲。`feedTtsPcm` 累积写入,`isFinal=true` 时整段编码下发后清空。 - * - * key = outputStreamId(OUT_UPLINK / OUT_DOWNLINK / OUT_SPEAKER 等) - */ - private val pcmBuffers = mutableMapOf() - /** 每条 leg 是否仍处于"等待首段触发"状态(路线 B 首段优先 flush 用)。 */ - private val firstSegmentPending = mutableMapOf() - private val bufferLock = Any() - private var encodeSeq = 0L - - /** 周期性 flush 巡检线程;start() 启动,stop() 关闭。 */ - private var flushWatcher: ScheduledExecutorService? = null - - /** 启动前置校验 + 进入翻译模式 */ + /** 启动前置校验 + 进入翻译模式。 */ fun start(): Result { if (!translationImpl.isInit) { return Result.failure(IllegalStateException("RCSP not init for ${device.address}")) @@ -285,172 +104,31 @@ class RcspTranslationRuntime( } if (!tempDir.exists()) tempDir.mkdirs() - upDecoder?.start() - downDecoder?.start() - stereoDecoder?.start() - translationImpl.addTranslationCallback(translationCallback) - Log.i(TAG, "[APP->SDK] enterMode addr=${device.address} mode=${mode.mode} type=${mode.audioType} sr=${mode.sampleRate} ch=${mode.channel} strategy=${mode.recordingStrategy} (waiting for onModeChange/onReceiveAudioData)") + Log.i( + TAG, + "[APP->SDK] enterMode addr=${device.address} mode=${mode.mode} type=${mode.audioType} " + + "sr=${mode.sampleRate} ch=${mode.channel} strategy=${mode.recordingStrategy} " + + "(SDK will drive bridge.startTranslating)" + ) translationImpl.enterMode(mode, translationCallback) - - // OPUS 模式启动周期性 flush 巡检:每 [PERIODIC_FLUSH_MS] 切一次缓冲。 - if (!isPcmMode) startPeriodicFlushWatcher() return Result.success(Unit) } - private fun startPeriodicFlushWatcher() { - flushWatcher?.shutdownNow() - val exec = Executors.newSingleThreadScheduledExecutor { r -> - Thread(r, "rcsp-tts-flush").apply { isDaemon = true } - } - flushWatcher = exec - exec.scheduleAtFixedRate( - { runCatching { periodicFlush() } }, - PERIODIC_FLUSH_MS, PERIODIC_FLUSH_MS, TimeUnit.MILLISECONDS, - ) - } - - /** 每秒巡检:每条 leg 的 buffer 只要非空就立即切走 + 编码 + 下发(不依赖任何外部信号)。 */ - private fun periodicFlush() { - val toEncode = mutableListOf>() - synchronized(bufferLock) { - for ((leg, buf) in pcmBuffers) { - if (buf.size() == 0) continue - val source = when (leg) { - TranslationStreams.OUT_UPLINK -> AudioData.SOURCE_E_SCO_UP_LINK - TranslationStreams.OUT_DOWNLINK -> AudioData.SOURCE_E_SCO_DOWN_LINK - else -> AudioData.SOURCE_PHONE_MIC - } - val out = buf.toByteArray() - buf.reset() - Log.d(TAG, "periodicFlush: leg=$leg bytes=${out.size}") - toEncode.add(Triple(source, leg, out)) - } - } - // encode 走 async,必须释放锁后做。 - for ((source, leg, pcm) in toEncode) { - encodeAndDispatchAsync(source, leg, pcm) - } - } - /** - * 把外部翻译服务回送的 PCM 注入回耳机。 + * 把外部翻译服务回送的 PCM 接力给 SDK。 * - * 策略: - * - PCM 模式:每次调用直接当作一段 AudioData 下发(SDK 内部按 blockMtu 分片)。 - * - OPUS 模式:纯时间驱动 flush —— 火山等端到端服务的段尾事件不可靠, - * 不能依赖 isFinal 或字节量阈值。改为: - * a) feedTtsPcm 仅追加 buffer,不主动 flush(除非 isFinal=true); - * b) [PERIODIC_FLUSH_MS] (1s) 巡检线程每秒切走 buffer 并整段编码下发; - * c) `isFinal=true` 仍短路立即 flush(兼容主动信号,不再依赖)。 - * 最大听感延迟 ≤ 1s,且短句 / 没有段尾事件的服务都不会卡死。 + * 行为完全委托给 [JieliAITranslationBridge.feedTtsPcm]: + * - PCM 模式:每次调用作为一段 [AudioData] 立即交付给 SDK; + * - OPUS 模式:累积,`isFinal=true` 时整段编码 → `onTranslateResult` 交给 SDK。 * - * @param outputStreamId 决定 source: - * [TranslationStreams.OUT_UPLINK] / [TranslationStreams.OUT_DOWNLINK] 用于通话翻译; - * 其它(speaker/localPlayback)由 ModeHandler 自己处理,不会落到这里。 - * @param isFinal 本帧是否为当前 utterance 的最后一帧;触发 buffer 残余立即 flush。 + * 调用方契约:**必须**在 utterance 末尾调一次 `isFinal=true`,否则音频会一直累积 + * 到 2MB 兜底上限才出。 */ - fun feedTtsPcm(outputStreamId: String, pcm: ByteArray, isFinal: Boolean): Boolean { - val source = when (outputStreamId) { - TranslationStreams.OUT_UPLINK -> AudioData.SOURCE_E_SCO_UP_LINK - TranslationStreams.OUT_DOWNLINK -> AudioData.SOURCE_E_SCO_DOWN_LINK - else -> AudioData.SOURCE_PHONE_MIC - } - if (isPcmMode) { - // PCM 直传:每次调用一段 AudioData。`isFinal` 在该模式下不影响下发节奏。 - writeBack(AudioData(source, Constants.AUDIO_TYPE_PCM, pcm)) - return true - } - // OPUS 模式:累积到 buffer,三条 flush 触发链(按优先级): - // 1. isFinal=true → 立即 flush,并把首段优先重置回 true(下一句新算) - // 2. 首段优先(路线 B) → 当前 leg 处于 firstSegmentPending=true 且 buffer - // 累积到 FIRST_SEGMENT_FLUSH_MS 字节量时立即 flush, - // 让"开口到对方听见"的延迟最小化(< 200ms) - // 3. 周期 flush(路线 A) → 由 [PERIODIC_FLUSH_MS] 巡检线程接管 - val pendingFlush: ByteArray? = synchronized(bufferLock) { - val buf = pcmBuffers.getOrPut(outputStreamId) { ByteArrayOutputStream() } - if (pcm.isNotEmpty()) { - if (buf.size() + pcm.size > PCM_BUFFER_HARD_LIMIT_BYTES) { - Log.w(TAG, "feedTtsPcm: leg=$outputStreamId pcm buffer hit hard limit ${PCM_BUFFER_HARD_LIMIT_BYTES}, dropping segment") - buf.reset() - return false - } - buf.write(pcm) - } - // 1) 段尾信号:立即下发尾段并把首段标志重置回 true。 - if (isFinal) { - firstSegmentPending[outputStreamId] = true - val out = buf.toByteArray() - buf.reset() - return@synchronized if (out.isEmpty()) null else out - } - // 2) 首段优先(路线 B):当前**实验中临时屏蔽**——为观察 500ms 周期 flush 的纯效果。 - // 屏蔽后所有 flush 走两条路:a) isFinal=true 短路;b) PERIODIC_FLUSH_MS 周期。 - // 实验完成后若需恢复,把下面 if (false) 改回 if (pending) 即可。 - @Suppress("ConstantConditionIf") - if (false) { - val pending = firstSegmentPending.getOrDefault(outputStreamId, true) - if (pending) { - val firstFlushBytes = (sampleRateHz * 2 * FIRST_SEGMENT_FLUSH_MS / 1000L).toInt() - if (buf.size() >= firstFlushBytes) { - firstSegmentPending[outputStreamId] = false - val out = buf.toByteArray() - buf.reset() - Log.d(TAG, "firstSegmentFlush: leg=$outputStreamId bytes=${out.size}") - return@synchronized if (out.isEmpty()) null else out - } - } - } - null - } - if (pendingFlush != null) { - encodeAndDispatchAsync(source, outputStreamId, pendingFlush) - } - return true - } - - private fun encodeAndDispatchAsync(source: Int, leg: String, pcmBytes: ByteArray) { - val seq = synchronized(bufferLock) { ++encodeSeq } - val pcmFile = File(tempDir, "tts_${source}_$seq.pcm") - val opusFile = File(tempDir, "tts_${source}_$seq.opus") - try { - pcmFile.writeBytes(pcmBytes) - } catch (e: Throwable) { - Log.w(TAG, "encodeAndDispatch leg=$leg seq=$seq write pcm file failed: ${e.message}") - return - } - // OpusManager.encodeFile(in, out, callback) 默认 OpusOption(带 head),与 demo 一致。 - val encoder = OpusManager() - encoder.encodeFile(pcmFile.absolutePath, opusFile.absolutePath, object : OnStateCallback { - override fun onStart() {} - override fun onComplete(path: String?) { - try { - val opusBytes = if (opusFile.exists()) opusFile.readBytes() else null - if (opusBytes == null || opusBytes.isEmpty()) { - Log.w(TAG, "encodeAndDispatch leg=$leg seq=$seq: empty opus output") - } else { - Log.d(TAG, "encodeAndDispatch leg=$leg seq=$seq pcm=${pcmBytes.size}B opus=${opusBytes.size}B") - writeBack(AudioData(source, Constants.AUDIO_TYPE_OPUS, opusBytes)) - } - } finally { - runCatching { encoder.release() } - runCatching { pcmFile.delete() } - runCatching { opusFile.delete() } - } - } - override fun onError(code: Int, message: String?) { - Log.w(TAG, "encodeAndDispatch leg=$leg seq=$seq encodeFile error code=$code msg=$message") - runCatching { encoder.release() } - runCatching { pcmFile.delete() } - runCatching { opusFile.delete() } - onError(code, "encodeFile $leg: $message") - } - }) - } + fun feedTtsPcm(outputStreamId: String, pcm: ByteArray, isFinal: Boolean): Boolean = + bridge.feedTtsPcm(outputStreamId, pcm, isFinal) fun stop() { - runCatching { flushWatcher?.shutdownNow() } - flushWatcher = null runCatching { Log.i(TAG, "[APP->SDK] exitMode addr=${device.address} mode=${mode.mode}") translationImpl.exitMode(object : OnRcspActionCallback { @@ -463,34 +141,8 @@ class RcspTranslationRuntime( }) } runCatching { translationImpl.removeTranslationCallback(translationCallback) } + // bridge 由 SDK 在 exitMode 后通过 stopTranslating 自动 release decoders; + // 这里不重复调,避免双重 stop。 runCatching { translationImpl.destroy() } - runCatching { upDecoder?.stop() } - runCatching { downDecoder?.stop() } - runCatching { stereoDecoder?.stop() } - runCatching { - uplinkWriter.clear() - downlinkWriter.clear() - phoneMicWriter.clear() - rawWriter.clear() - } - synchronized(bufferLock) { - pcmBuffers.values.forEach { it.reset() } - pcmBuffers.clear() - firstSegmentPending.clear() - } - } - - /** - * 把 OPUS / PCM 帧入队到对应 source 的 [WriteScheduler],由调度器串行发起 RCSP - * writeAudioData,避免并发写入触发 "Operation in progress"。 - */ - private fun writeBack(data: AudioData) { - val scheduler = when (data.source) { - AudioData.SOURCE_E_SCO_UP_LINK -> uplinkWriter - AudioData.SOURCE_E_SCO_DOWN_LINK -> downlinkWriter - AudioData.SOURCE_PHONE_MIC -> phoneMicWriter - else -> rawWriter - } - scheduler.enqueue(data) } } diff --git a/local_plugins/device_jieli/example/android/app/libs/jl_bluetooth_rcsp_V4.2.0_beta1_40255_20260324.aar b/local_plugins/device_jieli/example/android/app/libs/jl_bluetooth_rcsp_V4.2.0_beta1_40255_20260324.aar new file mode 100644 index 000000000..0ff72f1c6 Binary files /dev/null and b/local_plugins/device_jieli/example/android/app/libs/jl_bluetooth_rcsp_V4.2.0_beta1_40255_20260324.aar differ diff --git a/local_plugins/device_jieli/example/android/app/libs/jl_bluetooth_rcsp_V4.2.0_beta2_40214_20251224.aar b/local_plugins/device_jieli/example/android/app/libs/jl_bluetooth_rcsp_V4.2.0_beta2_40214_20251224.aar deleted file mode 100644 index 15eb8cf5e..000000000 Binary files a/local_plugins/device_jieli/example/android/app/libs/jl_bluetooth_rcsp_V4.2.0_beta2_40214_20251224.aar and /dev/null differ