Browse Source

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

newdev_shunjiawei
liwei1dao 9 months ago
parent
commit
af224bf73c
  1. 1
      lib/data/services/asr_service.dart
  2. 2
      lib/data/services/speech_impl/azure_asr_service.dart
  3. 14
      lib/data/services/speech_impl/azure_tts_service.dart
  4. 12
      lib/data/services/speech_impl/voice_clone_tts_service.dart
  5. 1
      lib/data/services/speech_impl/volcano_asr_api_service.dart
  6. 1
      lib/data/services/speech_impl/volcano_asr_service.dart
  7. 17
      lib/data/services/speech_impl/volcano_tts_service.dart
  8. 3
      lib/data/services/tts_service.dart
  9. 4
      lib/modules/agent/controllers/agent_controller.dart
  10. 34
      lib/modules/opus_test/controllers/opus_test_controller.dart
  11. 40
      lib/modules/translation/controllers/translation_controller.dart
  12. 2
      local_plugins/agent_service/android/src/main/kotlin/com/yunqiinnovation/agent_service/AgentService.kt
  13. 164
      local_plugins/agent_service/ios/agent_service/Sources/agent_service/AgentServiceImpl.swift
  14. 5
      local_plugins/agent_service/ios/agent_service/Sources/agent_service/AgentServicePlugin.swift
  15. 5
      local_plugins/agent_service/lib/agent_service.dart
  16. 21
      local_plugins/azure_speech/ios/azure_speech/Sources/azure_speech/AzureAsrHelper.swift
  17. 17
      local_plugins/azure_speech/ios/azure_speech/Sources/azure_speech/AzureSpeechPlugin.swift
  18. 1154
      local_plugins/azure_speech/ios/azure_speech/Sources/azure_speech/AzureTtsHelper.swift
  19. 210
      local_plugins/azure_speech/ios/azure_speech/Sources/tools/MicrophoneCapture.swift
  20. 261
      local_plugins/azure_speech/ios/azure_speech/Sources/tools/SimpleAudioReceiver.swift
  21. 2
      local_plugins/ble_service/android/src/main/kotlin/com/yunqiinnovation/ble_service/BleCommandSender.kt
  22. 4
      local_plugins/ble_service/android/src/main/kotlin/com/yunqiinnovation/ble_service/BleService.kt
  23. 80
      local_plugins/ble_service/ios/ble_service/Sources/ble_service/BleService.swift
  24. 2
      local_plugins/ble_service/ios/ble_service/Sources/ble_service/RecordingFile.swift
  25. 163
      local_plugins/ble_service/ios/ble_service/Sources/ble_service/SwiftBleServicePlugin.swift
  26. 91
      local_plugins/chat_api/ios/chat_api/Sources/chat_api/MCPClient.swift
  27. 549
      local_plugins/chat_storage/ios/chat_storage/Sources/chat_storage/ChatStorageHelper.swift
  28. 22
      local_plugins/jl_opus/ios/jl_opus/Package.swift
  29. 450
      local_plugins/jl_opus/ios/jl_opus/Sources/jl_opus/JlOpusPlugin.swift
  30. 4
      local_plugins/music_service/ios/music_service/Sources/music_service/MusicService.swift

1
lib/data/services/asr_service.dart

@ -22,6 +22,7 @@ abstract class AsrService {
Future<bool> startContinuousRecognition( Future<bool> startContinuousRecognition(
bool audioSourceType, { bool audioSourceType, {
bool isRemoveFirstPunctuation = true, bool isRemoveFirstPunctuation = true,
String mode = "normal",
}); });
/// 停止连续语音识别 /// 停止连续语音识别

2
lib/data/services/speech_impl/azure_asr_service.dart

@ -207,6 +207,7 @@ class AzureAsrService extends GetxService implements AsrService {
Future<bool> startContinuousRecognition( Future<bool> startContinuousRecognition(
bool audioSourceType, { bool audioSourceType, {
bool isRemoveFirstPunctuation = true, bool isRemoveFirstPunctuation = true,
String mode = "normal",
}) async { }) async {
// 防抖判断:短时间内重复调用直接拦截 // 防抖判断:短时间内重复调用直接拦截
final now = DateTime.now(); final now = DateTime.now();
@ -229,6 +230,7 @@ class AzureAsrService extends GetxService implements AsrService {
await _channel.invokeMethod('startContinuousRecognition', { await _channel.invokeMethod('startContinuousRecognition', {
'audioSourceType': audioSourceType, 'audioSourceType': audioSourceType,
'isRemoveFirstPunctuation': isRemoveFirstPunctuation, 'isRemoveFirstPunctuation': isRemoveFirstPunctuation,
'mode': mode,
}); });
if (!result) { if (!result) {

14
lib/data/services/speech_impl/azure_tts_service.dart

@ -182,6 +182,20 @@ class AzureTtsService extends GetxService implements TtsService {
} }
} }
@override
Future<bool> setTtsMode(String mod) async {
if (!_isInitialized) await initialize();
try {
final result = await _channel.invokeMethod('setTtsMode', {
'mod': mod,
});
return result;
} catch (e) {
Logger.error('设置TTS模式失败: ${e.toString()}');
return false;
}
}
@override @override
Future<bool> setVoiceFlocking(String speakerProfileId) async { Future<bool> setVoiceFlocking(String speakerProfileId) async {
final result = await _channel.invokeMethod('setVoiceFlocking', { final result = await _channel.invokeMethod('setVoiceFlocking', {

12
lib/data/services/speech_impl/voice_clone_tts_service.dart

@ -94,6 +94,18 @@ class VoiceCloneTtsService extends GetxService implements TtsService {
} }
} }
@override
Future<bool> setTtsMode(String mod) async {
if (!_isInitialized) await initialize();
try {
Logger.info('设置TTS模式: $mod');
return true;
} catch (e) {
Logger.error('设置TTS模式失败: ${e.toString()}');
return false;
}
}
@override @override
Future<bool> setVoiceFlocking(String speakerProfileId) async { Future<bool> setVoiceFlocking(String speakerProfileId) async {
return true; return true;

1
lib/data/services/speech_impl/volcano_asr_api_service.dart

@ -865,6 +865,7 @@ class VolcanoAsrApiService implements AsrService {
Future<bool> startContinuousRecognition( Future<bool> startContinuousRecognition(
bool audioSourceType, { bool audioSourceType, {
bool isRemoveFirstPunctuation = true, bool isRemoveFirstPunctuation = true,
String mode = "normal",
}) { }) {
// TODO: implement startContinuousRecognition // TODO: implement startContinuousRecognition
throw UnimplementedError(); throw UnimplementedError();

1
lib/data/services/speech_impl/volcano_asr_service.dart

@ -400,6 +400,7 @@ class VolcanoAsrService extends GetxService implements AsrService {
Future<bool> startContinuousRecognition( Future<bool> startContinuousRecognition(
bool audioSourceType, { bool audioSourceType, {
bool isRemoveFirstPunctuation = true, bool isRemoveFirstPunctuation = true,
String mode = "normal",
}) { }) {
// TODO: implement startContinuousRecognition // TODO: implement startContinuousRecognition
throw UnimplementedError(); throw UnimplementedError();

17
lib/data/services/speech_impl/volcano_tts_service.dart

@ -185,6 +185,23 @@ class VolcanoTtsService extends GetxService implements TtsService {
} }
} }
@override
Future<bool> setTtsMode(String mod) async {
if (!await _ensureInitialized()) return false;
try {
final result = await _channel.invokeMethod<bool>('setTtsMode', {
'mod': mod,
}) ??
false;
return result;
} catch (e) {
debugPrint('设置TTS模式失败:$e');
return false;
}
}
@override @override
Future<bool> setVoiceFlocking(String speakerProfileId) async { Future<bool> setVoiceFlocking(String speakerProfileId) async {
return true; return true;

3
lib/data/services/tts_service.dart

@ -19,6 +19,9 @@ abstract class TtsService {
/// 设置语音 /// 设置语音
Future<bool> setVoice(String voiceName); Future<bool> setVoice(String voiceName);
/// 设置TTS模式
Future<bool> setTtsMode(String mod);
/// 设置语音复刻 /// 设置语音复刻
Future<bool> setVoiceFlocking(String speakerProfileId); Future<bool> setVoiceFlocking(String speakerProfileId);

4
lib/modules/agent/controllers/agent_controller.dart

@ -1880,7 +1880,7 @@ class AgentController extends GetxController with WidgetsBindingObserver {
// 开始按住说话 // 开始按住说话
void startPushToTalk() async { void startPushToTalk() async {
Logger.i(TAG, '🎤 开始按住说话'); print('liwei--------- 🎤 开始按住说话');
// 只有在语音模式且按住说话模式下才能按住说话 // 只有在语音模式且按住说话模式下才能按住说话
if (isTextInputMode.value) return; if (isTextInputMode.value) return;
@ -1898,7 +1898,7 @@ class AgentController extends GetxController with WidgetsBindingObserver {
// 结束按住说话 // 结束按住说话
void endPushToTalk() async { void endPushToTalk() async {
Logger.i(TAG, '🎤 结束按住说话'); print('liwei----------🎤 结束按住说话');
// 如果不在按住说话状态,直接返回 // 如果不在按住说话状态,直接返回
if (!isPushToTalkActive.value) return; if (!isPushToTalkActive.value) return;

34
lib/modules/opus_test/controllers/opus_test_controller.dart

@ -4,6 +4,7 @@ import 'dart:typed_data';
import 'package:file_picker/file_picker.dart'; import 'package:file_picker/file_picker.dart';
import 'package:flutter/foundation.dart'; import 'package:flutter/foundation.dart';
import 'package:flutter/material.dart'; import 'package:flutter/material.dart';
import 'package:flutter/services.dart';
import 'package:get/get.dart'; import 'package:get/get.dart';
import 'package:path_provider/path_provider.dart'; import 'package:path_provider/path_provider.dart';
import 'package:just_audio/just_audio.dart'; import 'package:just_audio/just_audio.dart';
@ -42,8 +43,10 @@ class OpusTestController extends GetxController {
String? tempPcmPath; String? tempPcmPath;
// 杰理OPUS解码器 // 杰理OPUS解码器
late JlOpus jlOpus; JlOpus? jlOpus;
StreamSubscription? _eventSubscription; StreamSubscription? _eventSubscription;
static const MethodChannel _bleServiceChannel =
MethodChannel('com.yunqiinnovation.ble_service');
@override @override
void onInit() { void onInit() {
@ -56,16 +59,20 @@ class OpusTestController extends GetxController {
void onClose() { void onClose() {
player.dispose(); player.dispose();
_eventSubscription?.cancel(); _eventSubscription?.cancel();
jlOpus.dispose(); jlOpus?.dispose();
super.onClose(); super.onClose();
} }
// 初始化OPUS解码器 // 初始化OPUS解码器
Future<void> _initOpusDecoder() async { Future<void> _initOpusDecoder() async {
if (!Platform.isAndroid) {
return;
}
jlOpus = JlOpus(); jlOpus = JlOpus();
// 监听解码器事件 // 监听解码器事件
_eventSubscription = jlOpus.eventStream.listen((event) { _eventSubscription = jlOpus!.eventStream.listen((event) {
switch (event.event) { switch (event.event) {
case 'onStart': case 'onStart':
statusMessage.value = '开始${event.type == "file" ? "文件" : "流"}解码...'; statusMessage.value = '开始${event.type == "file" ? "文件" : "流"}解码...';
@ -86,12 +93,29 @@ class OpusTestController extends GetxController {
}); });
// 初始化OPUS解码器 // 初始化OPUS解码器
final initialized = await jlOpus.initOpusDecoder(); final initialized = await jlOpus!.initOpusDecoder();
if (!initialized) { if (!initialized) {
statusMessage.value = 'OPUS解码器初始化失败'; statusMessage.value = 'OPUS解码器初始化失败';
} }
} }
Future<String?> _decodeOpusFile(
String inPath, String outPath, OpusOption option) async {
if (Platform.isIOS) {
final result = await _bleServiceChannel.invokeMethod<String>(
'decodeOpusFile',
{
'inPath': inPath,
'outPath': outPath,
...option.toMap(),
},
);
return result;
}
return jlOpus?.decodeOpusFile(inPath, outPath, option);
}
// 请求必要权限 // 请求必要权限
Future<void> requestPermissions() async { Future<void> requestPermissions() async {
if (Platform.isAndroid) { if (Platform.isAndroid) {
@ -241,7 +265,7 @@ class OpusTestController extends GetxController {
); );
// 解码文件 // 解码文件
final pcmPath = await jlOpus.decodeOpusFile(inPath, outPath, option); final pcmPath = await _decodeOpusFile(inPath, outPath, option);
if (pcmPath == null) { if (pcmPath == null) {
statusMessage.value = '解码失败'; statusMessage.value = '解码失败';
isDecoding.value = false; isDecoding.value = false;

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

@ -344,6 +344,15 @@ class TranslationController extends GetxController with WidgetsBindingObserver {
} catch (_) {} } catch (_) {}
} }
//设置播报的模式
if (mode == 'simultaneous') {
await _ttsService.setTtsMode('phone_call');
} else if (mode == 'faceToFace') {
await _ttsService.setTtsMode('push_to_talk');
} else {
await _ttsService.setTtsMode('normal');
}
// 检查悬浮窗支持情况 // 检查悬浮窗支持情况
if (isFloatingWindowEnabled.value && !_isFloatingWindowSupportedMode()) { if (isFloatingWindowEnabled.value && !_isFloatingWindowSupportedMode()) {
// 当前模式不支持悬浮窗,自动禁用 // 当前模式不支持悬浮窗,自动禁用
@ -436,7 +445,8 @@ class TranslationController extends GetxController with WidgetsBindingObserver {
} }
await _asrService.startContinuousRecognition( await _asrService.startContinuousRecognition(
_audioSourceType, _audioSourceType,
isRemoveFirstPunctuation: false, isRemoveFirstPunctuation: true,
mode: 'push_to_talk',
); //开启识别 ); //开启识别
_startAsrActiveTracking(); //开启识别活动跟踪 _startAsrActiveTracking(); //开启识别活动跟踪
if (isRecording.value) { if (isRecording.value) {
@ -946,7 +956,17 @@ class TranslationController extends GetxController with WidgetsBindingObserver {
!bleManager.isCodecActive) { !bleManager.isCodecActive) {
return; return;
} }
await _startAsrService(); if (currentMode.value == 'simultaneous') {
Logger.info('liwei-------- 开始语音识别,模式:同声传译,ASR phone_call');
await _startAsrService('phone_call');
} else if (currentMode.value == 'faceToFace') {
Logger.info('iwei-------- 开始语音识别,模式:面对面,ASR push_to_talk');
await _startAsrService('push_to_talk');
} else {
Logger.info(
'iwei-------- 开始语音识别,模式:${currentMode.value},ASR 模式:normal');
await _startAsrService('normal');
}
await _finalizeRecognitionStart(); await _finalizeRecognitionStart();
} }
} catch (e) { } catch (e) {
@ -1045,8 +1065,20 @@ class TranslationController extends GetxController with WidgetsBindingObserver {
} }
/// 启动语音识别 /// 启动语音识别
Future<void> _startAsrService() async { Future<void> _startAsrService(String mode) async {
await _asrService.startContinuousRecognition(_audioSourceType); if (mode == 'push_to_talk') {
await _asrService.startContinuousRecognition(
_audioSourceType,
isRemoveFirstPunctuation: false,
mode: mode,
);
} else {
await _asrService.startContinuousRecognition(
_audioSourceType,
isRemoveFirstPunctuation: true,
mode: mode,
);
}
} }
/// 完成识别启动 /// 完成识别启动

2
local_plugins/agent_service/android/src/main/kotlin/com/yunqiinnovation/agent_service/AgentService.kt

@ -147,7 +147,7 @@ object AgentService : CoroutineScope {
// 空闲检测相关 // 空闲检测相关
private var idleCheckJob: Job? = null private var idleCheckJob: Job? = null
private val maxIdleSeconds = 10 // 最大空闲秒数(默认值) private val maxIdleSeconds = 5 // 最大空闲秒数(默认值)
// 语音识别模式 // 语音识别模式
private var currentRecognitionMode = "normal" // 当前识别模式:normal, ble_wakeup, phone_call, push_to_talk private var currentRecognitionMode = "normal" // 当前识别模式:normal, ble_wakeup, phone_call, push_to_talk

164
local_plugins/agent_service/ios/agent_service/Sources/agent_service/AgentServiceImpl.swift

@ -100,9 +100,10 @@ class AgentServiceImpl: NSObject {
private let timerQueue = DispatchQueue(label: "com.example.idleTimer.queue") private let timerQueue = DispatchQueue(label: "com.example.idleTimer.queue")
private var idleTimer: DispatchSourceTimer? private var idleTimer: DispatchSourceTimer?
private let maxIdleSeconds: TimeInterval = 10 private let maxIdleSeconds: TimeInterval = 5
var micCapture: MicrophoneCapture! var micCapture: MicrophoneCapture!
// 语音识别模式 // 语音识别模式
private var lastRecognitionMode: String = "normal"
public var currentRecognitionMode: String = "normal" public var currentRecognitionMode: String = "normal"
// 识别结果 // 识别结果
private var recognitionResult = "" private var recognitionResult = ""
@ -191,7 +192,7 @@ class AgentServiceImpl: NSObject {
""" """
} }
func initialize(config: [String: Any]) -> Bool { func initialize(config: [String: Any]) -> Bool {
if isInitialized { return true } // if isInitialized { return true }
// os_log("initialize: config=%{public}@", log: logger, type: .info, config) // os_log("initialize: config=%{public}@", log: logger, type: .info, config)
self.config = config self.config = config
@ -279,6 +280,8 @@ class AgentServiceImpl: NSObject {
self.xunfeiAccessKeySecret = xunfeiAccessKeySecret self.xunfeiAccessKeySecret = xunfeiAccessKeySecret
} }
//启动定位服务 //启动定位服务
NotificationCenter.default.post(name: .locationstartEvent, object: nil) NotificationCenter.default.post(name: .locationstartEvent, object: nil)
// addlocationmonitor() // addlocationmonitor()
@ -311,6 +314,8 @@ class AgentServiceImpl: NSObject {
LocationService.shared.getCurrentLocation() { location in LocationService.shared.getCurrentLocation() { location in
// self.locationInfo = location // self.locationInfo = location
} }
// currentRecognitionMode = "normal"
// checkAdnSetAudioSession()
return isInitialized return isInitialized
} }
@ -339,7 +344,6 @@ class AgentServiceImpl: NSObject {
sendError("初始化语音识别服务失败", code: "ASR_INIT_ERROR") sendError("初始化语音识别服务失败", code: "ASR_INIT_ERROR")
return false return false
} }
guard let ttsSuccess = azureTtsHelper?.initialize( guard let ttsSuccess = azureTtsHelper?.initialize(
ttsAppId: "", ttsAppId: "",
ttsAppToken: azureSpeechKey, ttsAppToken: azureSpeechKey,
@ -387,6 +391,12 @@ class AgentServiceImpl: NSObject {
return true return true
} }
private func checkAdnSetAudioSession() {
os_log("liwei------ 检查并设置音频会话", log: logger, type: .info)
let audioSourceType: AzureAsrHelper.AudioSourceType = BleService.shared.isConnected() ? .external : .microphone
azureAsrHelper?.performAudioStart(mod: currentRecognitionMode, audioSourceType: audioSourceType, audioDataCallback: nil)
}
private func startIdleCheckForMode(mode: String) { private func startIdleCheckForMode(mode: String) {
stopIdleCheck() stopIdleCheck()
guard isRecognizing else { return } guard isRecognizing else { return }
@ -394,24 +404,31 @@ class AgentServiceImpl: NSObject {
// 根据模式决定空闲检测策略 // 根据模式决定空闲检测策略
switch mode { switch mode {
case "phone_call": case "phone_call":
azureAsrHelper?.restoreOriginalAudioState()
// 通话模式:禁用空闲检测,保持持续激活 // 通话模式:禁用空闲检测,保持持续激活
os_log("通话模式:禁用空闲检测", log: logger, type: .info) os_log("通话模式:禁用空闲检测", log: logger, type: .info)
return return
case "ble_wakeup": case "ble_wakeup":
azureAsrHelper?.disableBluetoothAudio() if !isSpeaking {
azureAsrHelper?.disableBluetoothAudio()
}
// BLE唤醒模式:使用默认的空闲检测 // BLE唤醒模式:使用默认的空闲检测
startIdleCheck(idleSeconds: maxIdleSeconds) startIdleCheck(idleSeconds: maxIdleSeconds)
case "push_to_talk": case "push_to_talk":
azureAsrHelper?.disableBluetoothAudio() if !isSpeaking {
azureAsrHelper?.disableBluetoothAudio()
}
// 按住说话模式:使用默认的空闲检测 // 按住说话模式:使用默认的空闲检测
startIdleCheck(idleSeconds: maxIdleSeconds) startIdleCheck(idleSeconds: maxIdleSeconds)
case "normal": case "normal":
azureAsrHelper?.disableBluetoothAudio() if !isSpeaking {
azureAsrHelper?.disableBluetoothAudio()
}
// 普通模式:使用默认的空闲检测 // 普通模式:使用默认的空闲检测
startIdleCheck(idleSeconds: maxIdleSeconds) startIdleCheck(idleSeconds: maxIdleSeconds)
default: default:
azureAsrHelper?.disableBluetoothAudio() if !isSpeaking {
azureAsrHelper?.disableBluetoothAudio()
}
// 未知模式:使用默认的空闲检测 // 未知模式:使用默认的空闲检测
os_log("未知识别模式: %{public}@,使用默认空闲检测", log: logger, type: .error, mode) os_log("未知识别模式: %{public}@,使用默认空闲检测", log: logger, type: .error, mode)
startIdleCheck(idleSeconds: maxIdleSeconds) startIdleCheck(idleSeconds: maxIdleSeconds)
@ -425,7 +442,7 @@ class AgentServiceImpl: NSObject {
self.idleTimer = nil self.idleTimer = nil
} }
guard self.isRecognizing else { return } guard self.isRecognizing else { return }
os_log("启动空闲检测,超时时间: %.0f秒,模式: %{public}@", os_log("liwei--------启动空闲检测,超时时间: %.0f秒,模式: %{public}@",
log: self.logger, log: self.logger,
type: .info, type: .info,
idleSeconds, idleSeconds,
@ -436,7 +453,7 @@ class AgentServiceImpl: NSObject {
guard let self = self else { return } guard let self = self else { return }
DispatchQueue.main.async { [weak self] in DispatchQueue.main.async { [weak self] in
guard let self = self else { return } guard let self = self else { return }
os_log("空闲检测触发自动停止,模式: %{public}@", os_log("liwei--------空闲检测触发自动停止,模式: %{public}@",
log: self.logger, log: self.logger,
type: .info, type: .info,
self.currentRecognitionMode) self.currentRecognitionMode)
@ -472,7 +489,18 @@ class AgentServiceImpl: NSObject {
startIdleCheckForMode(mode: currentRecognitionMode) startIdleCheckForMode(mode: currentRecognitionMode)
} }
} }
func startRecognition(useBle: Bool = false, mode: String = "normal") -> Bool { func startRecognition(useBle: Bool = false, mode: String = "normal") -> Bool {
if !isInitialized {
sendError("服务未初始化", code: "NOT_INITIALIZED")
return false
}
if useBle, mode == "ble_wakeup", currentRecognitionMode == "ble_wakeup", (isRecognizing || isStartingRecognition) {
os_log("liwei-------- ble_wakeup 连续触发,识别中/启动中,忽略重复 startRecognition", log: logger, type: .info)
audioPlayer?.playStartSound()
return true
}
// 接收到唤醒信号,打开编码器 (设备侧) // 接收到唤醒信号,打开编码器 (设备侧)
print("ai启动语音\(useBle))") print("ai启动语音\(useBle))")
var openResult = false // 添加openResult变量定义 var openResult = false // 添加openResult变量定义
@ -481,18 +509,20 @@ class AgentServiceImpl: NSObject {
// 捕获openEncoder的返回值 // 捕获openEncoder的返回值
openResult = BleService.shared.openEncoder() // 正确捕获返回值 openResult = BleService.shared.openEncoder() // 正确捕获返回值
} }
if !isInitialized {
sendError("服务未初始化", code: "NOT_INITIALIZED")
return false
}
if !useBle { if !useBle {
guard AVAudioSession.sharedInstance().recordPermission == .granted else { guard AVAudioSession.sharedInstance().recordPermission == .granted else {
sendError("无麦克风权限", code: "PERMISSION_DENIED") sendError("无麦克风权限", code: "PERMISSION_DENIED")
return false return false
} }
} }
lastRecognitionMode = currentRecognitionMode
stopRecognition() stopRecognition()
if isAiStreaming {
isAiStreaming = false
_ = chatApiService?.cancelCurrentStream()
}
stopTts() stopTts()
@ -501,38 +531,33 @@ class AgentServiceImpl: NSObject {
// 如果音乐正在播放,则暂停 // 如果音乐正在播放,则暂停
if MusicService.shared.isPlaying { if MusicService.shared.isPlaying {
MusicService.shared.pause() MusicService.shared.pause()
} }
// 设置当前识别模式 // 设置当前识别模式
currentRecognitionMode = mode currentRecognitionMode = mode
os_log("设置语音识别模式: %{public}@", log: logger, type: .info, mode)
let audioSourceType: AzureAsrHelper.AudioSourceType = useBle ? .external : .microphone let audioSourceType: AzureAsrHelper.AudioSourceType = useBle ? .external : .microphone
if useBle {
// AudioSessionHub.shared.begin(.playback)
} else {
//AudioSessionHub.shared.begin(.voice)
}
// 新增:标记为"启动中",用于允许 stop 在启动未完成时也能生效 // 新增:标记为"启动中",用于允许 stop 在启动未完成时也能生效
isStartingRecognition = true isStartingRecognition = true
print("ai启动语音=\(audioSourceType)=\(useBle)") print("ai启动语音=\(audioSourceType)=\(useBle)")
guard let success = azureAsrHelper?.startContinuousRecognition( guard let success = azureAsrHelper?.startContinuousRecognition(
// callback: self, mod:self.currentRecognitionMode,
audioSourceType: audioSourceType, audioSourceType: audioSourceType,
isRemoveFirstPunctuation: false isRemoveFirstPunctuation: false
), success else { ), success else {
// 新增:启动失败时复位"启动中"状态 // 新增:启动失败时复位"启动中"状态
isStartingRecognition = false isStartingRecognition = false
sendError("启动语音识别失败", code: "RECOGNITION_START_ERROR") sendError("启动语音识别失败", code: "RECOGNITION_START_ERROR")
os_log("liwei--------------- 启动语音识别失败: %{public}@", log: logger, type: .info, mode)
return false return false
} }
if(useBle){ if(useBle){
isKeepResult = true isKeepResult = true
} }
os_log("liwei--------------- 启动语音识别: %{public}@", log: logger, type: .info, mode)
// 根据模式启动相应的空闲检测(此处 guard isRecognizing,会在 onSessionStarted 中再次启动) // 根据模式启动相应的空闲检测(此处 guard isRecognizing,会在 onSessionStarted 中再次启动)
startIdleCheckForMode(mode: mode) startIdleCheckForMode(mode: mode)
@ -560,29 +585,43 @@ class AgentServiceImpl: NSObject {
return true return true
} }
// 关键修复:如果正在结束的是通话模式,则立即将会话重置为 normal 状态
if self.currentRecognitionMode == "phone_call" {
os_log("通话模式结束,强制重置音频会话至 'normal' 状态", log: logger, type: .info)
self.currentRecognitionMode = "normal"
// 使用 normal 模式的配置来清理和重置音频会话
self.checkAdnSetAudioSession()
}
os_log("停止语音识别,当前模式: %{public}@", log: logger, type: .info, currentRecognitionMode) os_log("停止语音识别,当前模式: %{public}@", log: logger, type: .info, currentRecognitionMode)
let success = azureAsrHelper?.stopContinuousRecognition() ?? false let success = azureAsrHelper?.stopContinuousRecognition() ?? false
if success {
if !success && isStartingRecognition { do {
// let audioSession = AVAudioSession.sharedInstance()
// try audioSession.setCategory(.playback, mode: .spokenAudio, options: [.mixWithOthers, .allowBluetoothA2DP])
// try audioSession.setActive(true)
os_log("Audio session switched to playback successfully after recognition.", log: logger, type: .info)
} catch {
os_log("Failed to switch audio session to playback: %{public}@", log: logger, type: .error, error.localizedDescription)
}
} else if isStartingRecognition {
// 启动尚未完成,先记录一次待停止请求,onSessionStarted 到来后立即 stop // 启动尚未完成,先记录一次待停止请求,onSessionStarted 到来后立即 stop
stopRequestedDuringStart = true stopRequestedDuringStart = true
// 即使底层 stop 失败,但因为我们已经记录了停止请求,所以对调用方来说,可以认为是成功的
// 后续的 onSessionStarted 会处理真正的停止逻辑
} }
stopIdleCheck() stopIdleCheck()
// 只有之前在播放时才恢复播放 // 只有之前在播放时才恢复播放
if wasMusicPlayingBeforeRecognition { if wasMusicPlayingBeforeRecognition {
// QPlayAutoManager.shared.play()
wasMusicPlayingBeforeRecognition = false wasMusicPlayingBeforeRecognition = false
MusicService.shared.resume() MusicService.shared.resume()
} }
// 重置识别模式
// currentRecognitionMode = "normal"
// 启动中但已请求停止,也视为已接受停止请求 // 启动中但已请求停止,也视为已接受停止请求
return success || isStartingRecognition return success || stopRequestedDuringStart
} }
func disableBluetoothAudio() { func disableBluetoothAudio() {
azureAsrHelper?.disableBluetoothAudio() azureAsrHelper?.disableBluetoothAudio()
@ -603,7 +642,7 @@ class AgentServiceImpl: NSObject {
return true return true
} }
func speakText(sessionid sessionid:String,_ text: String) -> Bool { func speakText(sessionid: String, _ text: String) -> Bool {
if !isInitialized { if !isInitialized {
sendError("服务未初始化", code: "NOT_INITIALIZED") sendError("服务未初始化", code: "NOT_INITIALIZED")
return false return false
@ -612,8 +651,6 @@ class AgentServiceImpl: NSObject {
if text.isEmpty { if text.isEmpty {
return false return false
} }
_ = azureTtsHelper?.setSpeechParams(rate: ttsRatePercent, pitch: 0, volume: 100)
return azureTtsHelper?.speakOnce(sessionid:sessionid,text) ?? false return azureTtsHelper?.speakOnce(sessionid:sessionid,text) ?? false
} }
@ -675,8 +712,13 @@ class AgentServiceImpl: NSObject {
sendError("文本输入不能为空", code: "EMPTY_TEXT") sendError("文本输入不能为空", code: "EMPTY_TEXT")
return false return false
} }
stopIdleCheck()
os_log("开始处理文本输入", log: logger, type: .info) if isAiStreaming {
isAiStreaming = false
_ = chatApiService?.cancelCurrentStream()
}
stopTts()
processWithChatApiService(sessionid: sessionid, text: text, speakResponse: speakResponse) processWithChatApiService(sessionid: sessionid, text: text, speakResponse: speakResponse)
return true return true
} }
@ -686,7 +728,8 @@ class AgentServiceImpl: NSObject {
sendError("ChatAPI服务未初始化", code: "CHATAPI_NOT_INITIALIZED") sendError("ChatAPI服务未初始化", code: "CHATAPI_NOT_INITIALIZED")
return return
} }
_ = azureTtsHelper?.setSpeechParams(rate: ttsRatePercent, pitch: 0, volume: 100)
_ = azureTtsHelper?.setTtsMode(mod: currentRecognitionMode)
let userMessage = chatApiService.createUserMessage(content: text) let userMessage = chatApiService.createUserMessage(content: text)
processWithChatApiServiceInternal(sessionid: sessionid,userMessage: userMessage, displayText: text, speakResponse: speakResponse) processWithChatApiServiceInternal(sessionid: sessionid,userMessage: userMessage, displayText: text, speakResponse: speakResponse)
} }
@ -708,7 +751,7 @@ class AgentServiceImpl: NSObject {
stopTts(); stopTts();
isAiStreaming = true isAiStreaming = true
os_log("设置AI流式状态为true", log: logger, type: .info) // os_log("设置AI流式状态为true", log: logger, type: .info)
var messages: [[String: Any]] = [] var messages: [[String: Any]] = []
@ -725,17 +768,17 @@ class AgentServiceImpl: NSObject {
// print("用户画像数据: \(info)") // print("用户画像数据: \(info)")
let _systemPrompt = fillTemplate(systemPrompt, with: info) let _systemPrompt = fillTemplate(systemPrompt, with: info)
print("替换后的系统提示词: \(_systemPrompt)") // print("替换后的系统提示词: \(_systemPrompt)")
messages.append(chatApiService.createSystemMessage(content: _systemPrompt)) messages.append(chatApiService.createSystemMessage(content: _systemPrompt))
os_log("添加系统提示消息", log: logger, type: .info) // os_log("添加系统提示消息", log: logger, type: .info)
} }
messages.append(contentsOf: chatHistory) messages.append(contentsOf: chatHistory)
os_log("添加历史消息: %d条", log: logger, type: .info, chatHistory.count) // os_log("添加历史消息: %d条", log: logger, type: .info, chatHistory.count)
messages.append(userMessage) messages.append(userMessage)
os_log("添加用户消息", log: logger, type: .info) // os_log("添加用户消息", log: logger, type: .info)
// 保存用户消息 // 保存用户消息
self.chatStorageHelper.saveMessage( self.chatStorageHelper.saveMessage(
@ -1772,6 +1815,7 @@ class ChatApiStreamCallback: StreamCallback {
responseBuilder += token responseBuilder += token
if speakResponse && reply && broadcast{ if speakResponse && reply && broadcast{
// _ = agentService.azureTtsHelper?.setTtsMode(mod:agentService.currentRecognitionMode)
try agentService.azureTtsHelper?.speakStream(sessionid:sessionid,token) try agentService.azureTtsHelper?.speakStream(sessionid:sessionid,token)
// 在开始流式TTS时立即停止气泡音 // 在开始流式TTS时立即停止气泡音
} }
@ -1797,6 +1841,7 @@ class ChatApiStreamCallback: StreamCallback {
} }
if speakResponse && reply && broadcast && sessionid == agentService.currsessionId{ if speakResponse && reply && broadcast && sessionid == agentService.currsessionId{
// _ = agentService.azureTtsHelper?.setTtsMode(mod: agentService.currentRecognitionMode)
agentService.azureTtsHelper?.flushStream(sessionid:sessionid) agentService.azureTtsHelper?.flushStream(sessionid:sessionid)
} }
@ -1848,6 +1893,7 @@ class ChatApiStreamCallback: StreamCallback {
"message": message "message": message
]) ])
if(code == 2001){ if(code == 2001){
// _ = agentService.azureTtsHelper?.setTtsMode(mod: agentService.currentRecognitionMode)
agentService.azureTtsHelper?.speakStream(sessionid: sessionid,agentService.insufficientIntegralText) agentService.azureTtsHelper?.speakStream(sessionid: sessionid,agentService.insufficientIntegralText)
} }
agentService.isAiStreaming = false agentService.isAiStreaming = false
@ -1891,6 +1937,7 @@ class ChatApiStreamCallback: StreamCallback {
if (iscallingTool) { if (iscallingTool) {
os_log("收到函数调用: callingToolText:%{public}@", log: agentService.logger, type: .info, agentService.callingToolText) os_log("收到函数调用: callingToolText:%{public}@", log: agentService.logger, type: .info, agentService.callingToolText)
// _ = agentService.azureTtsHelper?.setTtsMode(mod: agentService.currentRecognitionMode)
agentService.azureTtsHelper?.speakStream(sessionid: sessionid,agentService.callingToolText) agentService.azureTtsHelper?.speakStream(sessionid: sessionid,agentService.callingToolText)
iscallingTool = false iscallingTool = false
} }
@ -1962,11 +2009,11 @@ extension AgentServiceImpl: TtsEventListener {
sendEvent(name: "tts_started", data: ["status": "started"]) sendEvent(name: "tts_started", data: ["status": "started"])
case .synthesisCompleted: case .synthesisCompleted:
isSpeaking = false
restartIdleCheck() restartIdleCheck()
sendEvent(name: "tts_completed", data: ["status": "completed"]) sendEvent(name: "tts_completed", data: ["status": "completed"])
case .playbackStarted: case .playbackStarted:
isSpeaking = true
restartIdleCheck() restartIdleCheck()
sendEvent(name: "playback_started", data: ["status": "playback_started"]) sendEvent(name: "playback_started", data: ["status": "playback_started"])
audioPlayer?.stopAwaitSound() audioPlayer?.stopAwaitSound()
@ -1975,6 +2022,7 @@ extension AgentServiceImpl: TtsEventListener {
BleService.shared.closeCodec() BleService.shared.closeCodec()
} }
case .playbackCompleted: case .playbackCompleted:
isSpeaking = false
restartIdleCheck() restartIdleCheck()
if(isInterrupt){//播报结束。打开解码器 if(isInterrupt){//播报结束。打开解码器
BleService.shared.openB1Encoder() BleService.shared.openB1Encoder()
@ -2087,10 +2135,12 @@ class AudioPlayer {
} }
func playStartSound() { func playStartSound() {
os_log("liwei 播放 playStartSound", log: self.logger, type: .debug)
playSound(named: "start", fileType: "mp3") playSound(named: "start", fileType: "mp3")
} }
func playAwaitSound() { func playAwaitSound() {
os_log("liwei 播放 playAwaitSound", log: self.logger, type: .debug)
// 若已处于等待音效播放状态,则直接返回,避免频繁启停 // 若已处于等待音效播放状态,则直接返回,避免频繁启停
if isAwaitSoundActive { return } if isAwaitSoundActive { return }
@ -2108,6 +2158,7 @@ class AudioPlayer {
} }
func playStopSound() { func playStopSound() {
os_log("liwei 播放 playStopSound", log: self.logger, type: .debug)
playSound(named: "stop", fileType: "mp3") playSound(named: "stop", fileType: "mp3")
} }
@ -2396,6 +2447,11 @@ extension AgentServiceImpl: AzureAsrHelper.ContinuousRecognizeCallback {
if (currentRecognitionMode == "ble_wakeup") { if (currentRecognitionMode == "ble_wakeup") {
audioPlayer?.playStopSound() audioPlayer?.playStopSound()
} }
if wasMusicPlayingBeforeRecognition {
wasMusicPlayingBeforeRecognition = false
MusicService.shared.resume()
}
} }
@ -2417,6 +2473,11 @@ extension AgentServiceImpl: AzureAsrHelper.ContinuousRecognizeCallback {
isStartingRecognition = false isStartingRecognition = false
stopRequestedDuringStart = false stopRequestedDuringStart = false
stopIdleCheck() stopIdleCheck()
if wasMusicPlayingBeforeRecognition {
wasMusicPlayingBeforeRecognition = false
MusicService.shared.resume()
}
} }
func onError(sessionid:String ,_ errorCode: Int, _ error: String) { func onError(sessionid:String ,_ errorCode: Int, _ error: String) {
@ -2427,6 +2488,11 @@ extension AgentServiceImpl: AzureAsrHelper.ContinuousRecognizeCallback {
isStartingRecognition = false isStartingRecognition = false
stopRequestedDuringStart = false stopRequestedDuringStart = false
stopIdleCheck() stopIdleCheck()
if wasMusicPlayingBeforeRecognition {
wasMusicPlayingBeforeRecognition = false
MusicService.shared.resume()
}
} }
// MARK: - isInterrupt 缓存方法 // MARK: - isInterrupt 缓存方法

5
local_plugins/agent_service/ios/agent_service/Sources/agent_service/AgentServicePlugin.swift

@ -74,7 +74,7 @@ public class AgentServicePlugin: NSObject, FlutterPlugin {
result(true) result(true)
return return
} }
print("startConversation") print("liwei----------AgentServicePlugin startConversation")
let mode = (call.arguments as? [String: Any])?["mode"] as? String ?? "normal" let mode = (call.arguments as? [String: Any])?["mode"] as? String ?? "normal"
impl.startRecognition(mode: mode) impl.startRecognition(mode: mode)
result(true) result(true)
@ -99,7 +99,7 @@ public class AgentServicePlugin: NSObject, FlutterPlugin {
result(true) result(true)
return return
} }
print("stopConversation") print("liwei-------AgentServicePlugin stopConversation")
// 关闭编码器 (设备侧) // 关闭编码器 (设备侧)
BleService.shared.closeCodec() BleService.shared.closeCodec()
impl.stopRecognition() impl.stopRecognition()
@ -121,6 +121,7 @@ public class AgentServicePlugin: NSObject, FlutterPlugin {
} }
let speakResponse = arguments["speakResponse"] as? Bool ?? false let speakResponse = arguments["speakResponse"] as? Bool ?? false
impl.currentRecognitionMode = "normal"
let success = impl.processTextInput(sessionid: sessionid,text, speakResponse: speakResponse) let success = impl.processTextInput(sessionid: sessionid,text, speakResponse: speakResponse)
result(success) result(success)

5
local_plugins/agent_service/lib/agent_service.dart

@ -265,8 +265,8 @@ class AgentService {
try { try {
// 在启动新对话前,确保之前的对话已完全停止 // 在启动新对话前,确保之前的对话已完全停止
// 这不仅可以清理状态,还能避免因快速切换导致的资源冲突 // 这不仅可以清理状态,还能避免因快速切换导致的资源冲突
await stopConversation(); // await stopConversation();
print('liwei--------- view startConversation');
final bool result = await _channel.invokeMethod('startConversation', { final bool result = await _channel.invokeMethod('startConversation', {
'mode': mode, 'mode': mode,
}); });
@ -337,6 +337,7 @@ class AgentService {
// 即使频繁调用,Native 层也能处理(我们已经修复了 Native 层的并发问题)。 // 即使频繁调用,Native 层也能处理(我们已经修复了 Native 层的并发问题)。
try { try {
print('liwei--------- view stopConversation');
final bool result = await _channel.invokeMethod('stopConversation'); final bool result = await _channel.invokeMethod('stopConversation');
return result; return result;
} on PlatformException catch (e) { } on PlatformException catch (e) {

21
local_plugins/azure_speech/ios/azure_speech/Sources/azure_speech/AzureAsrHelper.swift

@ -73,6 +73,7 @@ public class AzureAsrHelper: NSObject {
} }
public var audioSourceType = AudioSourceType.microphone public var audioSourceType = AudioSourceType.microphone
private var isPhoneCallMode = false
// 音频处理 // 音频处理
// private var externalAudioStream: ExternalAudioPullStream? // private var externalAudioStream: ExternalAudioPullStream?
@ -80,6 +81,10 @@ public class AzureAsrHelper: NSObject {
// 音频处理 // 音频处理
public var audioStream: SimpleAudioReceiver? public var audioStream: SimpleAudioReceiver?
public func asrProvider() -> String { public func asrProvider() -> String {
return useXunfei ? "xunfei" : "azure" return useXunfei ? "xunfei" : "azure"
} }
@ -264,6 +269,7 @@ public class AzureAsrHelper: NSObject {
/// 1) 新任务创建前,取消旧的未执行任务,确保只有“最后一次启动请求”会生效 /// 1) 新任务创建前,取消旧的未执行任务,确保只有“最后一次启动请求”会生效
/// 2) 任务执行前二次校验:必须是当前挂起任务且 audioStream?.isContinuousRecognitionActive 仍为 true 才执行启动 /// 2) 任务执行前二次校验:必须是当前挂起任务且 audioStream?.isContinuousRecognitionActive 仍为 true 才执行启动
private func startAudioRecordAsync( private func startAudioRecordAsync(
mod: String,
audioSourceType: AudioSourceType, audioSourceType: AudioSourceType,
audioDataCallback: SimpleAudioReceiver.AudioDataCallback? = nil, audioDataCallback: SimpleAudioReceiver.AudioDataCallback? = nil,
completion: @escaping (Bool) -> Void completion: @escaping (Bool) -> Void
@ -306,7 +312,7 @@ public class AzureAsrHelper: NSObject {
// 准备启动 // 准备启动
self.isAudioStarting = true self.isAudioStarting = true
self.pendingStopRequest = false // 启动前清空旧的停止请求 self.pendingStopRequest = false // 启动前清空旧的停止请求
let success = self.performAudioStart(audioSourceType: audioSourceType, audioDataCallback: audioDataCallback) let success = self.performAudioStart(mod: mod, audioSourceType: audioSourceType, audioDataCallback: audioDataCallback)
self.isAudioStarting = false self.isAudioStarting = false
self.isAudioStarted = success self.isAudioStarted = success
@ -338,13 +344,15 @@ public class AzureAsrHelper: NSObject {
/** /**
* 执行音频启动操作 * 执行音频启动操作
*/ */
private func performAudioStart( public func performAudioStart(
mod: String,
audioSourceType: AudioSourceType, audioSourceType: AudioSourceType,
audioDataCallback: SimpleAudioReceiver.AudioDataCallback? audioDataCallback: SimpleAudioReceiver.AudioDataCallback?
) -> Bool { ) -> Bool {
do { do {
if(!audioStream!.isContinuousRecognitionActive){ if(!audioStream!.isContinuousRecognitionActive){
audioStream?.startAudioRecord( audioStream?.startAudioRecord(
mod: mod,
audioSourceType: audioSourceType == .microphone ? .microphone : .external, audioSourceType: audioSourceType == .microphone ? .microphone : .external,
audioDataCallback: audioDataCallback audioDataCallback: audioDataCallback
) )
@ -411,6 +419,7 @@ public class AzureAsrHelper: NSObject {
* 若 stop 已发生则直接跳过 recognizer 的启动,避免“已停止但仍启动识别”的情况。 * 若 stop 已发生则直接跳过 recognizer 的启动,避免“已停止但仍启动识别”的情况。
*/ */
public func startContinuousRecognition( public func startContinuousRecognition(
mod: String = "normal",
audioSourceType: AudioSourceType = .microphone, audioSourceType: AudioSourceType = .microphone,
audioDataCallback: SimpleAudioReceiver.AudioDataCallback? = nil, audioDataCallback: SimpleAudioReceiver.AudioDataCallback? = nil,
isRemoveFirstPunctuation: Bool = true isRemoveFirstPunctuation: Bool = true
@ -468,7 +477,7 @@ public class AzureAsrHelper: NSObject {
} }
// 异步启动音频处理 // 异步启动音频处理
startAudioRecordAsync(audioSourceType: audioSourceType, audioDataCallback: proxyCallback) { [weak self] (success: Bool) in startAudioRecordAsync(mod: mod, audioSourceType: audioSourceType, audioDataCallback: proxyCallback) { [weak self] (success: Bool) in
guard let self = self else { return } guard let self = self else { return }
if success { if success {
@ -497,7 +506,7 @@ public class AzureAsrHelper: NSObject {
// self.audioStream?.isContinuousRecognitionActive = true // self.audioStream?.isContinuousRecognitionActive = true
// 异步启动音频处理 // 异步启动音频处理
startAudioRecordAsync(audioSourceType: audioSourceType, audioDataCallback: audioDataCallback) { [weak self] (success: Bool) in startAudioRecordAsync(mod: mod, audioSourceType: audioSourceType, audioDataCallback: audioDataCallback) { [weak self] (success: Bool) in
guard let self = self else { return } guard let self = self else { return }
if success { if success {
@ -768,6 +777,7 @@ public class AzureAsrHelper: NSObject {
if audioStream == nil { if audioStream == nil {
audioStream = SimpleAudioReceiver() audioStream = SimpleAudioReceiver()
audioStream?.initAudioRecord() // 确保调用初始化 audioStream?.initAudioRecord() // 确保调用初始化
} }
// 检查音频配置是否已存在 // 检查音频配置是否已存在
@ -992,7 +1002,7 @@ public class AzureAsrHelper: NSObject {
// 修复:添加类型转换 // 修复:添加类型转换
if(!audioStream!.isContinuousRecognitionActive){ if(!audioStream!.isContinuousRecognitionActive){
// 异步启动音频处理 // 异步启动音频处理
startAudioRecordAsync(audioSourceType: audioSourceType, audioDataCallback: audioDataCallback) { [weak self] (success: Bool) in startAudioRecordAsync(mod: "normal", audioSourceType: audioSourceType, audioDataCallback: audioDataCallback) { [weak self] (success: Bool) in
guard let self = self else { return } guard let self = self else { return }
} }
@ -1208,4 +1218,3 @@ private class AudioDataCallbackProxy: SimpleAudioReceiver.AudioDataCallback {

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

@ -382,6 +382,7 @@ private func sendAstEvent(_ event: [String: Any]) {
} }
let removeFirstPunctuation = args["isRemoveFirstPunctuation"] as? Bool ?? true let removeFirstPunctuation = args["isRemoveFirstPunctuation"] as? Bool ?? true
let mode = args["mode"] as? String ?? "normal"
// 根据参数确定音频源类型 // 根据参数确定音频源类型
let audioSourceType = isExternalActive ? let audioSourceType = isExternalActive ?
@ -392,6 +393,7 @@ private func sendAstEvent(_ event: [String: Any]) {
// 启动连续识别 // 启动连续识别
let success = azureAsrHelper.startContinuousRecognition( let success = azureAsrHelper.startContinuousRecognition(
mod: mode,
audioSourceType: audioSourceType, audioSourceType: audioSourceType,
isRemoveFirstPunctuation: removeFirstPunctuation isRemoveFirstPunctuation: removeFirstPunctuation
) )
@ -631,19 +633,14 @@ private func sendAstEvent(_ event: [String: Any]) {
let success = azureTtsHelper.setVoice(voiceName) let success = azureTtsHelper.setVoice(voiceName)
result(success) result(success)
case "setTtsMode":
case "setSpeechParams": guard let args = call.arguments as? [String: Any],
guard let args = call.arguments as? [String: Any] else { let mod = args["mod"] as? String else {
result(FlutterError(code: "INVALID_ARGUMENTS", message: "参数不能为空", details: nil)) result(FlutterError(code: "INVALID_ARGUMENTS", message: "参数不能为空", details: nil))
return return
} }
let success = azureTtsHelper.setTtsMode(mod: mod)
let rate = args["rate"] as? Int ?? 0
let pitch = args["pitch"] as? Int ?? 0
let volume = args["volume"] as? Int ?? 100
let success = azureTtsHelper.setSpeechParams(rate: rate, pitch: pitch, volume: volume)
result(success) result(success)
case "speakText": case "speakText":
guard let args = call.arguments as? [String: Any], guard let args = call.arguments as? [String: Any],
let sessionid = args["sessionid"] as? String, let sessionid = args["sessionid"] as? String,
@ -1119,7 +1116,7 @@ extension AzureSpeechPlugin: BleService.Callback {
} }
} }
print("分离音频数据 - 左声道: \(leftBuffer.count) bytes, 右声道: \(rightBuffer.count) bytes") // print("分离音频数据 - 左声道: \(leftBuffer.count) bytes, 右声道: \(rightBuffer.count) bytes")
// 安全解包版本 // 安全解包版本
guard let audioStream = azureAsrHelper.audioStream else { guard let audioStream = azureAsrHelper.audioStream else {

1154
local_plugins/azure_speech/ios/azure_speech/Sources/azure_speech/AzureTtsHelper.swift

File diff suppressed because it is too large

210
local_plugins/azure_speech/ios/azure_speech/Sources/tools/MicrophoneCapture.swift

@ -92,13 +92,7 @@ public class MicrophoneCapture: NSObject {
} }
do { do {
try audioSession.setPreferredSampleRate(sampleRate) try audioSession.setPreferredIOBufferDuration(0.02)
try audioSession.setPreferredIOBufferDuration(0.005) // 5ms缓冲
// 优先设置蓝牙音频路由
configureAudioRoute()
try audioSession.setActive(true)
print("音频会话配置成功") print("音频会话配置成功")
} catch { } catch {
print("音频会话配置失败: \(error.localizedDescription)") print("音频会话配置失败: \(error.localizedDescription)")
@ -153,11 +147,10 @@ public class MicrophoneCapture: NSObject {
// 取消扬声器强制输出,让音频通过蓝牙耳机输出 // 取消扬声器强制输出,让音频通过蓝牙耳机输出
do { do {
// 设置音频会话参数 // 设置音频会话参数
try audioSession.setCategory(.playAndRecord, try audioSession.setCategory(.playback,
mode: .videoChat, mode: .videoChat,
options: [.allowBluetoothA2DP, // 允许蓝牙耳机,不占用hfp链路 options: [.allowBluetoothA2DP, // 仅允许A2DP,不启用HFP
.mixWithOthers, .mixWithOthers])
.allowBluetooth]) // 添加音频优先级控制
try audioSession.overrideOutputAudioPort(.none) try audioSession.overrideOutputAudioPort(.none)
print("蓝牙模式:音频输出设置为蓝牙耳机") print("蓝牙模式:音频输出设置为蓝牙耳机")
@ -212,24 +205,55 @@ public class MicrophoneCapture: NSObject {
} }
// 检查麦克风权限 // 检查麦克风权限
let permission: AVAudioSession.RecordPermission
switch audioSession.recordPermission { switch audioSession.recordPermission {
case .granted: case .granted:
try setupAudioEngine() permission = .granted
case .denied: case .denied:
throw NSError(domain: "麦克风权限被拒绝", code: 0) permission = .denied
case .undetermined: case .undetermined:
let semaphore = DispatchSemaphore(value: 0)
var grantedResult = false
audioSession.requestRecordPermission { granted in audioSession.requestRecordPermission { granted in
if granted { grantedResult = granted
do { semaphore.signal()
try self.setupAudioEngine()
} catch {
print("音频引擎设置失败: \(error.localizedDescription)")
}
}
} }
let waitResult = semaphore.wait(timeout: .now() + 60)
if waitResult == .timedOut {
throw NSError(domain: "麦克风权限请求超时", code: 2)
}
permission = grantedResult ? .granted : .denied
@unknown default: @unknown default:
throw NSError(domain: "未知权限状态", code: 1) throw NSError(domain: "未知权限状态", code: 1)
} }
guard permission == .granted else {
throw NSError(domain: "麦克风权限被拒绝", code: 0)
}
do {
print("startCapture: 激活前 category=\(audioSession.category.rawValue) mode=\(audioSession.mode.rawValue) options=\(audioSession.categoryOptions)")
print("startCapture: 激活前 sampleRate=\(audioSession.sampleRate) ioBuffer=\(audioSession.ioBufferDuration)")
print("startCapture: 激活前 route inputs=\(audioSession.currentRoute.inputs.map { $0.portType.rawValue }) outputs=\(audioSession.currentRoute.outputs.map { $0.portType.rawValue })")
if let inputs = audioSession.availableInputs,
let builtInMic = inputs.first(where: { $0.portType == .builtInMic }) {
try? audioSession.setPreferredInput(builtInMic)
}
try audioSession.setActive(true)
print("startCapture: 激活后 category=\(audioSession.category.rawValue) mode=\(audioSession.mode.rawValue) options=\(audioSession.categoryOptions)")
print("startCapture: 激活后 sampleRate=\(audioSession.sampleRate) ioBuffer=\(audioSession.ioBufferDuration)")
print("startCapture: 激活后 route inputs=\(audioSession.currentRoute.inputs.map { $0.portType.rawValue }) outputs=\(audioSession.currentRoute.outputs.map { $0.portType.rawValue })")
} catch {
throw error
}
isCapturing = true
do {
try setupAudioEngine()
} catch {
isCapturing = false
cleanupAudioEngine()
throw error
}
} }
/** /**
@ -280,36 +304,54 @@ public class MicrophoneCapture: NSObject {
if let audioSession = audioSession { if let audioSession = audioSession {
sampleRate = audioSession.sampleRate sampleRate = audioSession.sampleRate
} }
// 设置音频格式
let inputFormat = audioInputNode.inputFormat(forBus: 0)
// iOS 13+ 启用语音处理
if #available(iOS 13.0, *) { if #available(iOS 13.0, *) {
try audioInputNode.setVoiceProcessingEnabled(true) let session = self.audioSession ?? AVAudioSession.sharedInstance()
let outputs = session.currentRoute.outputs
let hasBluetoothA2DPOutput = outputs.contains(where: { $0.portType == .bluetoothA2DP || $0.portType == .bluetoothLE })
let shouldEnableVoiceProcessing = (session.category == .playAndRecord) && (session.mode == .voiceChat || session.mode == .videoChat) && !hasBluetoothA2DPOutput
do {
try audioInputNode.setVoiceProcessingEnabled(shouldEnableVoiceProcessing)
} catch {
if let nsError = error as NSError? {
print("语音处理设置失败(enable=\(shouldEnableVoiceProcessing)): domain=\(nsError.domain) code=\(nsError.code) desc=\(nsError.localizedDescription)")
} else {
print("语音处理设置失败(enable=\(shouldEnableVoiceProcessing)): \(error.localizedDescription)")
}
}
} }
// 检查输入格式是否符合Microsoft要求 // 等待并获取有效的输入格式:在路由/类别切换瞬间可能出现 0Hz,直接 installTap 会触发系统断言崩溃
let needsConversion = !isMicrosoftCompatibleFormat(inputFormat) guard let inputFormat = waitForValidInputFormat(inputNode: audioInputNode, maxAttempts: 10, retryIntervalSeconds: 0.03) else {
let currentSession = AVAudioSession.sharedInstance()
// 需要格式转换 let inputs = currentSession.currentRoute.inputs.map { $0.portType.rawValue }.joined(separator: ",")
guard let targetFormat = audioFormat, let outputs = currentSession.currentRoute.outputs.map { $0.portType.rawValue }.joined(separator: ",")
let converter = AVAudioConverter(from: inputFormat, to: targetFormat) else { throw NSError(domain: "AudioSetup", code: 3, userInfo: [
NSLocalizedDescriptionKey: "输入音频格式无效(sampleRate/channelCount),可能处于路由或类别切换中",
"sessionCategory": currentSession.category.rawValue,
"sessionMode": currentSession.mode.rawValue,
"sessionSampleRate": currentSession.sampleRate,
"routeInputs": inputs,
"routeOutputs": outputs
])
}
// 需要格式转换(输出固定 16k/16bit/mono PCM,匹配 PushStream 默认/常用配置)
guard let targetFormat = audioFormat else {
throw NSError(domain: "AudioSetup", code: 2, userInfo: [NSLocalizedDescriptionKey: "目标音频格式为空"])
}
guard let converter = AVAudioConverter(from: inputFormat, to: targetFormat) else {
throw NSError(domain: "AudioSetup", code: 2, userInfo: [NSLocalizedDescriptionKey: "音频格式转换器创建失败"]) throw NSError(domain: "AudioSetup", code: 2, userInfo: [NSLocalizedDescriptionKey: "音频格式转换器创建失败"])
} }
print("输入格式不符合Microsoft要求,进行格式转换") print("输入格式: \(inputFormat.sampleRate)Hz, \(inputFormat.commonFormat.rawValue), ch=\(inputFormat.channelCount), interleaved=\(inputFormat.isInterleaved)")
print("输入格式: \(inputFormat.sampleRate)Hz, \(inputFormat.commonFormat.rawValue)") print("目标格式: \(targetFormat.sampleRate)Hz, \(targetFormat.commonFormat.rawValue), ch=\(targetFormat.channelCount), interleaved=\(targetFormat.isInterleaved)")
print("目标格式: \(targetFormat.sampleRate)Hz, \(targetFormat.commonFormat.rawValue)")
// 添加tap进行格式转换
audioInputNode.installTap(onBus: 0, audioInputNode.installTap(onBus: 0,
bufferSize: 1024, bufferSize: 1024,
format: inputFormat) { [weak self] buffer, when in format: inputFormat) { [weak self] buffer, _ in
guard let strongSelf = self else { return } guard let strongSelf = self else { return }
// 修复:添加安全检查,避免强制解包崩溃
guard let convertedBuffer = AVAudioPCMBuffer( guard let convertedBuffer = AVAudioPCMBuffer(
pcmFormat: targetFormat, pcmFormat: targetFormat,
frameCapacity: AVAudioFrameCount( frameCapacity: AVAudioFrameCount(
@ -319,36 +361,59 @@ public class MicrophoneCapture: NSObject {
print("创建转换缓冲区失败") print("创建转换缓冲区失败")
return return
} }
var error: NSError? var error: NSError?
// 执行音频格式转换
let status = converter.convert( let status = converter.convert(
to: convertedBuffer, to: convertedBuffer,
error: &error, error: &error,
withInputFrom: { inNumPackets, outStatus in withInputFrom: { _, outStatus in
outStatus.pointee = .haveData outStatus.pointee = .haveData
return buffer return buffer
} }
) )
// 转换成功且无错误
if status == .haveData, error == nil { if status == .haveData, error == nil {
let data = strongSelf.audioBufferToData(convertedBuffer) let data = strongSelf.audioBufferToData(convertedBuffer)
strongSelf.audioDataHandler?(data) strongSelf.audioDataHandler?(data)
} else if let error = error { } else if let error = error {
print("音频格式转换失败: \(error.localizedDescription)") print("音频格式转换失败: domain=\(error.domain) code=\(error.code) desc=\(error.localizedDescription)")
} }
} }
audioEngine.prepare()
// 启动引擎 // 启动引擎
do { do {
try audioEngine.start() try audioEngine.start()
isCapturing = true
} catch { } catch {
print("音频引擎启动失败: \(error.localizedDescription)") if let nsError = error as NSError? {
print("音频引擎启动失败: domain=\(nsError.domain) code=\(nsError.code) desc=\(nsError.localizedDescription)")
} else {
print("音频引擎启动失败: \(error.localizedDescription)")
}
throw error throw error
} }
} }
/**
* 等待并获取有效的输入音频格式
* - 参数:
* - inputNode: 输入节点(通常为 AVAudioEngine.inputNode)
* - maxAttempts: 最大重试次数
* - retryIntervalSeconds: 每次重试间隔(秒)
* - 返回值: 若在重试窗口内拿到有效格式则返回 AVAudioFormat,否则返回 nil
* - 异常: 无(内部不抛出异常)
*/
private func waitForValidInputFormat(inputNode: AVAudioInputNode, maxAttempts: Int, retryIntervalSeconds: TimeInterval) -> AVAudioFormat? {
for attempt in 1...maxAttempts {
let format = inputNode.inputFormat(forBus: 0)
if format.sampleRate > 0, format.channelCount > 0 {
return format
}
print("输入格式无效,等待重试 attempt=\(attempt) sampleRate=\(format.sampleRate) ch=\(format.channelCount)")
Thread.sleep(forTimeInterval: retryIntervalSeconds)
}
return nil
}
/** /**
* 安全清理音频引擎资源 * 安全清理音频引擎资源
@ -382,6 +447,15 @@ public class MicrophoneCapture: NSObject {
count: dataLength count: dataLength
) )
} }
// Handle 32-bit integer format
else if let int32Data = buffer.int32ChannelData {
let int32Buffer = int32Data.pointee
var int16Array = [Int16](repeating: 0, count: frameLength)
for i in 0..<frameLength {
int16Array[i] = Int16(max(Int32(Int16.min), min(Int32(Int16.max), int32Buffer[i] >> 16)))
}
return Data(bytes: int16Array, count: dataLength)
}
// Handle float format // Handle float format
else if let floatData = buffer.floatChannelData { else if let floatData = buffer.floatChannelData {
var int16Array = [Int16](repeating: 0, count: frameLength) var int16Array = [Int16](repeating: 0, count: frameLength)
@ -442,35 +516,7 @@ public class MicrophoneCapture: NSObject {
// 修复:使用安全的清理方法 // 修复:使用安全的清理方法
cleanupAudioEngine() cleanupAudioEngine()
// 安全地处理音频会话 // 安全地处理音频会话:这里不主动 deactivate,避免与后续播报会话切换打架导致卡顿/失败
guard let audioSession = audioSession else { return }
do {
try audioSession.setActive(false)
if self.hasBluetoothDevices {
print("耳机")
// 设置音频会话参数
try audioSession.setCategory(.playback,
mode: .videoChat,
options: [
.allowBluetoothA2DP,
.allowBluetooth
]) // 添加音频优先级控制
} else {
print("手机")
try audioSession.setCategory(.playback,
mode: .videoChat,
options: [
.mixWithOthers,
.defaultToSpeaker])
}
try audioSession.overrideOutputAudioPort(.none)
try audioSession.setActive(true)
} catch {
print("音频会话配置失败: \(error.localizedDescription)")
}
} }
/** /**
@ -487,4 +533,4 @@ public class MicrophoneCapture: NSObject {
} }
} }

261
local_plugins/azure_speech/ios/azure_speech/Sources/tools/SimpleAudioReceiver.swift

@ -54,6 +54,7 @@ public class SimpleAudioReceiver: NSObject {
private var currentRoute: AudioOutputRoute? private var currentRoute: AudioOutputRoute?
//public var onAudioData: ((Data) -> Void)? //public var onAudioData: ((Data) -> Void)?
public var recordfile: RecordFile? public var recordfile: RecordFile?
private var currentRecognitionMode = "normal" // 当前识别模式:normal, ble_wakeup, phone_call, push_to_talk
public var isHeadphones = true public var isHeadphones = true
/** /**
@ -119,29 +120,176 @@ public class SimpleAudioReceiver: NSObject {
print("AudioStream音频配置设置完成") print("AudioStream音频配置设置完成")
} }
/** /**
* 开始音频输入 * 开始音频输入
* - Parameter mod: 识别模式,如 "normal", "phone_call", "push_to_talk", "ble_wakeup"。
* - Parameter audioSourceType: 音频源类型(.microphone 或 .external)。
* - Parameter audioDataCallback: 音频数据回调。
*/ */
public func startAudioRecord(audioSourceType: AudioSourceType = .microphone, audioDataCallback: AudioDataCallback?) { public func startAudioRecord(mod: String = "normal", audioSourceType: AudioSourceType = .microphone, audioDataCallback: AudioDataCallback?) {
_isWriting=true self.currentRecognitionMode = mod
self.audioSourceType = audioSourceType; self.audioSourceType = audioSourceType
self.audioDataCallback = audioDataCallback self.audioDataCallback = audioDataCallback
print("startAudioRecord=audioSourceType\(audioSourceType)") self.isHeadphones = detectHeadphonesOutput()
switch audioSourceType {
case .microphone: _isWriting = true
runMicrophoneCapture()
case .external: print("liwei---------- startAudioRecord=mod:\(mod), audioSourceType:\(audioSourceType), isHeadphones:\(isHeadphones)")
runExternalCapture()
do {
try configureAudioSession()
let isHybridExternalMode = audioSourceType == .external && (currentRecognitionMode == "phone_call" || currentRecognitionMode == "push_to_talk")
if audioSourceType == .microphone || isHybridExternalMode {
runMicrophoneCapture()
} else { // True external mode
runExternalCapture()
}
} catch {
print("配置或启动音频捕获失败: \(error.localizedDescription)")
}
}
/**
* 根据当前模式和设备状态配置音频会话
*/
private func configureAudioSession() throws {
let audioSession = self.audioSession
var desiredCategory: AVAudioSession.Category
var desiredMode: AVAudioSession.Mode
var desiredOptions: AVAudioSession.CategoryOptions
let isHybridExternalMode = audioSourceType == .external && (currentRecognitionMode == "phone_call" || currentRecognitionMode == "push_to_talk")
if audioSourceType == .microphone || isHybridExternalMode {
// --- 麦克风输入逻辑 (包括混合模式) ---
if currentRecognitionMode == "phone_call" {
if isHeadphones {
desiredCategory = .playAndRecord
desiredMode = .default
desiredOptions = [.allowBluetoothA2DP]
} else {
desiredCategory = .playAndRecord
desiredMode = .voiceChat
desiredOptions = [.allowBluetooth, .defaultToSpeaker]
}
} else { // normal, push_to_talk, ble_wakeup 等:仅录音,不需要同时播报
desiredCategory = .record
desiredMode = .default
desiredOptions = []
}
// 如果在真实通话中,强制使用更适合的录音设置
if isTelephonyActive() {
desiredCategory = .record
desiredMode = .voiceChat
desiredOptions = [.allowBluetooth]
}
} else if audioSourceType == .external {
// --- 纯外部源 (仅播放) 逻辑 ---
desiredCategory = .playback
desiredMode = .spokenAudio
desiredOptions = isHeadphones ? [.allowBluetoothA2DP, .mixWithOthers] : [.mixWithOthers]
} else {
// 兜底,理论上不会执行
desiredCategory = .playAndRecord
desiredMode = .default
desiredOptions = []
}
// 仅在需要时更新会话,避免不必要的中断
if audioSession.category != desiredCategory || audioSession.mode != desiredMode || audioSession.categoryOptions != desiredOptions {
try audioSession.setCategory(desiredCategory, mode: desiredMode, options: desiredOptions)
}
// 激活会话
try audioSession.setActive(true, options: .notifyOthersOnDeactivation)
// --- 路由设置 ---
if audioSourceType == .microphone || isHybridExternalMode {
// 为需要麦克风的场景设置输入输出
if desiredCategory == .playAndRecord {
if isHeadphones {
try audioSession.overrideOutputAudioPort(.none) // 允许耳机/蓝牙输出
} else {
try audioSession.overrideOutputAudioPort(.speaker) // 强制外放
}
}
// 设置首选输入设备
if let inputs = audioSession.availableInputs {
// 针对电话和按住说话模式,在连接耳机时强制使用内置麦克风
if isHeadphones && (currentRecognitionMode == "phone_call" || currentRecognitionMode == "push_to_talk") {
if let builtInMic = inputs.first(where: { $0.portType == .builtInMic }) {
try audioSession.setPreferredInput(builtInMic)
print("耳机模式下,为 \(currentRecognitionMode) 模式强制设置内置麦克风")
}
} else if isHybridExternalMode {
// 混合模式: 强制使用内置麦克风
if let builtInMic = inputs.first(where: { $0.portType == .builtInMic }) {
try audioSession.setPreferredInput(builtInMic)
print("混合外部模式: 首选输入已设置为内置麦克风")
}
} else if let btInput = inputs.first(where: { $0.portType == .bluetoothHFP }), isHeadphones {
// 其他耳机模式: 优先使用蓝牙耳机麦克风
try audioSession.setPreferredInput(btInput)
print("标准耳机模式: 首选输入已设置为蓝牙麦克风")
} else if let builtInMic = inputs.first(where: { $0.portType == .builtInMic }) {
// 兜底或无耳机: 使用内置麦克风
try audioSession.setPreferredInput(builtInMic)
print("兜底: 首选输入已设置为内置麦克风")
}
} }
} else { // 纯外部源
// .overrideOutputAudioPort 仅对 .playAndRecord 有效,.playback 会触发 -50
}
}
/**
* 检测是否存在蓝牙HFP输入设备
* 参数:无
* 返回值:存在返回true,否则返回false
* 异常:无
*/
private func hasBluetoothHFPInput() -> Bool {
let session = AVAudioSession.sharedInstance()
if let inputs = session.availableInputs {
return inputs.contains { $0.portType == .bluetoothHFP }
}
// 兜底使用当前路由的输入
return session.currentRoute.inputs.contains { $0.portType == .bluetoothHFP }
}
/**
* 检测是否存在耳机/蓝牙输出设备(用于决定是否启用A2DP而非扬声器)
* 参数:无
* 返回值:存在返回true,否则返回false
* 异常:无
*/
private func detectHeadphonesOutput() -> Bool {
let session = AVAudioSession.sharedInstance()
let outputs = session.currentRoute.outputs
return outputs.contains { out in
out.portType == .bluetoothA2DP ||
out.portType == .bluetoothHFP ||
out.portType == .bluetoothLE ||
out.portType == .headphones ||
out.portType == .headsetMic
}
} }
/** /**
* 开启音频写入线程 * 开启音频写入线程
* 修复:使用userInitiated QoS避免优先级反转 * 参数:无
* 返回值:无
* 异常:无
*/ */
public func startWriteThread() { public func startWriteThread() {
// 使用userInitiated QoS匹配音频录制线程的优先级 writeThread = DispatchQueue(label: "audio.stream.writer", qos: .default)
writeThread = DispatchQueue(label: "audio.stream.writer", qos: .userInitiated)
writeThread?.async { [weak self] in writeThread?.async { [weak self] in
guard let self = self else { return } guard let self = self else { return }
@ -198,9 +346,6 @@ public class SimpleAudioReceiver: NSObject {
// 确保正在写入状态 // 确保正在写入状态
guard self._isWriting else { return } guard self._isWriting else { return }
// 根据通话状态进行安全的音频会话配置
try safeConfigureAudioSessionForMic()
// 启动麦克风采集 // 启动麦克风采集
try micCapture.startCapture() try micCapture.startCapture()
@ -225,74 +370,8 @@ public class SimpleAudioReceiver: NSObject {
*/ */
private func runExternalCapture() { private func runExternalCapture() {
if audioSourceType == .external { if audioSourceType == .external {
// do { // 停止麦克风采集,确保不走通话音道
// try audioSession.setCategory(.playback, micCapture.stopCapture()
// mode: .videoChat,
// options: [ .mixWithOthers,.allowBluetoothA2DP // 允许蓝牙耳机,不占用hfp链路
// ]) // 添加音频优先级控制
// try audioSession.overrideOutputAudioPort(.none)
// try audioSession.setActive(true)
// } catch {
// print("设置音频会话失败: \(error.localizedDescription)")
// }
//pushAudioData(data: Data())
micCapture.stopCapture()
}
}
/**
* 在通话/非通话场景下安全配置音频会话用于麦克风采集
* - 通话中:使用 .record + .voiceChat + [.allowBluetooth],避免 A2DP 和强制外放,兼容电话路由
* - 非通话:使用 .playAndRecord + .videoChat + [.allowBluetooth, .mixWithOthers, .defaultToSpeaker]
* - 失败兜底:退回到最简单的 .record + .default,确保录音尽可能可用
*/
private func safeConfigureAudioSessionForMic() throws {
let audioSession = self.audioSession
do {
// 计算目标配置
let desiredCategory: AVAudioSession.Category
let desiredMode: AVAudioSession.Mode
let desiredOptions: AVAudioSession.CategoryOptions
if isTelephonyActive() {
// 通话中:录音优先,避免 A2DP(仅播放),改用 HFP(.allowBluetooth)
desiredCategory = .record
desiredMode = .voiceChat
desiredOptions = [.allowBluetooth]
} else {
if(isHeadphones){
// 非通话:保持原有逻辑,但使用 .allowBluetooth 以支持录音(而不是 A2DP)
desiredCategory = .playAndRecord
desiredMode = .videoChat
desiredOptions = [.allowBluetooth, .mixWithOthers, .defaultToSpeaker]
//desiredOptions = [.mixWithOthers, .defaultToSpeaker]
}else{
desiredCategory = .playAndRecord
desiredMode = .videoChat
desiredOptions = [.mixWithOthers, .defaultToSpeaker]
}
}
// 仅在类别或模式变化时才调用 setCategory,避免触发不必要的 categoryChange
let needsSetCategory = (audioSession.category != desiredCategory) || (audioSession.mode != desiredMode)
if needsSetCategory {
try audioSession.setCategory(desiredCategory,
mode: desiredMode,
options: desiredOptions)
} else {
// 类别与模式一致时,也尽量避免重复 setCategory;如需确保 options,可在确有必要时再设置
// 如果你确实需要校正 options,可取消下一行注释,但这可能再次触发 categoryChange
// try audioSession.setCategory(desiredCategory, mode: desiredMode, options: desiredOptions)
}
// 保持激活;这里重复 setActive(true) 一般安全,但如果你观察到重复激活也触发路由回调,可在外层控制频率
try audioSession.setActive(true, options: .notifyOthersOnDeactivation)
} catch {
// 兜底:尽量保证能开始录音
print("音频会话配置失败,尝试兜底配置: \(error)")
try audioSession.setCategory(.record, mode: .default, options: [])
try audioSession.setActive(true, options: .notifyOthersOnDeactivation)
} }
} }
@ -314,14 +393,19 @@ public class SimpleAudioReceiver: NSObject {
* 继续麦克风捕获 * 继续麦克风捕获
*/ */
public func resumeRecord() { public func resumeRecord() {
guard _isWriting==false else { return } guard _isWriting==false else { return }
_isWriting = true _isWriting = true
// 启动录音引擎 // 启动录音引擎
do { do {
print("resumeRecord") print("resumeRecord")
try micCapture.startCapture() try configureAudioSession()
try micCapture.startCapture()
} catch { } catch {
os_log("Failed to start audio engine: %@", log: log, type: .error, error.localizedDescription) if let nsError = error as NSError? {
os_log("Failed to start audio engine: domain=%{public}@ code=%{public}d desc=%{public}@", log: log, type: .error, nsError.domain, nsError.code, nsError.localizedDescription)
} else {
os_log("Failed to start audio engine: %@", log: log, type: .error, error.localizedDescription)
}
} }
} }
@ -541,6 +625,3 @@ private class LinkedBlockingQueue<T> {
} }
} }
} }

2
local_plugins/ble_service/android/src/main/kotlin/com/yunqiinnovation/ble_service/BleCommandSender.kt

@ -456,7 +456,7 @@ class BleCommandSender {
fun processDeviceNotification(data: ByteArray) { fun processDeviceNotification(data: ByteArray) {
Log.d( Log.d(
TAG, TAG,
"收到设备主动上报: ${data.joinToString(", ") { "0x${(it.toInt() and 0xFF).toString(16)}" }}" "ble指令----------- 收到设备主动上报: ${data.joinToString(", ") { "0x${(it.toInt() and 0xFF).toString(16)}" }}"
) )
if (data.size < 3) { if (data.size < 3) {
Log.e(TAG, "设备主动上报数据格式错误:数据长度过短") Log.e(TAG, "设备主动上报数据格式错误:数据长度过短")

4
local_plugins/ble_service/android/src/main/kotlin/com/yunqiinnovation/ble_service/BleService.kt

@ -953,7 +953,7 @@ object BleService {
override fun onCharacteristicChanged(g: BluetoothGatt, c: BluetoothGattCharacteristic) { override fun onCharacteristicChanged(g: BluetoothGatt, c: BluetoothGattCharacteristic) {
val data = c.value ?: return val data = c.value ?: return
Log.i(TAG, "ble指令----------- 收到通知特征数据(命令和控制 uuid=${c.uuid}")
// 根据特征UUID区分处理 // 根据特征UUID区分处理
when (c.uuid) { when (c.uuid) {
// 音频特征数据 // 音频特征数据
@ -973,7 +973,7 @@ object BleService {
} }
// 通知特征数据(命令和控制) // 通知特征数据(命令和控制)
BleConst.NOTIFY_CHAR_UUID -> { BleConst.NOTIFY_CHAR_UUID -> {
Log.i(TAG, "收到通知特征数据(命令和控制)") // Log.i(TAG, "ble指令----------- 收到通知特征数据(命令和控制)")
// 判断数据类型 // 判断数据类型
when { when {
// 设备响应 (0xBB开头) // 设备响应 (0xBB开头)

80
local_plugins/ble_service/ios/ble_service/Sources/ble_service/BleService.swift

@ -55,6 +55,10 @@ public class BleService: NSObject {
// 添加录制文件管理器 // 添加录制文件管理器
private var recordingFile: RecordingFile? private var recordingFile: RecordingFile?
private var pendingStopOpusRecording = false
private var stopOpusRecordingWorkItem: DispatchWorkItem?
private var lastOpusPacketReceivedAt: TimeInterval = 0
private var isScanning = false private var isScanning = false
private var scanTimer: Timer? private var scanTimer: Timer?
@ -471,6 +475,7 @@ private var cmdReplyType: UInt8 = 0
// 停止音频发送定时器 // 停止音频发送定时器
stopAudioSendTimer() stopAudioSendTimer()
stopOpusRecording()
// 清空音频缓冲区和队列 // 清空音频缓冲区和队列
audioQueueLock.lock() audioQueueLock.lock()
@ -743,9 +748,7 @@ private var cmdReplyType: UInt8 = 0
public func closeCodec() -> Bool { public func closeCodec() -> Bool {
os_log("关闭编解码并停止录制...", log: logger, type: .info) os_log("关闭编解码并停止录制...", log: logger, type: .info)
// 停止录制
recordingFile?.closeFile()
recordingFile = nil
// 停止音频发送定时器 // 停止音频发送定时器
stopAudioSendTimer() stopAudioSendTimer()
@ -756,6 +759,9 @@ private var cmdReplyType: UInt8 = 0
audioDataQueue.removeAll() audioDataQueue.removeAll()
audioQueueLock.unlock() audioQueueLock.unlock()
requestStopOpusRecording()
// 不停止Opus解码流,只发送命令通知设备关闭编解码 // 不停止Opus解码流,只发送命令通知设备关闭编解码
let paramData = Data([BleConst.CODEC_CONTROL_CLOSE, BleConst.AUDIO_CHANNEL_RIGHT]) let paramData = Data([BleConst.CODEC_CONTROL_CLOSE, BleConst.AUDIO_CHANNEL_RIGHT])
return sendCommand(BleConst.CMD_CONTROL_CODEC, data: paramData) return sendCommand(BleConst.CMD_CONTROL_CODEC, data: paramData)
@ -828,7 +834,7 @@ private var cmdReplyType: UInt8 = 0
private func sendCommand(_ cmd: UInt8, data: Data = Data()) -> Bool { private func sendCommand(_ cmd: UInt8, data: Data = Data()) -> Bool {
guard let characteristic = writeCharacteristic, guard let characteristic = writeCharacteristic,
connectionState == BleConst.STATE_CONNECTED else { connectionState == BleConst.STATE_CONNECTED else {
os_log("发送命令失败: 设备未连接", log: logger, type: .error) os_log("liwei---- 发送命令失败: 设备未连接", log: logger, type: .error)
return false return false
} }
@ -856,7 +862,7 @@ private var cmdReplyType: UInt8 = 0
if !data.isEmpty { if !data.isEmpty {
dataHexString = data.map { String(format: "0x%02X", $0) }.joined(separator: ", ") dataHexString = data.map { String(format: "0x%02X", $0) }.joined(separator: ", ")
} }
os_log("发送命令: CMD=0x%02X, 长度=%d, 数据=[%@], CRC=0x%02X", os_log("liwei---- 发送命令: CMD=0x%02X, 长度=%d, 数据=[%@], CRC=0x%02X",
log: logger, type: .debug, cmd, cmdLength, dataHexString, crc) log: logger, type: .debug, cmd, cmdLength, dataHexString, crc)
@ -875,7 +881,7 @@ private var cmdReplyType: UInt8 = 0
return true return true
} catch { } catch {
os_log("发送命令异常: %@", log: logger, type: .error, error.localizedDescription) os_log("liwei---- 发送命令异常: %@", log: logger, type: .error, error.localizedDescription)
return false return false
} }
} }
@ -906,8 +912,67 @@ private var cmdReplyType: UInt8 = 0
return return
} }
lastOpusPacketReceivedAt = Date().timeIntervalSince1970
ensureOpusRecordingStarted(fileName: "耳机端").saveAudioData(data)
// 使用Opus处理器处理音频数据 // 使用Opus处理器处理音频数据
opusProcessor?.processAudioData(data) opusProcessor?.processAudioData(data)
if pendingStopOpusRecording {
requestStopOpusRecording()
}
}
/// 确保 Opus 文件已创建并处于可写状态
/// - Parameter fileName: 文件名前缀(不含时间戳与扩展名)
/// - Returns: 可用于追加写入的 RecordingFile 实例
/// - Throws: 不抛出;若创建文件失败将导致后续写入无效
private func ensureOpusRecordingStarted(fileName: String) -> RecordingFile {
let recorder = recordingFile ?? RecordingFile()
if recorder.fileName.isEmpty {
recorder.fileName = fileName
}
if !recorder.isRecording() {
recorder.createFile()
}
recordingFile = recorder
return recorder
}
/// 停止并释放 Opus 落盘资源
/// - Returns: 无
/// - Throws: 不抛出
private func stopOpusRecording() {
stopOpusRecordingWorkItem?.cancel()
stopOpusRecordingWorkItem = nil
pendingStopOpusRecording = false
recordingFile?.closeFile()
recordingFile = nil
}
/// 请求在短暂静默后停止 Opus 落盘,避免截断尾包导致文件不完整
/// - Parameters:
/// - graceSeconds: 等待静默的时间窗口
/// - Returns: 无
/// - Throws: 不抛出
private func requestStopOpusRecording(graceSeconds: TimeInterval = 1.0) {
pendingStopOpusRecording = true
stopOpusRecordingWorkItem?.cancel()
let workItem = DispatchWorkItem { [weak self] in
guard let self = self else { return }
let now = Date().timeIntervalSince1970
if now - self.lastOpusPacketReceivedAt >= graceSeconds {
self.stopOpusRecording()
} else {
self.requestStopOpusRecording(graceSeconds: graceSeconds)
}
}
stopOpusRecordingWorkItem = workItem
DispatchQueue.main.asyncAfter(deadline: .now() + graceSeconds, execute: workItem)
} }
/// 处理接收到的响应数据 /// 处理接收到的响应数据
@ -943,7 +1008,7 @@ private var cmdReplyType: UInt8 = 0
// 记录响应详情 // 记录响应详情
var hexString = data.map { String(format: "0x%02X", $0) }.joined(separator: ", ") var hexString = data.map { String(format: "0x%02X", $0) }.joined(separator: ", ")
os_log("收到设备响应: %@", log: logger, type: .debug, hexString) // os_log("liwei---- 收到设备响应: %@", log: logger, type: .debug, hexString)
// 检查数据长度(与 Android 统一失败判定保持一致) // 检查数据长度(与 Android 统一失败判定保持一致)
if data.count != Int(length) + 4 { // 帧头 + 命令 + 长度 + 数据 + CRC if data.count != Int(length) + 4 { // 帧头 + 命令 + 长度 + 数据 + CRC
@ -1786,6 +1851,7 @@ extension BleService: CBCentralManagerDelegate {
} }
updateConnectionState(BleConst.STATE_DISCONNECTED) updateConnectionState(BleConst.STATE_DISCONNECTED)
stopOpusRecording()
// 只有在非主动断开的情况下才重新连接 // 只有在非主动断开的情况下才重新连接
if !isManualDisconnect { if !isManualDisconnect {

2
local_plugins/ble_service/ios/ble_service/Sources/ble_service/RecordingFile.swift

@ -78,7 +78,7 @@ class RecordingFile {
do { do {
try self.fileHandle?.write(contentsOf: buffer) try self.fileHandle?.write(contentsOf: buffer)
self.totalBytesWritten += buffer.count self.totalBytesWritten += buffer.count
os_log("写入音频数据: %d bytes, 总计: %d bytes", log: self.logger, type: .debug, buffer.count, self.totalBytesWritten) // os_log("写入音频数据: %d bytes, 总计: %d bytes", log: self.logger, type: .debug, buffer.count, self.totalBytesWritten)
} catch { } catch {
os_log("写入音频数据失败: %@", log: self.logger, type: .error, error.localizedDescription) os_log("写入音频数据失败: %@", log: self.logger, type: .error, error.localizedDescription)
} }

163
local_plugins/ble_service/ios/ble_service/Sources/ble_service/SwiftBleServicePlugin.swift

@ -2,6 +2,7 @@ import Flutter
import UIKit import UIKit
import CoreBluetooth import CoreBluetooth
import os import os
import opus
@available(iOS 13.0, *) @available(iOS 13.0, *)
@objc(BleServicePlugin) @objc(BleServicePlugin)
@ -172,6 +173,40 @@ public class SwiftBleServicePlugin: NSObject, FlutterPlugin {
// 移除这个方法调用,让iOS系统自然处理权限 // 移除这个方法调用,让iOS系统自然处理权限
result(FlutterMethodNotImplemented) result(FlutterMethodNotImplemented)
case "decodeOpusFile":
guard let args = call.arguments as? [String: Any],
let inPath = args["inPath"] as? String,
let outPath = args["outPath"] as? String else {
result(FlutterError(code: "INVALID_ARGS", message: "缺少必要参数", details: nil))
return
}
let hasHeader = (args["hasHeader"] as? Bool) ?? false
let channel = Int32((args["channel"] as? Int) ?? 1)
let sampleRate = Int32((args["sampleRate"] as? Int) ?? 16000)
let packetSize = (args["packetSize"] as? Int) ?? 40
DispatchQueue.global(qos: .userInitiated).async {
do {
let output = try OpusFileDecoder.decodeFileToPcm(
inPath: inPath,
outPath: outPath,
hasHeader: hasHeader,
channel: channel,
sampleRate: sampleRate,
packetSize: packetSize
)
DispatchQueue.main.async {
result(output)
}
} catch {
DispatchQueue.main.async {
result(FlutterError(code: "DECODE_ERROR", message: "解码失败: \(error)", details: nil))
}
}
}
default: default:
result(FlutterMethodNotImplemented) result(FlutterMethodNotImplemented)
} }
@ -187,6 +222,134 @@ public class SwiftBleServicePlugin: NSObject, FlutterPlugin {
} }
} }
@available(iOS 13.0, *)
private enum OpusFileDecoder {
/// 将 Opus 文件(按固定 packetSize 切分的裸包流)解码为 PCM,并写入 outPath
/// - Parameters:
/// - inPath: 输入 Opus 文件路径
/// - outPath: 输出 PCM 文件路径
/// - hasHeader: 是否包含协议头(当前仅支持 false)
/// - channel: 通道数(1 或 2)
/// - sampleRate: 采样率
/// - packetSize: 固定包长(若不确定,会与 40/80 一起自动尝试)
/// - Returns: 输出文件路径(等于 outPath)
/// - Throws: 参数不合法、文件读写失败、或解码输出为空时抛出
static func decodeFileToPcm(
inPath: String,
outPath: String,
hasHeader: Bool,
channel: Int32,
sampleRate: Int32,
packetSize: Int
) throws -> String {
if hasHeader {
throw NSError(domain: "ble_service", code: -10, userInfo: [NSLocalizedDescriptionKey: "暂不支持带协议头的 Opus 文件解码"])
}
if channel != 1 && channel != 2 {
throw NSError(domain: "ble_service", code: -11, userInfo: [NSLocalizedDescriptionKey: "channel 仅支持 1 或 2"])
}
if sampleRate <= 0 {
throw NSError(domain: "ble_service", code: -12, userInfo: [NSLocalizedDescriptionKey: "sampleRate 必须大于 0"])
}
let inputData = try Data(contentsOf: URL(fileURLWithPath: inPath))
if inputData.isEmpty {
throw NSError(domain: "ble_service", code: -13, userInfo: [NSLocalizedDescriptionKey: "输入文件为空"])
}
let candidates = Array(Set([packetSize, 40, 80].filter { $0 > 0 }))
var best = Data()
var lastError: Error?
for size in candidates {
do {
let pcm = try decodeFixedPacketSize(
inputData: inputData,
packetSize: size,
sampleRate: sampleRate,
channels: channel
)
if pcm.count > best.count {
best = pcm
}
} catch {
lastError = error
}
}
if best.isEmpty {
throw lastError ?? NSError(domain: "ble_service", code: -14, userInfo: [NSLocalizedDescriptionKey: "解码输出为空"])
}
try best.write(to: URL(fileURLWithPath: outPath), options: .atomic)
return outPath
}
/// 按固定 packetSize 切分输入数据并逐包解码为 PCM
/// - Parameters:
/// - inputData: Opus 裸包流(无容器)
/// - packetSize: 固定包长
/// - sampleRate: 采样率
/// - channels: 通道数
/// - Returns: PCM 数据(16-bit little-endian)
/// - Throws: 解码器创建失败或输出为空时抛出
private static func decodeFixedPacketSize(
inputData: Data,
packetSize: Int,
sampleRate: Int32,
channels: Int32
) throws -> Data {
if packetSize <= 0 {
throw NSError(domain: "ble_service", code: -20, userInfo: [NSLocalizedDescriptionKey: "packetSize 必须大于 0"])
}
var errorCode: Int32 = 0
guard let decoder = opus_decoder_create(sampleRate, channels, &errorCode), errorCode == OPUS_OK else {
throw NSError(domain: "ble_service", code: Int(errorCode), userInfo: [NSLocalizedDescriptionKey: "Opus 解码器创建失败: \(errorCode)"])
}
defer { opus_decoder_destroy(decoder) }
let maxFrameSize: Int32 = 5760
var pcmBuffer = [opus_int16](repeating: 0, count: Int(maxFrameSize) * Int(channels))
var output = Data()
var offset = 0
while offset + packetSize <= inputData.count {
let packet = inputData.subdata(in: offset..<(offset + packetSize))
offset += packetSize
let decodedSamples: opus_int32 = packet.withUnsafeBytes { packetPtr in
guard let packetBase = packetPtr.bindMemory(to: UInt8.self).baseAddress else {
return opus_int32(-1)
}
return opus_decode(
decoder,
packetBase,
opus_int32(packet.count),
&pcmBuffer,
maxFrameSize,
0
)
}
if decodedSamples <= 0 {
continue
}
let bytesCount = Int(decodedSamples) * Int(channels) * MemoryLayout<opus_int16>.size
if bytesCount > 0 {
output.append(pcmBuffer.withUnsafeBytes { Data($0.prefix(bytesCount)) })
}
}
if output.isEmpty {
throw NSError(domain: "ble_service", code: -21, userInfo: [NSLocalizedDescriptionKey: "解码输出为空(packetSize=\(packetSize))"])
}
return output
}
}
// MARK: - BleService.Callback // MARK: - BleService.Callback
@available(iOS 13.0, *) @available(iOS 13.0, *)
extension SwiftBleServicePlugin: BleService.Callback { extension SwiftBleServicePlugin: BleService.Callback {

91
local_plugins/chat_api/ios/chat_api/Sources/chat_api/MCPClient.swift

@ -36,9 +36,11 @@ public class MCPSubClient {
// 使用官方MCP Swift SDK // 使用官方MCP Swift SDK
private var mcpClient: Client? private var mcpClient: Client?
private var transport: CustomSseClientTransport?
private var tools: [Tool] = [] private var tools: [Tool] = []
private var toolMaps: [[String: Any]] = [] private var toolMaps: [[String: Any]] = []
private var isConnectedFlag = false private var isConnectedFlag = false
private let toolsLock = NSLock()
// 重连配置 // 重连配置
private let maxRetryAttempts = 3 private let maxRetryAttempts = 3
@ -51,6 +53,10 @@ public class MCPSubClient {
// 工具调用配置 - 添加这两行 // 工具调用配置 - 添加这两行
private let maxToolCallRetries = 3 private let maxToolCallRetries = 3
private let toolCallTimeout: TimeInterval = 30.0 private let toolCallTimeout: TimeInterval = 30.0
private let keepAliveInterval: TimeInterval = 25.0
private let keepAliveTimeout: TimeInterval = 8.0
private var keepAliveTask: Task<Void, Never>?
// 连接状态锁 // 连接状态锁
private let connectionLock = NSLock() private let connectionLock = NSLock()
@ -89,6 +95,7 @@ public class MCPSubClient {
} }
} }
) )
self.transport = transport
// 3. 连接到服务器(添加超时) // 3. 连接到服务器(添加超时)
try await withTimeout(seconds: 10) { try await withTimeout(seconds: 10) {
@ -106,6 +113,7 @@ public class MCPSubClient {
isConnectedFlag = true isConnectedFlag = true
retryCount = 0 // 重置重试计数 retryCount = 0 // 重置重试计数
currentReconnectDelay = initialReconnectDelay // 重置延迟 currentReconnectDelay = initialReconnectDelay // 重置延迟
startKeepAlive()
return true return true
@ -116,6 +124,9 @@ public class MCPSubClient {
} }
private func processTools(_ toolList: [Tool]) { private func processTools(_ toolList: [Tool]) {
toolsLock.lock()
defer { toolsLock.unlock() }
tools.removeAll() tools.removeAll()
toolMaps.removeAll() toolMaps.removeAll()
@ -173,6 +184,13 @@ public class MCPSubClient {
/// 检查连接状态 /// 检查连接状态
/// 检查连接状态 /// 检查连接状态
public func checkConnection() async -> Bool { public func checkConnection() async -> Bool {
if isConnectedFlag, let transport {
let isActive = await transport.isConnectionActive()
if !isActive {
isConnectedFlag = false
}
}
if !isConnectedFlag { if !isConnectedFlag {
return await connect() return await connect()
} }
@ -339,14 +357,25 @@ public class MCPSubClient {
} }
public func getToolMaps() -> [[String: Any]] { public func getToolMaps() -> [[String: Any]] {
toolsLock.lock()
defer { toolsLock.unlock() }
return toolMaps return toolMaps
} }
public func containsTool(name: String) -> Bool { public func containsTool(name: String) -> Bool {
toolsLock.lock()
defer { toolsLock.unlock() }
return tools.contains { $0.name == name } return tools.contains { $0.name == name }
} }
public func close() async { public func close() async {
keepAliveTask?.cancel()
keepAliveTask = nil
reconnectTask?.cancel()
reconnectTask = nil
if let client = mcpClient { if let client = mcpClient {
await client.disconnect() await client.disconnect()
} }
@ -359,9 +388,9 @@ public class MCPSubClient {
if isConnectedFlag { if isConnectedFlag {
print("[\(serverId)] 连接已断开,更新状态") print("[\(serverId)] 连接已断开,更新状态")
isConnectedFlag = false isConnectedFlag = false
keepAliveTask?.cancel()
// 启动重连逻辑 keepAliveTask = nil
// startReconnection() startReconnection()
} }
} }
} }
@ -369,7 +398,9 @@ public class MCPSubClient {
/// 启动重连 /// 启动重连
private func startReconnection() { private func startReconnection() {
// 取消之前的重连任务 // 取消之前的重连任务
reconnectTask?.cancel() if reconnectTask != nil {
return
}
reconnectTask = Task { [weak self] in reconnectTask = Task { [weak self] in
await self?.performReconnection() await self?.performReconnection()
@ -378,6 +409,8 @@ public class MCPSubClient {
/// 执行重连 /// 执行重连
private func performReconnection() async { private func performReconnection() async {
defer { reconnectTask = nil }
while retryCount < maxRetryAttempts && !Task.isCancelled { while retryCount < maxRetryAttempts && !Task.isCancelled {
retryCount += 1 retryCount += 1
@ -403,6 +436,8 @@ public class MCPSubClient {
private func handleConnectionError(_ error: Error) async { private func handleConnectionError(_ error: Error) async {
print("[\(serverId)] 连接错误: \(error.localizedDescription)") print("[\(serverId)] 连接错误: \(error.localizedDescription)")
isConnectedFlag = false isConnectedFlag = false
keepAliveTask?.cancel()
keepAliveTask = nil
// 如果是MCP特定错误,可以进行特殊处理 // 如果是MCP特定错误,可以进行特殊处理
if let mcpError = error as? MCPError { if let mcpError = error as? MCPError {
@ -411,9 +446,49 @@ public class MCPSubClient {
} }
// 如果还有重试机会,启动重连 // 如果还有重试机会,启动重连
// if retryCount < maxRetryAttempts { startReconnection()
// startReconnection() }
// }
private func startKeepAlive() {
keepAliveTask?.cancel()
keepAliveTask = Task { [weak self] in
guard let self else { return }
while !Task.isCancelled {
try? await Task.sleep(for: .seconds(self.keepAliveInterval))
if Task.isCancelled {
return
}
if !self.isConnectedFlag {
continue
}
if let transport = self.transport {
let isActive = await transport.isConnectionActive()
if !isActive {
await self.handleConnectionLost()
continue
}
}
guard let client = self.mcpClient else {
await self.handleConnectionError(MCPError.internalError("MCP client not initialized"))
continue
}
do {
let (toolList, _) = try await withTimeout(seconds: self.keepAliveTimeout) {
try await client.listTools()
}
self.processTools(toolList)
} catch {
await self.handleConnectionError(error)
}
}
}
} }
/// 获取连接状态 /// 获取连接状态
@ -665,7 +740,7 @@ public class MCPClient {
} }
public func isConnected() -> Bool { public func isConnected() -> Bool {
return !subClients.isEmpty && subClients.values.allSatisfy { $0.getConnectionStatus() } return !subClients.isEmpty && subClients.values.contains { $0.getConnectionStatus() }
} }
public func disconnectAll() async { public func disconnectAll() async {

549
local_plugins/chat_storage/ios/chat_storage/Sources/chat_storage/ChatStorageHelper.swift

@ -15,6 +15,8 @@ public class ChatStorageHelper {
private let dbPath: String private let dbPath: String
private let logger = OSLog(subsystem: "com.yunqiinnovation.chat_storage", category: "ChatStorageHelper") private let logger = OSLog(subsystem: "com.yunqiinnovation.chat_storage", category: "ChatStorageHelper")
private let dbQueue = DispatchQueue(label: "com.yunqiinnovation.chat_storage.dbQueue")
private init() { private init() {
// 获取文档目录路径(不变) // 获取文档目录路径(不变)
let fileURL = try! FileManager.default let fileURL = try! FileManager.default
@ -23,39 +25,43 @@ public class ChatStorageHelper {
dbPath = fileURL.path dbPath = fileURL.path
// 打开数据库(不变) dbQueue.sync {
if sqlite3_open(dbPath, &db) != SQLITE_OK { // 打开数据库(不变)
let errmsg = String(cString: sqlite3_errmsg(db)!) if sqlite3_open(dbPath, &db) != SQLITE_OK {
os_log("无法打开数据库: %{public}@", log: logger, type: .error, errmsg) let errmsg = String(cString: sqlite3_errmsg(db)!)
return os_log("无法打开数据库: %{public}@", log: logger, type: .error, errmsg)
} return
}
// 关键修改:新增 UNIQUE (agent_id, session_id, sender) 联合唯一索引
let createTableString = """ // 关键修改:新增 UNIQUE (agent_id, session_id, sender) 联合唯一索引
CREATE TABLE IF NOT EXISTS messages ( let createTableString = """
id INTEGER PRIMARY KEY AUTOINCREMENT, CREATE TABLE IF NOT EXISTS messages (
agent_id TEXT NOT NULL, id INTEGER PRIMARY KEY AUTOINCREMENT,
session_id TEXT NOT NULL, agent_id TEXT NOT NULL,
message TEXT NOT NULL, session_id TEXT NOT NULL,
sender TEXT NOT NULL, message TEXT NOT NULL,
metadata TEXT, sender TEXT NOT NULL,
created_at INTEGER NOT NULL, metadata TEXT,
UNIQUE (agent_id, session_id, sender) ON CONFLICT REPLACE created_at INTEGER NOT NULL,
); UNIQUE (agent_id, session_id, sender) ON CONFLICT REPLACE
CREATE INDEX IF NOT EXISTS idx_agent_id ON messages (agent_id); );
CREATE INDEX IF NOT EXISTS idx_session_id ON messages (session_id); CREATE INDEX IF NOT EXISTS idx_agent_id ON messages (agent_id);
CREATE INDEX IF NOT EXISTS idx_created_at ON messages (created_at); CREATE INDEX IF NOT EXISTS idx_session_id ON messages (session_id);
""" CREATE INDEX IF NOT EXISTS idx_created_at ON messages (created_at);
"""
if sqlite3_exec(db, createTableString, nil, nil, nil) != SQLITE_OK {
let errmsg = String(cString: sqlite3_errmsg(db)!) if sqlite3_exec(db, createTableString, nil, nil, nil) != SQLITE_OK {
os_log("创建表失败: %{public}@", log: logger, type: .error, errmsg) let errmsg = String(cString: sqlite3_errmsg(db)!)
os_log("创建表失败: %{public}@", log: logger, type: .error, errmsg)
}
} }
} }
deinit { deinit {
if db != nil { dbQueue.sync {
sqlite3_close(db) if db != nil {
sqlite3_close(db)
}
} }
} }
@ -68,48 +74,48 @@ public class ChatStorageHelper {
* @return 插入的消息ID,失败则返回-1 * @return 插入的消息ID,失败则返回-1
*/ */
public func saveMessage(agentId: String, sessionId: String, message: String, sender: String, metadata: String?) -> Int64 { public func saveMessage(agentId: String, sessionId: String, message: String, sender: String, metadata: String?) -> Int64 {
// 关键修改:用 INSERT OR REPLACE 替换 INSERT,支持冲突时更新 var result: Int64 = -1
let insertStatementString = """ dbQueue.sync {
INSERT OR REPLACE INTO messages // 关键修改:用 INSERT OR REPLACE 替换 INSERT,支持冲突时更新
(agent_id, session_id, message, sender, metadata, created_at) let insertStatementString = """
VALUES (?, ?, ?, ?, ?, ?); INSERT OR REPLACE INTO messages
""" (agent_id, session_id, message, sender, metadata, created_at)
var insertStatement: OpaquePointer? VALUES (?, ?, ?, ?, ?, ?);
"""
if sqlite3_prepare_v2(db, insertStatementString, -1, &insertStatement, nil) == SQLITE_OK { var insertStatement: OpaquePointer?
// 绑定参数(逻辑不变,确保三个唯一字段正确传入)
sqlite3_bind_text(insertStatement, 1, (agentId as NSString).utf8String, -1, nil)
sqlite3_bind_text(insertStatement, 2, (sessionId as NSString).utf8String, -1, nil)
sqlite3_bind_text(insertStatement, 3, (message as NSString).utf8String, -1, nil)
sqlite3_bind_text(insertStatement, 4, (sender as NSString).utf8String, -1, nil)
if let metadata = metadata {
sqlite3_bind_text(insertStatement, 5, (metadata as NSString).utf8String, -1, nil)
} else {
sqlite3_bind_null(insertStatement, 5)
}
let currentTime = Int(Date().timeIntervalSince1970) if sqlite3_prepare_v2(db, insertStatementString, -1, &insertStatement, nil) == SQLITE_OK {
sqlite3_bind_int(insertStatement, 6, Int32(currentTime)) // 绑定参数(逻辑不变,确保三个唯一字段正确传入)
sqlite3_bind_text(insertStatement, 1, (agentId as NSString).utf8String, -1, nil)
// 执行语句(冲突时会自动替换,返回新的 rowid) sqlite3_bind_text(insertStatement, 2, (sessionId as NSString).utf8String, -1, nil)
if sqlite3_step(insertStatement) == SQLITE_DONE { sqlite3_bind_text(insertStatement, 3, (message as NSString).utf8String, -1, nil)
let id = sqlite3_last_insert_rowid(db) // 替换后返回新的 id(原 id 会被删除) sqlite3_bind_text(insertStatement, 4, (sender as NSString).utf8String, -1, nil)
if let metadata = metadata {
sqlite3_bind_text(insertStatement, 5, (metadata as NSString).utf8String, -1, nil)
} else {
sqlite3_bind_null(insertStatement, 5)
}
let currentTime = Int(Date().timeIntervalSince1970)
sqlite3_bind_int(insertStatement, 6, Int32(currentTime))
// 执行语句(冲突时会自动替换,返回新的 rowid)
if sqlite3_step(insertStatement) == SQLITE_DONE {
let id = sqlite3_last_insert_rowid(db) // 替换后返回新的 id(原 id 会被删除)
os_log("消息保存成功(新增/更新),id: %{public}lld", log: logger, type: .info, id)
result = id
} else {
let errmsg = String(cString: sqlite3_errmsg(db)!)
os_log("插入/更新消息失败: %{public}@", log: logger, type: .error, errmsg)
}
sqlite3_finalize(insertStatement) sqlite3_finalize(insertStatement)
os_log("消息保存成功(新增/更新),id: %{public}lld", log: logger, type: .info, id)
return id
} else { } else {
let errmsg = String(cString: sqlite3_errmsg(db)!) let errmsg = String(cString: sqlite3_errmsg(db)!)
os_log("插入/更新消息失败: %{public}@", log: logger, type: .error, errmsg) os_log("插入/更新消息语句准备失败: %{public}@", log: logger, type: .error, errmsg)
} }
sqlite3_finalize(insertStatement)
} else {
let errmsg = String(cString: sqlite3_errmsg(db)!)
os_log("插入/更新消息语句准备失败: %{public}@", log: logger, type: .error, errmsg)
} }
return result
return -1
} }
/** /**
@ -120,122 +126,106 @@ public class ChatStorageHelper {
* @return 消息列表的JSON字符串 * @return 消息列表的JSON字符串
*/ */
public func getMessages(agentId: String, page: Int, pageSize: Int) -> String { public func getMessages(agentId: String, page: Int, pageSize: Int) -> String {
let offset = (page - 1) * pageSize var result = "{\"messages\":[],\"page\":\(page),\"pageSize\":\(pageSize),\"totalCount\":0,\"totalPages\":0}"
var messagesArray: [[String: Any]] = [] dbQueue.sync {
let offset = (page - 1) * pageSize
// 首先获取总记录数 var messagesArray: [[String: Any]] = []
let countQuery = "SELECT COUNT(*) FROM messages WHERE agent_id = ?"
var countStatement: OpaquePointer? // 首先获取总记录数
var totalCount = 0 let countQuery = "SELECT COUNT(*) FROM messages WHERE agent_id = ?"
var countStatement: OpaquePointer?
if sqlite3_prepare_v2(db, countQuery, -1, &countStatement, nil) == SQLITE_OK { var totalCount = 0
sqlite3_bind_text(countStatement, 1, (agentId as NSString).utf8String, -1, nil)
if sqlite3_step(countStatement) == SQLITE_ROW { if sqlite3_prepare_v2(db, countQuery, -1, &countStatement, nil) == SQLITE_OK {
totalCount = Int(sqlite3_column_int(countStatement, 0)) sqlite3_bind_text(countStatement, 1, (agentId as NSString).utf8String, -1, nil)
if sqlite3_step(countStatement) == SQLITE_ROW {
totalCount = Int(sqlite3_column_int(countStatement, 0))
}
sqlite3_finalize(countStatement)
} }
sqlite3_finalize(countStatement)
}
// 计算总页数
let totalPages = totalCount == 0 ? 0 : (totalCount + pageSize - 1) / pageSize
let queryString = """
SELECT id, session_id, message, sender, metadata, created_at
FROM messages
WHERE agent_id = ?
ORDER BY created_at DESC
LIMIT ? OFFSET ?
"""
var queryStatement: OpaquePointer?
if sqlite3_prepare_v2(db, queryString, -1, &queryStatement, nil) == SQLITE_OK {
sqlite3_bind_text(queryStatement, 1, (agentId as NSString).utf8String, -1, nil)
sqlite3_bind_int(queryStatement, 2, Int32(pageSize))
sqlite3_bind_int(queryStatement, 3, Int32(offset))
while sqlite3_step(queryStatement) == SQLITE_ROW { // 计算总页数
let id = sqlite3_column_int(queryStatement, 0) let totalPages = totalCount == 0 ? 0 : (totalCount + pageSize - 1) / pageSize
let sessionIdPtr = sqlite3_column_text(queryStatement, 1) let queryString = """
let sessionId = sessionIdPtr != nil ? String(cString: sessionIdPtr!) : "" SELECT id, session_id, message, sender, metadata, created_at
FROM messages
let messagePtr = sqlite3_column_text(queryStatement, 2) WHERE agent_id = ?
let message = messagePtr != nil ? String(cString: messagePtr!) : "" ORDER BY created_at DESC
LIMIT ? OFFSET ?
let senderPtr = sqlite3_column_text(queryStatement, 3) """
let sender = senderPtr != nil ? String(cString: senderPtr!) : ""
var queryStatement: OpaquePointer?
let metadataPtr = sqlite3_column_text(queryStatement, 4)
let metadata = metadataPtr != nil ? String(cString: metadataPtr!) : nil if sqlite3_prepare_v2(db, queryString, -1, &queryStatement, nil) == SQLITE_OK {
sqlite3_bind_text(queryStatement, 1, (agentId as NSString).utf8String, -1, nil)
sqlite3_bind_int(queryStatement, 2, Int32(pageSize))
sqlite3_bind_int(queryStatement, 3, Int32(offset))
let createdAt = sqlite3_column_int(queryStatement, 5) while sqlite3_step(queryStatement) == SQLITE_ROW {
let id = sqlite3_column_int(queryStatement, 0)
let sessionIdPtr = sqlite3_column_text(queryStatement, 1)
let sessionId = sessionIdPtr != nil ? String(cString: sessionIdPtr!) : ""
let messagePtr = sqlite3_column_text(queryStatement, 2)
let message = messagePtr != nil ? String(cString: messagePtr!) : ""
let senderPtr = sqlite3_column_text(queryStatement, 3)
let sender = senderPtr != nil ? String(cString: senderPtr!) : ""
let metadataPtr = sqlite3_column_text(queryStatement, 4)
let metadata = metadataPtr != nil ? String(cString: metadataPtr!) : nil
let createdAt = sqlite3_column_int(queryStatement, 5)
// 将时间戳转换为ISO 8601格式的字符串
let date = Date(timeIntervalSince1970: TimeInterval(createdAt))
let formatter = ISO8601DateFormatter()
let timestamp = formatter.string(from: date)
os_log("getMessages: 读取历史记录: %{public}@", log: logger, type: .info, timestamp)
var messageDict: [String: Any] = [
"id": id,
"agentId": agentId, // 添加agentId字段
"sessionId": sessionId, // 添加sessionId字段
"message": message,
"sender": sender,
"timestamp": timestamp // 使用timestamp而不是created_at
]
if let metadata = metadata {
messageDict["metadata"] = metadata
}
messagesArray.append(messageDict)
}
// 将时间戳转换为ISO 8601格式的字符串 sqlite3_finalize(queryStatement)
let date = Date(timeIntervalSince1970: TimeInterval(createdAt))
let formatter = ISO8601DateFormatter()
let timestamp = formatter.string(from: date)
os_log("getMessages: 读取历史记录: %{public}@", log: logger, type: .info, timestamp)
var messageDict: [String: Any] = [ // 构建与Android版本一致的返回格式
"id": id, let resultDict: [String: Any] = [
"agentId": agentId, // 添加agentId字段 "messages": messagesArray,
"sessionId": sessionId, // 添加sessionId字段 "page": page,
"message": message, "pageSize": pageSize,
"sender": sender, "totalCount": totalCount,
"timestamp": timestamp // 使用timestamp而不是created_at "totalPages": totalPages
] ]
os_log("getMessages: 读取历史记录: %{public}@", log: logger, type: .error, messagesArray)
if let metadata = metadata { do {
messageDict["metadata"] = metadata let jsonData = try JSONSerialization.data(withJSONObject: resultDict, options: [])
} if let jsonString = String(data: jsonData, encoding: .utf8) {
result = jsonString
messagesArray.append(messageDict) }
} } catch {
os_log("getMessages: JSON转换失败: %{public}@", log: logger, type: .error, error.localizedDescription)
sqlite3_finalize(queryStatement)
// 构建与Android版本一致的返回格式
let result: [String: Any] = [
"messages": messagesArray,
"page": page,
"pageSize": pageSize,
"totalCount": totalCount,
"totalPages": totalPages
]
os_log("getMessages: 读取历史记录: %{public}@", log: logger, type: .error, messagesArray)
do {
let jsonData = try JSONSerialization.data(withJSONObject: result, options: [])
if let jsonString = String(data: jsonData, encoding: .utf8) {
return jsonString
} }
} catch { } else {
os_log("getMessages: JSON转换失败: %{public}@", log: logger, type: .error, error.localizedDescription) let errmsg = String(cString: sqlite3_errmsg(db)!)
} os_log("getMessages: SQL准备失败: %{public}@", log: logger, type: .error, errmsg)
} else {
let errmsg = String(cString: sqlite3_errmsg(db)!)
os_log("getMessages: SQL准备失败: %{public}@", log: logger, type: .error, errmsg)
}
// 返回空结果,但保持格式一致
let emptyResult: [String: Any] = [
"messages": [],
"page": page,
"pageSize": pageSize,
"totalCount": 0,
"totalPages": 0
]
do {
let jsonData = try JSONSerialization.data(withJSONObject: emptyResult, options: [])
if let jsonString = String(data: jsonData, encoding: .utf8) {
return jsonString
} }
} catch {
os_log("getMessages: 空结果JSON转换失败: %{public}@", log: logger, type: .error, error.localizedDescription)
} }
return result
return "{\"messages\":[],\"page\":\(page),\"pageSize\":\(pageSize),\"totalCount\":0,\"totalPages\":0}"
} }
/** /**
@ -245,56 +235,54 @@ public class ChatStorageHelper {
* @return 是否删除成功 * @return 是否删除成功
*/ */
public func deleteMessages(agentId: String?, messageIds: [Int]?) -> Bool { public func deleteMessages(agentId: String?, messageIds: [Int]?) -> Bool {
if let agentId = agentId { var result = false
let deleteString = "DELETE FROM messages WHERE agent_id = ?;" dbQueue.sync {
var deleteStatement: OpaquePointer? if let agentId = agentId {
let deleteString = "DELETE FROM messages WHERE agent_id = ?;"
if sqlite3_prepare_v2(db, deleteString, -1, &deleteStatement, nil) == SQLITE_OK { var deleteStatement: OpaquePointer?
sqlite3_bind_text(deleteStatement, 1, (agentId as NSString).utf8String, -1, nil)
if sqlite3_step(deleteStatement) == SQLITE_DONE { if sqlite3_prepare_v2(db, deleteString, -1, &deleteStatement, nil) == SQLITE_OK {
sqlite3_bind_text(deleteStatement, 1, (agentId as NSString).utf8String, -1, nil)
if sqlite3_step(deleteStatement) == SQLITE_DONE {
result = true
} else {
let errmsg = String(cString: sqlite3_errmsg(db)!)
os_log("删除会话消息失败: %{public}@", log: logger, type: .error, errmsg)
}
sqlite3_finalize(deleteStatement) sqlite3_finalize(deleteStatement)
return true
} else { } else {
let errmsg = String(cString: sqlite3_errmsg(db)!) let errmsg = String(cString: sqlite3_errmsg(db)!)
os_log("删除会话消息失败: %{public}@", log: logger, type: .error, errmsg) os_log("删除会话消息语句准备失败: %{public}@", log: logger, type: .error, errmsg)
}
sqlite3_finalize(deleteStatement)
} else {
let errmsg = String(cString: sqlite3_errmsg(db)!)
os_log("删除会话消息语句准备失败: %{public}@", log: logger, type: .error, errmsg)
}
} else if let messageIds = messageIds, !messageIds.isEmpty {
// 构建参数占位符
let placeholders = String(repeating: "?,", count: messageIds.count).dropLast()
let deleteString = "DELETE FROM messages WHERE id IN (\(placeholders));"
var deleteStatement: OpaquePointer?
if sqlite3_prepare_v2(db, deleteString, -1, &deleteStatement, nil) == SQLITE_OK {
for (index, id) in messageIds.enumerated() {
sqlite3_bind_int(deleteStatement, Int32(index + 1), Int32(id))
} }
} else if let messageIds = messageIds, !messageIds.isEmpty {
// 构建参数占位符
let placeholders = String(repeating: "?,", count: messageIds.count).dropLast()
let deleteString = "DELETE FROM messages WHERE id IN (\(placeholders));"
var deleteStatement: OpaquePointer?
if sqlite3_step(deleteStatement) == SQLITE_DONE { if sqlite3_prepare_v2(db, deleteString, -1, &deleteStatement, nil) == SQLITE_OK {
for (index, id) in messageIds.enumerated() {
sqlite3_bind_int(deleteStatement, Int32(index + 1), Int32(id))
}
if sqlite3_step(deleteStatement) == SQLITE_DONE {
result = true
} else {
let errmsg = String(cString: sqlite3_errmsg(db)!)
os_log("删除特定消息失败: %{public}@", log: logger, type: .error, errmsg)
}
sqlite3_finalize(deleteStatement) sqlite3_finalize(deleteStatement)
return true
} else { } else {
let errmsg = String(cString: sqlite3_errmsg(db)!) let errmsg = String(cString: sqlite3_errmsg(db)!)
os_log("删除特定消息失败: %{public}@", log: logger, type: .error, errmsg) os_log("删除特定消息语句准备失败: %{public}@", log: logger, type: .error, errmsg)
} }
sqlite3_finalize(deleteStatement)
} else { } else {
let errmsg = String(cString: sqlite3_errmsg(db)!) os_log("删除消息参数无效 - agentId和messageIds都为空", log: logger, type: .error)
os_log("删除特定消息语句准备失败: %{public}@", log: logger, type: .error, errmsg)
} }
} else {
os_log("删除消息参数无效 - agentId和messageIds都为空", log: logger, type: .error)
return false
} }
return result
return false
} }
/** /**
@ -302,15 +290,18 @@ public class ChatStorageHelper {
* @return 是否清空成功 * @return 是否清空成功
*/ */
public func clearDatabase() -> Bool { public func clearDatabase() -> Bool {
let deleteString = "DELETE FROM messages;" var result = false
dbQueue.sync {
if sqlite3_exec(db, deleteString, nil, nil, nil) == SQLITE_OK { let deleteString = "DELETE FROM messages;"
return true
} else { if sqlite3_exec(db, deleteString, nil, nil, nil) == SQLITE_OK {
let errmsg = String(cString: sqlite3_errmsg(db)!) result = true
os_log("清空数据库失败: %{public}@", log: logger, type: .error, errmsg) } else {
return false let errmsg = String(cString: sqlite3_errmsg(db)!)
os_log("清空数据库失败: %{public}@", log: logger, type: .error, errmsg)
}
} }
return result
} }
/** /**
@ -321,75 +312,75 @@ public class ChatStorageHelper {
*/ */
public func getRecentMessages(agentId: String, limit: Int) -> [[String: Any]] { public func getRecentMessages(agentId: String, limit: Int) -> [[String: Any]] {
var messages: [[String: Any]] = [] var messages: [[String: Any]] = []
dbQueue.sync {
// 首先检查数据库中是否有该会话的消息 // 首先检查数据库中是否有该会话的消息
let countQuery = "SELECT COUNT(*) FROM messages WHERE agent_id = ?" let countQuery = "SELECT COUNT(*) FROM messages WHERE agent_id = ?"
var countStatement: OpaquePointer? var countStatement: OpaquePointer?
if sqlite3_prepare_v2(db, countQuery, -1, &countStatement, nil) == SQLITE_OK { if sqlite3_prepare_v2(db, countQuery, -1, &countStatement, nil) == SQLITE_OK {
sqlite3_bind_text(countStatement, 1, (agentId as NSString).utf8String, -1, nil) sqlite3_bind_text(countStatement, 1, (agentId as NSString).utf8String, -1, nil)
if sqlite3_step(countStatement) == SQLITE_ROW { if sqlite3_step(countStatement) == SQLITE_ROW {
_ = Int(sqlite3_column_int(countStatement, 0)) _ = Int(sqlite3_column_int(countStatement, 0))
}
sqlite3_finalize(countStatement)
} }
sqlite3_finalize(countStatement)
}
// 构建查询语句 - 按时间倒序获取最近的N条,然后在结果中再按时间正序
let queryString = """
SELECT * FROM (
SELECT id, session_id, message, sender, metadata, created_at
FROM messages
WHERE agent_id = ?
ORDER BY created_at DESC
LIMIT ?
) tmp ORDER BY created_at ASC
"""
var queryStatement: OpaquePointer?
if sqlite3_prepare_v2(db, queryString, -1, &queryStatement, nil) == SQLITE_OK {
sqlite3_bind_text(queryStatement, 1, (agentId as NSString).utf8String, -1, nil)
sqlite3_bind_int(queryStatement, 2, Int32(limit))
while sqlite3_step(queryStatement) == SQLITE_ROW { // 构建查询语句 - 按时间倒序获取最近的N条,然后在结果中再按时间正序
let id = sqlite3_column_int(queryStatement, 0) let queryString = """
SELECT * FROM (
let sessionIdPtr = sqlite3_column_text(queryStatement, 1) SELECT id, session_id, message, sender, metadata, created_at
let sessionId = sessionIdPtr != nil ? String(cString: sessionIdPtr!) : "" FROM messages
WHERE agent_id = ?
ORDER BY created_at DESC
let messagePtr = sqlite3_column_text(queryStatement, 2) LIMIT ?
let message = messagePtr != nil ? String(cString: messagePtr!) : "" ) tmp ORDER BY created_at ASC
"""
let senderPtr = sqlite3_column_text(queryStatement, 3)
let sender = senderPtr != nil ? String(cString: senderPtr!) : "" var queryStatement: OpaquePointer?
let metadataPtr = sqlite3_column_text(queryStatement, 4) if sqlite3_prepare_v2(db, queryString, -1, &queryStatement, nil) == SQLITE_OK {
let metadata = metadataPtr != nil ? String(cString: metadataPtr!) : nil sqlite3_bind_text(queryStatement, 1, (agentId as NSString).utf8String, -1, nil)
sqlite3_bind_int(queryStatement, 2, Int32(limit))
let createdAt = sqlite3_column_int(queryStatement, 5)
var messageDict: [String: Any] = [
"id": id,
"sessionId":sessionId,
"message": message,
"sender": sender,
"created_at": createdAt
]
if let metadata = metadata { while sqlite3_step(queryStatement) == SQLITE_ROW {
messageDict["metadata"] = metadata let id = sqlite3_column_int(queryStatement, 0)
let sessionIdPtr = sqlite3_column_text(queryStatement, 1)
let sessionId = sessionIdPtr != nil ? String(cString: sessionIdPtr!) : ""
let messagePtr = sqlite3_column_text(queryStatement, 2)
let message = messagePtr != nil ? String(cString: messagePtr!) : ""
let senderPtr = sqlite3_column_text(queryStatement, 3)
let sender = senderPtr != nil ? String(cString: senderPtr!) : ""
let metadataPtr = sqlite3_column_text(queryStatement, 4)
let metadata = metadataPtr != nil ? String(cString: metadataPtr!) : nil
let createdAt = sqlite3_column_int(queryStatement, 5)
var messageDict: [String: Any] = [
"id": id,
"sessionId":sessionId,
"message": message,
"sender": sender,
"created_at": createdAt
]
if let metadata = metadata {
messageDict["metadata"] = metadata
}
messages.append(messageDict)
} }
messages.append(messageDict) sqlite3_finalize(queryStatement)
} else {
let errmsg = String(cString: sqlite3_errmsg(db)!)
os_log("查询最近消息失败: %{public}@", log: logger, type: .error, errmsg)
} }
sqlite3_finalize(queryStatement)
} else {
let errmsg = String(cString: sqlite3_errmsg(db)!)
os_log("查询最近消息失败: %{public}@", log: logger, type: .error, errmsg)
} }
return messages return messages
} }
} }

22
local_plugins/jl_opus/ios/jl_opus/Package.swift

@ -0,0 +1,22 @@
// swift-tools-version: 5.9
import PackageDescription
let package = Package(
name: "jl_opus",
platforms: [.iOS("16.0")],
products: [
.library(name: "jl-opus", targets: ["jl_opus"])
],
dependencies: [
.package(path: "../../../ble_service/ios/ble_service")
],
targets: [
.target(
name: "jl_opus",
dependencies: [
.product(name: "ble-service", package: "ble_service")
],
path: "Sources/jl_opus"
)
]
)

450
local_plugins/jl_opus/ios/jl_opus/Sources/jl_opus/JlOpusPlugin.swift

@ -0,0 +1,450 @@
import Flutter
import Foundation
import opus
@objc(JlOpusPlugin)
public class JlOpusPlugin: NSObject, FlutterPlugin, FlutterStreamHandler {
private var eventSink: FlutterEventSink?
private var isInitialized: Bool = false
private var isDecodingStream: Bool = false
private var streamHasHeader: Bool = false
private var streamChannels: Int32 = 1
private var streamSampleRate: Int32 = 16000
private var streamDecoder: OpaquePointer?
private let decodeQueue = DispatchQueue(label: "com.yunqiinnovation.jl_opus.decode", qos: .userInitiated)
/// Flutter 插件注册入口
/// - Parameter registrar: FlutterPluginRegistrar
/// - Returns: 无
/// - Throws: 不抛出
public static func register(with registrar: FlutterPluginRegistrar) {
let methodChannel = FlutterMethodChannel(
name: "com.yunqiinnovation.jl_opus",
binaryMessenger: registrar.messenger()
)
let eventChannel = FlutterEventChannel(
name: "com.yunqiinnovation.jl_opus/events",
binaryMessenger: registrar.messenger()
)
let instance = JlOpusPlugin()
registrar.addMethodCallDelegate(instance, channel: methodChannel)
eventChannel.setStreamHandler(instance)
}
/// 处理 Flutter 侧方法调用
/// - Parameters:
/// - call: FlutterMethodCall
/// - result: FlutterResult
/// - Returns: 无
/// - Throws: 不抛出(错误通过 FlutterError 返回)
public func handle(_ call: FlutterMethodCall, result: @escaping FlutterResult) {
switch call.method {
case "initOpusDecoder":
isInitialized = true
result(true)
case "decodeOpusFile":
decodeOpusFile(arguments: call.arguments, result: result)
case "startDecodeStream":
startDecodeStream(arguments: call.arguments, result: result)
case "stopDecodeStream":
stopDecodeStream(result: result)
case "writeAudioStream":
writeAudioStream(arguments: call.arguments, result: result)
case "isDecoding":
result(isDecodingStream)
case "dispose":
dispose(result: result)
default:
result(FlutterMethodNotImplemented)
}
}
/// 事件通道开始监听
/// - Parameters:
/// - arguments: Flutter 传入参数
/// - eventSink: 事件回调
/// - Returns: FlutterError?(通常为 nil)
/// - Throws: 不抛出
public func onListen(withArguments arguments: Any?, eventSink events: @escaping FlutterEventSink) -> FlutterError? {
eventSink = events
return nil
}
/// 事件通道取消监听
/// - Parameter arguments: Flutter 传入参数
/// - Returns: FlutterError?(通常为 nil)
/// - Throws: 不抛出
public func onCancel(withArguments arguments: Any?) -> FlutterError? {
eventSink = nil
return nil
}
/// 解码 Opus 文件为 PCM 文件(16-bit little-endian 原始 PCM)
/// - Parameters:
/// - arguments: Flutter 传入参数(需包含 inPath/outPath/hasHeader/channel/sampleRate/packetSize)
/// - result: FlutterResult,成功返回 outPath,失败返回 FlutterError
/// - Returns: 无
/// - Throws: 不直接抛出;内部错误会通过 FlutterError 与事件通道 onError 返回
private func decodeOpusFile(arguments: Any?, result: @escaping FlutterResult) {
guard isInitialized else {
result(FlutterError(code: "NOT_INITIALIZED", message: "Opus解码器未初始化", details: nil))
return
}
guard let args = arguments as? [String: Any],
let inPath = args["inPath"] as? String,
let outPath = args["outPath"] as? String else {
result(FlutterError(code: "INVALID_ARGS", message: "输入或输出路径不能为空", details: nil))
return
}
let hasHeader = (args["hasHeader"] as? Bool) ?? false
let channel = Int32((args["channel"] as? Int) ?? 1)
let sampleRate = Int32((args["sampleRate"] as? Int) ?? 16000)
let packetSize = (args["packetSize"] as? Int) ?? 40
sendEvent([
"event": "onStart",
"type": "file"
])
decodeQueue.async { [weak self] in
guard let self = self else { return }
do {
let inputUrl = URL(fileURLWithPath: inPath)
let outputUrl = URL(fileURLWithPath: outPath)
let inputData = try Data(contentsOf: inputUrl)
if inputData.isEmpty {
throw NSError(domain: "jl_opus", code: -2, userInfo: [NSLocalizedDescriptionKey: "输入文件为空"])
}
let packets = try self.splitPackets(inputData: inputData, hasHeader: hasHeader, packetSize: packetSize)
let pcmData = try self.decodePacketsToPcm(packets: packets, sampleRate: sampleRate, channels: channel)
try self.ensureParentDirectoryExists(for: outputUrl)
try pcmData.write(to: outputUrl, options: .atomic)
self.sendEvent([
"event": "onComplete",
"type": "file",
"filePath": outPath
])
DispatchQueue.main.async {
result(outPath)
}
} catch {
let message = error.localizedDescription
self.sendEvent([
"event": "onError",
"type": "file",
"code": -1,
"message": message
])
DispatchQueue.main.async {
result(FlutterError(code: "DECODE_ERROR", message: "解码出错: \(message)", details: nil))
}
}
}
}
/// 开始 Opus 数据流解码
/// - Parameters:
/// - arguments: Flutter 传入参数(需包含 hasHeader/channel/sampleRate)
/// - result: FlutterResult,成功返回 true,失败返回 FlutterError
/// - Returns: 无
/// - Throws: 不抛出;失败通过 FlutterError/事件 onError 返回
private func startDecodeStream(arguments: Any?, result: @escaping FlutterResult) {
guard isInitialized else {
result(FlutterError(code: "NOT_INITIALIZED", message: "Opus解码器未初始化", details: nil))
return
}
guard let args = arguments as? [String: Any] else {
result(FlutterError(code: "INVALID_ARGS", message: "缺少解码参数", details: nil))
return
}
let hasHeader = (args["hasHeader"] as? Bool) ?? false
let channel = Int32((args["channel"] as? Int) ?? 1)
let sampleRate = Int32((args["sampleRate"] as? Int) ?? 16000)
stopDecodeStreamInternal()
var error: Int32 = 0
guard let decoder = opus_decoder_create(sampleRate, channel, &error), error == OPUS_OK else {
sendEvent([
"event": "onError",
"type": "stream",
"code": Int(error),
"message": "Opus解码器创建失败: \(error)"
])
result(FlutterError(code: "DECODE_STREAM_ERROR", message: "数据流解码器创建失败", details: nil))
return
}
streamHasHeader = hasHeader
streamChannels = channel
streamSampleRate = sampleRate
streamDecoder = decoder
isDecodingStream = true
sendEvent([
"event": "onStart",
"type": "stream"
])
result(true)
}
/// 停止 Opus 数据流解码
/// - Parameter result: FlutterResult,已停止返回 true,未处于解码返回 false
/// - Returns: 无
/// - Throws: 不抛出
private func stopDecodeStream(result: @escaping FlutterResult) {
if isDecodingStream {
stopDecodeStreamInternal()
result(true)
} else {
result(false)
}
}
/// 停止并释放当前数据流解码器(内部方法)
/// - Returns: 无
/// - Throws: 不抛出
private func stopDecodeStreamInternal() {
if let decoder = streamDecoder {
opus_decoder_destroy(decoder)
}
streamDecoder = nil
isDecodingStream = false
}
/// 写入 Opus 数据到解码流,并通过事件通道返回 PCM 数据
/// - Parameters:
/// - arguments: Flutter 传入参数(需包含 data)
/// - result: FlutterResult,成功返回 true;不在解码状态返回 FlutterError
/// - Returns: 无
/// - Throws: 不直接抛出;解码错误通过事件 onError 返回
private func writeAudioStream(arguments: Any?, result: @escaping FlutterResult) {
guard isDecodingStream, let decoder = streamDecoder else {
result(FlutterError(code: "NOT_DECODING", message: "当前没有处于解码状态", details: nil))
return
}
guard let args = arguments as? [String: Any],
let data = args["data"] as? FlutterStandardTypedData else {
result(FlutterError(code: "INVALID_ARGS", message: "音频数据不能为空", details: nil))
return
}
let opusData = data.data
if opusData.isEmpty {
result(true)
return
}
let hasHeader = streamHasHeader
let channels = streamChannels
let sampleRate = streamSampleRate
decodeQueue.async { [weak self] in
guard let self = self else { return }
do {
let packets: [Data]
if hasHeader {
packets = try self.splitPacketsByLengthPrefix(inputData: opusData)
} else {
packets = [opusData]
}
for packet in packets {
let pcm = try self.decodeSinglePacketToPcm(decoder: decoder, packet: packet, sampleRate: sampleRate, channels: channels)
if !pcm.isEmpty {
self.sendEvent([
"event": "onDecodeStream",
"data": FlutterStandardTypedData(bytes: pcm)
])
}
}
} catch {
let message = error.localizedDescription
self.sendEvent([
"event": "onError",
"type": "stream",
"code": -1,
"message": message
])
}
}
result(true)
}
/// 释放插件资源
/// - Parameter result: FlutterResult
/// - Returns: 无
/// - Throws: 不抛出
private func dispose(result: @escaping FlutterResult) {
stopDecodeStreamInternal()
isInitialized = false
result(true)
}
/// 拆分输入文件数据为 Opus 包列表
/// - Parameters:
/// - inputData: 输入文件二进制数据
/// - hasHeader: 是否包含 2 字节小端长度前缀
/// - packetSize: 无协议头时按固定长度切包
/// - Returns: Opus 包数组
/// - Throws: 参数错误或解析失败时抛出 NSError
private func splitPackets(inputData: Data, hasHeader: Bool, packetSize: Int) throws -> [Data] {
if hasHeader {
return try splitPacketsByLengthPrefix(inputData: inputData)
}
guard packetSize > 0 else {
throw NSError(domain: "jl_opus", code: -5, userInfo: [NSLocalizedDescriptionKey: "packetSize 必须大于 0"])
}
var packets: [Data] = []
var offset = 0
while offset < inputData.count {
let end = min(offset + packetSize, inputData.count)
packets.append(inputData.subdata(in: offset..<end))
offset = end
}
if packets.isEmpty {
throw NSError(domain: "jl_opus", code: -6, userInfo: [NSLocalizedDescriptionKey: "未解析到任何Opus数据包"])
}
return packets
}
/// 按 2 字节小端长度前缀拆分 Opus 包
/// - Parameter inputData: 输入数据
/// - Returns: Opus 包数组
/// - Throws: 协议头解析失败时抛出 NSError
private func splitPacketsByLengthPrefix(inputData: Data) throws -> [Data] {
var packets: [Data] = []
var cursor = 0
while cursor + 2 <= inputData.count {
let lo = Int(inputData[cursor])
let hi = Int(inputData[cursor + 1]) << 8
let length = lo | hi
cursor += 2
if length <= 0 || cursor + length > inputData.count {
throw NSError(domain: "jl_opus", code: -3, userInfo: [NSLocalizedDescriptionKey: "协议头解析失败"])
}
packets.append(inputData.subdata(in: cursor..<(cursor + length)))
cursor += length
}
if packets.isEmpty {
throw NSError(domain: "jl_opus", code: -4, userInfo: [NSLocalizedDescriptionKey: "未解析到任何Opus数据包"])
}
return packets
}
/// 将 Opus 包列表解码为 PCM 数据
/// - Parameters:
/// - packets: Opus 包数组
/// - sampleRate: 采样率
/// - channels: 通道数
/// - Returns: PCM 数据(16-bit little-endian)
/// - Throws: 解码器创建失败或任意包解码失败时抛出 NSError
private func decodePacketsToPcm(packets: [Data], sampleRate: Int32, channels: Int32) throws -> Data {
var error: Int32 = 0
guard let decoder = opus_decoder_create(sampleRate, channels, &error), error == OPUS_OK else {
throw NSError(domain: "jl_opus", code: Int(error), userInfo: [NSLocalizedDescriptionKey: "Opus解码器创建失败: \(error)"])
}
defer { opus_decoder_destroy(decoder) }
var pcmOutput = Data()
for packet in packets {
let pcmChunk = try decodeSinglePacketToPcm(decoder: decoder, packet: packet, sampleRate: sampleRate, channels: channels)
pcmOutput.append(pcmChunk)
}
return pcmOutput
}
/// 解码单个 Opus 包为 PCM
/// - Parameters:
/// - decoder: Opus 解码器指针
/// - packet: 单个 Opus 包
/// - sampleRate: 采样率
/// - channels: 通道数
/// - Returns: PCM 数据(16-bit little-endian)
/// - Throws: 解码失败时抛出 NSError
private func decodeSinglePacketToPcm(decoder: OpaquePointer, packet: Data, sampleRate: Int32, channels: Int32) throws -> Data {
let maxSamplesPerChannel = Int(sampleRate * 120 / 1000)
var pcmBuffer = [opus_int16](repeating: 0, count: maxSamplesPerChannel * Int(channels))
let decodedSamples: Int32 = packet.withUnsafeBytes { packetPtr in
guard let inPtr = packetPtr.bindMemory(to: UInt8.self).baseAddress else {
return OPUS_BAD_ARG
}
return opus_decode(
decoder,
inPtr,
Int32(packet.count),
&pcmBuffer,
Int32(maxSamplesPerChannel),
0
)
}
if decodedSamples < 0 {
throw NSError(domain: "jl_opus", code: Int(decodedSamples), userInfo: [NSLocalizedDescriptionKey: "Opus解码失败: \(decodedSamples)"])
}
let sampleCount = Int(decodedSamples) * Int(channels)
return pcmBuffer.withUnsafeBytes { rawPtr in
Data(rawPtr.bindMemory(to: UInt8.self).prefix(sampleCount * MemoryLayout<opus_int16>.size))
}
}
/// 确保输出路径的父目录存在
/// - Parameter url: 输出文件 URL
/// - Returns: 无
/// - Throws: 创建目录失败时抛出错误
private func ensureParentDirectoryExists(for url: URL) throws {
let dir = url.deletingLastPathComponent()
var isDirectory: ObjCBool = false
if FileManager.default.fileExists(atPath: dir.path, isDirectory: &isDirectory), isDirectory.boolValue {
return
}
try FileManager.default.createDirectory(at: dir, withIntermediateDirectories: true, attributes: nil)
}
/// 通过事件通道向 Flutter 发送事件
/// - Parameter event: 事件字典
/// - Returns: 无
/// - Throws: 不抛出
private func sendEvent(_ event: [String: Any]) {
guard let sink = eventSink else { return }
DispatchQueue.main.async {
sink(event)
}
}
}

4
local_plugins/music_service/ios/music_service/Sources/music_service/MusicService.swift

@ -41,10 +41,10 @@ public class MusicService: NSObject {
private override init() { private override init() {
super.init() super.init()
setupRemoteCommandCenter() setupRemoteCommandCenter()
setupAudioSession() // setupAudioSession()
setupNotificationObservers() setupNotificationObservers()
// 监听通话状态(需主线程队列) // 监听通话状态(需主线程队列)
// callObserver.setDelegate(self, queue: DispatchQueue.main) //callObserver.setDelegate(self, queue: DispatchQueue.main)
} }
public func initialize(serviceUrl: String, userToken: String) { public func initialize(serviceUrl: String, userToken: String) {

Loading…
Cancel
Save