Browse Source

上传新版本的通话翻译代码

weicu
liwei1dao 5 months ago
parent
commit
f35bce68d8
  1. 8
      lib/data/services/ble_manager.dart
  2. 111
      local_plugins/azure_speech/android/src/main/kotlin/com/yunqiinnovation/azure_speech/AstCallbacks.kt
  3. 27
      local_plugins/azure_speech/android/src/main/kotlin/com/yunqiinnovation/azure_speech/DoubaoE2ETranslateHelper.kt
  4. BIN
      local_plugins/device_jieli/android/libs/jl_bluetooth_rcsp_V4.2.0_beta2_40214_20251224.aar
  5. 2
      local_plugins/device_jieli/android/src/main/kotlin/com/jielihome/jielihome/feature/record/JieliDeviceRecordPort.kt
  6. 281
      local_plugins/device_jieli/android/src/main/kotlin/com/jielihome/jielihome/feature/translation/runtime/JieliAITranslationBridge.kt
  7. 17
      local_plugins/device_jieli/android/src/main/kotlin/com/jielihome/jielihome/feature/translation/runtime/NoOpAITranslationApi.kt
  8. 466
      local_plugins/device_jieli/android/src/main/kotlin/com/jielihome/jielihome/feature/translation/runtime/RcspTranslationRuntime.kt
  9. BIN
      local_plugins/device_jieli/example/android/app/libs/jl_bluetooth_rcsp_V4.2.0_beta1_40255_20260324.aar
  10. BIN
      local_plugins/device_jieli/example/android/app/libs/jl_bluetooth_rcsp_V4.2.0_beta2_40214_20251224.aar

8
lib/data/services/ble_manager.dart

@ -265,8 +265,12 @@ class BleManager extends GetxService {
_eventSub?.cancel(); _eventSub?.cancel();
_eventSub = Jielihome.instance.events.listen( _eventSub = Jielihome.instance.events.listen(
(event) { (event) {
// 总开关日志:每一个事件都打一次,便于诊断 native → dart 通道是否畅通 // 总开关日志:每一个事件都打一次,便于诊断 native → dart 通道是否畅通。
print('[BLE] [EVT] ${event.runtimeType}'); // 通话翻译 PCM 上推 (TranslationAudioEvent) 是高频事件(20ms 一帧),
// 日志量大且无诊断价值,单独跳过;其它事件保留 runtimeType 打印。
if (event is! TranslationAudioEvent) {
print('[BLE] [EVT] ${event.runtimeType}');
}
if (event is AdapterStatusEvent) { if (event is AdapterStatusEvent) {
print('[BLE] [EVT] 适配器状态 enabled=${event.enabled} hasBle=${event.hasBle}'); print('[BLE] [EVT] 适配器状态 enabled=${event.enabled} hasBle=${event.hasBle}');
} else if (event is ScanStatusEvent) { } else if (event is ScanStatusEvent) {

111
local_plugins/azure_speech/android/src/main/kotlin/com/yunqiinnovation/azure_speech/AstCallbacks.kt

@ -209,6 +209,9 @@ class AzureAstCallback(
"error" to error "error" to error
) )
) )
// 与豆包 / 百炼对齐:合成失败时也要 markEnd,让下游 bridge 把已累积的半段 PCM 立即刷出,
// 否则会粘到下一段 utterance 的 isFinal 才一起下发。
audioWriter?.markEnd()
} }
override fun onSynthesisProgress( override fun onSynthesisProgress(
@ -459,6 +462,9 @@ class IflytekAstCallback(
"error" to error "error" to error
) )
) )
// 与豆包 / 百炼对齐:合成失败时也要 markEnd,让下游 bridge 把已累积的半段 PCM 立即刷出,
// 否则会粘到下一段 utterance 的 isFinal 才一起下发。
audioWriter?.markEnd()
} }
override fun onSynthesisProgress( override fun onSynthesisProgress(
@ -536,7 +542,35 @@ class DoubaoAstCallback(
private val tag = "DoubaoAstCallback" private val tag = "DoubaoAstCallback"
private val doubaoFinalSourceTextCache: MutableMap<String, String> = mutableMapOf() private val doubaoFinalSourceTextCache: MutableMap<String, String> = 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<String, SessionStat> = java.util.concurrent.ConcurrentHashMap()
private fun statOf(sessionId: String): SessionStat =
sessionStats.getOrPut(sessionId) { SessionStat() }
override fun onSessionStarted(sessionId: String) { override fun onSessionStarted(sessionId: String) {
sessionStats[sessionId] = SessionStat()
FileLogger.i(tag, "[E2E-DIAG] [$serviceId/$direction] onSessionStarted sid=$sessionId")
eventSender.send( eventSender.send(
mapOf( mapOf(
"type" to "serviceInitialized", "type" to "serviceInitialized",
@ -548,6 +582,8 @@ class DoubaoAstCallback(
} }
override fun onPartialSourceText(sessionId: String, text: String) { override fun onPartialSourceText(sessionId: String, text: String) {
val s = statOf(sessionId)
s.partialSourceTextCount++
eventSender.send( eventSender.send(
mapOf( mapOf(
"type" to "recognizing", "type" to "recognizing",
@ -561,6 +597,12 @@ class DoubaoAstCallback(
} }
override fun onFinalSourceText(sessionId: String, finalText: String) { 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" val key = "$serviceId:$sessionId"
doubaoFinalSourceTextCache[key] = finalText doubaoFinalSourceTextCache[key] = finalText
eventSender.send( eventSender.send(
@ -576,6 +618,8 @@ class DoubaoAstCallback(
} }
override fun onPartialText(sessionId: String, text: String) { override fun onPartialText(sessionId: String, text: String) {
val s = statOf(sessionId)
s.partialTextCount++
eventSender.send( eventSender.send(
mapOf( mapOf(
"type" to "translatedInterim", "type" to "translatedInterim",
@ -590,6 +634,23 @@ class DoubaoAstCallback(
} }
override fun onPartialAudio(sessionId: String, data: ByteArray) { 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) audioWriter?.write(data)
} }
@ -598,6 +659,17 @@ class DoubaoAstCallback(
finalText: String, finalText: String,
finalAudio: ByteArray 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( eventSender.send(
mapOf( mapOf(
"type" to "translated", "type" to "translated",
@ -615,9 +687,18 @@ class DoubaoAstCallback(
// } // }
// 火山豆包段尾:触发下游 RCSP runtime 立即整段下发。 // 火山豆包段尾:触发下游 RCSP runtime 立即整段下发。
audioWriter?.markEnd() audioWriter?.markEnd()
sessionStats.remove(sessionId)
} }
override fun onSessionError(sessionId: String, code: Int, message: String) { 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( eventSender.send(
mapOf( mapOf(
"type" to "error", "type" to "error",
@ -630,9 +711,20 @@ class DoubaoAstCallback(
) )
// 错误也要把当前 buffer 释放掉,避免残留。 // 错误也要把当前 buffer 释放掉,避免残留。
audioWriter?.markEnd() audioWriter?.markEnd()
sessionStats.remove(sessionId)
} }
override fun onFinalTranslatedText(sessionId: String, finalText: String) { 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 key = "$serviceId:$sessionId"
val original = doubaoFinalSourceTextCache.remove(key) ?: "" val original = doubaoFinalSourceTextCache.remove(key) ?: ""
eventSender.send( 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
}
} }
/** /**

27
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) 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() private var conf: Config = Config()
@ -264,6 +282,15 @@ class DoubaoE2ETranslateHelper(
return 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()) { if (data.isNotEmpty()) {
recvAudio.write(data) recvAudio.write(data)
Log.d(TAG, "onMessage: partial audio appended size=${data.size}") Log.d(TAG, "onMessage: partial audio appended size=${data.size}")

BIN
local_plugins/device_jieli/android/libs/jl_bluetooth_rcsp_V4.2.0_beta2_40214_20251224.aar

Binary file not shown.

2
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 val addr = device?.address
if (impl != null) { if (impl != null) {
runCatching { 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<Int> { impl.exitMode(object : OnRcspActionCallback<Int> {
override fun onSuccess(d: BluetoothDevice?, t: Int?) { override fun onSuccess(d: BluetoothDevice?, t: Int?) {
Log.i(TAG, "[APP->SDK] exitMode onSuccess addr=${d?.address} t=$t") Log.i(TAG, "[APP->SDK] exitMode onSuccess addr=${d?.address} t=$t")

281
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<String, ByteArrayOutputStream>()
private val bufferLock = Any()
private val encodeSeq = AtomicLong(0)
/** 调试统计 */
@Volatile private var rxFirstLogged = false
private val rxFirstPerSource = java.util.concurrent.ConcurrentHashMap<Int, Boolean>()
// ─── 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
}
}

17
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 import com.jieli.bluetooth.interfaces.rcsp.translation.IAITranslationApi
/** /**
* SDK 的 [TranslationImpl] 构造函数要求一个 [IAITranslationApi]。 * SDK 的 `TranslationImpl` 构造函数要求一个 [IAITranslationApi]。本类是给**纯录音通路**
* 我们走的是「外部翻译服务」路线,SDK 自带的 AI 流程不启用,所以这里给个空实现。 * 占位用的空实现:上层只想拿原始 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 { internal class NoOpAITranslationApi : IAITranslationApi {
override fun isWorking(): Boolean = false override fun isWorking(): Boolean = false

466
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 package com.jielihome.jielihome.feature.translation.runtime
import android.bluetooth.BluetoothDevice import android.bluetooth.BluetoothDevice
import android.util.Log
import com.jieli.bluetooth.bean.translation.AudioData import com.jieli.bluetooth.bean.translation.AudioData
import com.jieli.bluetooth.bean.translation.TranslationMode import com.jieli.bluetooth.bean.translation.TranslationMode
import com.jieli.bluetooth.constant.Constants 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.impl.rcsp.translation.TranslationImpl
import com.jieli.bluetooth.interfaces.rcsp.callback.OnRcspActionCallback import com.jieli.bluetooth.interfaces.rcsp.callback.OnRcspActionCallback
import com.jieli.bluetooth.interfaces.rcsp.translation.TranslationCallback 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.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 翻译模式运行时。 * RCSP 翻译模式运行时(call translation 路径专用)。
*
* audioType:
* - [Constants.AUDIO_TYPE_OPUS](默认):上行 OPUS → 解码 PCM;下行 PCM → **整段** 编码 OPUS 写回
* - [Constants.AUDIO_TYPE_PCM]:上下行直接走 PCM,不经编解码
* *
* # 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 把每帧当成"独立的一段 * 之前的 runtime 大约 500 行:自己用 [TranslationCallback.onReceiveAudioData] 拿上行 OPUS、
* TTS 起点",解码器频繁 reset → 杂音 / 断续。 * 自己 OPUS 编码 TTS、自己维护 `WriteScheduler` 串行队列调 [TranslationImpl.writeAudioData]。
* 实测耳机端解码器频繁 reset 导致杂音 / 断续。
* *
* 现在按 demo:[feedTtsPcm] 把 PCM 累积到 per-leg 缓冲;上层在每段 utterance 末尾把 * 现在按官方 demo `AITranslationImpl` 的标准做法:
* `isFinal=true` 透传过来;runtime 这一刻: * 1. 构造 [TranslationImpl] 时传入真正的 [JieliAITranslationBridge],SDK 通过它把上行
* 1. 把累积 PCM 写到临时 .pcm 文件; * OPUS 推给我们解码;
* 2. `OpusManager.encodeFile(pcmFile, opusFile, callback)` 离线整段编码(带 head 的默认 OpusOption); * 2. 下行 TTS PCM 累积成段后通过 [AITranslationCallback.onTranslateResult] 交给 SDK,
* 3. 整段读出 → **一个** [AudioData] → [WriteScheduler] 单次入队下发。 * 由 SDK 内部完成 cmd=52 切包 / 写时序 / 缓冲水位([TranslationImpl.PushDataWrapper])。
* *
* # 文件命名 * # 上下行 source 字段语义
* 临时文件位于 [tempDir],按 `tts_<source>_<seq>.<pcm|opus>` 命名;编码完成异步删除。
*
* # source 字段写回方向
* - 通话翻译给「对端听」 → AudioData.source = SOURCE_E_SCO_UP_LINK * - 通话翻译给「对端听」 → AudioData.source = SOURCE_E_SCO_UP_LINK
* - 通话翻译给「本机听」 → AudioData.source = SOURCE_E_SCO_DOWN_LINK * - 通话翻译给「本机听」 → AudioData.source = SOURCE_E_SCO_DOWN_LINK
* - 录音/音视频/面对面 → AudioData.source = SOURCE_PHONE_MIC(SDK 按 mode 自分发) * - 录音/音视频/面对面 → AudioData.source = SOURCE_PHONE_MIC(SDK 按 mode 自分发)
@ -56,90 +43,44 @@ class RcspTranslationRuntime(
/** 解码(或直传 PCM)后的音频上行;source 为 SDK 原值 */ /** 解码(或直传 PCM)后的音频上行;source 为 SDK 原值 */
private val onPcm: (source: Int, pcm: ByteArray) -> Unit, private val onPcm: (source: Int, pcm: ByteArray) -> Unit,
private val onError: (code: Int, msg: String?) -> 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 opusPacketSize: Int = if (mode.channel == 2) 80 else 200,
) { ) {
private val translationImpl = TranslationImpl(btManager, NoOpAITranslationApi(), device) companion object {
private val isPcmMode = mode.audioType == Constants.AUDIO_TYPE_PCM private const val TAG = "RcspTranslationRuntime"
private val sampleRateHz = mode.sampleRate.takeIf { it > 0 } ?: 16000 }
/** 入栈解码器:仅 OPUS 模式下创建 */ /** SDK 的 AI hook 实现:上行解码、下行 TTS 接力都在这里。 */
private val upDecoder: OpusStreamDecoder? = if (isPcmMode) null else OpusStreamDecoder( private val bridge = JieliAITranslationBridge(
channel = 1, mode = mode,
packetSize = opusPacketSize, tempDir = tempDir,
sampleRate = sampleRateHz, onPcm = onPcm,
onPcm = { pcm -> onPcm(AudioData.SOURCE_E_SCO_UP_LINK, pcm) }, onError = onError,
onError = { c, m -> onError(c, "upDecoder: $m") }, opusPacketSize = opusPacketSize,
) )
private val downDecoder: OpusStreamDecoder? = private val translationImpl = TranslationImpl(btManager, bridge, device)
if (!isPcmMode && mode.mode == TranslationMode.MODE_CALL_TRANSLATION) OpusStreamDecoder( private val isPcmMode = mode.audioType == Constants.AUDIO_TYPE_PCM
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<Int, Boolean>()
/**
* 仅订阅 mode 变化事件用于日志 / 异常上报;不在这里消费音频
* (音频走 [bridge] 的 `IAITranslationApi.writeAudio`,避免双路重复消费)。
*/
private val translationCallback = object : TranslationCallback { private val translationCallback = object : TranslationCallback {
override fun onModeChange(d: BluetoothDevice, m: TranslationMode) { 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) { override fun onReceiveAudioData(d: BluetoothDevice, data: AudioData) {
rxCount++ // 音频走 bridge.writeAudio;这里不消费,避免与 bridge 双路重复解码。
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")
}
} }
override fun onError(d: BluetoothDevice, code: Int, msg: String) { 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<Boolean>) -> Unit,
) {
private val queue = ArrayDeque<AudioData>(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<Boolean> {
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<String, ByteArrayOutputStream>()
/** 每条 leg 是否仍处于"等待首段触发"状态(路线 B 首段优先 flush 用)。 */
private val firstSegmentPending = mutableMapOf<String, Boolean>()
private val bufferLock = Any()
private var encodeSeq = 0L
/** 周期性 flush 巡检线程;start() 启动,stop() 关闭。 */
private var flushWatcher: ScheduledExecutorService? = null
/** 启动前置校验 + 进入翻译模式 */
fun start(): Result<Unit> { fun start(): Result<Unit> {
if (!translationImpl.isInit) { if (!translationImpl.isInit) {
return Result.failure(IllegalStateException("RCSP not init for ${device.address}")) return Result.failure(IllegalStateException("RCSP not init for ${device.address}"))
@ -285,172 +104,31 @@ class RcspTranslationRuntime(
} }
if (!tempDir.exists()) tempDir.mkdirs() if (!tempDir.exists()) tempDir.mkdirs()
upDecoder?.start()
downDecoder?.start()
stereoDecoder?.start()
translationImpl.addTranslationCallback(translationCallback) 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) translationImpl.enterMode(mode, translationCallback)
// OPUS 模式启动周期性 flush 巡检:每 [PERIODIC_FLUSH_MS] 切一次缓冲。
if (!isPcmMode) startPeriodicFlushWatcher()
return Result.success(Unit) 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<Triple<Int, String, ByteArray>>()
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。
* *
* 策略: * 行为完全委托给 [JieliAITranslationBridge.feedTtsPcm]:
* - PCM 模式:每次调用直接当作一段 AudioData 下发(SDK 内部按 blockMtu 分片)。 * - PCM 模式:每次调用作为一段 [AudioData] 立即交付给 SDK;
* - OPUS 模式:纯时间驱动 flush —— 火山等端到端服务的段尾事件不可靠, * - OPUS 模式:累积,`isFinal=true` 时整段编码 → `onTranslateResult` 交给 SDK。
* 不能依赖 isFinal 或字节量阈值。改为:
* a) feedTtsPcm 仅追加 buffer,不主动 flush(除非 isFinal=true);
* b) [PERIODIC_FLUSH_MS] (1s) 巡检线程每秒切走 buffer 并整段编码下发;
* c) `isFinal=true` 仍短路立即 flush(兼容主动信号,不再依赖)。
* 最大听感延迟 ≤ 1s,且短句 / 没有段尾事件的服务都不会卡死。
* *
* @param outputStreamId 决定 source: * 调用方契约:**必须**在 utterance 末尾调一次 `isFinal=true`,否则音频会一直累积
* [TranslationStreams.OUT_UPLINK] / [TranslationStreams.OUT_DOWNLINK] 用于通话翻译; * 到 2MB 兜底上限才出。
* 其它(speaker/localPlayback)由 ModeHandler 自己处理,不会落到这里。
* @param isFinal 本帧是否为当前 utterance 的最后一帧;触发 buffer 残余立即 flush。
*/ */
fun feedTtsPcm(outputStreamId: String, pcm: ByteArray, isFinal: Boolean): Boolean { fun feedTtsPcm(outputStreamId: String, pcm: ByteArray, isFinal: Boolean): Boolean =
val source = when (outputStreamId) { bridge.feedTtsPcm(outputStreamId, pcm, isFinal)
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 stop() { fun stop() {
runCatching { flushWatcher?.shutdownNow() }
flushWatcher = null
runCatching { runCatching {
Log.i(TAG, "[APP->SDK] exitMode addr=${device.address} mode=${mode.mode}") Log.i(TAG, "[APP->SDK] exitMode addr=${device.address} mode=${mode.mode}")
translationImpl.exitMode(object : OnRcspActionCallback<Int> { translationImpl.exitMode(object : OnRcspActionCallback<Int> {
@ -463,34 +141,8 @@ class RcspTranslationRuntime(
}) })
} }
runCatching { translationImpl.removeTranslationCallback(translationCallback) } runCatching { translationImpl.removeTranslationCallback(translationCallback) }
// bridge 由 SDK 在 exitMode 后通过 stopTranslating 自动 release decoders;
// 这里不重复调,避免双重 stop。
runCatching { translationImpl.destroy() } 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)
} }
} }

BIN
local_plugins/device_jieli/example/android/app/libs/jl_bluetooth_rcsp_V4.2.0_beta1_40255_20260324.aar

Binary file not shown.

BIN
local_plugins/device_jieli/example/android/app/libs/jl_bluetooth_rcsp_V4.2.0_beta2_40214_20251224.aar

Binary file not shown.
Loading…
Cancel
Save