Browse Source

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

weicu
liwei1dao 9 months ago
parent
commit
400d3d5f09
  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. 140
      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. 1114
      local_plugins/azure_speech/ios/azure_speech/Sources/azure_speech/AzureTtsHelper.swift
  19. 186
      local_plugins/azure_speech/ios/azure_speech/Sources/tools/MicrophoneCapture.swift
  20. 247
      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. 69
      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(
bool audioSourceType, {
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(
bool audioSourceType, {
bool isRemoveFirstPunctuation = true,
String mode = "normal",
}) async {
// 防抖判断:短时间内重复调用直接拦截
final now = DateTime.now();
@ -229,6 +230,7 @@ class AzureAsrService extends GetxService implements AsrService {
await _channel.invokeMethod('startContinuousRecognition', {
'audioSourceType': audioSourceType,
'isRemoveFirstPunctuation': isRemoveFirstPunctuation,
'mode': mode,
});
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
Future<bool> setVoiceFlocking(String speakerProfileId) async {
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
Future<bool> setVoiceFlocking(String speakerProfileId) async {
return true;

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

@ -865,6 +865,7 @@ class VolcanoAsrApiService implements AsrService {
Future<bool> startContinuousRecognition(
bool audioSourceType, {
bool isRemoveFirstPunctuation = true,
String mode = "normal",
}) {
// TODO: implement startContinuousRecognition
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(
bool audioSourceType, {
bool isRemoveFirstPunctuation = true,
String mode = "normal",
}) {
// TODO: implement startContinuousRecognition
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
Future<bool> setVoiceFlocking(String speakerProfileId) async {
return true;

3
lib/data/services/tts_service.dart

@ -19,6 +19,9 @@ abstract class TtsService {
/// 设置语音
Future<bool> setVoice(String voiceName);
/// 设置TTS模式
Future<bool> setTtsMode(String mod);
/// 设置语音复刻
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 {
Logger.i(TAG, '🎤 开始按住说话');
print('liwei--------- 🎤 开始按住说话');
// 只有在语音模式且按住说话模式下才能按住说话
if (isTextInputMode.value) return;
@ -1898,7 +1898,7 @@ class AgentController extends GetxController with WidgetsBindingObserver {
// 结束按住说话
void endPushToTalk() async {
Logger.i(TAG, '🎤 结束按住说话');
print('liwei----------🎤 结束按住说话');
// 如果不在按住说话状态,直接返回
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:flutter/foundation.dart';
import 'package:flutter/material.dart';
import 'package:flutter/services.dart';
import 'package:get/get.dart';
import 'package:path_provider/path_provider.dart';
import 'package:just_audio/just_audio.dart';
@ -42,8 +43,10 @@ class OpusTestController extends GetxController {
String? tempPcmPath;
// 杰理OPUS解码器
late JlOpus jlOpus;
JlOpus? jlOpus;
StreamSubscription? _eventSubscription;
static const MethodChannel _bleServiceChannel =
MethodChannel('com.yunqiinnovation.ble_service');
@override
void onInit() {
@ -56,16 +59,20 @@ class OpusTestController extends GetxController {
void onClose() {
player.dispose();
_eventSubscription?.cancel();
jlOpus.dispose();
jlOpus?.dispose();
super.onClose();
}
// 初始化OPUS解码器
Future<void> _initOpusDecoder() async {
if (!Platform.isAndroid) {
return;
}
jlOpus = JlOpus();
// 监听解码器事件
_eventSubscription = jlOpus.eventStream.listen((event) {
_eventSubscription = jlOpus!.eventStream.listen((event) {
switch (event.event) {
case 'onStart':
statusMessage.value = '开始${event.type == "file" ? "文件" : "流"}解码...';
@ -86,12 +93,29 @@ class OpusTestController extends GetxController {
});
// 初始化OPUS解码器
final initialized = await jlOpus.initOpusDecoder();
final initialized = await jlOpus!.initOpusDecoder();
if (!initialized) {
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 {
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) {
statusMessage.value = '解码失败';
isDecoding.value = false;

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

@ -344,6 +344,15 @@ class TranslationController extends GetxController with WidgetsBindingObserver {
} 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()) {
// 当前模式不支持悬浮窗,自动禁用
@ -436,7 +445,8 @@ class TranslationController extends GetxController with WidgetsBindingObserver {
}
await _asrService.startContinuousRecognition(
_audioSourceType,
isRemoveFirstPunctuation: false,
isRemoveFirstPunctuation: true,
mode: 'push_to_talk',
); //开启识别
_startAsrActiveTracking(); //开启识别活动跟踪
if (isRecording.value) {
@ -946,7 +956,17 @@ class TranslationController extends GetxController with WidgetsBindingObserver {
!bleManager.isCodecActive) {
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();
}
} catch (e) {
@ -1045,8 +1065,20 @@ class TranslationController extends GetxController with WidgetsBindingObserver {
}
/// 启动语音识别
Future<void> _startAsrService() async {
await _asrService.startContinuousRecognition(_audioSourceType);
Future<void> _startAsrService(String mode) async {
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 val maxIdleSeconds = 10 // 最大空闲秒数(默认值)
private val maxIdleSeconds = 5 // 最大空闲秒数(默认值)
// 语音识别模式
private var currentRecognitionMode = "normal" // 当前识别模式:normal, ble_wakeup, phone_call, push_to_talk

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

5
local_plugins/agent_service/lib/agent_service.dart

@ -265,8 +265,8 @@ class AgentService {
try {
// 在启动新对话前,确保之前的对话已完全停止
// 这不仅可以清理状态,还能避免因快速切换导致的资源冲突
await stopConversation();
// await stopConversation();
print('liwei--------- view startConversation');
final bool result = await _channel.invokeMethod('startConversation', {
'mode': mode,
});
@ -337,6 +337,7 @@ class AgentService {
// 即使频繁调用,Native 层也能处理(我们已经修复了 Native 层的并发问题)。
try {
print('liwei--------- view stopConversation');
final bool result = await _channel.invokeMethod('stopConversation');
return result;
} 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
private var isPhoneCallMode = false
// 音频处理
// private var externalAudioStream: ExternalAudioPullStream?
@ -80,6 +81,10 @@ public class AzureAsrHelper: NSObject {
// 音频处理
public var audioStream: SimpleAudioReceiver?
public func asrProvider() -> String {
return useXunfei ? "xunfei" : "azure"
}
@ -264,6 +269,7 @@ public class AzureAsrHelper: NSObject {
/// 1) 新任务创建前,取消旧的未执行任务,确保只有“最后一次启动请求”会生效
/// 2) 任务执行前二次校验:必须是当前挂起任务且 audioStream?.isContinuousRecognitionActive 仍为 true 才执行启动
private func startAudioRecordAsync(
mod: String,
audioSourceType: AudioSourceType,
audioDataCallback: SimpleAudioReceiver.AudioDataCallback? = nil,
completion: @escaping (Bool) -> Void
@ -306,7 +312,7 @@ public class AzureAsrHelper: NSObject {
// 准备启动
self.isAudioStarting = true
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.isAudioStarted = success
@ -338,13 +344,15 @@ public class AzureAsrHelper: NSObject {
/**
* 执行音频启动操作
*/
private func performAudioStart(
public func performAudioStart(
mod: String,
audioSourceType: AudioSourceType,
audioDataCallback: SimpleAudioReceiver.AudioDataCallback?
) -> Bool {
do {
if(!audioStream!.isContinuousRecognitionActive){
audioStream?.startAudioRecord(
mod: mod,
audioSourceType: audioSourceType == .microphone ? .microphone : .external,
audioDataCallback: audioDataCallback
)
@ -411,6 +419,7 @@ public class AzureAsrHelper: NSObject {
* 若 stop 已发生则直接跳过 recognizer 的启动,避免“已停止但仍启动识别”的情况。
*/
public func startContinuousRecognition(
mod: String = "normal",
audioSourceType: AudioSourceType = .microphone,
audioDataCallback: SimpleAudioReceiver.AudioDataCallback? = nil,
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 }
if success {
@ -497,7 +506,7 @@ public class AzureAsrHelper: NSObject {
// 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 }
if success {
@ -768,6 +777,7 @@ public class AzureAsrHelper: NSObject {
if audioStream == nil {
audioStream = SimpleAudioReceiver()
audioStream?.initAudioRecord() // 确保调用初始化
}
// 检查音频配置是否已存在
@ -992,7 +1002,7 @@ public class AzureAsrHelper: NSObject {
// 修复:添加类型转换
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 }
}
@ -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 mode = args["mode"] as? String ?? "normal"
// 根据参数确定音频源类型
let audioSourceType = isExternalActive ?
@ -392,6 +393,7 @@ private func sendAstEvent(_ event: [String: Any]) {
// 启动连续识别
let success = azureAsrHelper.startContinuousRecognition(
mod: mode,
audioSourceType: audioSourceType,
isRemoveFirstPunctuation: removeFirstPunctuation
)
@ -631,19 +633,14 @@ private func sendAstEvent(_ event: [String: Any]) {
let success = azureTtsHelper.setVoice(voiceName)
result(success)
case "setSpeechParams":
guard let args = call.arguments as? [String: Any] else {
case "setTtsMode":
guard let args = call.arguments as? [String: Any],
let mod = args["mod"] as? String else {
result(FlutterError(code: "INVALID_ARGUMENTS", message: "参数不能为空", details: nil))
return
}
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)
let success = azureTtsHelper.setTtsMode(mod: mod)
result(success)
case "speakText":
guard let args = call.arguments as? [String: Any],
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 {

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

File diff suppressed because it is too large

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

@ -92,13 +92,7 @@ public class MicrophoneCapture: NSObject {
}
do {
try audioSession.setPreferredSampleRate(sampleRate)
try audioSession.setPreferredIOBufferDuration(0.005) // 5ms缓冲
// 优先设置蓝牙音频路由
configureAudioRoute()
try audioSession.setActive(true)
try audioSession.setPreferredIOBufferDuration(0.02)
print("音频会话配置成功")
} catch {
print("音频会话配置失败: \(error.localizedDescription)")
@ -153,11 +147,10 @@ public class MicrophoneCapture: NSObject {
// 取消扬声器强制输出,让音频通过蓝牙耳机输出
do {
// 设置音频会话参数
try audioSession.setCategory(.playAndRecord,
try audioSession.setCategory(.playback,
mode: .videoChat,
options: [.allowBluetoothA2DP, // 允许蓝牙耳机,不占用hfp链路
.mixWithOthers,
.allowBluetooth]) // 添加音频优先级控制
options: [.allowBluetoothA2DP, // 仅允许A2DP,不启用HFP
.mixWithOthers])
try audioSession.overrideOutputAudioPort(.none)
print("蓝牙模式:音频输出设置为蓝牙耳机")
@ -212,24 +205,55 @@ public class MicrophoneCapture: NSObject {
}
// 检查麦克风权限
let permission: AVAudioSession.RecordPermission
switch audioSession.recordPermission {
case .granted:
try setupAudioEngine()
permission = .granted
case .denied:
throw NSError(domain: "麦克风权限被拒绝", code: 0)
permission = .denied
case .undetermined:
let semaphore = DispatchSemaphore(value: 0)
var grantedResult = false
audioSession.requestRecordPermission { granted in
if granted {
do {
try self.setupAudioEngine()
} catch {
print("音频引擎设置失败: \(error.localizedDescription)")
}
grantedResult = granted
semaphore.signal()
}
let waitResult = semaphore.wait(timeout: .now() + 60)
if waitResult == .timedOut {
throw NSError(domain: "麦克风权限请求超时", code: 2)
}
permission = grantedResult ? .granted : .denied
@unknown default:
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
}
}
/**
@ -281,35 +305,53 @@ public class MicrophoneCapture: NSObject {
sampleRate = audioSession.sampleRate
}
// 设置音频格式
let inputFormat = audioInputNode.inputFormat(forBus: 0)
// iOS 13+ 启用语音处理
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要求
let needsConversion = !isMicrosoftCompatibleFormat(inputFormat)
// 等待并获取有效的输入格式:在路由/类别切换瞬间可能出现 0Hz,直接 installTap 会触发系统断言崩溃
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: ",")
let outputs = currentSession.currentRoute.outputs.map { $0.portType.rawValue }.joined(separator: ",")
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
])
}
// 需要格式转换
guard let targetFormat = audioFormat,
let converter = AVAudioConverter(from: inputFormat, to: targetFormat) else {
// 需要格式转换(输出固定 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: "音频格式转换器创建失败"])
}
print("输入格式不符合Microsoft要求,进行格式转换")
print("输入格式: \(inputFormat.sampleRate)Hz, \(inputFormat.commonFormat.rawValue)")
print("目标格式: \(targetFormat.sampleRate)Hz, \(targetFormat.commonFormat.rawValue)")
print("输入格式: \(inputFormat.sampleRate)Hz, \(inputFormat.commonFormat.rawValue), ch=\(inputFormat.channelCount), interleaved=\(inputFormat.isInterleaved)")
print("目标格式: \(targetFormat.sampleRate)Hz, \(targetFormat.commonFormat.rawValue), ch=\(targetFormat.channelCount), interleaved=\(targetFormat.isInterleaved)")
// 添加tap进行格式转换
audioInputNode.installTap(onBus: 0,
bufferSize: 1024,
format: inputFormat) { [weak self] buffer, when in
format: inputFormat) { [weak self] buffer, _ in
guard let strongSelf = self else { return }
// 修复:添加安全检查,避免强制解包崩溃
guard let convertedBuffer = AVAudioPCMBuffer(
pcmFormat: targetFormat,
frameCapacity: AVAudioFrameCount(
@ -321,35 +363,58 @@ public class MicrophoneCapture: NSObject {
}
var error: NSError?
// 执行音频格式转换
let status = converter.convert(
to: convertedBuffer,
error: &error,
withInputFrom: { inNumPackets, outStatus in
withInputFrom: { _, outStatus in
outStatus.pointee = .haveData
return buffer
}
)
// 转换成功且无错误
if status == .haveData, error == nil {
let data = strongSelf.audioBufferToData(convertedBuffer)
strongSelf.audioDataHandler?(data)
} else if let error = error {
print("音频格式转换失败: \(error.localizedDescription)")
print("音频格式转换失败: domain=\(error.domain) code=\(error.code) desc=\(error.localizedDescription)")
}
}
audioEngine.prepare()
// 启动引擎
do {
try audioEngine.start()
isCapturing = true
} catch {
if let nsError = error as NSError? {
print("音频引擎启动失败: domain=\(nsError.domain) code=\(nsError.code) desc=\(nsError.localizedDescription)")
} else {
print("音频引擎启动失败: \(error.localizedDescription)")
}
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
)
}
// 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
else if let floatData = buffer.floatChannelData {
var int16Array = [Int16](repeating: 0, count: frameLength)
@ -442,35 +516,7 @@ public class MicrophoneCapture: NSObject {
// 修复:使用安全的清理方法
cleanupAudioEngine()
// 安全地处理音频会话
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)")
}
// 安全地处理音频会话:这里不主动 deactivate,避免与后续播报会话切换打架导致卡顿/失败
}
/**

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

@ -54,6 +54,7 @@ public class SimpleAudioReceiver: NSObject {
private var currentRoute: AudioOutputRoute?
//public var onAudioData: ((Data) -> Void)?
public var recordfile: RecordFile?
private var currentRecognitionMode = "normal" // 当前识别模式:normal, ble_wakeup, phone_call, push_to_talk
public var isHeadphones = true
/**
@ -119,29 +120,176 @@ public class SimpleAudioReceiver: NSObject {
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?) {
_isWriting=true
self.audioSourceType = audioSourceType;
public func startAudioRecord(mod: String = "normal", audioSourceType: AudioSourceType = .microphone, audioDataCallback: AudioDataCallback?) {
self.currentRecognitionMode = mod
self.audioSourceType = audioSourceType
self.audioDataCallback = audioDataCallback
print("startAudioRecord=audioSourceType\(audioSourceType)")
switch audioSourceType {
case .microphone:
self.isHeadphones = detectHeadphonesOutput()
_isWriting = true
print("liwei---------- startAudioRecord=mod:\(mod), audioSourceType:\(audioSourceType), isHeadphones:\(isHeadphones)")
do {
try configureAudioSession()
let isHybridExternalMode = audioSourceType == .external && (currentRecognitionMode == "phone_call" || currentRecognitionMode == "push_to_talk")
if audioSourceType == .microphone || isHybridExternalMode {
runMicrophoneCapture()
case .external:
} 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() {
// 使用userInitiated QoS匹配音频录制线程的优先级
writeThread = DispatchQueue(label: "audio.stream.writer", qos: .userInitiated)
writeThread = DispatchQueue(label: "audio.stream.writer", qos: .default)
writeThread?.async { [weak self] in
guard let self = self else { return }
@ -198,9 +346,6 @@ public class SimpleAudioReceiver: NSObject {
// 确保正在写入状态
guard self._isWriting else { return }
// 根据通话状态进行安全的音频会话配置
try safeConfigureAudioSessionForMic()
// 启动麦克风采集
try micCapture.startCapture()
@ -225,74 +370,8 @@ public class SimpleAudioReceiver: NSObject {
*/
private func runExternalCapture() {
if audioSourceType == .external {
// do {
// try audioSession.setCategory(.playback,
// 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)
}
}
@ -319,11 +398,16 @@ public class SimpleAudioReceiver: NSObject {
// 启动录音引擎
do {
print("resumeRecord")
try configureAudioSession()
try micCapture.startCapture()
} catch {
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)
}
}
}
// MARK: - 协议定义
@ -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) {
Log.d(
TAG,
"收到设备主动上报: ${data.joinToString(", ") { "0x${(it.toInt() and 0xFF).toString(16)}" }}"
"ble指令----------- 收到设备主动上报: ${data.joinToString(", ") { "0x${(it.toInt() and 0xFF).toString(16)}" }}"
)
if (data.size < 3) {
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) {
val data = c.value ?: return
Log.i(TAG, "ble指令----------- 收到通知特征数据(命令和控制 uuid=${c.uuid}")
// 根据特征UUID区分处理
when (c.uuid) {
// 音频特征数据
@ -973,7 +973,7 @@ object BleService {
}
// 通知特征数据(命令和控制)
BleConst.NOTIFY_CHAR_UUID -> {
Log.i(TAG, "收到通知特征数据(命令和控制)")
// Log.i(TAG, "ble指令----------- 收到通知特征数据(命令和控制)")
// 判断数据类型
when {
// 设备响应 (0xBB开头)

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

@ -56,6 +56,10 @@ public class BleService: NSObject {
// 添加录制文件管理器
private var recordingFile: RecordingFile?
private var pendingStopOpusRecording = false
private var stopOpusRecordingWorkItem: DispatchWorkItem?
private var lastOpusPacketReceivedAt: TimeInterval = 0
private var isScanning = false
private var scanTimer: Timer?
@ -471,6 +475,7 @@ private var cmdReplyType: UInt8 = 0
// 停止音频发送定时器
stopAudioSendTimer()
stopOpusRecording()
// 清空音频缓冲区和队列
audioQueueLock.lock()
@ -743,9 +748,7 @@ private var cmdReplyType: UInt8 = 0
public func closeCodec() -> Bool {
os_log("关闭编解码并停止录制...", log: logger, type: .info)
// 停止录制
recordingFile?.closeFile()
recordingFile = nil
// 停止音频发送定时器
stopAudioSendTimer()
@ -756,6 +759,9 @@ private var cmdReplyType: UInt8 = 0
audioDataQueue.removeAll()
audioQueueLock.unlock()
requestStopOpusRecording()
// 不停止Opus解码流,只发送命令通知设备关闭编解码
let paramData = Data([BleConst.CODEC_CONTROL_CLOSE, BleConst.AUDIO_CHANNEL_RIGHT])
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 {
guard let characteristic = writeCharacteristic,
connectionState == BleConst.STATE_CONNECTED else {
os_log("发送命令失败: 设备未连接", log: logger, type: .error)
os_log("liwei---- 发送命令失败: 设备未连接", log: logger, type: .error)
return false
}
@ -856,7 +862,7 @@ private var cmdReplyType: UInt8 = 0
if !data.isEmpty {
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)
@ -875,7 +881,7 @@ private var cmdReplyType: UInt8 = 0
return true
} catch {
os_log("发送命令异常: %@", log: logger, type: .error, error.localizedDescription)
os_log("liwei---- 发送命令异常: %@", log: logger, type: .error, error.localizedDescription)
return false
}
}
@ -906,8 +912,67 @@ private var cmdReplyType: UInt8 = 0
return
}
lastOpusPacketReceivedAt = Date().timeIntervalSince1970
ensureOpusRecordingStarted(fileName: "耳机端").saveAudioData(data)
// 使用Opus处理器处理音频数据
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: ", ")
os_log("收到设备响应: %@", log: logger, type: .debug, hexString)
// os_log("liwei---- 收到设备响应: %@", log: logger, type: .debug, hexString)
// 检查数据长度(与 Android 统一失败判定保持一致)
if data.count != Int(length) + 4 { // 帧头 + 命令 + 长度 + 数据 + CRC
@ -1786,6 +1851,7 @@ extension BleService: CBCentralManagerDelegate {
}
updateConnectionState(BleConst.STATE_DISCONNECTED)
stopOpusRecording()
// 只有在非主动断开的情况下才重新连接
if !isManualDisconnect {

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

@ -78,7 +78,7 @@ class RecordingFile {
do {
try self.fileHandle?.write(contentsOf: buffer)
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 {
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 CoreBluetooth
import os
import opus
@available(iOS 13.0, *)
@objc(BleServicePlugin)
@ -172,6 +173,40 @@ public class SwiftBleServicePlugin: NSObject, FlutterPlugin {
// 移除这个方法调用,让iOS系统自然处理权限
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:
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
@available(iOS 13.0, *)
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
private var mcpClient: Client?
private var transport: CustomSseClientTransport?
private var tools: [Tool] = []
private var toolMaps: [[String: Any]] = []
private var isConnectedFlag = false
private let toolsLock = NSLock()
// 重连配置
private let maxRetryAttempts = 3
@ -52,6 +54,10 @@ public class MCPSubClient {
private let maxToolCallRetries = 3
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()
@ -89,6 +95,7 @@ public class MCPSubClient {
}
}
)
self.transport = transport
// 3. 连接到服务器(添加超时)
try await withTimeout(seconds: 10) {
@ -106,6 +113,7 @@ public class MCPSubClient {
isConnectedFlag = true
retryCount = 0 // 重置重试计数
currentReconnectDelay = initialReconnectDelay // 重置延迟
startKeepAlive()
return true
@ -116,6 +124,9 @@ public class MCPSubClient {
}
private func processTools(_ toolList: [Tool]) {
toolsLock.lock()
defer { toolsLock.unlock() }
tools.removeAll()
toolMaps.removeAll()
@ -173,6 +184,13 @@ public class MCPSubClient {
/// 检查连接状态
/// 检查连接状态
public func checkConnection() async -> Bool {
if isConnectedFlag, let transport {
let isActive = await transport.isConnectionActive()
if !isActive {
isConnectedFlag = false
}
}
if !isConnectedFlag {
return await connect()
}
@ -339,14 +357,25 @@ public class MCPSubClient {
}
public func getToolMaps() -> [[String: Any]] {
toolsLock.lock()
defer { toolsLock.unlock() }
return toolMaps
}
public func containsTool(name: String) -> Bool {
toolsLock.lock()
defer { toolsLock.unlock() }
return tools.contains { $0.name == name }
}
public func close() async {
keepAliveTask?.cancel()
keepAliveTask = nil
reconnectTask?.cancel()
reconnectTask = nil
if let client = mcpClient {
await client.disconnect()
}
@ -359,9 +388,9 @@ public class MCPSubClient {
if isConnectedFlag {
print("[\(serverId)] 连接已断开,更新状态")
isConnectedFlag = false
// 启动重连逻辑
// startReconnection()
keepAliveTask?.cancel()
keepAliveTask = nil
startReconnection()
}
}
}
@ -369,7 +398,9 @@ public class MCPSubClient {
/// 启动重连
private func startReconnection() {
// 取消之前的重连任务
reconnectTask?.cancel()
if reconnectTask != nil {
return
}
reconnectTask = Task { [weak self] in
await self?.performReconnection()
@ -378,6 +409,8 @@ public class MCPSubClient {
/// 执行重连
private func performReconnection() async {
defer { reconnectTask = nil }
while retryCount < maxRetryAttempts && !Task.isCancelled {
retryCount += 1
@ -403,6 +436,8 @@ public class MCPSubClient {
private func handleConnectionError(_ error: Error) async {
print("[\(serverId)] 连接错误: \(error.localizedDescription)")
isConnectedFlag = false
keepAliveTask?.cancel()
keepAliveTask = nil
// 如果是MCP特定错误,可以进行特殊处理
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 {
return !subClients.isEmpty && subClients.values.allSatisfy { $0.getConnectionStatus() }
return !subClients.isEmpty && subClients.values.contains { $0.getConnectionStatus() }
}
public func disconnectAll() async {

69
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 logger = OSLog(subsystem: "com.yunqiinnovation.chat_storage", category: "ChatStorageHelper")
private let dbQueue = DispatchQueue(label: "com.yunqiinnovation.chat_storage.dbQueue")
private init() {
// 获取文档目录路径(不变)
let fileURL = try! FileManager.default
@ -23,6 +25,7 @@ public class ChatStorageHelper {
dbPath = fileURL.path
dbQueue.sync {
// 打开数据库(不变)
if sqlite3_open(dbPath, &db) != SQLITE_OK {
let errmsg = String(cString: sqlite3_errmsg(db)!)
@ -52,12 +55,15 @@ public class ChatStorageHelper {
os_log("创建表失败: %{public}@", log: logger, type: .error, errmsg)
}
}
}
deinit {
dbQueue.sync {
if db != nil {
sqlite3_close(db)
}
}
}
/**
* 保存消息到数据库
@ -68,6 +74,8 @@ public class ChatStorageHelper {
* @return 插入的消息ID,失败则返回-1
*/
public func saveMessage(agentId: String, sessionId: String, message: String, sender: String, metadata: String?) -> Int64 {
var result: Int64 = -1
dbQueue.sync {
// 关键修改:用 INSERT OR REPLACE 替换 INSERT,支持冲突时更新
let insertStatementString = """
INSERT OR REPLACE INTO messages
@ -95,21 +103,19 @@ public class ChatStorageHelper {
// 执行语句(冲突时会自动替换,返回新的 rowid)
if sqlite3_step(insertStatement) == SQLITE_DONE {
let id = sqlite3_last_insert_rowid(db) // 替换后返回新的 id(原 id 会被删除)
sqlite3_finalize(insertStatement)
os_log("消息保存成功(新增/更新),id: %{public}lld", log: logger, type: .info, id)
return id
result = id
} else {
let errmsg = String(cString: sqlite3_errmsg(db)!)
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 -1
}
return result
}
/**
@ -120,6 +126,8 @@ public class ChatStorageHelper {
* @return 消息列表的JSON字符串
*/
public func getMessages(agentId: String, page: Int, pageSize: Int) -> String {
var result = "{\"messages\":[],\"page\":\(page),\"pageSize\":\(pageSize),\"totalCount\":0,\"totalPages\":0}"
dbQueue.sync {
let offset = (page - 1) * pageSize
var messagesArray: [[String: Any]] = []
@ -196,7 +204,7 @@ public class ChatStorageHelper {
sqlite3_finalize(queryStatement)
// 构建与Android版本一致的返回格式
let result: [String: Any] = [
let resultDict: [String: Any] = [
"messages": messagesArray,
"page": page,
"pageSize": pageSize,
@ -205,9 +213,9 @@ public class ChatStorageHelper {
]
os_log("getMessages: 读取历史记录: %{public}@", log: logger, type: .error, messagesArray)
do {
let jsonData = try JSONSerialization.data(withJSONObject: result, options: [])
let jsonData = try JSONSerialization.data(withJSONObject: resultDict, options: [])
if let jsonString = String(data: jsonData, encoding: .utf8) {
return jsonString
result = jsonString
}
} catch {
os_log("getMessages: JSON转换失败: %{public}@", log: logger, type: .error, error.localizedDescription)
@ -216,26 +224,8 @@ public class ChatStorageHelper {
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 "{\"messages\":[],\"page\":\(page),\"pageSize\":\(pageSize),\"totalCount\":0,\"totalPages\":0}"
return result
}
/**
@ -245,6 +235,8 @@ public class ChatStorageHelper {
* @return 是否删除成功
*/
public func deleteMessages(agentId: String?, messageIds: [Int]?) -> Bool {
var result = false
dbQueue.sync {
if let agentId = agentId {
let deleteString = "DELETE FROM messages WHERE agent_id = ?;"
var deleteStatement: OpaquePointer?
@ -253,8 +245,7 @@ public class ChatStorageHelper {
sqlite3_bind_text(deleteStatement, 1, (agentId as NSString).utf8String, -1, nil)
if sqlite3_step(deleteStatement) == SQLITE_DONE {
sqlite3_finalize(deleteStatement)
return true
result = true
} else {
let errmsg = String(cString: sqlite3_errmsg(db)!)
os_log("删除会话消息失败: %{public}@", log: logger, type: .error, errmsg)
@ -277,13 +268,11 @@ public class ChatStorageHelper {
}
if sqlite3_step(deleteStatement) == SQLITE_DONE {
sqlite3_finalize(deleteStatement)
return true
result = true
} else {
let errmsg = String(cString: sqlite3_errmsg(db)!)
os_log("删除特定消息失败: %{public}@", log: logger, type: .error, errmsg)
}
sqlite3_finalize(deleteStatement)
} else {
let errmsg = String(cString: sqlite3_errmsg(db)!)
@ -291,10 +280,9 @@ public class ChatStorageHelper {
}
} else {
os_log("删除消息参数无效 - agentId和messageIds都为空", log: logger, type: .error)
return false
}
return false
}
return result
}
/**
@ -302,16 +290,19 @@ public class ChatStorageHelper {
* @return 是否清空成功
*/
public func clearDatabase() -> Bool {
var result = false
dbQueue.sync {
let deleteString = "DELETE FROM messages;"
if sqlite3_exec(db, deleteString, nil, nil, nil) == SQLITE_OK {
return true
result = true
} else {
let errmsg = String(cString: sqlite3_errmsg(db)!)
os_log("清空数据库失败: %{public}@", log: logger, type: .error, errmsg)
return false
}
}
return result
}
/**
* 获取指定会话的最近N条消息
@ -321,7 +312,7 @@ public class ChatStorageHelper {
*/
public func getRecentMessages(agentId: String, limit: Int) -> [[String: Any]] {
var messages: [[String: Any]] = []
dbQueue.sync {
// 首先检查数据库中是否有该会话的消息
let countQuery = "SELECT COUNT(*) FROM messages WHERE agent_id = ?"
var countStatement: OpaquePointer?
@ -389,7 +380,7 @@ public class ChatStorageHelper {
let errmsg = String(cString: sqlite3_errmsg(db)!)
os_log("查询最近消息失败: %{public}@", log: logger, type: .error, errmsg)
}
}
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() {
super.init()
setupRemoteCommandCenter()
setupAudioSession()
// setupAudioSession()
setupNotificationObservers()
// 监听通话状态(需主线程队列)
// callObserver.setDelegate(self, queue: DispatchQueue.main)
//callObserver.setDelegate(self, queue: DispatchQueue.main)
}
public func initialize(serviceUrl: String, userToken: String) {

Loading…
Cancel
Save