Browse Source

Merge branch 'new_dev' of https://github.com/deepcloud2048/deep_voice into 重写ble

weicu
liwei1dao 1 year ago
parent
commit
2b9fc14bc2
  1. 2
      lib/core/translations/language/zh_cn.dart
  2. 65
      lib/modules/agent/controllers/agent_controller.dart
  3. 29
      local_plugins/agent_service/android/src/main/kotlin/com/yunqiinnovation/agent_service/AgentService.kt
  4. 63
      local_plugins/agent_service/ios/agent_service/Sources/agent_service/AgentServiceImpl.swift
  5. 2
      local_plugins/azure_speech/android/src/main/kotlin/com/yunqiinnovation/azure_speech/AzureTtsHelper.kt
  6. 2
      local_plugins/ble_service/ios/ble_service/Sources/ble_service/SwiftBleServicePlugin.swift
  7. 47
      local_plugins/chat_api/android/src/main/kotlin/com/yunqiinnovation/chat_api/ChatApiService.kt
  8. 83
      local_plugins/chat_api/android/src/main/kotlin/com/yunqiinnovation/chat_api/QQMusicSystemFunctionHandler.kt
  9. 202
      local_plugins/chat_api/ios/chat_api/Sources/chat_api/ChatApiService.swift
  10. 20
      local_plugins/chat_api/ios/chat_api/Sources/chat_api/MCPClient.swift
  11. 179
      local_plugins/chat_api/ios/chat_api/Sources/chat_api/NetworkStateMonitor.swift

2
lib/core/translations/language/zh_cn.dart

@ -511,7 +511,7 @@ const Map<String, String> zhCN = {
"loadMoreMessages": "加载更多消息...",
"loadFailedClickRetry": "加载失败,点击重试",
"swipeUpForMore": "向上滑动查看更多",
"callingTool": "正在为您调用工具查询相关信息",
"callingTool": "正在调用工具!",
// 预设提示词
"presetMusic": "来点音乐", // 来点音乐

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

@ -722,23 +722,25 @@ class AgentController extends GetxController with WidgetsBindingObserver {
break;
case AgentServiceEventType.error:
//验证错误码
final sessionid = event.data['sessionid'] ?? '';
final message = event.data['message'] ?? '';
// //验证错误码
if (event.data['code'] == 1000) {
isListening.value = false;
isSpeaking.value = false;
isProcessing.value = false;
isImageProcessing.value = false;
if (sessionid != "") {
final index = messages.lastIndexWhere(
(msg) => msg.sessionid == sessionid && !msg.isUser);
if (index >= 0) {
messages.removeAt(index);
messages.refresh();
}
}
Get.snackbar('error'.tr, '${event.data['message']}');
}
// 移除临时的识别消息
final index =
messages.lastIndexWhere((msg) => msg.isRecognizing && msg.isUser);
if (index >= 0) {
messages.removeAt(index);
// messages.refresh();
}
// Get.snackbar('error'.tr, '${event.data['message']}');
Logger.e(TAG, '代理服务错误: ${event.data['message']}');
break;
@ -2088,7 +2090,7 @@ class AgentController extends GetxController with WidgetsBindingObserver {
color: Colors.grey.withOpacity(0.1),
borderRadius: BorderRadius.circular(8),
),
child: const Icon(Icons.navigation, color: Colors.grey),
child: const Icon(Icons.map, color: Colors.orange),
),
title: const Text('苹果地图'),
subtitle: const Text('使用苹果地图导航'),
@ -2097,6 +2099,23 @@ class AgentController extends GetxController with WidgetsBindingObserver {
openAppleMap(destination, mode);
},
),
ListTile(
leading: Container(
padding: const EdgeInsets.all(8),
decoration: BoxDecoration(
color: Colors.grey.withOpacity(0.1),
borderRadius: BorderRadius.circular(8),
),
child: const Icon(Icons.navigation, color: Colors.grey),
),
title: const Text('内置导航'),
subtitle: const Text('使用内置地图导航'),
onTap: () {
Get.back();
openInternalGaodeMap(
start, destination, endlatitude, endlongitude, mode);
},
),
// 底部安全区域
SizedBox(height: MediaQuery.of(Get.context!).padding.bottom + 20),
],
@ -2489,6 +2508,8 @@ class AgentController extends GetxController with WidgetsBindingObserver {
// 播放成功后,如果歌曲不在当前列表中,保持列表显示状态不变
// Get.snackbar('开始播放', song.name);
_qqmusicManager.syncCurrentPlayInfo();
} else {
Get.snackbar('播放失败', result['error'] ?? '无法播放歌曲');
}
return;
}
@ -2509,6 +2530,8 @@ class AgentController extends GetxController with WidgetsBindingObserver {
// 播放成功后,如果歌曲不在当前列表中,保持列表显示状态不变
// Get.snackbar('开始播放', song.name);
_qqmusicManager.syncCurrentPlayInfo();
} else {
Get.snackbar('播放失败', result['error'] ?? '无法播放歌曲');
}
}
}
@ -2518,6 +2541,24 @@ class AgentController extends GetxController with WidgetsBindingObserver {
}
}
/// 打开内置高德地图进行导航
Future<bool> openInternalGaodeMap(
String start,
String destination,
String endlatitude,
String endlongitude,
String mode // 'driving', 'transit', 'walking', 'cycling'
) async {
try {
_navigationManager.open(start, "$endlongitude,$endlatitude", mode);
return true;
} catch (e) {
print('打开内置高德地图失败: $e');
Get.snackbar('错误', '打开地图时发生错误');
return false;
}
}
/// 打开Google地图进行导航
Future<bool> openGoogleMap(String daddr, String mode) async {
try {

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

@ -39,6 +39,7 @@ import java.util.Date
import android.app.PendingIntent
import android.app.NotificationManager
import android.app.NotificationChannel
import android.os.Message
import androidx.core.app.NotificationCompat
@ -721,6 +722,7 @@ object AgentService : CoroutineScope {
_isRecognitionActive.set(false)
stopIdleCheck()
BleService.closeCodec()
audioPlayer?.stopAudio()
Log.d(TAG, "语音识别错误: $error")
sendEvent(
"error", mapOf(
@ -1071,20 +1073,23 @@ object AgentService : CoroutineScope {
}
}
override fun onError(sessionid: String, e: Exception) {
override fun onError(sessionid: String, code:Int,message: String) {
try {
Log.e(TAG, "AI处理出错", e)
sendEvent(
"error", mapOf(
"sessionid" to sessionid,
"code" to "AI_ERROR",
"message" to e.message.toString()
if (sessionid == currsessionId) {
audioPlayer?.stopAudio()
Log.e(TAG, "AI处理出错 $message")
// 优先使用 ChatApiException.code;否则根据底层异常类型推断错误码
sendEvent(
"error", mapOf(
"sessionid" to sessionid,
"code" to code,
"message" to message
)
)
)
// 标记AI流式输出已完成
_isAiStreaming.set(false)
currentAiJob = null
// 标记AI流式输出已完成
_isAiStreaming.set(false)
currentAiJob = null
}
} catch (e: Exception) {
Log.e(TAG, "liwei--------------- AI Call onError 异常", e)
}

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

@ -62,6 +62,7 @@ class AgentServiceImpl: NSObject {
private let agentId = "default_agent"
//工具调用提示
public var callingTool = ""
private var navigationMode = ""
private var isInitialized: Bool = false
private var isRecognizing: Bool = false
@ -168,6 +169,10 @@ class AgentServiceImpl: NSObject {
self.config = config
if let navigationMode = config["navigationMode"] as? String {
self.navigationMode = navigationMode
}
if let callingTool = config["callingTool"] as? String {
self.callingTool = callingTool
}
@ -218,6 +223,7 @@ class AgentServiceImpl: NSObject {
sendError("Azure语音区域缺失", code: "CONFIG_ERROR")
return false
}
//启动定位服务
NotificationCenter.default.post(name: .locationstartEvent, object: nil)
// addlocationmonitor()
@ -1023,7 +1029,7 @@ private func jsonToString(_ json: [String: Any]) -> String? {
// 应用在前台时的原有逻辑
// 根据导航模式选择不同的导航方式
switch mode {
switch self.navigationMode {
case "internal_gaode":
// 使用内置高德导航
// startInternalNavigation(start: start, destination: destination,
@ -1708,27 +1714,29 @@ class ChatApiStreamCallback: StreamCallback {
}
func onError(_ sessionid:String,_ error: Error) {
func onError(_ sessionid:String,_ code: Int,_ message:String) {
do {
guard let agentService = try agentService else { return }
if sessionid == agentService.currsessionId {
// 停止等待音效
agentService.audioPlayer?.stopAwaitSound()
}
os_log("ChatAPI错误: %{public}@", log: agentService.logger, type: .error, error.localizedDescription)
if (sessionid != agentService.currsessionId) {
return
}
// 停止等待音效
agentService.audioPlayer?.stopAwaitSound()
try agentService.sendEvent(name: "error", data: [
"sessionid":sessionid,
"code": "AI_ERROR",
"message": error.localizedDescription
"code": code,
"message": message
])
agentService.isAiStreaming = false
os_log("因错误设置AI流式状态为false", log: agentService.logger, type: .info)
os_log(
"liwei--------------- AI Call onError code: %{public}d message: %{public}@ ",
log: agentService.logger,
type: .info,
code, // 直接传 Int 类型
message // 假设 message 是 String 类型
)
}catch{
print("liwei--------------- AI Call onError 异常: \(error)")
}
@ -1768,7 +1776,7 @@ class ChatApiStreamCallback: StreamCallback {
}catch{
print("liwei--------------- AI Call onFunctionCall 异常: \(error)")
}
}
func onFunctionCallResult(_ sessionid:String,_ functionCall: [String: Any], _ functionCallResult: [String: Any]) {
@ -1968,16 +1976,23 @@ class AudioPlayer {
}
func stopAwaitSound() {
// 停止音频播放
awaitPlayer?.stop()
awaitPlayer = nil
// 停止超时定时器
stopAwaitTimeout()
os_log("停止等待音效播放", log: logger, type: .debug)
// 确保在主线程执行音频停止操作
DispatchQueue.main.async { [weak self] in
guard let self = self else { return }
// 停止音频播放
if let player = self.awaitPlayer {
if player.isPlaying {
player.stop()
}
self.awaitPlayer = nil
}
// 停止超时定时器
self.stopAwaitTimeout()
os_log("停止等待音效播放", log: self.logger, type: .debug)
}
}
private func playSound(named: String, fileType: String = "mp3") {
@ -2182,7 +2197,7 @@ extension AgentServiceImpl: AzureAsrHelper.ContinuousRecognizeCallback {
}
func onError(sessionid:String ,_ errorCode: Int, _ error: String) {
let data: [String: Any] = ["sessionid":sessionid,"message": error.isEmpty ? "未知错误" : error]
let data: [String: Any] = ["sessionid":sessionid,"code":errorCode, "message": error.isEmpty ? "未知错误" : error]
sendEvent(name: "error", data: data)
isRecognizing = false
// 新增:出错时复位"启动中/待停止"状态

2
local_plugins/azure_speech/android/src/main/kotlin/com/yunqiinnovation/azure_speech/AzureTtsHelper.kt

@ -671,8 +671,8 @@ class AzureTtsHelper(private val context: Context) : ITtsService {
try {
synthesizer?.StopSpeakingAsync()
synthesizer?.close()
streamBuffer.clear()
isSpeaking = false
return true
} catch (e: Exception) {

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

@ -141,7 +141,7 @@ public class SwiftBleServicePlugin: NSObject, FlutterPlugin {
case "getPairedMacAddress":
let pairedUUID = BleService.shared.getPairedMacAddress()
result(nil)
result(pairedUUID)
case "clearAssociations":
BleService.shared.clearAssociations()
result(true)

47
local_plugins/chat_api/android/src/main/kotlin/com/yunqiinnovation/chat_api/ChatApiService.kt

@ -29,8 +29,7 @@ import java.util.Collections
/**
* ChatAPI服务异常
*/
class ChatApiException(message: String) : Exception("ChatApiException: $message")
class ChatApiException(message: String, val code: String? = null, cause: Throwable? = null) : Exception("ChatApiException: $message", cause)
/**
* 流式回调接口
* 与 iOS 版本 StreamCallback 协议保持完全一致
@ -53,7 +52,7 @@ interface StreamCallback {
/**
* 出现错误
*/
fun onError(sessionid: String,error: Exception)
fun onError(sessionid: String,code:Int,message:String)
/**
* 函数调用 - 兼容JSONObject格式
@ -354,13 +353,20 @@ class ChatApiService(private val context: android.content.Context? = null) : Cor
if (!isInitialized || apiKey.isEmpty() || openAI == null) {
Log.e("ChatApiService", "ChatAPI服务未初始化,无法发送消息")
try {
getSessionCallback(sessionid)?.onError(sessionid,ChatApiException("ChatAPI服务未初始化"))
getSessionCallback(sessionid)?.onError(sessionid,1002,"ChatAPI服务未初始化")
} catch (e: Exception) {
Log.e(TAG, "onError回调异常: ${e.message}", e)
}
return
}
// 在开始流式请求前检查网络状态
if (!isNetworkAvailable()) {
Log.w("ChatApiService", "[Session: $sessionid] 网络不可用,直接返回网络错误")
getSessionCallback(sessionid)?.onError(sessionid,1000,"网络不可用")
return
}
// 重置状态
currentMessages = messages
toolCalls.clear()
@ -418,7 +424,7 @@ class ChatApiService(private val context: android.content.Context? = null) : Cor
Log.e("ChatApiService", "创建ChatCompletionRequest或调用chatCompletions失败: ${e.message}", e)
if (sessionid == currSessionId) {
try {
getSessionCallback(sessionid)?.onError(sessionid,ChatApiException("流式请求失败: ${e.message}"))
getSessionCallback(sessionid)?.onError(sessionid,1000,"流式请求失败: ${e.message}")
} catch (ex: Exception) {
Log.e(TAG, "onError回调异常: ${ex.message}", ex)
}
@ -477,14 +483,6 @@ class ChatApiService(private val context: android.content.Context? = null) : Cor
}
}
}
//else{
// delta.content?.let { content ->
// Log.d("ChatApiService", "liwei------------ [Session: $sessionid] 中间过程不输出 $content")
// }
// }
// Log.d(TAG, "liwei-------------------------开始AI 对话 7-4")
// 收集工具调用信息
delta.toolCalls?.forEach { toolCall ->
@ -559,7 +557,8 @@ class ChatApiService(private val context: android.content.Context? = null) : Cor
} catch (e: Exception) {
if (sessionid == currSessionId) {
try {
getSessionCallback(sessionid)?.onError(sessionid,ChatApiException("流式请求失败: ${e.message}"))
Log.d(TAG, "AI聊天异常 Session $sessionid 错误类型:${e::class.simpleName} 错误完整类型:${e::class.qualifiedName}")
getSessionCallback(sessionid)?.onError(sessionid,1000,"流式请求失败: ${e.message}")
} catch (ex: Exception) {
Log.e(TAG, "onError回调异常: ${ex.message}", ex)
}
@ -632,13 +631,13 @@ class ChatApiService(private val context: android.content.Context? = null) : Cor
}
// 调用MCP工具
val toolResult = try {
withTimeout(60000) { // 60秒超时
withTimeout(10000) { // 100秒超时
// 再次检查会话状态
if (sessionid != currSessionId) {
throw CancellationException("Session cancelled")
}
_mcpClient?.callTool(functionName, arguments)
}
}
} catch (e: TimeoutCancellationException) {
Log.w("ChatApiService", "[Session: $sessionid] MCP工具调用超时: $functionName")
mapOf(
@ -708,7 +707,8 @@ class ChatApiService(private val context: android.content.Context? = null) : Cor
convertMapToJsonObject(result)
)
}else{
getSessionCallback(sessionid)?.onError(sessionid, ChatApiException(errorText))
Log.d(TAG, "[Session: $sessionid] 工具调用失败 $errorText")
//getSessionCallback(sessionid)?.onError(sessionid, ChatApiException(errorText))
}
} catch (e: Exception) {
Log.e(TAG, "onFunctionCallResult回调异常: ${e.message}", e)
@ -798,6 +798,19 @@ class ChatApiService(private val context: android.content.Context? = null) : Cor
sendMessageStream(sessionid,fullMessages)
}
// 添加网络检查方法
private fun isNetworkAvailable(): Boolean {
return try {
val connectivityManager = context?.getSystemService(android.content.Context.CONNECTIVITY_SERVICE) as? android.net.ConnectivityManager
val activeNetwork = connectivityManager?.activeNetworkInfo
activeNetwork?.isConnectedOrConnecting == true
} catch (e: Exception) {
Log.w("ChatApiService", "检查网络状态失败: ${e.message}")
true // 如果检查失败,假设网络可用,让后续的网络请求来处理
}
}
/**
* 取消当前流式请求
*/

83
local_plugins/chat_api/android/src/main/kotlin/com/yunqiinnovation/chat_api/QQMusicSystemFunctionHandler.kt

@ -6,6 +6,7 @@ import kotlinx.coroutines.delay
import com.yunqiinnovation.qq_music.QQMusicSingleton
import kotlin.coroutines.suspendCoroutine
import kotlin.coroutines.resume
import kotlinx.coroutines.CompletableDeferred
/**
* QQ音乐系统功能处理器
* 负责注册QQ音乐相关的MCP函数
@ -211,7 +212,8 @@ private class QQMusicPlaySongsHandler(private val context: Context?) : FunctionH
val index = (arguments["index"] as? Number)?.toInt() ?: 0
val qqMusicSingleton = QQMusicSingleton.getInstance()
var result = ""
// 创建一个CompletableDeferred用于等待回调结果
val deferred = CompletableDeferred<String>()
qqMusicSingleton.playSongs(songMids, index) { playResult ->
result = if (playResult.isSuccess) {
"{\"broadcast\": false, \"success\": true, \"message\": \"开始播放歌曲列表\"}"
@ -219,11 +221,12 @@ private class QQMusicPlaySongsHandler(private val context: Context?) : FunctionH
val error = playResult.exceptionOrNull()
"{\"success\": false, \"message\": \"播放失败: ${error?.message}\"}"
}
// 回调完成,唤醒挂起的协程
deferred.complete(result)
}
// 等待回调完成 十秒
delay(10000)
result.ifEmpty { "{\"nocard\": false,\"broadcast\": true, \"success\": false, \"message\": \"播放超时\"}" }
// 挂起等待回调结果(无超时,直到回调触发)
deferred.await()
} catch (e: Exception) {
Log.e("QQMusicPlaySongsHandler", "QQ音乐播放歌曲失败", e)
"{\"nocard\": true, \"success\": false, \"message\": \"播放异常: ${e.message}\"}"
@ -236,7 +239,8 @@ private class QQMusicPlayHandler(private val context: Context?) : FunctionHandle
return try {
val qqMusicSingleton = QQMusicSingleton.getInstance()
var result = ""
// 创建一个CompletableDeferred用于等待回调结果
val deferred = CompletableDeferred<String>()
qqMusicSingleton.playMusic { playResult ->
result = if (playResult.isSuccess) {
"{\"broadcast\": false, \"success\": true, \"message\": \"开始播放\"}"
@ -244,10 +248,12 @@ private class QQMusicPlayHandler(private val context: Context?) : FunctionHandle
val error = playResult.exceptionOrNull()
"{\"success\": false, \"message\": \"播放失败: ${error?.message}\"}"
}
// 回调完成,唤醒挂起的协程
deferred.complete(result)
}
delay(10000)
result.ifEmpty { "{\"success\": false, \"message\": \"播放超时\"}" }
// 挂起等待回调结果(无超时,直到回调触发)
deferred.await()
} catch (e: Exception) {
Log.e("QQMusicPlayHandler", "QQ音乐播放失败", e)
"{\"success\": false, \"message\": \"播放异常: ${e.message}\"}"
@ -260,7 +266,8 @@ private class QQMusicPauseHandler(private val context: Context?) : FunctionHandl
return try {
val qqMusicSingleton = QQMusicSingleton.getInstance()
var result = ""
// 创建一个CompletableDeferred用于等待回调结果
val deferred = CompletableDeferred<String>()
qqMusicSingleton.pauseMusic { pauseResult ->
result = if (pauseResult.isSuccess) {
"{\"broadcast\": false, \"success\": true, \"message\": \"暂停播放\"}"
@ -268,10 +275,12 @@ private class QQMusicPauseHandler(private val context: Context?) : FunctionHandl
val error = pauseResult.exceptionOrNull()
"{\"success\": false, \"message\": \"暂停失败: ${error?.message}\"}"
}
// 回调完成,唤醒挂起的协程
deferred.complete(result)
}
delay(2000)
result.ifEmpty { "{\"success\": false, \"message\": \"暂停超时\"}" }
// 挂起等待回调结果(无超时,直到回调触发)
deferred.await()
} catch (e: Exception) {
Log.e("QQMusicPauseHandler", "QQ音乐暂停失败", e)
"{\"success\": false, \"message\": \"暂停异常: ${e.message}\"}"
@ -284,7 +293,8 @@ private class QQMusicResumeHandler(private val context: Context?) : FunctionHand
return try {
val qqMusicSingleton = QQMusicSingleton.getInstance()
var result = ""
// 创建一个CompletableDeferred用于等待回调结果
val deferred = CompletableDeferred<String>()
qqMusicSingleton.resumeMusic { resumeResult ->
result = if (resumeResult.isSuccess) {
"{\"broadcast\": false, \"success\": true, \"message\": \"恢复播放\"}"
@ -292,10 +302,12 @@ private class QQMusicResumeHandler(private val context: Context?) : FunctionHand
val error = resumeResult.exceptionOrNull()
"{\"success\": false, \"message\": \"恢复播放失败: ${error?.message}\"}"
}
// 回调完成,唤醒挂起的协程
deferred.complete(result)
}
delay(2000)
result.ifEmpty { "{\"success\": false, \"message\": \"恢复播放超时\"}" }
// 挂起等待回调结果(无超时,直到回调触发)
deferred.await()
} catch (e: Exception) {
Log.e("QQMusicResumeHandler", "QQ音乐恢复播放失败", e)
"{\"success\": false, \"message\": \"恢复播放异常: ${e.message}\"}"
@ -308,7 +320,8 @@ private class QQMusicStopHandler(private val context: Context?) : FunctionHandle
return try {
val qqMusicSingleton = QQMusicSingleton.getInstance()
var result = ""
// 创建一个CompletableDeferred用于等待回调结果
val deferred = CompletableDeferred<String>()
qqMusicSingleton.stopMusic { stopResult ->
result = if (stopResult.isSuccess) {
"{\"broadcast\": false, \"success\": true, \"message\": \"停止播放\"}"
@ -316,10 +329,12 @@ private class QQMusicStopHandler(private val context: Context?) : FunctionHandle
val error = stopResult.exceptionOrNull()
"{\"success\": false, \"message\": \"停止失败: ${error?.message}\"}"
}
// 回调完成,唤醒挂起的协程
deferred.complete(result)
}
delay(2000)
result.ifEmpty { "{\"success\": false, \"message\": \"停止超时\"}" }
// 挂起等待回调结果(无超时,直到回调触发)
deferred.await()
} catch (e: Exception) {
Log.e("QQMusicStopHandler", "QQ音乐停止失败", e)
"{\"success\": false, \"message\": \"停止异常: ${e.message}\"}"
@ -332,7 +347,8 @@ private class QQMusicNextHandler(private val context: Context?) : FunctionHandle
return try {
val qqMusicSingleton = QQMusicSingleton.getInstance()
var result = ""
// 创建一个CompletableDeferred用于等待回调结果
val deferred = CompletableDeferred<String>()
qqMusicSingleton.skipToNext { nextResult ->
result = if (nextResult.isSuccess) {
"{\"broadcast\": false, \"success\": true, \"message\": \"切换到下一首\"}"
@ -340,10 +356,12 @@ private class QQMusicNextHandler(private val context: Context?) : FunctionHandle
val error = nextResult.exceptionOrNull()
"{\"success\": false, \"message\": \"切换失败: ${error?.message}\"}"
}
// 回调完成,唤醒挂起的协程
deferred.complete(result)
}
delay(2000)
result.ifEmpty { "{\"success\": false, \"message\": \"切换超时\"}" }
// 挂起等待回调结果(无超时,直到回调触发)
deferred.await()
} catch (e: Exception) {
Log.e("QQMusicNextHandler", "QQ音乐下一首失败", e)
"{\"success\": false, \"message\": \"切换异常: ${e.message}\"}"
@ -356,7 +374,8 @@ private class QQMusicPreviousHandler(private val context: Context?) : FunctionHa
return try {
val qqMusicSingleton = QQMusicSingleton.getInstance()
var result = ""
// 创建一个CompletableDeferred用于等待回调结果
val deferred = CompletableDeferred<String>()
qqMusicSingleton.skipToPrevious { previousResult ->
result = if (previousResult.isSuccess) {
"{\"broadcast\": false, \"success\": true, \"message\": \"切换到上一首\"}"
@ -364,10 +383,12 @@ private class QQMusicPreviousHandler(private val context: Context?) : FunctionHa
val error = previousResult.exceptionOrNull()
"{\"success\": false, \"message\": \"切换失败: ${error?.message}\"}"
}
// 回调完成,唤醒挂起的协程
deferred.complete(result)
}
delay(2000)
result.ifEmpty { "{\"success\": false, \"message\": \"切换超时\"}" }
// 挂起等待回调结果(无超时,直到回调触发)
deferred.await()
} catch (e: Exception) {
Log.e("QQMusicPreviousHandler", "QQ音乐上一首失败", e)
"{\"success\": false, \"message\": \"切换异常: ${e.message}\"}"
@ -380,7 +401,8 @@ private class QQMusicGetPlaybackStateHandler(private val context: Context?) : Fu
return try {
val qqMusicSingleton = QQMusicSingleton.getInstance()
var result = ""
// 创建一个CompletableDeferred用于等待回调结果
val deferred = CompletableDeferred<String>()
qqMusicSingleton.getPlaybackState { stateResult ->
result = if (stateResult.isSuccess) {
val state = stateResult.getOrNull() ?: emptyMap()
@ -392,10 +414,12 @@ private class QQMusicGetPlaybackStateHandler(private val context: Context?) : Fu
val error = stateResult.exceptionOrNull()
"{\"success\": false, \"message\": \"获取播放状态失败: ${error?.message}\"}"
}
// 回调完成,唤醒挂起的协程
deferred.complete(result)
}
delay(2000)
result.ifEmpty { "{\"success\": false, \"message\": \"获取播放状态超时\"}" }
// 挂起等待回调结果(无超时,直到回调触发)
deferred.await()
} catch (e: Exception) {
Log.e("CahtApi", "获取QQ音乐播放状态失败", e)
"{\"success\": false, \"message\": \"获取播放状态异常: ${e.message}\"}"
@ -408,7 +432,8 @@ private class QQMusicGetCurrentSongHandler(private val context: Context?) : Func
return try {
val qqMusicSingleton = QQMusicSingleton.getInstance()
var result = ""
// 创建一个CompletableDeferred用于等待回调结果
val deferred = CompletableDeferred<String>()
qqMusicSingleton.getCurrentSong { songResult ->
result = if (songResult.isSuccess) {
val song = songResult.getOrNull() ?: emptyMap()
@ -420,10 +445,12 @@ private class QQMusicGetCurrentSongHandler(private val context: Context?) : Func
val error = songResult.exceptionOrNull()
"{\"success\": false, \"message\": \"获取当前歌曲失败: ${error?.message}\"}"
}
// 回调完成,唤醒挂起的协程
deferred.complete(result)
}
delay(2000)
result.ifEmpty { "{\"success\": false, \"message\": \"获取当前歌曲超时\"}" }
// 挂起等待回调结果(无超时,直到回调触发)
deferred.await()
} catch (e: Exception) {
Log.e("QQMusicLoginHandler", "获取QQ音乐当前歌曲失败", e)
"{\"success\": false, \"message\": \"获取当前歌曲异常: ${e.message}\"}"

202
local_plugins/chat_api/ios/chat_api/Sources/chat_api/ChatApiService.swift

@ -20,7 +20,7 @@ public protocol StreamCallback {
func onUsage(_ sessionid: String,_ prompt_tokens: Int?,_ completion_tokens: Int?,_ total_tokens: Int?)
func onToken(_ sessionId:String,_ token: String)
func onComplete(_ sessionId:String)
func onError(_ sessionId:String,_ error: Error)
func onError(_ sessionId:String,_ code:Int,_ message:String)
func onFunctionCall(_ sessionId:String,_ functionCall: [String: Any])
func onFunctionCallResult(_ sessionId:String,_ functionCall: [String: Any], _ functionCallResult: [String: Any])
}
@ -59,9 +59,16 @@ public class ChatApiService: NSObject {
private var currentMessages: [[String: Any]] = []
private var toolCalls: [Int: ToolCallInfo] = [:]
private var currSessionId = ""
// 网络监控 - 延迟初始化
private var networkMonitor: NetworkStateMonitor!
// MARK: - 初始化
public override init() {
super.init()
// 在这里初始化网络监控
networkMonitor = NetworkStateMonitor()
networkMonitor.initialize(listener: self)
}
// MARK: - 公共方法
@ -274,8 +281,12 @@ public func sendMessageStream(_ sessionId: String, messages: [[String: Any]]) {
currSessionId = sessionId
guard isInitialized && !apiKey.isEmpty, let openAI = openAI else {
let error = ChatApiException("ChatAPI服务未初始化")
getSessionCallback(sessionId)?.onError(sessionId, error)
getSessionCallback(sessionId)?.onError(sessionId,1002,"服务未初始化")
return
}
//网络检测
if (!checkNetworkStatus()) {
getSessionCallback(sessionId)?.onError(sessionId,1000,"当前网络不可用")
return
}
@ -398,9 +409,9 @@ public func sendMessageStream(_ sessionId: String, messages: [[String: Any]]) {
} catch {
if sessionId == self.currSessionId {
let errorMessage = "流式请求失败: \(error.localizedDescription)"
let chatApiError = ChatApiException(errorMessage)
self.getSessionCallback(sessionId)?.onError(sessionId, chatApiError)
// let errorMessage = "流式请求失败: \(error.localizedDescription)"
// let chatApiError = ChatApiException(errorMessage)
self.getSessionCallback(sessionId)?.onError(sessionId, 1001,"流式请求失败: \(error.localizedDescription)")
self.clearCurrentSession()
}
}
@ -433,7 +444,6 @@ private func roughTokenCount(text: String) -> Int {
}
/// 处理工具调用
/// 处理工具调用
private func processToolCalls(_ sessionId:String) async -> Bool {
// 验证工具调用集合不为空
if toolCalls.isEmpty {
@ -467,8 +477,6 @@ private func processToolCalls(_ sessionId:String) async -> Bool {
// 转换为JSON对象格式
let jsonFunctionCall = convertMapToJsonObject(functionCall)
getSessionCallback(sessionId)?.onFunctionCall(sessionId, jsonFunctionCall)
// 通知上层工具调用事件
// getSessionCallback(sessionId)?.onFunctionCall(sessionId, functionCall)
// 在后台处理工具调用
Task { [weak self] in
@ -479,72 +487,86 @@ private func processToolCalls(_ sessionId:String) async -> Bool {
// 解析参数
let args = try self.parseJsonArguments(firstToolCall.arguments)
// 通过MCP客户端处理工具调用
// 通过MCP客户端处理工具调用,添加8秒超时
var result: [String: Any] = [:]
var isError = false
var errorStr = ""
if let client = self.mcpClient {
// 调用MCP工具
if let toolResult = await client.callTool(name: firstToolCall.name, arguments: args) {
print("[Session: \(sessionId)] AI调用工具结果: \(firstToolCall.name), 参数: \(args), 结果: \(toolResult)")
// 统一结果格式
if toolResult["isError"] as? Bool == true {
if let content = toolResult["content"] as? [[String: Any]],
let firstContent = content.first,
let errorText = firstContent["text"] as? String {
errorStr = errorText
result = ["context": errorText]
// 检查网络状态
if !self.checkNetworkStatus() {
print("[Session: \(sessionId)] 网络连接不可用,工具调用失败")
result = ["context": "网络连接不可用,请检查网络设置"]
isError = true
} else {
// 使用8秒超时调用MCP工具
do {
let toolResult = try await withTimeout(seconds: 8.0) {
return await client.callTool(name: firstToolCall.name, arguments: args)
}
if let toolResult = toolResult {
print("[Session: \(sessionId)] AI调用工具结果: \(firstToolCall.name), 参数: \(args), 结果: \(toolResult)")
// 统一结果格式
if toolResult["isError"] as? Bool == true {
if let content = toolResult["content"] as? [[String: Any]],
let firstContent = content.first,
let errorText = firstContent["text"] as? String {
errorStr = errorText
result = ["context": errorText]
} else {
result = ["context": "Tool execution failed"]
}
} else if let context = toolResult["context"] {
// 本地函数结果
result = ["context": context]
} else if let content = toolResult["content"] as? [[String: Any]],
let firstContent = content.first,
let text = firstContent["text"] as? String {
// MCP工具结果
result = ["context": text]
} else {
print("[Session: \(sessionId)] MCP工具调用返回无法解析的结果")
result = ["context": "Tool call failed"]
}
} else {
result = ["context": "Tool execution failed"]
print("[Session: \(sessionId)] MCP工具调用返回nil")
result = ["context": "Tool call failed"]
}
} else if let context = toolResult["context"] {
// 本地函数结果
result = ["context": context]
} else if let content = toolResult["content"] as? [[String: Any]],
let firstContent = content.first,
let text = firstContent["text"] as? String {
// MCP工具结果
result = ["context": text]
} else {
print("[Session: \(sessionId)] MCP工具调用返回无法解析的结果")
result = ["context": "Tool call failed"]
} catch is TimeoutError {
print("[Session: \(sessionId)] 工具调用超时: \(firstToolCall.name)")
result = ["context": "工具调用超时,请重试"]
isError = true
} catch {
print("[Session: \(sessionId)] 工具调用异常: \(error.localizedDescription)")
result = ["context": "工具调用失败: \(error.localizedDescription)"]
isError = true
}
} else {
print("[Session: \(sessionId)] MCP工具调用返回nil")
result = ["context": "Tool call failed"]
}
} else {
// 工具不存在
result = ["context": "Tool not found: \(firstToolCall.name)"]
}
// 再次检查会话是否仍然有效
if sessionId == self.currSessionId {
// if (firstToolCall.name != "set_user_profile_field"){
// 处理结果
if (!isError){
self.getSessionCallback(sessionId)?.onFunctionCallResult(
sessionId,
functionCall,
result
)
}else{
self.getSessionCallback(sessionId)?.onError(
sessionId,
ChatApiException(errorStr)
)
}
// 将结果发送回OpenAI继续对话
await self.sendFunctionCallResultInternal(
sessionId: sessionId,
messages: self.currentMessages,
functionCall: functionCall,
functionResult: self.jsonToString(result) ?? "{}"
// 处理结果
if (!isError){
self.getSessionCallback(sessionId)?.onFunctionCallResult(
sessionId,
functionCall,
result
)
// }else{
// self.getSessionCallback(sessionId)?.onComplete(sessionId)
// self.clearCurrentSession()
// }
}
// 将结果发送回OpenAI继续对话
await self.sendFunctionCallResultInternal(
sessionId: sessionId,
messages: self.currentMessages,
functionCall: functionCall,
functionResult: self.jsonToString(result) ?? "{}"
)
}
}
} catch {
@ -995,6 +1017,65 @@ private func convertMapToJsonObject(_ map: [String: Any]) -> [String: Any] {
return chatTools
}
/**
* 检查网络状态
* @return 网络是否可用
*/
private func checkNetworkStatus() -> Bool {
// 直接使用 NetworkStateMonitor 的 checkNetworkStatus 方法
let isAvailable = networkMonitor.checkNetworkStatus()
return isAvailable
}
/// 分析网络错误并返回错误代码和消息
private func analyzeNetworkError(_ error: Error) -> (code: Int, message: String) {
let errorDescription = error.localizedDescription.lowercased()
// 使用 NetworkStateMonitor 的网络错误检测
if networkMonitor.isNetworkRelatedError(reason: errorDescription, errorDetails: "") {
if errorDescription.contains("timeout") {
return (1004, "网络请求超时")
} else if errorDescription.contains("connection") {
if errorDescription.contains("abort") {
return (1005, "网络连接中断")
} else if errorDescription.contains("reset") {
return (1006, "网络连接重置")
} else {
return (1007, "网络连接失败")
}
} else if errorDescription.contains("dns") || errorDescription.contains("host") {
return (1008, "域名解析失败")
} else if errorDescription.contains("ssl") || errorDescription.contains("tls") {
return (1009, "SSL/TLS 连接失败")
} else {
return (1010, "网络错误: \(error.localizedDescription)")
}
}
// 非网络错误
return (1003, "流式请求失败: \(error.localizedDescription)")
}
}
// MARK: - NetworkStateListener
extension ChatApiService: NetworkStateMonitor.NetworkStateListener {
public func onNetworkAvailable() {
print("[ChatApiService] 网络连接恢复")
// 可以在这里通知上层网络恢复
}
public func onNetworkLost() {
print("[ChatApiService] 网络连接丢失")
// 网络丢失时中止当前请求
if !currSessionId.isEmpty {
getSessionCallback(currSessionId)?.onError(currSessionId, 1001, "网络连接丢失")
abortCurrentSession()
}
}
}
/// 本地函数处理器
@ -1009,3 +1090,4 @@ private class LocalFunctionHandler: FunctionHandler {
return "LOCAL_FUNCTION:\(functionName)"
}
}

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

@ -450,7 +450,7 @@ public class MCPClient {
private func initializeSystemFunctions() {
do {
let handler = SystemFunctionHandler()
let handler: SystemFunctionHandler = SystemFunctionHandler()
handler.registerAllFunctions(client: self)
} catch {
@ -460,8 +460,8 @@ public class MCPClient {
private func initializeMusicSystemFunction() {
do {
// 初始化QQ音乐功能
let qqMusicHandler = MusicSystemFunctionHandler()
qqMusicHandler.registerAllFunctions(client: self)
let musicHandler: MusicSystemFunctionHandler = MusicSystemFunctionHandler()
musicHandler.registerAllFunctions(client: self)
} catch {
// 静默处理错误
}
@ -469,7 +469,7 @@ public class MCPClient {
private func initializeQQMusicSystemFunction() {
do {
// 初始化QQ音乐功能
let qqMusicHandler = QQMusicSystemFunctionHandler()
let qqMusicHandler: QQMusicSystemFunctionHandler = QQMusicSystemFunctionHandler()
qqMusicHandler.registerAllFunctions(client: self)
} catch {
// 静默处理错误
@ -479,26 +479,26 @@ public class MCPClient {
public func connectToSSE(mcpConfigJson: String) async throws -> Bool {
// 清除现有连接
for client in subClients.values {
for client: MCPSubClient in subClients.values {
await client.close()
}
subClients.removeAll()
guard let data = mcpConfigJson.data(using: .utf8) else {
guard let data: Data = mcpConfigJson.data(using: .utf8) else {
return false
}
guard let config = try? JSONSerialization.jsonObject(with: data) as? [String: Any] else {
guard let config: [String : Any] = try? JSONSerialization.jsonObject(with: data) as? [String: Any] else {
return false
}
guard let mcpServers = config["mcpServers"] as? [String: Any] else {
guard let mcpServers: [String : Any] = config["mcpServers"] as? [String: Any] else {
return false
}
var connectedCount = 0
for (serverId, serverConfig) in mcpServers {
guard let configDict = serverConfig as? [String: Any] else {
for (serverId: String, serverConfig) in mcpServers {
guard let configDict: [String : Any] = serverConfig as? [String: Any] else {
continue
}

179
local_plugins/chat_api/ios/chat_api/Sources/chat_api/NetworkStateMonitor.swift

@ -0,0 +1,179 @@
import Foundation
import SystemConfiguration
import Network
/// 网络状态监听器(iOS)
public final class NetworkStateMonitor {
private let tag = "NetworkStateMonitor"
// iOS 12+
private var monitor: NWPathMonitor?
private let monitorQueue = DispatchQueue(label: "NetworkStateMonitor.queue")
// iOS 10–11 回退
private var reachability: SCNetworkReachability?
private var reachabilityQueue: DispatchQueue?
private(set) var isNetworkAvailable: Bool = true
private weak var networkStateListener: NetworkStateListener?
/// 网络状态监听接口
public protocol NetworkStateListener: AnyObject {
func onNetworkAvailable()
func onNetworkLost()
}
public init() {}
/// 初始化网络状态监听
public func initialize(listener: NetworkStateListener? = nil) {
networkStateListener = listener
initNetworkMonitoring()
}
/// 检查当前网络状态
/// - Returns: true 表示网络可用
public func checkNetworkStatus() -> Bool {
if #available(iOS 12.0, *) {
// NWPathMonitor 只在回调里更新状态,因此这里返回缓存的 isNetworkAvailable
return isNetworkAvailable
} else {
// 使用 Reachability 即时检测
var zeroAddress = sockaddr_in()
zeroAddress.sin_len = UInt8(MemoryLayout<sockaddr_in>.size)
zeroAddress.sin_family = sa_family_t(AF_INET)
let ref = withUnsafePointer(to: &zeroAddress) {
$0.withMemoryRebound(to: sockaddr.self, capacity: 1) {
SCNetworkReachabilityCreateWithAddress(nil, $0)
}
}
guard let reachRef = ref else { return false }
var flags = SCNetworkReachabilityFlags()
if !SCNetworkReachabilityGetFlags(reachRef, &flags) { return false }
return Self.isReachable(flags: flags)
}
}
/// 释放资源
public func dispose() {
if #available(iOS 12.0, *) {
monitor?.cancel()
monitor = nil
} else {
if let reachability = reachability {
SCNetworkReachabilitySetDispatchQueue(reachability, nil)
SCNetworkReachabilitySetCallback(reachability, nil, nil)
}
reachability = nil
reachabilityQueue = nil
}
networkStateListener = nil
}
/// 判断是否为网络相关错误(关键字简单匹配)
public func isNetworkRelatedError(reason: String, errorDetails: String) -> Bool {
let networkErrorKeywords = [
"network","connection","timeout","unreachable",
"dns","socket","ssl","tls","certificate",
"网络","连接","超时","不可达"
]
let combined = (reason + " " + errorDetails).lowercased()
return networkErrorKeywords.contains { combined.contains($0.lowercased()) }
}
}
// MARK: - Private
private extension NetworkStateMonitor {
func initNetworkMonitoring() {
if #available(iOS 12.0, *) {
let monitor = NWPathMonitor()
self.monitor = monitor
monitor.pathUpdateHandler = { [weak self] path in
guard let self = self else { return }
let hasInternet = (path.status == .satisfied)
// 只有状态变化时才回调
if hasInternet && !self.isNetworkAvailable {
self.isNetworkAvailable = true
self.networkStateListener?.onNetworkAvailable()
Self.log(self.tag, "网络连接可用")
} else if !hasInternet && self.isNetworkAvailable {
self.isNetworkAvailable = false
self.networkStateListener?.onNetworkLost()
Self.log(self.tag, "网络连接丢失")
} else {
Self.log(self.tag, "网络能力变化,是否有网络: \(hasInternet)")
}
}
monitor.start(queue: monitorQueue)
} else {
// iOS 10–11 使用 Reachability 回退
var zeroAddress = sockaddr_in()
zeroAddress.sin_len = UInt8(MemoryLayout<sockaddr_in>.size)
zeroAddress.sin_family = sa_family_t(AF_INET)
reachability = withUnsafePointer(to: &zeroAddress) {
$0.withMemoryRebound(to: sockaddr.self, capacity: 1) {
SCNetworkReachabilityCreateWithAddress(nil, $0)
}
}
guard let reachability = reachability else {
Self.log(tag, "初始化网络监听失败: Reachability 创建失败")
return
}
let queue = DispatchQueue(label: "NetworkStateMonitor.reachability")
reachabilityQueue = queue
var context = SCNetworkReachabilityContext(
version: 0,
info: UnsafeMutableRawPointer(Unmanaged.passUnretained(self).toOpaque()),
retain: nil,
release: nil,
copyDescription: nil
)
let callback: SCNetworkReachabilityCallBack = { (_, flags, info) in
guard let info = info else { return }
let monitor = Unmanaged<NetworkStateMonitor>.fromOpaque(info).takeUnretainedValue()
let hasInternet = NetworkStateMonitor.isReachable(flags: flags)
if hasInternet && !monitor.isNetworkAvailable {
monitor.isNetworkAvailable = true
monitor.networkStateListener?.onNetworkAvailable()
NetworkStateMonitor.log(monitor.tag, "网络连接可用")
} else if !hasInternet && monitor.isNetworkAvailable {
monitor.isNetworkAvailable = false
monitor.networkStateListener?.onNetworkLost()
NetworkStateMonitor.log(monitor.tag, "网络连接丢失")
} else {
NetworkStateMonitor.log(monitor.tag, "网络能力变化,是否有网络: \(hasInternet)")
}
}
if SCNetworkReachabilitySetCallback(reachability, callback, &context),
SCNetworkReachabilitySetDispatchQueue(reachability, queue) {
// 触发一次初始查询
var flags = SCNetworkReachabilityFlags()
if SCNetworkReachabilityGetFlags(reachability, &flags) {
self.isNetworkAvailable = Self.isReachable(flags: flags)
}
} else {
Self.log(tag, "初始化网络监听失败: 设置回调/队列失败")
SCNetworkReachabilitySetCallback(reachability, nil, nil)
SCNetworkReachabilitySetDispatchQueue(reachability, nil)
self.reachability = nil
}
}
}
static func isReachable(flags: SCNetworkReachabilityFlags) -> Bool {
// 经典 Reachability 判定
let reachable = flags.contains(.reachable)
let requiresConnection = flags.contains(.connectionRequired)
let canConnectAutomatically = flags.contains(.connectionOnTraffic) || flags.contains(.connectionOnDemand)
let canConnectWithoutUser = canConnectAutomatically && !flags.contains(.interventionRequired)
return reachable && (!requiresConnection || canConnectWithoutUser)
}
static func log(_ tag: String, _ message: String) {
#if DEBUG
print("[\(tag)] \(message)")
#endif
}
}
Loading…
Cancel
Save