Browse Source

Merge branch 'new_dev' of https://github.com/deepcloud2048/deep_voice into liwei

newdev_shunjiawei
liwei1dao 6 months ago
parent
commit
1687d9555d
  1. 48
      lib/modules/translation/controllers/translation_controller.dart
  2. 86
      local_plugins/azure_speech/ios/azure_speech/Sources/azure_speech/AliyunBailianE2EHelper.swift
  3. 6
      local_plugins/azure_speech/ios/azure_speech/Sources/azure_speech/AzureSpeechPlugin.swift
  4. 18
      local_plugins/azure_speech/ios/azure_speech/Sources/azure_speech/DoubaoE2ETranslateHelper.swift
  5. 2
      pubspec.yaml

48
lib/modules/translation/controllers/translation_controller.dart

@ -159,6 +159,8 @@ class TranslationController extends GetxController with WidgetsBindingObserver {
// ==================== 重新初始化守卫 ==================== // ==================== 重新初始化守卫 ====================
bool _isReinitializing = false; // 语言切换重新初始化期间为 true,抑制旧会话的错误事件 bool _isReinitializing = false; // 语言切换重新初始化期间为 true,抑制旧会话的错误事件
// 初始化/切换 provider 完成后的宽限期,用于过滤底层异步 dispose 旧 session 时产生的残留错误事件(如 1012/1011)
DateTime? _astErrorGraceUntil;
// 语言切换互斥:连续切换语言时串行执行,避免并发的 _reinitializeAsrService 导致 // 语言切换互斥:连续切换语言时串行执行,避免并发的 _reinitializeAsrService 导致
// 原生 AST 已启动但 isRecognizing 仍为 false 的状态不一致。 // 原生 AST 已启动但 isRecognizing 仍为 false 的状态不一致。
@ -940,6 +942,11 @@ class TranslationController extends GetxController with WidgetsBindingObserver {
}); });
Logger.info('[STS] 已订阅 astStream, subscription=${_astEventSubscription.hashCode}'); Logger.info('[STS] 已订阅 astStream, subscription=${_astEventSubscription.hashCode}');
// 初始化完成后给 2s 宽限期,过滤掉底层异步 dispose 旧 provider session 时产生的
// cancelled(1012) / Socket not connected(1011) 残留事件,避免被误判为新 session 错误
_astErrorGraceUntil =
DateTime.now().add(const Duration(milliseconds: 2000));
Logger.info('通话模式语音翻译服务初始化完成'); Logger.info('通话模式语音翻译服务初始化完成');
} catch (e) { } catch (e) {
Logger.error('通话模式语音翻译服务初始化失败: ${e.toString()}'); Logger.error('通话模式语音翻译服务初始化失败: ${e.toString()}');
@ -1101,7 +1108,8 @@ class TranslationController extends GetxController with WidgetsBindingObserver {
Logger.i('Translation', '开始语音识别,模式:面对面,ASR push_to_talk'); Logger.i('Translation', '开始语音识别,模式:面对面,ASR push_to_talk');
await _startAsrService('push_to_talk'); await _startAsrService('push_to_talk');
} else if (currentMode.value == 'call') { } else if (currentMode.value == 'call') {
Logger.i('Translation', '通话模式:AST端到端服务已启动,跳过Azure ASR'); Logger.i('Translation', '通话模式:启动 AST 端到端服务');
await _astService.startContinuousTranslation();
} else { } else {
Logger.i('Translation', Logger.i('Translation',
'开始语音识别,模式:${currentMode.value},ASR 模式:normal'); '开始语音识别,模式:${currentMode.value},ASR 模式:normal');
@ -1177,11 +1185,23 @@ class TranslationController extends GetxController with WidgetsBindingObserver {
/// 配置通话模式(只做 BLE 配置,AST 已在 _initializeCallModeTranslationService 中初始化) /// 配置通话模式(只做 BLE 配置,AST 已在 _initializeCallModeTranslationService 中初始化)
Future<void> _configureCallMode() async { Future<void> _configureCallMode() async {
isPreparing.value = true; isPreparing.value = true;
await bleManager.openA2DPDecoder(); // 最多重试 3 轮 A2DP 解码器打开,每轮等待 2 秒(约 13 次 150ms 轮询)
int attempts = 0; // 部分耳机首次收到 CMD_CONTROL_CODEC 时需要重放命令才能真正切到 A2DP 流式模式
while (!bleManager.isCodecActive && attempts < 10) { const int maxRounds = 3;
await Future.delayed(const Duration(milliseconds: 150)); const int attemptsPerRound = 13;
attempts++; for (int round = 0; round < maxRounds; round++) {
Logger.info('[CALL-CFG] 打开 A2DP 解码器 (round=${round + 1}/$maxRounds)');
await bleManager.openA2DPDecoder();
int attempts = 0;
while (!bleManager.isCodecActive && attempts < attemptsPerRound) {
await Future.delayed(const Duration(milliseconds: 150));
attempts++;
}
if (bleManager.isCodecActive) {
Logger.info('[CALL-CFG] A2DP 解码器就绪 (round=${round + 1}, attempts=$attempts)');
break;
}
Logger.warning('[CALL-CFG] A2DP 解码器未激活,重试 (round=${round + 1})');
} }
isPreparing.value = false; isPreparing.value = false;
if (!bleManager.isCodecActive) { if (!bleManager.isCodecActive) {
@ -1521,6 +1541,12 @@ class TranslationController extends GetxController with WidgetsBindingObserver {
Logger.info('[STS] 正在重新初始化,忽略旧会话错误'); Logger.info('[STS] 正在重新初始化,忽略旧会话错误');
break; break;
} }
// 初始化刚完成的宽限期内,底层仍在异步回调旧 session 的 1011/1012,忽略
if (_astErrorGraceUntil != null &&
DateTime.now().isBefore(_astErrorGraceUntil!)) {
Logger.info('[STS] 初始化宽限期内,忽略旧会话错误: ${event.error}');
break;
}
// 非重新初始化期间的错误:终止通话服务并还原状态 // 非重新初始化期间的错误:终止通话服务并还原状态
_terminateCallOnError('AST 错误: ${event.error}'); _terminateCallOnError('AST 错误: ${event.error}');
break; break;
@ -1531,6 +1557,11 @@ class TranslationController extends GetxController with WidgetsBindingObserver {
Logger.info('[STS] 正在重新初始化,忽略旧会话取消'); Logger.info('[STS] 正在重新初始化,忽略旧会话取消');
break; break;
} }
if (_astErrorGraceUntil != null &&
DateTime.now().isBefore(_astErrorGraceUntil!)) {
Logger.info('[STS] 初始化宽限期内,忽略旧会话取消: ${event.error}');
break;
}
_terminateCallOnError('AST 取消: ${event.error}'); _terminateCallOnError('AST 取消: ${event.error}');
break; break;
@ -2222,6 +2253,11 @@ class TranslationController extends GetxController with WidgetsBindingObserver {
} catch (e) { } catch (e) {
Logger.error('重新初始化ASR服务失败: ${e.toString()}'); Logger.error('重新初始化ASR服务失败: ${e.toString()}');
} finally { } finally {
// 初始化完成后设置一个宽限期,过滤底层异步 dispose 旧 session 时产生的
// cancelled(1012)/Socket not connected(1011) 等残留错误事件,避免它们把新会话
// 误终止掉。2 秒对 urlSession didCompleteWithError 的回调窗口足够。
_astErrorGraceUntil =
DateTime.now().add(const Duration(milliseconds: 2000));
_isReinitializing = false; _isReinitializing = false;
} }
} }

86
local_plugins/azure_speech/ios/azure_speech/Sources/azure_speech/AliyunBailianE2EHelper.swift

@ -50,6 +50,11 @@ class AliyunBailianE2EHelper: NSObject, URLSessionWebSocketDelegate {
private var audioChunkBuffer = Data() private var audioChunkBuffer = Data()
private var isStarted = false private var isStarted = false
private var openContinuation: CheckedContinuation<Bool, Never>?
// isStarted 成立前的音频缓冲,didOpen 后再冲刷,避免首包被丢
private var pendingAudioChunks: [Data] = []
private let pendingLock = NSLock()
private let maxPendingChunks = 50
deinit { deinit {
os_log("AliyunBailianE2EHelper deinit", log: log, type: .info) os_log("AliyunBailianE2EHelper deinit", log: log, type: .info)
@ -67,6 +72,15 @@ class AliyunBailianE2EHelper: NSObject, URLSessionWebSocketDelegate {
urlSession = nil urlSession = nil
webSocket = nil webSocket = nil
isStarted = false isStarted = false
// 抛弃上一次未完成的等待,避免泄漏
if let cont = openContinuation {
openContinuation = nil
cont.resume(returning: false)
}
// 清空上一会话遗留的待发音频
pendingLock.lock()
pendingAudioChunks.removeAll()
pendingLock.unlock()
conf = config conf = config
callback = cb callback = cb
@ -77,7 +91,27 @@ class AliyunBailianE2EHelper: NSObject, URLSessionWebSocketDelegate {
urlSession = URLSession(configuration: configuration, delegate: self, delegateQueue: OperationQueue.main) urlSession = URLSession(configuration: configuration, delegate: self, delegateQueue: OperationQueue.main)
return startContinuousConversation() // 真正等到 WebSocket didOpen(isStarted=true)再返回,防止后续 pushAudioData 被丢
let kicked = startContinuousConversation()
if !kicked { return false }
let opened = await withCheckedContinuation { (cont: CheckedContinuation<Bool, Never>) in
// 如已 open(极少见),立即放行
if isStarted {
cont.resume(returning: true)
return
}
openContinuation = cont
// 超时兜底:10s 未 open 就放行避免 Flutter 侧死等
DispatchQueue.main.asyncAfter(deadline: .now() + 10.0) { [weak self] in
guard let self = self else { return }
if let pending = self.openContinuation {
self.openContinuation = nil
os_log("initialize: WebSocket didOpen 超时,返回 false", log: self.log, type: .error)
pending.resume(returning: false)
}
}
}
return opened
} }
/** /**
@ -317,8 +351,17 @@ class AliyunBailianE2EHelper: NSObject, URLSessionWebSocketDelegate {
* 推送音频 PCM 数据(Base64 编码) * 推送音频 PCM 数据(Base64 编码)
*/ */
func pushAudioData(_ data: Data) -> Bool { func pushAudioData(_ data: Data) -> Bool {
guard isStarted, let ws = webSocket else { // isStarted 还没就绪时,缓冲音频(最多 maxPendingChunks 包),didOpen 后冲刷
os_log("pushAudioData ignored: isStarted=%{public}d, ws is nil", log: log, type: .info, isStarted) if !isStarted || webSocket == nil {
pendingLock.lock()
if pendingAudioChunks.count >= maxPendingChunks {
pendingAudioChunks.removeFirst()
}
pendingAudioChunks.append(data)
pendingLock.unlock()
return true
}
guard let ws = webSocket else {
return false return false
} }
os_log("pushAudioData: %{public}d bytes", log: log, type: .info, data.count) os_log("pushAudioData: %{public}d bytes", log: log, type: .info, data.count)
@ -350,7 +393,11 @@ class AliyunBailianE2EHelper: NSObject, URLSessionWebSocketDelegate {
if let ws = webSocket { if let ws = webSocket {
ws.cancel(with: .normalClosure, reason: "User stopped".data(using: .utf8)) ws.cancel(with: .normalClosure, reason: "User stopped".data(using: .utf8))
} }
webSocket = nil
isStarted = false isStarted = false
pendingLock.lock()
pendingAudioChunks.removeAll()
pendingLock.unlock()
return true return true
} }
@ -367,6 +414,13 @@ class AliyunBailianE2EHelper: NSObject, URLSessionWebSocketDelegate {
isStarted = false isStarted = false
urlSession?.invalidateAndCancel() urlSession?.invalidateAndCancel()
urlSession = nil urlSession = nil
if let cont = openContinuation {
openContinuation = nil
cont.resume(returning: false)
}
pendingLock.lock()
pendingAudioChunks.removeAll()
pendingLock.unlock()
} }
/** /**
@ -380,6 +434,22 @@ class AliyunBailianE2EHelper: NSObject, URLSessionWebSocketDelegate {
isStarted = true // 在 session.update 发送后才允许推送音频,与 Android 行为对齐 isStarted = true // 在 session.update 发送后才允许推送音频,与 Android 行为对齐
callback?.onSessionStarted(sessionId: sessionId) callback?.onSessionStarted(sessionId: sessionId)
receiveLoop() receiveLoop()
// 冲刷缓冲的音频
pendingLock.lock()
let pending = pendingAudioChunks
pendingAudioChunks.removeAll()
pendingLock.unlock()
if !pending.isEmpty {
os_log("Flushing %{public}d buffered audio chunks", log: log, type: .info, pending.count)
for chunk in pending {
_ = pushAudioData(chunk)
}
}
// 唤醒 initialize 中等待 didOpen 的协程
if let cont = openContinuation {
openContinuation = nil
cont.resume(returning: true)
}
} }
/** /**
@ -419,6 +489,11 @@ class AliyunBailianE2EHelper: NSObject, URLSessionWebSocketDelegate {
let reasonStr = String(data: reason ?? Data(), encoding: .utf8) ?? "" let reasonStr = String(data: reason ?? Data(), encoding: .utf8) ?? ""
os_log("WebSocket didClose code=%{public}d reason=%{public}@", log: log, type: .info, closeCode.rawValue, reasonStr) os_log("WebSocket didClose code=%{public}d reason=%{public}@", log: log, type: .info, closeCode.rawValue, reasonStr)
isStarted = false isStarted = false
// 如果 initialize 还在等 didOpen 就收到 close,直接返回 false
if let cont = openContinuation {
openContinuation = nil
cont.resume(returning: false)
}
callback?.onSessionFinished(sessionId: sessionId, finalText: fullTextBuffer, finalAudio: fullAudioBuffer) callback?.onSessionFinished(sessionId: sessionId, finalText: fullTextBuffer, finalAudio: fullAudioBuffer)
} }
@ -433,6 +508,11 @@ class AliyunBailianE2EHelper: NSObject, URLSessionWebSocketDelegate {
callback?.onSessionError(sessionId: sessionId, code: 1012, message: e.localizedDescription) callback?.onSessionError(sessionId: sessionId, code: 1012, message: e.localizedDescription)
} }
isStarted = false isStarted = false
// 如果 initialize 还在等 didOpen 就失败,立即返回 false
if let cont = openContinuation {
openContinuation = nil
cont.resume(returning: false)
}
} }
} }

6
local_plugins/azure_speech/ios/azure_speech/Sources/azure_speech/AzureSpeechPlugin.swift

@ -970,6 +970,7 @@ private func sendAudioDataEvent(_ event: [String: Any]) {
os_log("[AST-INIT] Azure: initializing B service...", log: self.ctLog, type: .info) os_log("[AST-INIT] Azure: initializing B service...", log: self.ctLog, type: .info)
await azureAstHelperB.initialize(azureConfig: azureConfigI, translationConfig: translationConfigI, serviceConfig: svcB, callback: callbackB) await azureAstHelperB.initialize(azureConfig: azureConfigI, translationConfig: translationConfigI, serviceConfig: svcB, callback: callbackB)
os_log("[AST-INIT] Azure: A+B both initialized", log: self.ctLog, type: .info) os_log("[AST-INIT] Azure: A+B both initialized", log: self.ctLog, type: .info)
DispatchQueue.main.async { result(true) }
} }
} else if currentAstProvider == "volcano" { } else if currentAstProvider == "volcano" {
// 初始化豆包 AST(双路) // 初始化豆包 AST(双路)
@ -996,6 +997,7 @@ private func sendAudioDataEvent(_ event: [String: Any]) {
os_log("[AST-INIT] Volcano: initializing B service...", log: self.ctLog, type: .info) os_log("[AST-INIT] Volcano: initializing B service...", log: self.ctLog, type: .info)
let rB = await doubaoAstHelperB.initialize(config: cfgB, cb: cbB) let rB = await doubaoAstHelperB.initialize(config: cfgB, cb: cbB)
os_log("[AST-INIT] Volcano: B service result=%d. Both services ready.", log: self.ctLog, type: .info, rB ? 1 : 0) os_log("[AST-INIT] Volcano: B service result=%d. Both services ready.", log: self.ctLog, type: .info, rB ? 1 : 0)
DispatchQueue.main.async { result(true) }
} }
} else if currentAstProvider == "alibaba" { } else if currentAstProvider == "alibaba" {
// 初始化阿里百炼(双路) // 初始化阿里百炼(双路)
@ -1030,9 +1032,11 @@ private func sendAudioDataEvent(_ event: [String: Any]) {
os_log("[AST-INIT] Alibaba: A service result=%d, initializing B service...", log: self.ctLog, type: .info, rA ? 1 : 0) os_log("[AST-INIT] Alibaba: A service result=%d, initializing B service...", log: self.ctLog, type: .info, rA ? 1 : 0)
let rB = await bailianAstHelperB.initialize(config: cfgB, cb: cbB) let rB = await bailianAstHelperB.initialize(config: cfgB, cb: cbB)
os_log("[AST-INIT] Alibaba: B service result=%d. Both services ready.", log: self.ctLog, type: .info, rB ? 1 : 0) os_log("[AST-INIT] Alibaba: B service result=%d. Both services ready.", log: self.ctLog, type: .info, rB ? 1 : 0)
DispatchQueue.main.async { result(true) }
} }
} else {
result(true)
} }
result(true)
default: default:
result(FlutterMethodNotImplemented) result(FlutterMethodNotImplemented)

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

@ -215,14 +215,22 @@ class DoubaoE2ETranslateHelper: NSObject, URLSessionWebSocketDelegate {
} }
/** /**
* 停止会话并发送结束请求(Protobuf) * 停止会话并发送结束请求(Protobuf),并取消 WebSocket,以便下次 start 时重新建连
*/ */
func stopContinuousTranslation() -> Bool { func stopContinuousTranslation() -> Bool {
os_log("Stopping continuous translation", log: log, type: .info) os_log("Stopping continuous translation", log: log, type: .info)
guard let ws = webSocket else { return false } stopKeepAliveTimer()
let req = makeFinishRequest(sessionId: sessionId) if let ws = webSocket,
guard let msg = try? req.serializedData() else { return false } let msg = try? makeFinishRequest(sessionId: sessionId).serializedData() {
ws.send(.data(msg)) { _ in } ws.send(.data(msg)) { _ in }
ws.cancel(with: .normalClosure, reason: "stop".data(using: .utf8))
}
webSocket = nil
isStarted = false
lock.lock()
isWebSocketConnected = false
pendingAudioData.removeAll()
lock.unlock()
return true return true
} }

2
pubspec.yaml

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

Loading…
Cancel
Save