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 690e45dab..cb577021f 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 @@ -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) diff --git a/local_plugins/azure_speech/ios/azure_speech/Sources/azure_speech/DoubaoE2ETranslateHelper.swift b/local_plugins/azure_speech/ios/azure_speech/Sources/azure_speech/DoubaoE2ETranslateHelper.swift index 82896fa45..64f944d65 100644 --- a/local_plugins/azure_speech/ios/azure_speech/Sources/azure_speech/DoubaoE2ETranslateHelper.swift +++ b/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,16 +103,15 @@ 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 isStarted = false lock.lock() isWebSocketConnected = false lock.unlock() - + DispatchQueue.main.asyncAfter(deadline: .now() + 0.2) { [weak self] in guard let self = self else { return } if !self.isStarted { @@ -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) - } - // 2. 否则(连接中、ws 为空、或者刚断开)-> 缓冲 - else { - // os_log("Buffering audio data, size: %d", log: log, type: .debug, data.count) - pendingAudioData.append(data) - lock.unlock() - return true + if !data.isEmpty { + audioSendBufferLock.lock() + audioSendBuffer.append(data) + audioSendBufferLock.unlock() } + 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) - } - /** * 接收消息循环 */ diff --git a/pubspec.yaml b/pubspec.yaml index 59dc30c73..50f3980b4 100644 --- a/pubspec.yaml +++ b/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"