Browse Source

上传通话翻译静默优化翻译

newdev_shunjiawei
liwei1dao 6 months ago
parent
commit
35c25a330f
  1. 81
      local_plugins/azure_speech/android/src/main/kotlin/com/yunqiinnovation/azure_speech/DoubaoE2ETranslateHelper.kt
  2. 160
      local_plugins/azure_speech/ios/azure_speech/Sources/azure_speech/DoubaoE2ETranslateHelper.swift
  3. 2
      pubspec.yaml

81
local_plugins/azure_speech/android/src/main/kotlin/com/yunqiinnovation/azure_speech/DoubaoE2ETranslateHelper.kt

@ -71,6 +71,17 @@ class DoubaoE2ETranslateHelper(
private val isStarted = AtomicBoolean(false)
private val scope = CoroutineScope(Dispatchers.IO + SupervisorJob())
// 上行音频 pump:
// 会话建立后启动一个固定 20ms 节拍的协程,持续从 audioSendBuffer 取 640B 一帧上行;
// 取不到时自动补静音帧,保证服务端看到的是稳定实时流,不会触发 AudioSendSlow。
// 16kHz / 16bit / mono,20ms = 16000 * 0.02 * 2 = 640 字节
private val pumpFrameBytes = 640
private val pumpFrameIntervalMs = 20L
private val silentFrame = ByteArray(pumpFrameBytes)
private val audioSendBufferLock = Any()
private var audioSendBuffer = ByteArray(0)
private var pumpJob: Job? = null
/**
* 初始化助手,设置配置与回调。
*/
@ -99,6 +110,7 @@ class DoubaoE2ETranslateHelper(
*/
private fun restartSession() {
Log.i(TAG, "restartSession: performing auto-restart...")
stopAudioPump()
try { webSocket?.close(1000, "restarting") } catch (_: Exception) {}
webSocket = null
isStarted.set(false)
@ -161,6 +173,7 @@ class DoubaoE2ETranslateHelper(
ws.send(ByteString.of(*startReq.toByteArray()))
Log.d(TAG, "onOpen: StartSession sent")
isStarted.set(true)
startAudioPump()
callback?.onSessionStarted(sessionId)
}
@ -289,6 +302,7 @@ class DoubaoE2ETranslateHelper(
*/
override fun onClosed(ws: WebSocket, code: Int, reason: String) {
Log.d(TAG, "onClosed: code=${code} reason=${reason}")
stopAudioPump()
try { recvAudio.close() } catch (_: Exception) {}
isStarted.set(false)
}
@ -336,25 +350,81 @@ class DoubaoE2ETranslateHelper(
*/
private var pushCount = 0L
/**
* 只负责把外部音频塞进上行缓冲区,真正的发送由 pump 按 20ms 节拍统一调度。
*/
fun pushAudioData(data: ByteArray): Boolean {
pushCount++
if (pushCount % 100 == 1L) {
Log.d(TAG, "pushAudioData: size=${data.size} isStarted=${isStarted.get()} wsIsNull=${webSocket==null} count=$pushCount")
}
val ws = webSocket ?: return false
if (!isStarted.get()) return false
synchronized(sessionLock) {
val req = makeChunkRequest(sessionId, data)
val ok = ws.send(ByteString.of(*req.toByteArray()))
return ok
if (data.isNotEmpty()) {
synchronized(audioSendBufferLock) {
audioSendBuffer += data
}
}
return true
}
/**
* 从上行缓冲区取一帧 20ms 数据;不够则返回静音帧。
*/
private fun nextUplinkFrame(): ByteArray {
synchronized(audioSendBufferLock) {
if (audioSendBuffer.size < pumpFrameBytes) return silentFrame
val frame = audioSendBuffer.copyOfRange(0, pumpFrameBytes)
audioSendBuffer = audioSendBuffer.copyOfRange(pumpFrameBytes, audioSendBuffer.size)
return frame
}
}
/**
* 启动固定 20ms 节拍的上行 pump,使用绝对时间避免漂移。
*/
private fun startAudioPump() {
stopAudioPump()
pumpJob = scope.launch {
var nextTick = System.currentTimeMillis()
while (isActive && isStarted.get()) {
val ws = webSocket
if (ws == null) {
delay(pumpFrameIntervalMs)
continue
}
val frame = nextUplinkFrame()
try {
synchronized(sessionLock) {
val req = makeChunkRequest(sessionId, frame)
ws.send(ByteString.of(*req.toByteArray()))
}
} catch (e: Exception) {
Log.e(TAG, "audioPump: send failed", e)
}
nextTick += pumpFrameIntervalMs
val wait = nextTick - System.currentTimeMillis()
if (wait > 0) {
delay(wait)
} else {
// 落后超过一帧时,重对齐节拍,避免补发积压
nextTick = System.currentTimeMillis()
}
}
}
}
private fun stopAudioPump() {
pumpJob?.cancel()
pumpJob = null
synchronized(audioSendBufferLock) { audioSendBuffer = ByteArray(0) }
}
/**
* 停止会话,发送 FinishSession,并等待服务端返回。
*/
fun stopContinuousTranslation(): Boolean {
Log.d(TAG, "stopContinuousTranslation: wsIsNull=${webSocket==null}")
stopAudioPump()
val ws = webSocket ?: return false
synchronized(sessionLock) {
val finishReq = makeFinishRequest(sessionId)
@ -369,6 +439,7 @@ class DoubaoE2ETranslateHelper(
*/
fun dispose() {
Log.d(TAG, "dispose: begin")
stopAudioPump()
try { webSocket?.close(1000, "dispose") } catch (_: Exception) {}
webSocket = null
isStarted.set(false)

160
local_plugins/azure_speech/ios/azure_speech/Sources/azure_speech/DoubaoE2ETranslateHelper.swift

@ -56,13 +56,18 @@ class DoubaoE2ETranslateHelper: NSObject, URLSessionWebSocketDelegate {
private let lock = NSLock()
private var isWebSocketConnected = false
private var pendingAudioData: [Data] = []
private var keepAliveTimer: Timer?
private var lastAudioSentAt = Date.distantPast
private let keepAliveIntervalSec: TimeInterval = 1.0
private let keepAliveIdleThresholdSec: TimeInterval = 2.5
private let keepAliveSilentChunk = Data(repeating: 0, count: 640)
private var hasSentInitialWarmupPacket = false
// 上行音频 pump:
// 会话建立后启动一个 20ms 固定节拍的定时器,持续从 audioSendBuffer 取 640B 一帧上行;
// 取不到时自动补静音帧,保证服务端看到的是稳定实时流,不会触发 AudioSendSlow。
// 16kHz / 16bit / mono,20ms = 640 字节
private let pumpFrameBytes = 640
private let pumpFrameIntervalMs: Int = 20
private let silentFrame = Data(repeating: 0, count: 640)
private let audioSendBufferLock = NSLock()
private var audioSendBuffer = Data()
private var pumpTimer: DispatchSourceTimer?
private let pumpQueue = DispatchQueue(label: "com.azure.speech.DoubaoE2E.pump", qos: .userInitiated)
/**
* 初始化助手(异步)
@ -98,8 +103,7 @@ class DoubaoE2ETranslateHelper: NSObject, URLSessionWebSocketDelegate {
*/
private func restartSession() {
os_log("Restarting session due to timeout...", log: log, type: .info)
stopKeepAliveTimer()
hasSentInitialWarmupPacket = false
stopAudioPump()
webSocket?.cancel(with: .normalClosure, reason: "restarting".data(using: .utf8))
webSocket = nil
@ -135,13 +139,15 @@ class DoubaoE2ETranslateHelper: NSObject, URLSessionWebSocketDelegate {
recvAudio.removeAll()
recvText.removeAll()
recvSourceText.removeAll()
hasSentInitialWarmupPacket = false
lock.lock()
isWebSocketConnected = false
pendingAudioData.removeAll()
lock.unlock()
audioSendBufferLock.lock()
audioSendBuffer.removeAll(keepingCapacity: true)
audioSendBufferLock.unlock()
var request = URLRequest(url: url)
request.timeoutInterval = 60 // 握手超时可以设置,但连接成功后不应受此限制
@ -170,31 +176,33 @@ class DoubaoE2ETranslateHelper: NSObject, URLSessionWebSocketDelegate {
}
/**
* 推送音频数据(PCM/WAV)
* 只负责把外部音频塞进上行缓冲区,真正的发送由 pump 按 20ms 节拍统一调度。
*/
func pushAudioData(_ data: Data) -> Bool {
// 先检查 isStarted,避免未启动就调用
guard isStarted else {
os_log("Push audio failed: not started", log: log, type: .error)
return false
}
// 尝试获取 webSocket,如果为空,说明还在连接中,应该进入缓冲逻辑
// 如果已经连接成功但 webSocket 还是 nil,那就是异常状态
lock.lock()
// 1. 如果已连接且有 ws 实例 -> 直接发
if isWebSocketConnected, let ws = webSocket {
lock.unlock()
return sendAudioData(ws, data: data)
if !data.isEmpty {
audioSendBufferLock.lock()
audioSendBuffer.append(data)
audioSendBufferLock.unlock()
}
// 2. 否则(连接中、ws 为空、或者刚断开)-> 缓冲
else {
// os_log("Buffering audio data, size: %d", log: log, type: .debug, data.count)
pendingAudioData.append(data)
lock.unlock()
return true
return true
}
/**
* 从上行缓冲区取一帧 20ms 数据;不够则返回静音帧。
*/
private func nextUplinkFrame() -> Data {
audioSendBufferLock.lock()
defer { audioSendBufferLock.unlock() }
if audioSendBuffer.count < pumpFrameBytes {
return silentFrame
}
let frame = audioSendBuffer.prefix(pumpFrameBytes)
audioSendBuffer.removeFirst(pumpFrameBytes)
return Data(frame)
}
/**
@ -207,19 +215,48 @@ class DoubaoE2ETranslateHelper: NSObject, URLSessionWebSocketDelegate {
ws.send(.data(msg)) { [weak self] error in
if let e = error {
os_log("发送音频失败: %{public}@", log: self?.log ?? .default, type: .error, e.localizedDescription)
} else {
self?.markAudioSentNow()
}
}
return true
}
/**
* 启动固定 20ms 节拍的上行 pump
*/
private func startAudioPump() {
stopAudioPump()
let timer = DispatchSource.makeTimerSource(queue: pumpQueue)
timer.schedule(deadline: .now() + .milliseconds(pumpFrameIntervalMs),
repeating: .milliseconds(pumpFrameIntervalMs),
leeway: .milliseconds(2))
timer.setEventHandler { [weak self] in
guard let self = self else { return }
guard self.isStarted, let ws = self.webSocket else { return }
self.lock.lock()
let connected = self.isWebSocketConnected
self.lock.unlock()
guard connected else { return }
let frame = self.nextUplinkFrame()
_ = self.sendAudioData(ws, data: frame)
}
pumpTimer = timer
timer.resume()
}
private func stopAudioPump() {
pumpTimer?.cancel()
pumpTimer = nil
audioSendBufferLock.lock()
audioSendBuffer.removeAll(keepingCapacity: true)
audioSendBufferLock.unlock()
}
/**
* 停止会话并发送结束请求(Protobuf),并取消 WebSocket,以便下次 start 时重新建连
*/
func stopContinuousTranslation() -> Bool {
os_log("Stopping continuous translation", log: log, type: .info)
stopKeepAliveTimer()
stopAudioPump()
if let ws = webSocket,
let msg = try? makeFinishRequest(sessionId: sessionId).serializedData() {
ws.send(.data(msg)) { _ in }
@ -229,7 +266,6 @@ class DoubaoE2ETranslateHelper: NSObject, URLSessionWebSocketDelegate {
isStarted = false
lock.lock()
isWebSocketConnected = false
pendingAudioData.removeAll()
lock.unlock()
return true
}
@ -239,13 +275,12 @@ class DoubaoE2ETranslateHelper: NSObject, URLSessionWebSocketDelegate {
*/
func dispose() {
os_log("Disposing resources", log: log, type: .info)
stopKeepAliveTimer()
stopAudioPump()
do { try webSocket?.cancel(with: .normalClosure, reason: "dispose".data(using: .utf8)) } catch { }
webSocket = nil
isStarted = false
lock.lock()
isWebSocketConnected = false
pendingAudioData.removeAll()
lock.unlock()
urlSession?.invalidateAndCancel()
urlSession = nil
@ -256,8 +291,6 @@ class DoubaoE2ETranslateHelper: NSObject, URLSessionWebSocketDelegate {
*/
func urlSession(_ session: URLSession, webSocketTask: URLSessionWebSocketTask, didOpenWithProtocol protocol: String?) {
os_log("WebSocket didOpen", log: log, type: .info)
markAudioSentNow()
startKeepAliveTimer()
let startReq = makeStartRequest(sessionId: sessionId)
if let data = try? startReq.serializedData() {
os_log("Sending StartSession (protobuf), size=%{public}d", log: log, type: .info, data.count)
@ -267,24 +300,15 @@ class DoubaoE2ETranslateHelper: NSObject, URLSessionWebSocketDelegate {
os_log("StartSession 发送失败: %{public}@", log: self?.log ?? .default, type: .error, e.localizedDescription)
} else {
os_log("StartSession sent successfully", log: self?.log ?? .default, type: .info)
self?.sendInitialWarmupPacketIfNeeded(ws: webSocketTask)
}
}
}
lock.lock()
isWebSocketConnected = true
let pending = pendingAudioData
pendingAudioData.removeAll()
lock.unlock()
if !pending.isEmpty {
os_log("Flushing %d buffered audio chunks", log: log, type: .info, pending.count)
}
for data in pending {
_ = sendAudioData(webSocketTask, data: data)
}
startAudioPump()
callback?.onSessionStarted(sessionId: sessionId)
receiveLoop()
@ -296,7 +320,7 @@ class DoubaoE2ETranslateHelper: NSObject, URLSessionWebSocketDelegate {
func urlSession(_ session: URLSession, webSocketTask: URLSessionWebSocketTask, didCloseWith closeCode: URLSessionWebSocketTask.CloseCode, reason: Data?) {
let reasonStr = String(data: reason ?? Data(), encoding: .utf8) ?? ""
os_log("WebSocket didClose, code: %d, reason: %{public}@", log: log, type: .info, closeCode.rawValue, reasonStr)
stopKeepAliveTimer()
stopAudioPump()
lock.lock()
isWebSocketConnected = false
lock.unlock()
@ -310,7 +334,7 @@ class DoubaoE2ETranslateHelper: NSObject, URLSessionWebSocketDelegate {
if let e = error {
os_log("WebSocket task failed: %{public}@", log: log, type: .error, e.localizedDescription)
callback?.onSessionError(sessionId: sessionId, code: 1012, message: e.localizedDescription)
stopKeepAliveTimer()
stopAudioPump()
lock.lock()
isWebSocketConnected = false
@ -319,46 +343,6 @@ class DoubaoE2ETranslateHelper: NSObject, URLSessionWebSocketDelegate {
}
}
private func startKeepAliveTimer() {
stopKeepAliveTimer()
keepAliveTimer = Timer.scheduledTimer(withTimeInterval: keepAliveIntervalSec, repeats: true) { [weak self] _ in
self?.sendKeepAliveIfNeeded()
}
if let timer = keepAliveTimer {
RunLoop.main.add(timer, forMode: .common)
}
}
private func stopKeepAliveTimer() {
keepAliveTimer?.invalidate()
keepAliveTimer = nil
}
private func markAudioSentNow() {
lastAudioSentAt = Date()
}
private func sendKeepAliveIfNeeded() {
lock.lock()
let canSend = isStarted && isWebSocketConnected
let ws = webSocket
lock.unlock()
guard canSend, let ws = ws else { return }
let idle = Date().timeIntervalSince(lastAudioSentAt)
if idle < keepAliveIdleThresholdSec { return }
os_log("Sending keepalive silent chunk, idle=%.2fs", log: log, type: .debug, idle)
_ = sendAudioData(ws, data: keepAliveSilentChunk)
}
private func sendInitialWarmupPacketIfNeeded(ws: URLSessionWebSocketTask) {
guard !hasSentInitialWarmupPacket else { return }
hasSentInitialWarmupPacket = true
os_log("Sending initial warmup packet to avoid first-packet timeout", log: log, type: .info)
_ = sendAudioData(ws, data: keepAliveSilentChunk)
}
/**
* 接收消息循环
*/

2
pubspec.yaml

@ -1,7 +1,7 @@
name: voitrans
description: "Voitrans - AI Voice Assistant."
publish_to: "none"
version: 1.0.28+103
version: 1.0.29+103
environment:
sdk: ">=3.3.0 <4.0.0"

Loading…
Cancel
Save