|
|
|
@ -60,12 +60,48 @@ class JieliAITranslationBridge( |
|
|
|
/** 解码 / 编码失败的统一出口。 */ |
|
|
|
private val onError: (code: Int, msg: String?) -> Unit, |
|
|
|
private val opusPacketSize: Int = if (mode.channel == 2) 80 else 200, |
|
|
|
/** |
|
|
|
* 调试用:整句 TTS 落地目录。给定后每段交给 SDK 的整句音频会落成 .wav + .opus, |
|
|
|
* 用于人工听辨"丢句"发生在 app(我们累积/编码就缺)还是 SDK(我们给全了但耳机没播)。 |
|
|
|
* null = 不落地。建议传外部目录(context.getExternalFilesDir)方便 adb pull。 |
|
|
|
*/ |
|
|
|
private val dumpDir: File? = null, |
|
|
|
) : 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 |
|
|
|
|
|
|
|
/** |
|
|
|
* 每段末尾补的静音时长(ms)。修"每句尾音不完整": |
|
|
|
* - 防 OPUS 编码丢掉最后不足一帧的真实尾音(丢的改成丢这段静音); |
|
|
|
* - 防设备切到下一段时把当前段尾巴咬掉(咬掉的也是静音)。 |
|
|
|
* 100ms 在通话里听不出来,但能保住尾音。觉得还缺就加大。0 = 关闭。 |
|
|
|
*/ |
|
|
|
private const val TTS_TAIL_PAD_MS = 100 |
|
|
|
|
|
|
|
// ─── 下行 TTS 段合并(coalesce)───────────────────────────────────── |
|
|
|
// 背景:Doubao E2E 的 TTS 按"译文小句"切,每句一个 TTSSentenceEnd→isFinal, |
|
|
|
// 段都很碎(1~2.5s)。碎段各自独立 OPUS 发给 SDK → 设备端解码器频繁 reset、 |
|
|
|
// 段间留空隙 → 对端听着不连贯。合并:把连续到达的几段 PCM 拼成一整段再发。 |
|
|
|
/** 合并开关。true = 把"间隔 < COALESCE_IDLE_MS 的连续碎段"在 per-leg buffer 里 |
|
|
|
* 无缝拼接(PCM 层拼接,无 OPUS 边界)成一整段再下发;A 合 A、B 合 B,绝不跨腿。 |
|
|
|
* false = 老行为(每个 isFinal 立刻下发)。 |
|
|
|
* 目的:Doubao 把一句译文切成多个 1~2.5s 碎段,各自独立 OPUS 会让设备解码器频繁 reset、 |
|
|
|
* 段间留空隙 → 对端听着细碎不连贯。合并后段间无 OPUS 边界、更连贯,并顺带降低推流次数。 |
|
|
|
* 注意:启用后才激活下面 scheduleIdleFlush 的空闲兜底(关时那段是永不执行的死代码)。 |
|
|
|
* 有了兜底,isFinal(TTSSentenceEnd 段尾)延迟/丢失时最多卡 COALESCE_IDLE_MS 就下发, |
|
|
|
* 不会再一直卡到对端说下一句才补播。 */ |
|
|
|
private const val COALESCE_ENABLED = true |
|
|
|
/** 累计音频达到该时长(ms)就立刻下发:连续说话(段间隔一直 <IDLE)时的最大起播延迟上限, |
|
|
|
* 防止一直接话把 buffer 越滚越大、延迟失控。= "连贯优先"档。 */ |
|
|
|
private const val COALESCE_TARGET_MS = 3000 |
|
|
|
/** 距上一段语音空闲超过该时长(ms)即认为"这一组说完了",把已累积音频整段下发。 |
|
|
|
* = "合并窗口":下一段在该时长内到达就拼进当前组,超过就把当前组单独发。 |
|
|
|
* 也是一组说完后的最大尾延迟。2000ms = 把间隔 2 秒内到达的碎段合并成一个。 */ |
|
|
|
private const val COALESCE_IDLE_MS = 2000L |
|
|
|
|
|
|
|
} |
|
|
|
|
|
|
|
private val isPcmMode = mode.audioType == Constants.AUDIO_TYPE_PCM |
|
|
|
@ -103,7 +139,27 @@ class JieliAITranslationBridge( |
|
|
|
private val bufferLock = Any() |
|
|
|
private val encodeSeq = AtomicLong(0) |
|
|
|
|
|
|
|
/** 调试统计 */ |
|
|
|
/** 段合并的去抖定时器(per-leg):空闲超时后把累积的 PCM 整段下发。 */ |
|
|
|
private val coalesceExecutor = |
|
|
|
java.util.concurrent.Executors.newSingleThreadScheduledExecutor { r -> |
|
|
|
Thread(r, "tts-coalesce").apply { isDaemon = true } |
|
|
|
} |
|
|
|
private val idleFutures = |
|
|
|
java.util.concurrent.ConcurrentHashMap<String, java.util.concurrent.ScheduledFuture<*>>() |
|
|
|
|
|
|
|
/** 调试落地:整句 TTS 音频目录(懒创建)。null = 关闭。 */ |
|
|
|
private val resolvedDumpDir: File? = dumpDir?.also { runCatching { if (!it.exists()) it.mkdirs() } } |
|
|
|
|
|
|
|
/** |
|
|
|
* 整通会话级 per-leg PCM 落地(disk-backed,避免长通话 OOM)。 |
|
|
|
* key = 方向标签(翻译-左 / 翻译-右)。每段 flush 出来的 PCM 追加到对应文件, |
|
|
|
* stopTranslating 时统一导出 翻译-左/右 的 .wav + .opus 各一份。 |
|
|
|
*/ |
|
|
|
private val sessionPcmStreams = mutableMapOf<String, java.io.FileOutputStream>() |
|
|
|
private val sessionPcmFiles = mutableMapOf<String, File>() |
|
|
|
private val sessionDumpLock = Any() |
|
|
|
|
|
|
|
/** 首帧确认日志用:标记总首帧 / 各 source 首帧是否已打过日志。 */ |
|
|
|
@Volatile private var rxFirstLogged = false |
|
|
|
private val rxFirstPerSource = java.util.concurrent.ConcurrentHashMap<Int, Boolean>() |
|
|
|
|
|
|
|
@ -113,6 +169,13 @@ class JieliAITranslationBridge( |
|
|
|
|
|
|
|
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}") |
|
|
|
resolvedDumpDir?.let { dir -> |
|
|
|
Log.i(TAG, "[TTS-DUMP] enabled dir=${dir.absolutePath}") |
|
|
|
// 清理上次异常退出(未走 stopTranslating)残留的会话级临时 PCM |
|
|
|
runCatching { |
|
|
|
dir.listFiles { f -> f.name.endsWith(".session.pcm") }?.forEach { it.delete() } |
|
|
|
} |
|
|
|
} |
|
|
|
sdkCallback = callback |
|
|
|
if (!tempDir.exists()) tempDir.mkdirs() |
|
|
|
upDecoder?.start() |
|
|
|
@ -125,6 +188,9 @@ class JieliAITranslationBridge( |
|
|
|
override fun stopTranslating() { |
|
|
|
Log.i(TAG, "[SDK->APP] stopTranslating") |
|
|
|
working = false |
|
|
|
// 取消所有未触发的合并定时器(utterance 已被打断,剩余直接丢)。 |
|
|
|
idleFutures.values.forEach { it.cancel(false) } |
|
|
|
idleFutures.clear() |
|
|
|
runCatching { upDecoder?.stop() } |
|
|
|
runCatching { downDecoder?.stop() } |
|
|
|
runCatching { stereoDecoder?.stop() } |
|
|
|
@ -133,6 +199,8 @@ class JieliAITranslationBridge( |
|
|
|
pcmBuffers.values.forEach { it.reset() } |
|
|
|
pcmBuffers.clear() |
|
|
|
} |
|
|
|
// 会话结束:把整通累积的 per-leg PCM 导出成 翻译-左/右 的 wav + opus 各一份。 |
|
|
|
finalizeSessionDumps() |
|
|
|
sdkCallback?.onStop(0, "stopTranslating") |
|
|
|
sdkCallback = null |
|
|
|
} |
|
|
|
@ -190,30 +258,73 @@ class JieliAITranslationBridge( |
|
|
|
postTranslationResult(cb, source, Constants.AUDIO_TYPE_PCM, pcm) |
|
|
|
return true |
|
|
|
} |
|
|
|
// OPUS 模式:累积 + isFinal/超限触发整段编码 |
|
|
|
val pendingFlush: ByteArray? = synchronized(bufferLock) { |
|
|
|
|
|
|
|
// 合并目标字节(按当前采样率换算):累计达到就立刻下发。 |
|
|
|
val targetBytes = COALESCE_TARGET_MS * (sampleRateHz / 1000 * 2).coerceAtLeast(1) |
|
|
|
var pendingFlush: ByteArray? = null |
|
|
|
var armIdle = false |
|
|
|
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 |
|
|
|
|
|
|
|
when { |
|
|
|
// 硬上限兜底:无论如何先放出去,防止内存爆。 |
|
|
|
buf.size() > PCM_BUFFER_HARD_LIMIT_BYTES -> { |
|
|
|
Log.w(TAG, "feedTtsPcm leg=$outputStreamId hit hard limit, force flush") |
|
|
|
pendingFlush = buf.toByteArray(); buf.reset() |
|
|
|
} |
|
|
|
// 不合并:老行为,isFinal 立即下发。 |
|
|
|
!COALESCE_ENABLED && isFinal -> { |
|
|
|
val out = buf.toByteArray(); buf.reset() |
|
|
|
if (out.isNotEmpty()) pendingFlush = out |
|
|
|
} |
|
|
|
// 合并:累计够长了立刻发,否则等空闲去抖(期待后续段拼进来)。 |
|
|
|
COALESCE_ENABLED && isFinal -> { |
|
|
|
if (buf.size() >= targetBytes) { |
|
|
|
val out = buf.toByteArray(); buf.reset() |
|
|
|
if (out.isNotEmpty()) pendingFlush = out |
|
|
|
} else if (buf.size() > 0) { |
|
|
|
armIdle = true |
|
|
|
} |
|
|
|
} |
|
|
|
} |
|
|
|
} |
|
|
|
// 正在收音频 = 没空闲,取消挂着的空闲定时器,避免一段 burst 中途被 flush 切断。 |
|
|
|
if (pcm.isNotEmpty()) cancelIdleFlush(outputStreamId) |
|
|
|
|
|
|
|
if (pendingFlush != null) { |
|
|
|
encodeAndDeliverAsync(cb, source, outputStreamId, pendingFlush) |
|
|
|
cancelIdleFlush(outputStreamId) |
|
|
|
encodeAndDeliverAsync(cb, source, outputStreamId, pendingFlush!!) |
|
|
|
} else if (armIdle) { |
|
|
|
scheduleIdleFlush(cb, source, outputStreamId) |
|
|
|
} |
|
|
|
return true |
|
|
|
} |
|
|
|
|
|
|
|
/** 重置/启动 per-leg 空闲定时器:空闲 [COALESCE_IDLE_MS] 后把累积 PCM 整段下发。 */ |
|
|
|
private fun scheduleIdleFlush(cb: AITranslationCallback, source: Int, leg: String) { |
|
|
|
idleFutures.remove(leg)?.cancel(false) |
|
|
|
val task = Runnable { |
|
|
|
val out: ByteArray? = synchronized(bufferLock) { |
|
|
|
val buf = pcmBuffers[leg] |
|
|
|
if (buf != null && buf.size() > 0) buf.toByteArray().also { buf.reset() } else null |
|
|
|
} |
|
|
|
if (out != null && working) { |
|
|
|
encodeAndDeliverAsync(cb, source, leg, out) |
|
|
|
} |
|
|
|
} |
|
|
|
val future = runCatching { |
|
|
|
coalesceExecutor.schedule(task, COALESCE_IDLE_MS, java.util.concurrent.TimeUnit.MILLISECONDS) |
|
|
|
}.getOrNull() |
|
|
|
if (future != null) idleFutures[leg] = future |
|
|
|
} |
|
|
|
|
|
|
|
private fun cancelIdleFlush(leg: String) { |
|
|
|
idleFutures.remove(leg)?.cancel(false) |
|
|
|
} |
|
|
|
|
|
|
|
private fun postTranslationResult( |
|
|
|
cb: AITranslationCallback, |
|
|
|
source: Int, |
|
|
|
@ -237,8 +348,14 @@ class JieliAITranslationBridge( |
|
|
|
val seq = encodeSeq.incrementAndGet() |
|
|
|
val pcmFile = File(tempDir, "tts_${source}_$seq.pcm") |
|
|
|
val opusFile = File(tempDir, "tts_${source}_$seq.opus") |
|
|
|
// 调试落地:把本段 PCM 追加进"整通会话级"per-leg 累积(翻译-左/右 各一份), |
|
|
|
// 会话结束(stopTranslating)时统一导出 wav + opus,不再每段一个文件。 |
|
|
|
appendSessionPcm(leg, pcmBytes) |
|
|
|
// 末尾补静音再编码:保住尾音(防编码丢尾帧 / 防设备切段咬尾)。 |
|
|
|
val padBytes = TTS_TAIL_PAD_MS * (sampleRateHz / 1000 * 2).coerceAtLeast(1) |
|
|
|
val toEncode = if (padBytes > 0) pcmBytes + ByteArray(padBytes) else pcmBytes |
|
|
|
try { |
|
|
|
pcmFile.writeBytes(pcmBytes) |
|
|
|
pcmFile.writeBytes(toEncode) |
|
|
|
} catch (e: Throwable) { |
|
|
|
Log.w(TAG, "encodeAndDeliver leg=$leg seq=$seq write pcm failed: ${e.message}") |
|
|
|
return |
|
|
|
@ -249,19 +366,16 @@ class JieliAITranslationBridge( |
|
|
|
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 { |
|
|
|
val opusBytes = runCatching { if (opusFile.exists()) opusFile.readBytes() else null }.getOrNull() |
|
|
|
runCatching { encoder.release() } |
|
|
|
runCatching { pcmFile.delete() } |
|
|
|
runCatching { opusFile.delete() } |
|
|
|
|
|
|
|
if (opusBytes == null || opusBytes.isEmpty()) { |
|
|
|
Log.w(TAG, "encodeAndDeliver leg=$leg seq=$seq: empty opus output") |
|
|
|
return |
|
|
|
} |
|
|
|
postTranslationResult(cb, source, Constants.AUDIO_TYPE_OPUS, opusBytes) |
|
|
|
} |
|
|
|
override fun onError(code: Int, message: String?) { |
|
|
|
Log.w(TAG, "encodeAndDeliver leg=$leg seq=$seq encodeFile error code=$code msg=$message") |
|
|
|
@ -278,4 +392,113 @@ class JieliAITranslationBridge( |
|
|
|
TranslationStreams.OUT_DOWNLINK -> AudioData.SOURCE_E_SCO_DOWN_LINK |
|
|
|
else -> AudioData.SOURCE_PHONE_MIC |
|
|
|
} |
|
|
|
|
|
|
|
// ─── 调试落地 ───────────────────────────────────────────────────────── |
|
|
|
|
|
|
|
/** leg → 中文方向标签,用于文件名:上行=翻译-左,下行=翻译-右,其它=翻译-mic。 */ |
|
|
|
private fun legLabel(leg: String): String = when (leg) { |
|
|
|
TranslationStreams.OUT_UPLINK -> "翻译-左" |
|
|
|
TranslationStreams.OUT_DOWNLINK -> "翻译-右" |
|
|
|
else -> "翻译-mic" |
|
|
|
} |
|
|
|
|
|
|
|
/** 把本段 PCM 追加进对应方向的会话级临时 PCM 文件(disk-backed)。 */ |
|
|
|
private fun appendSessionPcm(leg: String, pcm: ByteArray) { |
|
|
|
val dir = resolvedDumpDir ?: return |
|
|
|
if (pcm.isEmpty()) return |
|
|
|
val label = legLabel(leg) |
|
|
|
synchronized(sessionDumpLock) { |
|
|
|
val os = sessionPcmStreams[label] ?: runCatching { |
|
|
|
val f = File(dir, "$label.session.pcm") |
|
|
|
runCatching { f.delete() } |
|
|
|
java.io.FileOutputStream(f, false).also { |
|
|
|
sessionPcmStreams[label] = it |
|
|
|
sessionPcmFiles[label] = f |
|
|
|
} |
|
|
|
}.getOrNull() ?: return |
|
|
|
runCatching { os.write(pcm) } |
|
|
|
.onFailure { Log.w(TAG, "[TTS-DUMP] append $label failed: ${it.message}") } |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
/** |
|
|
|
* 会话结束:关闭各方向 PCM 流,把整通累积的 PCM 导出成 |
|
|
|
* 翻译-左/右 的 .wav(编码前原始)+ .opus(编码后)各一份。 |
|
|
|
*/ |
|
|
|
private fun finalizeSessionDumps() { |
|
|
|
val dir = resolvedDumpDir ?: return |
|
|
|
val files: Map<String, File> |
|
|
|
synchronized(sessionDumpLock) { |
|
|
|
sessionPcmStreams.values.forEach { runCatching { it.flush(); it.close() } } |
|
|
|
sessionPcmStreams.clear() |
|
|
|
files = sessionPcmFiles.toMap() |
|
|
|
sessionPcmFiles.clear() |
|
|
|
} |
|
|
|
files.forEach { (label, pcmFile) -> |
|
|
|
if (!pcmFile.exists() || pcmFile.length() == 0L) { |
|
|
|
runCatching { pcmFile.delete() } |
|
|
|
return@forEach |
|
|
|
} |
|
|
|
// 1) WAV(同步):标准 44 字节头 + 整通 PCM |
|
|
|
runCatching { |
|
|
|
val wav = File(dir, "$label.wav") |
|
|
|
writeWavFromPcmFile(wav, pcmFile) |
|
|
|
val ms = pcmFile.length() / (sampleRateHz / 1000 * 2).coerceAtLeast(1) |
|
|
|
Log.i(TAG, "[TTS-DUMP] session wav $label pcm=${pcmFile.length()}B dur=${ms}ms → ${wav.absolutePath}") |
|
|
|
}.onFailure { Log.w(TAG, "[TTS-DUMP] session wav $label failed: ${it.message}") } |
|
|
|
// 2) OPUS(异步编码):整通 PCM 编一份;完成后删临时 PCM |
|
|
|
val opusOut = File(dir, "$label.opus") |
|
|
|
val enc = OpusManager() |
|
|
|
runCatching { |
|
|
|
enc.encodeFile(pcmFile.absolutePath, opusOut.absolutePath, object : OnStateCallback { |
|
|
|
override fun onStart() {} |
|
|
|
override fun onComplete(path: String?) { |
|
|
|
runCatching { enc.release() } |
|
|
|
runCatching { pcmFile.delete() } |
|
|
|
Log.i(TAG, "[TTS-DUMP] session opus $label → ${opusOut.absolutePath}") |
|
|
|
} |
|
|
|
override fun onError(code: Int, message: String?) { |
|
|
|
runCatching { enc.release() } |
|
|
|
runCatching { pcmFile.delete() } |
|
|
|
Log.w(TAG, "[TTS-DUMP] session opus $label encode error code=$code msg=$message") |
|
|
|
} |
|
|
|
}) |
|
|
|
}.onFailure { |
|
|
|
runCatching { enc.release() } |
|
|
|
runCatching { pcmFile.delete() } |
|
|
|
Log.w(TAG, "[TTS-DUMP] session opus $label start failed: ${it.message}") |
|
|
|
} |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
/** 流式写 WAV:先写 44 字节头(按 pcm 文件大小),再把 PCM 文件内容拷进去(不占内存)。 */ |
|
|
|
private fun writeWavFromPcmFile(wavFile: File, pcmFile: File) { |
|
|
|
val pcmLen = pcmFile.length().toInt() |
|
|
|
wavFile.outputStream().use { os -> |
|
|
|
os.write(wavHeader(pcmLen, sampleRateHz, 1)) |
|
|
|
pcmFile.inputStream().use { it.copyTo(os) } |
|
|
|
} |
|
|
|
} |
|
|
|
|
|
|
|
/** 标准 44 字节 PCM WAV 头(16bit)。 */ |
|
|
|
private fun wavHeader(pcmSize: Int, sampleRate: Int, channels: Int): ByteArray { |
|
|
|
val bitsPerSample = 16 |
|
|
|
val byteRate = sampleRate * channels * bitsPerSample / 8 |
|
|
|
val blockAlign = channels * bitsPerSample / 8 |
|
|
|
val header = java.nio.ByteBuffer.allocate(44).order(java.nio.ByteOrder.LITTLE_ENDIAN) |
|
|
|
header.put("RIFF".toByteArray(Charsets.US_ASCII)) |
|
|
|
header.putInt(36 + pcmSize) |
|
|
|
header.put("WAVE".toByteArray(Charsets.US_ASCII)) |
|
|
|
header.put("fmt ".toByteArray(Charsets.US_ASCII)) |
|
|
|
header.putInt(16) // fmt chunk size (PCM) |
|
|
|
header.putShort(1.toShort()) // audio format = PCM |
|
|
|
header.putShort(channels.toShort()) |
|
|
|
header.putInt(sampleRate) |
|
|
|
header.putInt(byteRate) |
|
|
|
header.putShort(blockAlign.toShort()) |
|
|
|
header.putShort(bitsPerSample.toShort()) |
|
|
|
header.put("data".toByteArray(Charsets.US_ASCII)) |
|
|
|
header.putInt(pcmSize) |
|
|
|
return header.array() |
|
|
|
} |
|
|
|
} |
|
|
|
|