From c7ab1fe5f39ff4bcc6b8b6609fbb4740040f0fcf Mon Sep 17 00:00:00 2001 From: wolfplus Date: Sat, 7 Jun 2025 13:08:15 +0100 Subject: [PATCH 1/3] add --- .../agent_service/AgentService.kt | 3 +- local_plugins/chat_api/android/README.md | 70 --- .../chat_api/android/build.gradle.kts | 12 +- .../yunqiinnovation/chat_api/ChatApiPlugin.kt | 291 ++++------ .../chat_api/ChatApiService.kt | 536 ++++++++++++++++-- .../chat_api/CustomSseClientTransport.kt | 271 +++++++++ .../com/yunqiinnovation/chat_api/MCPClient.kt | 360 ++++++++++++ .../yunqiinnovation/chat_api/MCPSubClient.kt | 329 +++++++++++ .../chat_api/SystemFunctionHandler.kt | 471 +++++++++++++++ .../com/yunqiinnovation/chat_api/Utils.kt | 57 ++ 10 files changed, 2111 insertions(+), 289 deletions(-) delete mode 100644 local_plugins/chat_api/android/README.md create mode 100644 local_plugins/chat_api/android/src/main/kotlin/com/yunqiinnovation/chat_api/CustomSseClientTransport.kt create mode 100644 local_plugins/chat_api/android/src/main/kotlin/com/yunqiinnovation/chat_api/MCPClient.kt create mode 100644 local_plugins/chat_api/android/src/main/kotlin/com/yunqiinnovation/chat_api/MCPSubClient.kt create mode 100644 local_plugins/chat_api/android/src/main/kotlin/com/yunqiinnovation/chat_api/SystemFunctionHandler.kt create mode 100644 local_plugins/chat_api/android/src/main/kotlin/com/yunqiinnovation/chat_api/Utils.kt diff --git a/local_plugins/agent_service/android/src/main/kotlin/com/yunqiinnovation/agent_service/AgentService.kt b/local_plugins/agent_service/android/src/main/kotlin/com/yunqiinnovation/agent_service/AgentService.kt index d9f358d34..d560c4684 100644 --- a/local_plugins/agent_service/android/src/main/kotlin/com/yunqiinnovation/agent_service/AgentService.kt +++ b/local_plugins/agent_service/android/src/main/kotlin/com/yunqiinnovation/agent_service/AgentService.kt @@ -152,7 +152,8 @@ object AgentService : CoroutineScope { config["openaiApiKey"]?.toString() ?: "", config["openaiBaseUrl"]?.toString() ?: "", config["openaiModel"]?.toString() ?: "", - config["mcpServer"]?.toString() ?: "" + "" + //config["mcpServer"]?.toString() ?: "" ) // 加载最近的聊天记录 diff --git a/local_plugins/chat_api/android/README.md b/local_plugins/chat_api/android/README.md deleted file mode 100644 index 99c528b5f..000000000 --- a/local_plugins/chat_api/android/README.md +++ /dev/null @@ -1,70 +0,0 @@ -# ChatAPI Android 实现 - -基于 [openai-kotlin 4.0.1](https://github.com/aallam/openai-kotlin/tree/4.0.1) 的 ChatAPI 插件 Android 实现。 - -## 特性 - -- ✅ 基于 openai-kotlin 4.0.1 的 OpenAI API 客户端 -- ✅ 与 iOS 版本接口完全一致 -- ✅ 支持流式和非流式对话 -- ✅ 支持多模态内容(文本和图像识别) -- ✅ 工具调用架构预留(MCP 功能暂时留空) -- ✅ 协程支持和异步处理 -- ✅ 完整的错误处理和取消机制 - -## 依赖 - -- openai-kotlin 4.0.1 -- ktor-client-okhttp 2.3.2 -- kotlinx-coroutines -- gson - -## 主要文件 - -- `ChatApiPlugin.kt` - Flutter 插件桥接层 -- `ChatApiService.kt` - 核心 ChatAPI 服务实现 -- `StreamCallback.kt` - 流式回调接口 - -## 与 iOS 版本的一致性 - -Android 实现完全遵循 iOS 版本的接口定义: - -1. **初始化方法**: `initialize(apiKey, baseUrl, model, mcpServer)` -2. **消息创建**: `createUserMessage()`, `createAssistantMessage()` -3. **消息发送**: `sendMessage()` (非流式), `sendMessageStream()` (流式) -4. **函数注册**: `registerFunction()` -5. **MCP 相关**: 接口保留,功能暂时留空 -6. **流式事件**: 完全一致的事件类型和回调 - -## MCP 功能状态 - -MCP (Model Context Protocol) 相关功能在 Android 版本中暂时留空,包括: - -- `initializeMcpClient()` - 返回 false -- `isMcpInitialized()` - 返回 false -- `handleMcpToolCall()` - 返回占位符消息 -- 工具调用会通知上层但不会实际执行 - -## 使用示例 - -```kotlin -// 初始化 -chatApiService.initialize( - apiKey = "your-api-key", - baseUrl = "https://api.openai.com/v1/", - model = "gpt-3.5-turbo", - mcpServer = "" -) - -// 流式对话 -chatApiService.sendMessageStream(listOf( - mapOf("role" to "user", "content" to "Hello!") -)) -``` - -## 注意事项 - -1. 需要网络权限 `INTERNET` 和 `ACCESS_NETWORK_STATE` -2. 所有网络请求在后台线程执行 -3. 流式回调在主线程触发 -4. 支持取消正在进行的请求 \ No newline at end of file diff --git a/local_plugins/chat_api/android/build.gradle.kts b/local_plugins/chat_api/android/build.gradle.kts index 1eb929ecf..f9f8ddefc 100644 --- a/local_plugins/chat_api/android/build.gradle.kts +++ b/local_plugins/chat_api/android/build.gradle.kts @@ -42,7 +42,17 @@ dependencies { // OpenAI Kotlin 4.0.1 implementation("com.aallam.openai:openai-client:4.0.1") - implementation("io.ktor:ktor-client-okhttp:2.3.2") + + // MCP Kotlin SDK - 0.5.0官方版本,包含所有必要的传输层支持 + implementation("io.modelcontextprotocol:kotlin-sdk:0.5.0") + + // Ktor HTTP Client with OkHttp engine - 升级到3.1.2版本 + implementation("io.ktor:ktor-client-core:3.1.2") + implementation("io.ktor:ktor-client-okhttp:3.1.2") + + // OkHttp + implementation("com.squareup.okhttp3:okhttp:4.11.0") + implementation("com.squareup.okhttp3:logging-interceptor:4.11.0") // JSON处理 implementation("org.jetbrains.kotlinx:kotlinx-serialization-json:1.6.0") diff --git a/local_plugins/chat_api/android/src/main/kotlin/com/yunqiinnovation/chat_api/ChatApiPlugin.kt b/local_plugins/chat_api/android/src/main/kotlin/com/yunqiinnovation/chat_api/ChatApiPlugin.kt index bc7801951..c369deac9 100644 --- a/local_plugins/chat_api/android/src/main/kotlin/com/yunqiinnovation/chat_api/ChatApiPlugin.kt +++ b/local_plugins/chat_api/android/src/main/kotlin/com/yunqiinnovation/chat_api/ChatApiPlugin.kt @@ -1,7 +1,6 @@ package com.yunqiinnovation.chat_api -import android.os.Handler -import android.os.Looper +import androidx.annotation.NonNull import io.flutter.embedding.engine.plugins.FlutterPlugin import io.flutter.plugin.common.EventChannel import io.flutter.plugin.common.MethodCall @@ -9,6 +8,7 @@ import io.flutter.plugin.common.MethodChannel import io.flutter.plugin.common.MethodChannel.MethodCallHandler import io.flutter.plugin.common.MethodChannel.Result import kotlinx.coroutines.* +import android.util.Log /** * ChatApiPlugin @@ -17,215 +17,144 @@ import kotlinx.coroutines.* * 与 iOS 版本接口完全一致 */ class ChatApiPlugin : FlutterPlugin, MethodCallHandler, EventChannel.StreamHandler { - private lateinit var methodChannel: MethodChannel + private lateinit var channel: MethodChannel private lateinit var eventChannel: EventChannel private var eventSink: EventChannel.EventSink? = null - private val chatApiService = ChatApiService() - private val mainHandler = Handler(Looper.getMainLooper()) + private var chatApiService: ChatApiService? = null - // 协程作用域 + // 插件协程作用域 private val pluginScope = CoroutineScope(Dispatchers.Main + SupervisorJob()) - override fun onAttachedToEngine(flutterPluginBinding: FlutterPlugin.FlutterPluginBinding) { - methodChannel = MethodChannel(flutterPluginBinding.binaryMessenger, "com.yunqiinnovation.chat_api/methods") - methodChannel.setMethodCallHandler(this) + override fun onAttachedToEngine(@NonNull flutterPluginBinding: FlutterPlugin.FlutterPluginBinding) { + channel = MethodChannel(flutterPluginBinding.binaryMessenger, "chat_api") + channel.setMethodCallHandler(this) eventChannel = EventChannel(flutterPluginBinding.binaryMessenger, "com.yunqiinnovation.chat_api/events") eventChannel.setStreamHandler(this) + // 初始化ChatAPI服务 + chatApiService = ChatApiService(flutterPluginBinding.applicationContext) + // 设置流式回调 - chatApiService.setStreamCallback(object : com.yunqiinnovation.chat_api.StreamCallback { + chatApiService?.setStreamCallback(object : StreamCallback { override fun onToken(token: String) { - sendEvent("token", token) + channel.invokeMethod("onToken", token) } override fun onComplete() { - sendEvent("complete", null) + channel.invokeMethod("onComplete", null) } override fun onError(error: Exception) { - sendEvent("error", error.message ?: "未知错误") + channel.invokeMethod("onError", error.message) } override fun onFunctionCall(functionCall: org.json.JSONObject) { - try { - val jsonString = functionCall.toString() - sendEvent("functionCall", jsonString) - } catch (e: Exception) { - sendEvent("error", "Failed to serialize function call: ${e.message}") - } + channel.invokeMethod("onFunctionCall", functionCall.toString()) } override fun onFunctionCallResult(functionCall: org.json.JSONObject, functionCallResult: org.json.JSONObject) { - try { - val meta = org.json.JSONObject().apply { - put("functionCall", functionCall) - put("functionCallResult", functionCallResult) - } - val metaString = meta.toString() - val context = functionCallResult.optString("context", "") - - sendEvent("functionCall", context, metaString) - } catch (e: Exception) { - sendEvent("error", "Failed to serialize function call result: ${e.message}") - } + channel.invokeMethod("onFunctionCallResult", mapOf( + "functionCall" to functionCall.toString(), + "functionCallResult" to functionCallResult.toString() + )) } }) } - override fun onDetachedFromEngine(binding: FlutterPlugin.FlutterPluginBinding) { - methodChannel.setMethodCallHandler(null) + override fun onDetachedFromEngine(@NonNull binding: FlutterPlugin.FlutterPluginBinding) { + channel.setMethodCallHandler(null) eventChannel.setStreamHandler(null) + + // 清理资源 pluginScope.cancel() + chatApiService?.closeMcpClient() + chatApiService = null } - override fun onMethodCall(call: MethodCall, result: Result) { + override fun onMethodCall(@NonNull call: MethodCall, @NonNull result: Result) { + Log.d("ChatApiPlugin", "onMethodCall: ${call.method}") when (call.method) { - "initialize" -> handleInitialize(call, result) - "createUserMessage" -> handleCreateUserMessage(call, result) - "createAssistantMessage" -> handleCreateAssistantMessage(call, result) - "sendMessage" -> handleSendMessage(call, result) - "sendMessageStream" -> handleSendMessageStream(call, result) - "sendFunctionCallResult" -> handleSendFunctionCallResult(call, result) - "cancelCurrentStream" -> handleCancelCurrentStream(result) - "registerFunction" -> handleRegisterFunction(call, result) - "initializeMcpClient" -> handleInitializeMcpClient(call, result) - "isMcpInitialized" -> handleIsMcpInitialized(result) - "closeMcpClient" -> handleCloseMcpClient(result) - "handleMcpToolCall" -> handleMcpToolCall(call, result) - else -> result.notImplemented() - } - } - - // MARK: - Method Handlers - - private fun handleInitialize(call: MethodCall, result: Result) { - val apiKey = call.argument("apiKey") - if (apiKey.isNullOrEmpty()) { - result.error("INVALID_ARGS", "缺少必要参数", null) - return - } - - val baseUrl = call.argument("baseUrl") ?: "" - val model = call.argument("model") ?: "" - val mcpServer = call.argument("mcpServer") ?: "" - - val success = chatApiService.initialize(apiKey, baseUrl, model, mcpServer) - result.success(success) - } - - private fun handleCreateUserMessage(call: MethodCall, result: Result) { - val content = call.argument("content") - if (content.isNullOrEmpty()) { - result.error("INVALID_ARGS", "缺少必要参数", null) - return - } - - val message = chatApiService.createUserMessage(content) - result.success(message) - } - - private fun handleCreateAssistantMessage(call: MethodCall, result: Result) { - val content = call.argument("content") - if (content.isNullOrEmpty()) { - result.error("INVALID_ARGS", "缺少必要参数", null) - return - } - - val message = chatApiService.createAssistantMessage(content) - result.success(message) - } - - private fun handleSendMessage(call: MethodCall, result: Result) { - @Suppress("UNCHECKED_CAST") - val messages = call.argument>>("messages") - if (messages == null) { - result.error("INVALID_ARGS", "缺少必要参数", null) - return - } - - pluginScope.launch { - try { - val response = chatApiService.sendMessage(messages) - mainHandler.post { - result.success(response) + "initialize" -> { + val apiKey = call.argument("apiKey") ?: "" + val baseUrl = call.argument("baseUrl") ?: "" + val model = call.argument("model") ?: "" + val mcpServer = call.argument("mcpServer") ?: "" + val success = chatApiService?.initialize(apiKey, baseUrl, model, mcpServer) ?: false + result.success(success) + } + "chatCompletionStream" -> { + val messages = call.argument>>("messages") ?: emptyList() + val tool = call.argument("tool") ?: false + pluginScope.launch { + chatApiService?.chatCompletionStream(messages, tool) + } + result.success(true) + } + "cancelChatStream" -> { + chatApiService?.cancelChatStream() + result.success(true) + } + "processImage" -> { + val imagePath = call.argument("imagePath") ?: "" + val prompt = call.argument("prompt") ?: "" + val maxWidth = call.argument("maxWidth") ?: 2048.0 + val detail = call.argument("detail") ?: "auto" + val base64Image = chatApiService?.processImage(imagePath, prompt, maxWidth, detail) + result.success(base64Image) + } + "initializeMcpClient" -> { + val serverUrl = call.argument("serverUrl") ?: "" + val success = chatApiService?.initializeMcpClient(serverUrl) ?: false + result.success(success) + } + "isMcpInitialized" -> { + val initialized = chatApiService?.isMcpInitialized() ?: false + result.success(initialized) + } + "closeMcpClient" -> { + chatApiService?.closeMcpClient() + result.success(true) + } + "getToolMaps" -> { + val mcpClient = chatApiService?.mcpClient + if (mcpClient != null) { + val toolMaps = mcpClient.getToolMaps() + result.success(toolMaps) + } else { + result.success(emptyList>()) } - } catch (e: Exception) { - mainHandler.post { - result.error("SEND_ERROR", e.message, null) + } + "hasToolWithName" -> { + val name = call.argument("name") ?: "" + val mcpClient = chatApiService?.mcpClient + val hasTool = mcpClient?.hasToolWithName(name) ?: false + result.success(hasTool) + } + "getToolType" -> { + val name = call.argument("name") ?: "" + val mcpClient = chatApiService?.mcpClient + val toolType = mcpClient?.getToolType(name) + result.success(toolType?.name) + } + "registerFunction" -> { + val name = call.argument("name") ?: "" + val description = call.argument("description") ?: "" + val parametersJson = call.argument("parameters") ?: "{}" + val parameters = try { + org.json.JSONObject(parametersJson) + } catch (e: Exception) { + org.json.JSONObject() } + val success = chatApiService?.registerFunction(name, description, parameters) ?: false + result.success(success) + } + else -> { + result.notImplemented() } } } - - private fun handleSendMessageStream(call: MethodCall, result: Result) { - @Suppress("UNCHECKED_CAST") - val messages = call.argument>>("messages") - if (messages == null) { - result.error("INVALID_ARGS", "缺少必要参数", null) - return - } - - chatApiService.sendMessageStream(messages) - result.success(true) - } - - private fun handleSendFunctionCallResult(call: MethodCall, result: Result) { - result.error("DEPRECATED", "sendFunctionCallResult已废弃,工具调用结果现在自动处理", null) - } - - private fun handleCancelCurrentStream(result: Result) { - val success = chatApiService.cancelCurrentStream() - result.success(success) - } - - private fun handleRegisterFunction(call: MethodCall, result: Result) { - val name = call.argument("name") - val description = call.argument("description") - @Suppress("UNCHECKED_CAST") - val parameters = call.argument>("parameters") - - if (name.isNullOrEmpty() || description.isNullOrEmpty() || parameters == null) { - result.error("INVALID_ARGS", "缺少必要参数", null) - return - } - - val success = chatApiService.registerFunction(name, description, parameters) - result.success(success) - } - - private fun handleInitializeMcpClient(call: MethodCall, result: Result) { - val serverUrl = call.argument("serverUrl") - if (serverUrl.isNullOrEmpty()) { - result.error("INVALID_ARGS", "缺少必要参数", null) - return - } - - // MCP功能暂时留空 - result.success(false) - } - - private fun handleIsMcpInitialized(result: Result) { - // MCP功能暂时留空 - result.success(false) - } - - private fun handleCloseMcpClient(result: Result) { - // MCP功能暂时留空 - result.success(true) - } - - private fun handleMcpToolCall(call: MethodCall, result: Result) { - val functionCallJson = call.argument("functionCall") - if (functionCallJson.isNullOrEmpty()) { - result.error("INVALID_ARGS", "缺少必要参数", null) - return - } - - // MCP功能暂时留空 - result.success("MCP功能暂未实现") - } // MARK: - Event Stream Handler @@ -238,15 +167,13 @@ class ChatApiPlugin : FlutterPlugin, MethodCallHandler, EventChannel.StreamHandl } private fun sendEvent(type: String, content: Any?, meta: String? = null) { - mainHandler.post { - eventSink?.let { sink -> - val event = mutableMapOf( - "type" to type - ) - content?.let { event["content"] = it } - meta?.let { event["meta"] = it } - sink.success(event) - } + eventSink?.let { sink -> + val event = mutableMapOf( + "type" to type + ) + content?.let { event["content"] = it } + meta?.let { event["meta"] = it } + sink.success(event) } } } \ No newline at end of file diff --git a/local_plugins/chat_api/android/src/main/kotlin/com/yunqiinnovation/chat_api/ChatApiService.kt b/local_plugins/chat_api/android/src/main/kotlin/com/yunqiinnovation/chat_api/ChatApiService.kt index dcf8b01bb..d54ce8298 100644 --- a/local_plugins/chat_api/android/src/main/kotlin/com/yunqiinnovation/chat_api/ChatApiService.kt +++ b/local_plugins/chat_api/android/src/main/kotlin/com/yunqiinnovation/chat_api/ChatApiService.kt @@ -20,6 +20,9 @@ import kotlin.math.min import kotlin.math.sqrt import kotlin.time.Duration.Companion.seconds import android.util.Log +import org.json.JSONObject +import android.os.Handler +import android.os.Looper /** * ChatAPI服务异常 @@ -83,9 +86,21 @@ class ChatApiService(private val context: android.content.Context? = null) : Cor private var visionModel = "gpt-4-vision-preview" private var isInitialized = false + // Handler for main thread + private val mainHandler = Handler(Looper.getMainLooper()) + // OpenAI 客户端 private var openAI: OpenAI? = null + // MCP 客户端 + private var _mcpClient: MCPClient? = null + + // 公开的MCP客户端访问器 + val mcpClient: MCPClient? + get() = _mcpClient + + private var mcpConfigJson: String? = null + // 流式请求相关 private var currentStreamJob: Job? = null private var streamCallback: StreamCallback? = null @@ -109,6 +124,13 @@ class ChatApiService(private val context: android.content.Context? = null) : Cor * 初始化ChatAPI服务 */ fun initialize(apiKey: String, baseUrl: String, model: String, mcpServer: String): Boolean { + Log.d("ChatApiService", "=== ChatApiService 初始化开始 ===") + Log.d("ChatApiService", "参数检查:") + Log.d("ChatApiService", " - apiKey: ${if (apiKey.isNotEmpty()) "已提供(${apiKey.length}字符)" else "未提供"}") + Log.d("ChatApiService", " - baseUrl: '$baseUrl'") + Log.d("ChatApiService", " - model: '$model'") + Log.d("ChatApiService", " - mcpServer: '$mcpServer'") + this.apiKey = apiKey if (baseUrl.isNotEmpty()) { this.baseUrl = baseUrl @@ -134,24 +156,33 @@ class ChatApiService(private val context: android.content.Context? = null) : Cor Log.d("ChatApiService", "处理后 baseUrl: $processedBaseUrl") return try { + Log.d("ChatApiService", "开始创建OpenAI配置...") // 创建OpenAI配置 val config = OpenAIConfig( token = apiKey, timeout = Timeout(socket = 60.seconds), host = OpenAIHost(baseUrl = processedBaseUrl) ) + Log.d("ChatApiService", "OpenAI配置创建成功") + Log.d("ChatApiService", "正在创建OpenAI客户端...") openAI = OpenAI(config) + Log.d("ChatApiService", "OpenAI客户端创建成功") - // 初始化MCP客户端 (暂时留空,但保持接口一致) - if (mcpServer.isNotEmpty()) { - // MCP功能暂时留空,但记录服务器地址以便后续实现 - // initializeMcpClient(mcpServer) + // 初始化MCP客户端 + Log.d("ChatApiService", "准备初始化MCP客户端...") + // 保存MCP配置以便后续使用 + mcpConfigJson = mcpServer + // 异步初始化MCP客户端 + launch { + initializeMcpClient(mcpServer) } isInitialized = apiKey.isNotEmpty() + Log.d("ChatApiService", "ChatApiService初始化完成,isInitialized: $isInitialized") true } catch (e: Exception) { + Log.e("ChatApiService", "ChatApiService初始化失败: ${e.message}", e) false } } @@ -200,12 +231,80 @@ class ChatApiService(private val context: android.content.Context? = null) : Cor // 转换消息格式 val chatMessages = convertToChatMessages(messages) + // 获取MCP工具列表 + val tools = mutableListOf() + mcpClient?.getToolMaps()?.forEach { toolMap -> + try { + Log.d("ChatApiService", "处理工具映射: $toolMap") + + val type = toolMap["type"] as? String + if (type == "function") { + @Suppress("UNCHECKED_CAST") + val functionMap = toolMap["function"] as? Map + if (functionMap == null) { + Log.e("ChatApiService", "工具function映射为null") + return@forEach + } + + val name = functionMap["name"] as? String + if (name == null) { + Log.e("ChatApiService", "工具name为null") + return@forEach + } + + val description = functionMap["description"] as? String ?: "" + val parametersMap = functionMap["parameters"] as? Map + + Log.d("ChatApiService", "工具信息 - 名称: $name, 描述: $description") + Log.d("ChatApiService", "参数映射: $parametersMap") + + if (parametersMap == null) { + Log.e("ChatApiService", "工具 $name 的parameters为null") + return@forEach + } + + // 验证parametersMap的基本结构 + if (!parametersMap.containsKey("type")) { + Log.e("ChatApiService", "工具 $name 的parameters缺少type字段") + return@forEach + } + + val parametersJson = gson.toJson(parametersMap) + Log.d("ChatApiService", "参数JSON: $parametersJson") + + try { + val parameters = com.aallam.openai.api.core.Parameters.fromJsonString(parametersJson) + Log.d("ChatApiService", "成功创建Parameters对象") + + tools.add(Tool.function( + name = name, + description = description, + parameters = parameters + )) + + Log.d("ChatApiService", "成功添加工具: $name") + + } catch (e: Exception) { + Log.e("ChatApiService", "创建Parameters对象失败 for 工具 $name: ${e.message}", e) + Log.e("ChatApiService", "失败的参数JSON: $parametersJson") + // 跳过这个工具,继续处理其他工具 + } + } else { + Log.d("ChatApiService", "跳过非function类型的工具: $type") + } + } catch (e: Exception) { + Log.e("ChatApiService", "处理工具映射时出错: ${e.message}", e) + Log.e("ChatApiService", "出错的工具映射: $toolMap") + } + } + // 构建请求 val chatCompletionRequest = ChatCompletionRequest( model = ModelId(currentModel), messages = chatMessages, maxTokens = 2000, - temperature = 0.7 + temperature = 0.7, + tools = if (tools.isNotEmpty()) tools else null ) val result = openAI!!.chatCompletion(chatCompletionRequest) @@ -245,7 +344,14 @@ class ChatApiService(private val context: android.content.Context? = null) : Cor * 发送消息(流式输出) */ fun sendMessageStream(messages: List>) { + Log.d("ChatApiService", "=== sendMessageStream 开始 ===") + Log.d("ChatApiService", "isInitialized: $isInitialized") + Log.d("ChatApiService", "apiKey isEmpty: ${apiKey.isEmpty()}") + Log.d("ChatApiService", "openAI is null: ${openAI == null}") + Log.d("ChatApiService", "messages size: ${messages.size}") + if (!isInitialized || apiKey.isEmpty() || openAI == null) { + Log.e("ChatApiService", "ChatAPI服务未初始化,无法发送消息") streamCallback?.onError(ChatApiException("ChatAPI服务未初始化")) return } @@ -257,47 +363,227 @@ class ChatApiService(private val context: android.content.Context? = null) : Cor // 检查是否包含图片,决定使用哪个模型 val currentModel = if (containsImage(messages)) visionModel else model + Log.d("ChatApiService", "选择的模型: $currentModel (默认模型: $model, 视觉模型: $visionModel)") currentStreamJob = launch { try { + Log.e("ChatApiService", "=== 开始流式请求 ===") + Log.e("ChatApiService", "消息数量: ${messages.size}") + Log.e("ChatApiService", "当前模型: $currentModel") + Log.e("ChatApiService", "MCP客户端状态: ${mcpClient != null}") + // 转换消息格式 - val chatMessages = convertToChatMessages(messages) + Log.e("ChatApiService", "开始转换消息格式...") + val chatMessages = try { + convertToChatMessages(messages) + } catch (e: Exception) { + Log.e("ChatApiService", "ERROR: 转换消息格式失败: ${e.message}", e) + throw e + } + Log.e("ChatApiService", "转换后的聊天消息数量: ${chatMessages.size}") + + // 验证转换后的消息 + chatMessages.forEachIndexed { index, message -> + Log.e("ChatApiService", "消息 $index: role=${message.role}, content类型=${message.content?.javaClass?.simpleName}") + } + + // 获取MCP工具列表 + Log.e("ChatApiService", "开始获取MCP工具列表...") + val tools = mutableListOf() + + val toolMaps = mcpClient?.getToolMaps() + Log.e("ChatApiService", "获取到的工具映射数量: ${toolMaps?.size ?: 0}") + + toolMaps?.forEach { toolMap -> + try { + Log.e("ChatApiService", "处理工具映射: $toolMap") + + val type = toolMap["type"] as? String + if (type == "function") { + @Suppress("UNCHECKED_CAST") + val functionMap = toolMap["function"] as? Map + if (functionMap == null) { + Log.e("ChatApiService", "ERROR: 工具function映射为null") + return@forEach + } + + val name = functionMap["name"] as? String + if (name == null) { + Log.e("ChatApiService", "ERROR: 工具name为null") + return@forEach + } + + val description = functionMap["description"] as? String ?: "" + val parametersMap = functionMap["parameters"] as? Map + + Log.e("ChatApiService", "工具信息 - 名称: $name, 描述: $description") + Log.e("ChatApiService", "参数映射: $parametersMap") + + if (parametersMap == null) { + Log.e("ChatApiService", "ERROR: 工具 $name 的parameters为null") + return@forEach + } + + // 验证parametersMap的基本结构 + if (!parametersMap.containsKey("type")) { + Log.e("ChatApiService", "ERROR: 工具 $name 的parameters缺少type字段") + return@forEach + } + + val parametersJson = gson.toJson(parametersMap) + Log.e("ChatApiService", "参数JSON: $parametersJson") + + try { + Log.e("ChatApiService", "正在调用 Parameters.fromJsonString...") + Log.e("ChatApiService", "参数JSON长度: ${parametersJson.length}") + Log.e("ChatApiService", "参数JSON内容预览: ${parametersJson.take(200)}...") + + if (parametersJson.isBlank()) { + Log.e("ChatApiService", "ERROR: parametersJson为空或空白") + return@forEach + } + + val parameters = com.aallam.openai.api.core.Parameters.fromJsonString(parametersJson) + Log.e("ChatApiService", "SUCCESS: 成功创建Parameters对象") + + Log.e("ChatApiService", "正在创建Tool.function...") + Log.e("ChatApiService", "工具参数 - name: '$name' (isEmpty: ${name.isEmpty()})") + Log.e("ChatApiService", "工具参数 - description: '$description' (isEmpty: ${description.isEmpty()})") + Log.e("ChatApiService", "工具参数 - parameters: ${parameters != null}") + + if (name.isEmpty()) { + Log.e("ChatApiService", "ERROR: 工具名称为空,跳过") + return@forEach + } + + val tool = Tool.function( + name = name, + description = description, + parameters = parameters + ) + tools.add(tool) + + Log.e("ChatApiService", "SUCCESS: 成功添加工具: $name") + + } catch (e: Exception) { + Log.e("ChatApiService", "ERROR: 创建Parameters对象失败 for 工具 $name: ${e.message}", e) + Log.e("ChatApiService", "失败的参数JSON: $parametersJson") + Log.e("ChatApiService", "异常堆栈: ${e.stackTraceToString()}") + // 跳过这个工具,继续处理其他工具 + } + } else { + Log.e("ChatApiService", "跳过非function类型的工具: $type") + } + } catch (e: Exception) { + Log.e("ChatApiService", "ERROR: 处理工具映射时出错: ${e.message}", e) + Log.e("ChatApiService", "出错的工具映射: $toolMap") + } + } + + Log.e("ChatApiService", "最终工具列表大小: ${tools.size}") // 构建请求 - val chatCompletionRequest = ChatCompletionRequest( - model = ModelId(currentModel), - messages = chatMessages, - maxTokens = 2000, - temperature = 0.7 - ) + Log.e("ChatApiService", "正在构建ChatCompletionRequest...") + Log.e("ChatApiService", "构建参数检查:") + Log.e("ChatApiService", " - currentModel: '$currentModel' (isEmpty: ${currentModel.isEmpty()})") + Log.e("ChatApiService", " - chatMessages size: ${chatMessages.size}") + Log.e("ChatApiService", " - tools size: ${tools.size}") + Log.e("ChatApiService", " - openAI对象: ${openAI != null}") + + if (currentModel.isEmpty()) { + Log.e("ChatApiService", "ERROR: 模型名称为空") + throw IllegalArgumentException("模型名称不能为空") + } + + if (chatMessages.isEmpty()) { + Log.e("ChatApiService", "ERROR: 消息列表为空") + throw IllegalArgumentException("消息列表不能为空") + } - val chatsFlow = openAI!!.chatCompletions(chatCompletionRequest) + val chatsFlow = try { + val chatCompletionRequest = ChatCompletionRequest( + model = ModelId(currentModel), + messages = chatMessages, + maxTokens = 2000, + temperature = 0.7, + tools = if (tools.isNotEmpty()) tools else null + ) + Log.e("ChatApiService", "SUCCESS: ChatCompletionRequest创建成功") + + if (openAI == null) { + Log.e("ChatApiService", "ERROR: openAI对象为null") + throw IllegalStateException("OpenAI客户端未初始化") + } + + Log.e("ChatApiService", "正在调用 openAI.chatCompletions...") + val flow = openAI!!.chatCompletions(chatCompletionRequest) + Log.e("ChatApiService", "SUCCESS: 获得chatsFlow,开始收集...") + flow + } catch (e: Exception) { + Log.e("ChatApiService", "ERROR: 创建ChatCompletionRequest或调用chatCompletions失败: ${e.message}", e) + throw e + } chatsFlow.collect { result -> if (isCanceled) return@collect - val choice = result.choices.firstOrNull() ?: return@collect - val delta = choice.delta ?: return@collect + Log.e("ChatApiService", "收到流式响应,结果类型: ${result.javaClass.simpleName}") + Log.e("ChatApiService", "result.choices大小: ${result.choices?.size ?: 0}") + + val choice = result.choices?.firstOrNull() + if (choice == null) { + Log.e("ChatApiService", "ERROR: choice为null") + return@collect + } + + Log.e("ChatApiService", "choice存在,delta类型: ${choice.delta?.javaClass?.simpleName ?: "null"}") + + val delta = choice.delta + if (delta == null) { + Log.e("ChatApiService", "ERROR: delta为null") + return@collect + } + + Log.e("ChatApiService", "delta存在,content: ${delta.content}, toolCalls: ${delta.toolCalls?.size ?: 0}") // 处理普通文本内容 delta.content?.let { content -> + Log.e("ChatApiService", "收到文本内容: $content") streamCallback?.onToken(content) } // 收集工具调用信息 delta.toolCalls?.forEach { toolCall -> + Log.e("ChatApiService", "处理工具调用: index=${toolCall.index}, id=${toolCall.id}") val index = toolCall.index // 创建或获取现有的工具调用信息 val toolCallInfo = toolCalls.getOrPut(index) { ToolCallInfo() } - // 更新ID - toolCall.id?.let { toolCallInfo.id = it.toString() } + // 安全处理工具调用ID + try { + toolCall.id?.let { id -> + toolCallInfo.id = id.toString() + Log.e("ChatApiService", "更新工具调用ID: $id") + } + } catch (e: Exception) { + Log.e("ChatApiService", "处理工具调用ID异常: ${e.message}") + } - // 更新函数信息 - toolCall.function?.let { function -> - function.name?.let { toolCallInfo.name = it } - function.arguments?.let { toolCallInfo.arguments += it } + // 安全处理函数信息 + try { + toolCall.function?.let { function -> + function.name?.let { name -> + toolCallInfo.name = name + Log.e("ChatApiService", "更新工具名称: $name") + } + function.arguments?.let { args -> + toolCallInfo.arguments += args + Log.e("ChatApiService", "添加参数片段: $args") + } + } + } catch (e: Exception) { + Log.e("ChatApiService", "处理工具调用函数信息异常: ${e.message}") } } } @@ -322,7 +608,16 @@ class ChatApiService(private val context: android.content.Context? = null) : Cor * 处理工具调用 */ private suspend fun processToolCalls(): Boolean { - val firstToolCall = toolCalls.values.firstOrNull { it.isValid() } ?: return false + Log.d("ChatApiService", "开始处理工具调用,工具调用数量: ${toolCalls.size}") + + val firstToolCall = toolCalls.values.firstOrNull { it.isValid() } + if (firstToolCall == null) { + Log.d("ChatApiService", "没有有效的工具调用") + return false + } + + Log.d("ChatApiService", "第一个有效的工具调用: name=${firstToolCall.name}, id=${firstToolCall.id}") + Log.d("ChatApiService", "工具调用参数: ${firstToolCall.arguments}") // 创建函数调用字典 val functionCall = mapOf( @@ -331,6 +626,8 @@ class ChatApiService(private val context: android.content.Context? = null) : Cor "id" to firstToolCall.id ) + Log.d("ChatApiService", "创建的函数调用字典: $functionCall") + // 通知上层工具调用事件 streamCallback?.onFunctionCall(convertMapToJsonObject(functionCall)) @@ -338,8 +635,57 @@ class ChatApiService(private val context: android.content.Context? = null) : Cor launch { try { if (!isCanceled) { - // 这里暂时返回占位符结果,实际MCP功能留空 - val result = mapOf("context" to "MCP功能暂未实现") + // 通过MCP客户端处理工具调用 + val functionName = firstToolCall.name + val argumentsJson = firstToolCall.arguments + + Log.d("ChatApiService", "准备调用MCP工具: $functionName") + Log.d("ChatApiService", "参数JSON: $argumentsJson") + + val result = if (_mcpClient?.hasToolWithName(functionName) == true) { + Log.d("ChatApiService", "MCP客户端中找到工具: $functionName") + + // 解析参数 + val arguments = _mcpClient?.parseJsonArguments(argumentsJson) ?: emptyMap() + Log.d("ChatApiService", "解析后的参数: $arguments") + + // 调用MCP工具 + Log.d("ChatApiService", "开始调用MCP工具...") + val toolResult = _mcpClient?.callTool(functionName, arguments) + Log.d("ChatApiService", "MCP工具调用结果: $toolResult") + + // 处理结果 + if (toolResult != null) { + if (toolResult["isError"] == true) { + Log.w("ChatApiService", "MCP工具调用返回错误") + // 处理错误情况 + val content = toolResult["content"] as? List<*> + val firstContent = content?.firstOrNull() as? Map<*, *> + val errorText = firstContent?.get("text") as? String ?: "Tool execution failed" + mapOf("context" to errorText) + } else if (toolResult.containsKey("context")) { + Log.d("ChatApiService", "本地函数结果") + // 本地函数结果 + toolResult + } else { + Log.d("ChatApiService", "MCP工具结果") + // MCP工具结果 + val content = toolResult["content"] as? List<*> + val firstContent = content?.firstOrNull() as? Map<*, *> + val text = firstContent?.get("text") as? String ?: "" + mapOf("context" to text) + } + } else { + Log.e("ChatApiService", "MCP工具调用返回null") + mapOf("context" to "Tool call failed") + } + } else { + Log.w("ChatApiService", "MCP客户端中未找到工具: $functionName") + // 工具不存在 + mapOf("context" to "Tool not found: $functionName") + } + + Log.d("ChatApiService", "最终处理结果: $result") if (!isCanceled) { // 处理结果 @@ -357,6 +703,7 @@ class ChatApiService(private val context: android.content.Context? = null) : Cor } } } catch (e: Exception) { + Log.e("ChatApiService", "工具调用处理过程中出错: ${e.message}", e) if (!isCanceled) { val errorMessage = "工具调用处理失败: ${e.message}" sendFunctionCallResultInternal( @@ -435,40 +782,135 @@ class ChatApiService(private val context: android.content.Context? = null) : Cor /** * 注册函数 */ - fun registerFunction(name: String, description: String, parameters: Map): Boolean { - // 暂时返回true,实际功能留空,但保持与iOS版本接口一致 - return true + fun registerFunction(name: String, description: String, parameters: JSONObject): Boolean { + // 将JSONObject转换为Map + val parametersMap = convertJsonObjectToMap(parameters) + return registerFunction(name, description, parametersMap) + } + + /** + * 注册函数 (内部版本,接收Map参数) + */ + private fun registerFunction(name: String, description: String, parameters: Map): Boolean { + return try { + if (_mcpClient == null) { + _mcpClient = MCPClient(context) + } + + // 创建一个本地函数处理器 + val handler = LocalFunctionHandler(name) + + // 注册到MCP客户端 + _mcpClient!!.registerLocalFunction( + name = name, + description = description, + parameters = parameters, + handler = handler + ) + + true + } catch (e: Exception) { + Log.e("ChatApiService", "Failed to register function: $name", e) + false + } } /** * 初始化MCP客户端 */ fun initializeMcpClient(serverUrl: String): Boolean { - // MCP功能暂时留空,但保持与iOS版本接口一致 - return false + if (_mcpClient == null) { + _mcpClient = MCPClient(context) + } + + // 直接使用类的CoroutineScope启动协程 + launch { + try { + val result = _mcpClient?.connectToSSE(serverUrl) ?: false + Log.d("ChatApiService", "MCP客户端初始化${if (result) "成功" else "失败"}") + } catch (e: Exception) { + Log.e("ChatApiService", "MCP客户端初始化失败: ${e.message}", e) + } + } + + return true // 立即返回,实际连接在后台进行 } /** * MCP客户端是否已初始化 */ fun isMcpInitialized(): Boolean { - // MCP功能暂时留空,但保持与iOS版本接口一致 - return false + return _mcpClient?.isConnected() ?: false } /** * 关闭MCP客户端 */ fun closeMcpClient() { - // MCP功能暂时留空,但保持与iOS版本接口一致 + runBlocking { + _mcpClient?.disconnectAll() + _mcpClient = null + } } /** * 处理MCP工具调用 + * + * @param functionCall 函数调用JSON对象,必须包含name和arguments字段 + * @return 工具调用结果 + */ + suspend fun handleMcpToolCall(functionCall: JSONObject): String { + if (_mcpClient == null || !isMcpInitialized()) { + return "MCP客户端未初始化" + } + + try { + // 获取函数名称 + val name = functionCall.getString("name") + + // 获取参数 + val argumentsJson = functionCall.getString("arguments") + val arguments = _mcpClient?.parseJsonArguments(argumentsJson) ?: mapOf() + + // 调用工具 + val result = _mcpClient?.callTool(name, arguments) + + // 返回工具调用结果 + return result?.get("context") as? String ?: "工具调用失败" + } catch (e: Exception) { + Log.e("ChatApiService", "处理MCP工具调用失败: ${e.message}", e) + return "处理MCP工具调用失败: ${e.message}" + } + } + + /** + * 获取所有可用的工具定义 + */ + fun getToolDefinitions(): List> { + return _mcpClient?.getToolMaps() ?: emptyList() + } + + /** + * 处理图片 + */ + fun processImage(imagePath: String, prompt: String, maxWidth: Double, detail: String): String? { + return fileToBase64(imagePath, (maxWidth * 2).toInt()) // 简化处理,使用宽度的两倍作为最大KB数 + } + + /** + * 取消聊天流 */ - suspend fun handleMcpToolCall(functionCallJson: String): String { - // MCP功能暂时留空,但保持与iOS版本接口一致 - return "MCP功能暂未实现" + fun cancelChatStream() { + cancelCurrentStream() + } + + /** + * 聊天完成流式接口 + * 与iOS版本保持一致的接口 + */ + fun chatCompletionStream(messages: List>, tool: Boolean = false) { + // 直接调用sendMessageStream,因为该方法已经处理了工具调用 + sendMessageStream(messages) } // MARK: - 工具方法 @@ -796,4 +1238,28 @@ class ChatApiService(private val context: android.content.Context? = null) : Cor return jsonObject } + + /** + * 本地函数处理器 + * 用于处理在Flutter端定义的函数 + */ + private inner class LocalFunctionHandler( + private val functionName: String + ) : FunctionHandler { + override suspend fun handle(arguments: Map): String { + // 将函数调用通知到Flutter端 + val functionCall = JSONObject().apply { + put("name", functionName) + put("arguments", JSONObject(arguments)) + } + + // 通知Flutter端处理函数调用 + mainHandler.post { + streamCallback?.onFunctionCall(functionCall) + } + + // 返回一个标记,表示函数已被调用 + return "Function '$functionName' called with arguments: $arguments" + } + } } \ No newline at end of file diff --git a/local_plugins/chat_api/android/src/main/kotlin/com/yunqiinnovation/chat_api/CustomSseClientTransport.kt b/local_plugins/chat_api/android/src/main/kotlin/com/yunqiinnovation/chat_api/CustomSseClientTransport.kt new file mode 100644 index 000000000..945dde9f3 --- /dev/null +++ b/local_plugins/chat_api/android/src/main/kotlin/com/yunqiinnovation/chat_api/CustomSseClientTransport.kt @@ -0,0 +1,271 @@ +package com.yunqiinnovation.chat_api + +import android.util.Log +import io.ktor.client.* +import io.ktor.client.plugins.sse.* +import io.ktor.client.request.* +import io.ktor.client.statement.* +import io.ktor.http.* +import io.modelcontextprotocol.kotlin.sdk.JSONRPCMessage +import io.modelcontextprotocol.kotlin.sdk.shared.AbstractTransport +import kotlinx.coroutines.* +import kotlinx.serialization.encodeToString +import kotlinx.serialization.json.Json +import kotlinx.serialization.decodeFromString +import kotlin.properties.Delegates +import kotlin.time.Duration +import java.util.concurrent.atomic.AtomicBoolean + +/** + * 自定义SSE客户端传输层,解决官方SseClientTransport的URL路径问题 + * 直接使用原始URL,不自动添加/sse后缀 + */ +class CustomSseClientTransport( + private val client: HttpClient, + private val urlString: String?, + private val reconnectionTime: Duration? = null, + private val requestBuilder: HttpRequestBuilder.() -> Unit = {}, +) : AbstractTransport() { + private val TAG = "CustomSseClientTransport" + + private val scope by lazy { + CoroutineScope(session.coroutineContext + SupervisorJob()) + } + + private val initialized = AtomicBoolean(false) + private var session: ClientSSESession by Delegates.notNull() + private val endpoint = CompletableDeferred() + + private var job: Job? = null + + // 创建JSON解析器 + private val json = Json { + ignoreUnknownKeys = true + isLenient = true + coerceInputValues = true + encodeDefaults = true + explicitNulls = false + } + + // URL解析结果 + private var hostPart: String = "" + private var pathPart: String = "" + private var queryParams: Map = emptyMap() + + /** + * 解析URL,分离主机、路径和查询参数 + */ + private fun parseUrl(url: String): Triple> { + return try { + val params = mutableMapOf() + + var processedUrl = url.trim() + if (!processedUrl.startsWith("http://") && !processedUrl.startsWith("https://")) { + processedUrl = "https://$processedUrl" + } + + val urlObj = java.net.URL(processedUrl) + + // 解析查询参数 + if (urlObj.query != null) { + urlObj.query.split("&").forEach { param -> + val parts = param.split("=", limit = 2) + if (parts.size == 2) { + params[parts[0]] = parts[1] + } + } + } + + // 构建主机部分URL + val port = if (urlObj.port == -1) "" else ":${urlObj.port}" + val hostUrl = "${urlObj.protocol}://${urlObj.host}$port" + + // 路径部分 + val path = urlObj.path + + Triple(hostUrl, path, params) + } catch (e: Exception) { + Log.e(TAG, "解析URL失败: $url, ${e.message}") + Triple(url, "", emptyMap()) + } + } + + /** + * 收集SSE事件 + */ + private suspend fun collectEvents() { + job = scope.launch(CoroutineName("CustomSseMcpClientTransport.collect#${hashCode()}")) { + session.incoming.collect { event -> + when (event.event) { + "error" -> { + val e = IllegalStateException("SSE error: ${event.data}") + Log.e(TAG, "SSE错误: ${event.data}") + _onError(e) + throw e + } + + "open" -> { + Log.d(TAG, "SSE连接已打开") + } + + "endpoint" -> { + try { + val eventData = event.data ?: "" + Log.d(TAG, "收到endpoint事件: $eventData") + + // 构建完整的端点URL + val fullEndpoint = if (eventData.contains(hostPart)) { + eventData + } else if (eventData.startsWith("/")) { + "$hostPart$eventData" + } else { + eventData + } + + // 添加查询参数 + val endpointWithParams = if (queryParams.isNotEmpty()) { + if (fullEndpoint.contains("?")) { + val queryString = queryParams.entries.joinToString("&") { "${it.key}=${it.value}" } + "$fullEndpoint&$queryString" + } else { + val queryString = queryParams.entries.joinToString("&") { "${it.key}=${it.value}" } + "$fullEndpoint?$queryString" + } + } else { + fullEndpoint + } + + Log.d(TAG, "最终消息端点: $endpointWithParams") + endpoint.complete(endpointWithParams) + } catch (e: Exception) { + Log.e(TAG, "处理endpoint事件失败: ${e.message}", e) + _onError(e) + close() + error(e) + } + } + + else -> { + try { + val data = event.data + if (data != null) { + try { + val message = json.decodeFromString(data) + _onMessage(message) + } catch (e: Exception) { + Log.e(TAG, "解析JSON-RPC消息失败: ${e.message}", e) + _onError(e) + } + } + } catch (e: Exception) { + Log.e(TAG, "处理事件失败: ${e.message}", e) + _onError(e) + } + } + } + } + } + } + + /** + * 启动传输层 + */ + override suspend fun start() { + if (!initialized.compareAndSet(false, true)) { + Log.e(TAG, "传输层已经启动,不能重复启动") + error("CustomSseClientTransport already started!") + } + + // 解析URL + if (urlString != null) { + val urlInfo = parseUrl(urlString) + hostPart = urlInfo.first + pathPart = urlInfo.second + queryParams = urlInfo.third + + Log.d(TAG, "原始URL: $urlString") + Log.d(TAG, "主机部分: $hostPart") + Log.d(TAG, "路径部分: $pathPart") + Log.d(TAG, "查询参数: $queryParams") + } + + // 创建SSE会话 - 直接使用原始URL + session = urlString?.let { + val sseConnectUrl = if (queryParams.isNotEmpty()) { + if (pathPart.contains("?")) { + "$hostPart$pathPart" + } else { + val queryString = queryParams.entries.joinToString("&") { "${it.key}=${it.value}" } + "$hostPart$pathPart?$queryString" + } + } else { + "$hostPart$pathPart" + } + + Log.d(TAG, "SSE连接URL: $sseConnectUrl") + + client.sseSession( + urlString = sseConnectUrl, + reconnectionTime = reconnectionTime, + block = requestBuilder, + ) + } ?: client.sseSession( + reconnectionTime = reconnectionTime, + block = requestBuilder, + ) + + // 收集SSE事件 + collectEvents() + + // 等待endpoint就绪 + endpoint.await() + Log.d(TAG, "传输层启动完成,消息端点已就绪") + } + + /** + * 发送消息 + */ + @OptIn(ExperimentalCoroutinesApi::class) + override suspend fun send(message: JSONRPCMessage) { + if (!endpoint.isCompleted) { + Log.e(TAG, "发送失败: 未连接") + error("Not connected") + } + + try { + val messageEndpoint = endpoint.getCompleted() + Log.d(TAG, "发送消息到: $messageEndpoint") + + val jsonString = json.encodeToString(message) + + val response = client.post(messageEndpoint) { + headers.append(HttpHeaders.ContentType, ContentType.Application.Json.toString()) + setBody(jsonString) + } + + if (!response.status.isSuccess()) { + val text = response.bodyAsText() + Log.e(TAG, "发送消息失败: HTTP ${response.status}, $text") + error("Error POSTing to endpoint (HTTP ${response.status}): $text") + } + } catch (e: Exception) { + Log.e(TAG, "发送消息异常: ${e.message}", e) + _onError(e) + throw e + } + } + + /** + * 关闭传输层 + */ + override suspend fun close() { + if (!initialized.get()) { + Log.e(TAG, "关闭失败: 传输层未初始化") + error("CustomSseClientTransport is not initialized!") + } + session.cancel() + _onClose() + job?.cancelAndJoin() + Log.d(TAG, "传输层已关闭") + } +} \ No newline at end of file diff --git a/local_plugins/chat_api/android/src/main/kotlin/com/yunqiinnovation/chat_api/MCPClient.kt b/local_plugins/chat_api/android/src/main/kotlin/com/yunqiinnovation/chat_api/MCPClient.kt new file mode 100644 index 000000000..9aadf3d43 --- /dev/null +++ b/local_plugins/chat_api/android/src/main/kotlin/com/yunqiinnovation/chat_api/MCPClient.kt @@ -0,0 +1,360 @@ +package com.yunqiinnovation.chat_api + +import android.content.Context +import android.util.Log +import kotlinx.coroutines.* +import org.json.JSONObject +import org.json.JSONArray +import io.ktor.client.* + +/** + * 工具类型枚举 + */ +enum class ToolType { + LOCAL_FUNCTION, // 本地函数 + MCP_TOOL // MCP工具 +} + +/** + * 函数处理器接口 + */ +interface FunctionHandler { + /** + * 处理函数调用 + * @param arguments 函数参数,Map格式 + * @return 函数执行结果,字符串 + */ + suspend fun handle(arguments: Map): String +} + +/** + * MCP客户端 + * 与 iOS 版本 MCPClient 功能对等 + */ +class MCPClient(private val context: Context? = null) : AutoCloseable { + + companion object { + private const val TAG = "MCPClient" + } + + // 本地函数Map,函数名 -> 处理器 + private val localFunctions = mutableMapOf() + + // 本地函数定义Map,函数名 -> 定义 + private val localFunctionDefs = mutableMapOf>() + + // 子客户端列表,每个连接一个MCP服务器 + private val subClients = mutableMapOf() + + // 是否已连接 + private var isConnectedFlag = false + + // 共享的HttpClient,用于所有子客户端 + private val sharedHttpClient by lazy { createSslTrustAllClient() } + + init { + initializeSystemFunctions() + } + + /** + * 初始化系统函数 + */ + private fun initializeSystemFunctions() { + try { + val handler = SystemFunctionHandler(context) + handler.registerAllFunctions(this) + } catch (e: Exception) { + Log.w(TAG, "Failed to initialize system functions", e) + } + } + + /** + * 连接到SSE服务器 + * 直接接收完整的JSON配置字符串 + * + * @param mcpConfigJson 包含mcpServers字段的JSON配置字符串 + * @return 是否连接成功 + */ + suspend fun connectToSSE(mcpConfigJson: String): Boolean { + // 清除现有连接 + closeAllConnections() + + return try { + val config = JSONObject(mcpConfigJson) + val mcpServers = config.optJSONObject("mcpServers") ?: return false + + var connectedCount = 0 + val serverIds = mcpServers.keys() + + while (serverIds.hasNext()) { + val serverId = serverIds.next() + val serverConfig = mcpServers.optJSONObject(serverId) ?: continue + + val url = serverConfig.optString("url", "") + if (url.isEmpty()) continue + + val subClient = MCPSubClient(serverId, url, sharedHttpClient) + + if (subClient.connect()) { + subClients[serverId] = subClient + connectedCount++ + Log.d(TAG, "Connected to MCP server: $serverId") + } else { + Log.w(TAG, "Failed to connect to MCP server: $serverId") + } + } + + isConnectedFlag = connectedCount > 0 + Log.d(TAG, "Connected to $connectedCount MCP servers") + connectedCount > 0 + + } catch (e: Exception) { + Log.e(TAG, "Failed to connect to SSE", e) + false + } + } + + /** + * 关闭所有连接 + */ + private fun closeAllConnections() { + subClients.forEach { (serverId, client) -> + try { + client.close() + Log.d(TAG, "已关闭子客户端 [$serverId]") + } catch (e: Exception) { + Log.e(TAG, "关闭子客户端 [$serverId] 失败: ${e.message}") + } + } + subClients.clear() + isConnectedFlag = false + } + + /** + * 注册函数(简单版本,用于兼容) + */ + fun registerFunction(name: String, handler: FunctionHandler) { + localFunctions[name] = handler + } + + /** + * 注册本地函数 + */ + fun registerLocalFunction( + name: String, + description: String, + parameters: Any, + handler: FunctionHandler + ): Boolean { + try { + // 将参数统一转换为Map格式 + val parametersMap: Map = when (parameters) { + is Map<*, *> -> { + @Suppress("UNCHECKED_CAST") + parameters as Map + } + is JSONObject -> { + convertJsonObjectToMap(parameters) + } + else -> { + Log.e(TAG, "参数类型不支持: ${parameters.javaClass.name}") + return false + } + } + + // 检查参数是否包含必要字段 + if (!parametersMap.containsKey("type") || (parametersMap["type"] != "object")) { + Log.e(TAG, "参数必须是object类型") + return false + } + + localFunctions[name] = handler + + val functionDef = mapOf( + "name" to name, + "description" to description, + "parameters" to parametersMap + ) + localFunctionDefs[name] = functionDef + + Log.d(TAG, "Registered local function: $name") + return true + } catch (e: Exception) { + Log.e(TAG, "注册本地函数失败: ${e.message}", e) + return false + } + } + + /** + * 注销本地函数 + */ + fun unregisterLocalFunction(name: String): Boolean { + val removed = localFunctions.remove(name) != null + if (removed) { + localFunctionDefs.remove(name) + Log.d(TAG, "Unregistered local function: $name") + } + return removed + } + + /** + * 获取工具映射列表 + */ + fun getToolMaps(): List> { + val allToolMaps = mutableListOf>() + + // 添加本地函数 + for (functionDef in localFunctionDefs.values) { + allToolMaps.add(mapOf( + "type" to "function", + "function" to functionDef + )) + } + + // 添加MCP工具 + for (client in subClients.values) { + allToolMaps.addAll(client.getToolMaps()) + } + + return allToolMaps + } + + /** + * 获取工具类型 + */ + fun getToolType(name: String): ToolType? { + if (localFunctions.containsKey(name)) { + return ToolType.LOCAL_FUNCTION + } + + if (subClients.values.any { it.containsTool(name) }) { + return ToolType.MCP_TOOL + } + + return null + } + + /** + * 调用工具 + */ + suspend fun callTool(name: String, arguments: Map): Map? { + // 首先检查本地函数 + localFunctions[name]?.let { handler -> + return try { + val result = handler.handle(arguments) + mapOf("context" to result) + } catch (e: Exception) { + mapOf( + "content" to listOf(mapOf( + "type" to "text", + "text" to "Function call failed: ${e.message}" + )), + "isError" to true + ) + } + } + + // 然后检查MCP工具 + for (client in subClients.values) { + if (client.containsTool(name)) { + return client.callTool(name, arguments) + } + } + + // 工具未找到 + return mapOf( + "content" to listOf(mapOf( + "type" to "text", + "text" to "Tool not found: $name" + )), + "isError" to true + ) + } + + /** + * 检查是否有指定名称的工具 + */ + fun hasToolWithName(name: String): Boolean { + return localFunctions.containsKey(name) || + subClients.values.any { it.containsTool(name) } + } + + /** + * 解析JSON参数 + */ + fun parseJsonArguments(json: String): Map { + return try { + val jsonObject = JSONObject(json) + convertJsonObjectToMap(jsonObject) + } catch (e: Exception) { + Log.w(TAG, "Failed to parse JSON arguments", e) + emptyMap() + } + } + + /** + * 是否已连接 + */ + fun isConnected(): Boolean { + return isConnectedFlag || localFunctions.isNotEmpty() + } + + /** + * 断开所有连接 + */ + suspend fun disconnectAll() { + closeAllConnections() + Log.d(TAG, "Disconnected all MCP clients") + } + + /** + * 关闭连接 + */ + override fun close() { + runBlocking { + closeAllConnections() + localFunctions.clear() + localFunctionDefs.clear() + Log.d(TAG, "已关闭MCP客户端") + } + } + + /** + * 将JSONObject转换为Map + */ + private fun convertJsonObjectToMap(jsonObject: JSONObject): Map { + val map = mutableMapOf() + + val keys = jsonObject.keys() + while (keys.hasNext()) { + val key = keys.next() + val value = jsonObject.get(key) + + map[key] = when (value) { + is JSONObject -> convertJsonObjectToMap(value) + is JSONArray -> convertJsonArrayToList(value) + else -> value + } + } + + return map + } + + /** + * 将JSONArray转换为List + */ + private fun convertJsonArrayToList(jsonArray: JSONArray): List { + val list = mutableListOf() + + for (i in 0 until jsonArray.length()) { + val value = jsonArray.get(i) + list.add(when (value) { + is JSONObject -> convertJsonObjectToMap(value) + is JSONArray -> convertJsonArrayToList(value) + else -> value + }) + } + + return list + } +} \ No newline at end of file diff --git a/local_plugins/chat_api/android/src/main/kotlin/com/yunqiinnovation/chat_api/MCPSubClient.kt b/local_plugins/chat_api/android/src/main/kotlin/com/yunqiinnovation/chat_api/MCPSubClient.kt new file mode 100644 index 000000000..eb3a7d4f4 --- /dev/null +++ b/local_plugins/chat_api/android/src/main/kotlin/com/yunqiinnovation/chat_api/MCPSubClient.kt @@ -0,0 +1,329 @@ +package com.yunqiinnovation.chat_api + +import android.util.Log +import io.ktor.client.* +import io.modelcontextprotocol.kotlin.sdk.* +import io.modelcontextprotocol.kotlin.sdk.client.* +import io.modelcontextprotocol.kotlin.sdk.shared.* +import io.modelcontextprotocol.kotlin.sdk.CallToolRequest +import io.modelcontextprotocol.kotlin.sdk.Tool +import kotlinx.coroutines.* +import kotlinx.coroutines.sync.Mutex +import kotlinx.coroutines.sync.withLock +import kotlinx.serialization.json.* +import kotlin.collections.mutableMapOf +import kotlin.collections.mutableListOf + +/** + * MCP子客户端,管理单个MCP服务器的连接 + * 使用MCP SDK 0.5.0的官方API实现真实连接 + */ +class MCPSubClient( + private val serverId: String, + private val serverUrl: String, + private val httpClient: HttpClient? = null +) : AutoCloseable { + + companion object { + private const val TAG = "MCPSubClient" + } + + private val scope = CoroutineScope(Dispatchers.IO + SupervisorJob()) + private val connectionMutex = Mutex() + private var mcpClient: Client? = null + private var isConnected = false + private var availableTools = mutableListOf() + + /** + * 连接到MCP服务器 + */ + suspend fun connect(): Boolean = connectionMutex.withLock { + if (isConnected) return true + + return try { + Log.d(TAG, "[$serverId] 开始连接到MCP服务器: $serverUrl") + + // 创建MCP客户端实例 + val client = Client( + clientInfo = Implementation( + name = "deep-voice-chat-api", + version = "1.0.0" + ) + ) + + // 根据URL类型选择传输方式 + val transport = when { + serverUrl.startsWith("http://") || serverUrl.startsWith("https://") -> { + // SSE传输 - 使用自定义的CustomSseClientTransport + Log.d(TAG, "[$serverId] 创建自定义SSE传输,URL: $serverUrl") + val mcpHttpClient = httpClient ?: createMcpHttpClient() + CustomSseClientTransport( + client = mcpHttpClient, + urlString = serverUrl + ) + } + else -> { + Log.e(TAG, "[$serverId] 不支持的服务器URL格式: $serverUrl") + return false + } + } + + // 连接到服务器 + client.connect(transport) + + // 获取可用工具列表 + try { + val toolsResult = client.listTools() + if (toolsResult != null) { + availableTools.clear() + availableTools.addAll(toolsResult.tools) + Log.d(TAG, "[$serverId] 获取到 ${availableTools.size} 个工具") + + // 打印工具信息 + availableTools.forEach { tool -> + Log.d(TAG, "[$serverId] 工具: ${tool.name} - ${tool.description}") + } + } + } catch (e: Exception) { + Log.w(TAG, "[$serverId] 获取工具列表失败: ${e.message}") + // 即使获取工具失败,连接也可能是成功的 + } + + mcpClient = client + isConnected = true + Log.d(TAG, "[$serverId] MCP连接成功") + true + + } catch (e: Exception) { + Log.e(TAG, "[$serverId] MCP连接失败: ${e.message}", e) + false + } + } + + /** + * 检查是否包含指定工具 + */ + fun containsTool(name: String): Boolean { + return availableTools.any { it.name == name } + } + + /** + * 获取工具映射列表 + */ + fun getToolMaps(): List> { + Log.d(TAG, "[$serverId] 开始获取工具映射,工具数量: ${availableTools.size}") + + val toolMaps = availableTools.map { tool -> + Log.d(TAG, "[$serverId] 处理工具: ${tool.name}") + Log.d(TAG, "[$serverId] 工具描述: ${tool.description}") + Log.d(TAG, "[$serverId] 输入Schema: ${tool.inputSchema}") + + val parametersMap = tool.inputSchema?.let { inputSchema -> + convertInputSchemaToMap(inputSchema) + } ?: mapOf( + "type" to "object", + "properties" to emptyMap(), + "required" to emptyList() + ) + + Log.d(TAG, "[$serverId] 转换后的参数映射: $parametersMap") + + val toolMap = mapOf( + "type" to "function", + "function" to mapOf( + "name" to tool.name, + "description" to (tool.description ?: ""), + "parameters" to parametersMap + ) + ) + + Log.d(TAG, "[$serverId] 最终工具映射: $toolMap") + toolMap + } + + Log.d(TAG, "[$serverId] 完成工具映射生成,返回 ${toolMaps.size} 个工具") + return toolMaps + } + + /** + * 调用MCP工具 + */ + suspend fun callTool(name: String, arguments: Map): Map? { + val client = mcpClient ?: return null + + return try { + Log.d(TAG, "[$serverId] 调用工具: $name, 参数: $arguments") + + // 创建工具调用请求 - 将Map转换为JsonObject + val argumentsJson = kotlinx.serialization.json.buildJsonObject { + arguments.forEach { (key, value) -> + when (value) { + is String -> put(key, value) + is Number -> put(key, kotlinx.serialization.json.JsonPrimitive(value)) + is Boolean -> put(key, value) + else -> put(key, value.toString()) + } + } + } + + val request = CallToolRequest( + name = name, + arguments = argumentsJson + ) + + // 调用工具 + val result = client.callTool(request) + + result?.let { callResult -> + Log.d(TAG, "[$serverId] 工具调用结果: ${callResult.content.size} 个内容项") + + // 将结果转换为统一格式 + val contentList = callResult.content.map { contentItem -> + // 根据不同的内容类型处理 + mapOf( + "type" to "text", + "text" to (contentItem.toString()) + ) + } + + mapOf( + "content" to contentList, + "isError" to (callResult.isError ?: false) + ) + } + + } catch (e: Exception) { + Log.e(TAG, "[$serverId] 调用工具 '$name' 失败: ${e.message}", e) + mapOf( + "content" to listOf(mapOf( + "type" to "text", + "text" to "Tool call failed: ${e.message}" + )), + "isError" to true + ) + } + } + + /** + * 将Tool.Input转换为Map格式,供OpenAI使用 + */ + private fun convertInputSchemaToMap(inputSchema: Tool.Input): Map { + Log.d(TAG, "开始转换Input Schema: $inputSchema") + + val properties = mutableMapOf() + val required = mutableListOf() + + // 处理properties + inputSchema.properties?.let { propsJsonObject -> + Log.d(TAG, "处理properties: $propsJsonObject") + for ((key, value) in propsJsonObject) { + Log.d(TAG, "处理属性: $key = $value (${value::class.java.simpleName})") + when (value) { + is JsonPrimitive -> { + if (value.isString) { + properties[key] = mapOf("type" to value.content) + } else { + properties[key] = value.content + } + } + is JsonObject -> { + properties[key] = convertJsonObjectToMap(value) + } + else -> { + Log.w(TAG, "未知的属性值类型: ${value::class.java.simpleName}") + properties[key] = value.toString() + } + } + } + } ?: Log.d(TAG, "properties为null") + + // 处理required + inputSchema.required?.let { requiredList -> + Log.d(TAG, "处理required: $requiredList") + required.addAll(requiredList) + } ?: Log.d(TAG, "required为null") + + val result = mapOf( + "type" to "object", + "properties" to properties, + "required" to required + ) + + Log.d(TAG, "转换后的Schema Map: $result") + return result + } + + /** + * 将JsonObject转换为Map + */ + private fun convertJsonObjectToMap(jsonObject: JsonObject): Map { + val map = mutableMapOf() + + for ((key, value) in jsonObject) { + map[key] = when (value) { + is JsonPrimitive -> { + when { + value.isString -> value.content + value.booleanOrNull != null -> value.boolean + value.longOrNull != null -> value.long + value.doubleOrNull != null -> value.double + else -> value.toString() + } + } + is JsonObject -> convertJsonObjectToMap(value) + is JsonArray -> value.map { element -> + when (element) { + is JsonPrimitive -> element.content + is JsonObject -> convertJsonObjectToMap(element) + else -> element.toString() + } + } + else -> value.toString() + } + } + + return map + } + + /** + * 刷新工具列表 + */ + suspend fun refreshTools(): Boolean { + val client = mcpClient ?: return false + + return try { + val toolsResult = client.listTools() + if (toolsResult != null) { + availableTools.clear() + availableTools.addAll(toolsResult.tools) + Log.d(TAG, "[$serverId] 刷新工具列表成功,共 ${availableTools.size} 个工具") + true + } else { + false + } + } catch (e: Exception) { + Log.e(TAG, "[$serverId] 刷新工具列表失败: ${e.message}", e) + false + } + } + + /** + * 关闭连接 + */ + override fun close() { + scope.launch { + connectionMutex.withLock { + try { + mcpClient?.close() + mcpClient = null + isConnected = false + availableTools.clear() + Log.d(TAG, "[$serverId] MCP连接已关闭") + } catch (e: Exception) { + Log.e(TAG, "[$serverId] 关闭MCP连接时出错: ${e.message}", e) + } + } + } + scope.cancel() + } +} \ No newline at end of file diff --git a/local_plugins/chat_api/android/src/main/kotlin/com/yunqiinnovation/chat_api/SystemFunctionHandler.kt b/local_plugins/chat_api/android/src/main/kotlin/com/yunqiinnovation/chat_api/SystemFunctionHandler.kt new file mode 100644 index 000000000..93a019bcb --- /dev/null +++ b/local_plugins/chat_api/android/src/main/kotlin/com/yunqiinnovation/chat_api/SystemFunctionHandler.kt @@ -0,0 +1,471 @@ +package com.yunqiinnovation.chat_api + +import android.content.Context +import android.content.Intent +import android.content.pm.PackageManager +import android.location.LocationManager +import android.net.Uri +import android.os.Build +import android.provider.CalendarContract +import android.telephony.SmsManager +import androidx.core.content.ContextCompat +import java.text.SimpleDateFormat +import java.util.* + +/** + * 系统函数处理器 + * 提供内置的系统函数,与 iOS 版本保持一致 + */ +class SystemFunctionHandler(private val context: Context? = null) { + + /** + * 注册所有系统函数 + */ + fun registerAllFunctions(client: MCPClient) { + // 注册退出交互函数 + client.registerLocalFunction( + name = "exit_interaction", + description = "结束当前交互", + parameters = mapOf( + "type" to "object", + "properties" to emptyMap(), + "required" to emptyList() + ), + handler = ExitInteractionHandler(context) + ) + + // 注册翻译模式函数 + client.registerLocalFunction( + name = "enter_translation_mode", + description = "用户请求进入实时翻译模式时,启动实时翻译功能", + parameters = mapOf( + "type" to "object", + "properties" to emptyMap(), + "required" to emptyList() + ), + handler = TranslationModeHandler(context) + ) + + // 注册发送短信函数 + client.registerLocalFunction( + name = "send_text_message", + description = "发送短信", + parameters = mapOf( + "type" to "object", + "properties" to mapOf( + "contact" to mapOf( + "type" to "string", + "description" to "联系人姓名或电话号码" + ), + "message" to mapOf( + "type" to "string", + "description" to "短信内容" + ) + ), + "required" to listOf("contact", "message") + ), + handler = SendTextMessageHandler(context) + ) + + // 注册拨打电话函数 + client.registerLocalFunction( + name = "make_phone_call", + description = "拨打电话", + parameters = mapOf( + "type" to "object", + "properties" to mapOf( + "contact" to mapOf( + "type" to "string", + "description" to "联系人姓名或电话号码" + ) + ), + "required" to listOf("contact") + ), + handler = MakePhoneCallHandler(context) + ) + + // 注册设置提醒函数 + client.registerLocalFunction( + name = "set_reminder", + description = "设置提醒事项", + parameters = mapOf( + "type" to "object", + "properties" to mapOf( + "title" to mapOf( + "type" to "string", + "description" to "提醒标题" + ), + "content" to mapOf( + "type" to "string", + "description" to "提醒内容" + ), + "time" to mapOf( + "type" to "string", + "description" to "提醒时间,格式为'yyyy-MM-dd HH:mm',如'2023-12-31 14:30'" + ) + ), + "required" to listOf("title", "time") + ), + handler = SetReminderHandler(context) + ) + + // 注册获取当前时间函数 + client.registerLocalFunction( + name = "get_current_time", + description = "获取当前日期和时间", + parameters = mapOf( + "type" to "object", + "properties" to mapOf( + "format" to mapOf( + "type" to "string", + "description" to "时间格式,可选,默认为标准格式" + ) + ), + "required" to emptyList() + ), + handler = GetCurrentTimeHandler() + ) + + // 注册获取当前位置函数 + client.registerLocalFunction( + name = "get_current_location", + description = "获取当前地理位置", + parameters = mapOf( + "type" to "object", + "properties" to emptyMap(), + "required" to emptyList() + ), + handler = GetCurrentLocationHandler(context) + ) + + // 注册媒体播放功能 + client.registerLocalFunction( + name = "media_play", + description = "播放媒体", + parameters = mapOf( + "type" to "object", + "properties" to emptyMap(), + "required" to emptyList() + ), + handler = MediaPlayHandler(context) + ) + + // 注册媒体暂停功能 + client.registerLocalFunction( + name = "media_pause", + description = "暂停媒体播放", + parameters = mapOf( + "type" to "object", + "properties" to emptyMap(), + "required" to emptyList() + ), + handler = MediaPauseHandler(context) + ) + + // 注册媒体上一首功能 + client.registerLocalFunction( + name = "media_previous", + description = "播放上一首", + parameters = mapOf( + "type" to "object", + "properties" to emptyMap(), + "required" to emptyList() + ), + handler = MediaPreviousHandler(context) + ) + + // 注册媒体下一首功能 + client.registerLocalFunction( + name = "media_next", + description = "播放下一首", + parameters = mapOf( + "type" to "object", + "properties" to emptyMap(), + "required" to emptyList() + ), + handler = MediaNextHandler(context) + ) + + // 注册打开录音机功能 + client.registerLocalFunction( + name = "open_recorder", + description = "打开系统录音机并开始录音", + parameters = mapOf( + "type" to "object", + "properties" to emptyMap(), + "required" to emptyList() + ), + handler = OpenRecorderHandler(context) + ) + } +} + +/** + * 退出交互处理器 + */ +private class ExitInteractionHandler(private val context: Context?) : FunctionHandler { + override suspend fun handle(arguments: Map): String { + // 发送广播通知退出交互 + context?.sendBroadcast(Intent("com.yunqiinnovation.deepsound.EXIT_INTERACTION")) + return "{\"result\": \"已结束当前交互\"}" + } +} + +/** + * 翻译模式处理器 + */ +private class TranslationModeHandler(private val context: Context?) : FunctionHandler { + override suspend fun handle(arguments: Map): String { + // 发送广播通知进入翻译模式 + context?.sendBroadcast(Intent("com.yunqiinnovation.deepsound.ENTER_TRANSLATION_MODE")) + return "{\"result\": \"已进入翻译模式\"}" + } +} + +/** + * 发送短信处理器 + */ +private class SendTextMessageHandler(private val context: Context?) : FunctionHandler { + override suspend fun handle(arguments: Map): String { + val contact = arguments["contact"] as? String + val message = arguments["message"] as? String + + if (contact == null || message == null) { + return "{\"result\": \"缺少必要参数\"}" + } + + return try { + val smsIntent = Intent(Intent.ACTION_SENDTO).apply { + data = Uri.parse("smsto:$contact") + putExtra("sms_body", message) + flags = Intent.FLAG_ACTIVITY_NEW_TASK + } + + if (context?.packageManager?.queryIntentActivities(smsIntent, 0)?.isNotEmpty() == true) { + context.startActivity(smsIntent) + "{\"result\": \"已打开短信应用\"}" + } else { + "{\"result\": \"无法打开短信应用\"}" + } + } catch (e: Exception) { + "{\"result\": \"发送短信失败:${e.message}\"}" + } + } +} + +/** + * 拨打电话处理器 + */ +private class MakePhoneCallHandler(private val context: Context?) : FunctionHandler { + override suspend fun handle(arguments: Map): String { + val contact = arguments["contact"] as? String + + if (contact == null) { + return "{\"result\": \"缺少联系人参数\"}" + } + + return try { + val callIntent = Intent(Intent.ACTION_DIAL).apply { + data = Uri.parse("tel:$contact") + flags = Intent.FLAG_ACTIVITY_NEW_TASK + } + + if (context?.packageManager?.queryIntentActivities(callIntent, 0)?.isNotEmpty() == true) { + context.startActivity(callIntent) + "{\"result\": \"已发起电话呼叫\"}" + } else { + "{\"result\": \"无法拨打电话\"}" + } + } catch (e: Exception) { + "{\"result\": \"拨打电话失败:${e.message}\"}" + } + } +} + +/** + * 设置提醒处理器 + */ +private class SetReminderHandler(private val context: Context?) : FunctionHandler { + override suspend fun handle(arguments: Map): String { + val title = arguments["title"] as? String + val content = arguments["content"] as? String ?: "" + val timeString = arguments["time"] as? String + + if (title == null || timeString == null) { + return "{\"result\": \"缺少必要参数\"}" + } + + return try { + // 解析时间 + val formatter = SimpleDateFormat("yyyy-MM-dd HH:mm", Locale.getDefault()) + val date = formatter.parse(timeString) + + if (date == null) { + return "{\"result\": \"时间格式错误\"}" + } + + // 创建日历事件 + val calendarIntent = Intent(Intent.ACTION_INSERT).apply { + data = CalendarContract.Events.CONTENT_URI + putExtra(CalendarContract.Events.TITLE, title) + putExtra(CalendarContract.Events.DESCRIPTION, content) + putExtra(CalendarContract.EXTRA_EVENT_BEGIN_TIME, date.time) + putExtra(CalendarContract.EXTRA_EVENT_END_TIME, date.time + 60 * 60 * 1000) // 默认1小时 + putExtra(CalendarContract.Events.HAS_ALARM, 1) + flags = Intent.FLAG_ACTIVITY_NEW_TASK + } + + if (context?.packageManager?.queryIntentActivities(calendarIntent, 0)?.isNotEmpty() == true) { + context.startActivity(calendarIntent) + "{\"result\": \"提醒设置成功\"}" + } else { + "{\"result\": \"无法打开日历应用\"}" + } + } catch (e: Exception) { + "{\"result\": \"设置提醒失败:${e.message}\"}" + } + } +} + +/** + * 获取当前时间处理器 + */ +private class GetCurrentTimeHandler : FunctionHandler { + override suspend fun handle(arguments: Map): String { + val format = arguments["format"] as? String + + return try { + val formatter = if (!format.isNullOrEmpty()) { + SimpleDateFormat(format, Locale.getDefault()) + } else { + SimpleDateFormat("yyyy年MM月dd日 HH:mm:ss", Locale.CHINA) + } + + val currentTime = formatter.format(Date()) + "{\"result\": \"$currentTime\", \"time\": \"$currentTime\"}" + } catch (e: Exception) { + "{\"result\": \"获取时间失败:${e.message}\"}" + } + } +} + +/** + * 获取当前位置处理器 + */ +private class GetCurrentLocationHandler(private val context: Context?) : FunctionHandler { + override suspend fun handle(arguments: Map): String { + if (context == null) { + return "{\"result\": \"上下文未初始化\", \"success\": false}" + } + + // 检查位置权限 + val hasPermission = if (Build.VERSION.SDK_INT >= Build.VERSION_CODES.M) { + ContextCompat.checkSelfPermission( + context, + android.Manifest.permission.ACCESS_FINE_LOCATION + ) == PackageManager.PERMISSION_GRANTED + } else { + true + } + + if (!hasPermission) { + return "{\"result\": \"位置权限被拒绝。请前往设置中启用位置权限\", \"success\": false}" + } + + // 检查位置服务是否启用 + val locationManager = context.getSystemService(Context.LOCATION_SERVICE) as? LocationManager + val isLocationEnabled = locationManager?.isProviderEnabled(LocationManager.GPS_PROVIDER) == true || + locationManager?.isProviderEnabled(LocationManager.NETWORK_PROVIDER) == true + + if (!isLocationEnabled) { + return "{\"result\": \"位置服务未启用。请前往设置中启用位置服务\", \"success\": false}" + } + + // 注意:实际的位置获取需要异步处理,这里只返回提示信息 + // 在实际应用中应该使用 LocationCallback 或 Coroutines 来获取实时位置 + return "{\"result\": \"需要通过位置服务获取当前位置\", \"success\": true}" + } +} + +/** + * 媒体播放处理器 + */ +private class MediaPlayHandler(private val context: Context?) : FunctionHandler { + override suspend fun handle(arguments: Map): String { + // 发送媒体播放广播 + context?.sendBroadcast(Intent("com.yunqiinnovation.deepsound.MEDIA_PLAY")) + return "{\"result\": \"已开始播放媒体\"}" + } +} + +/** + * 媒体暂停处理器 + */ +private class MediaPauseHandler(private val context: Context?) : FunctionHandler { + override suspend fun handle(arguments: Map): String { + // 发送媒体暂停广播 + context?.sendBroadcast(Intent("com.yunqiinnovation.deepsound.MEDIA_PAUSE")) + return "{\"result\": \"已暂停媒体播放\"}" + } +} + +/** + * 媒体上一首处理器 + */ +private class MediaPreviousHandler(private val context: Context?) : FunctionHandler { + override suspend fun handle(arguments: Map): String { + // 发送切换上一首广播 + context?.sendBroadcast(Intent("com.yunqiinnovation.deepsound.MEDIA_PREVIOUS")) + return "{\"result\": \"已切换到上一首\"}" + } +} + +/** + * 媒体下一首处理器 + */ +private class MediaNextHandler(private val context: Context?) : FunctionHandler { + override suspend fun handle(arguments: Map): String { + // 发送切换下一首广播 + context?.sendBroadcast(Intent("com.yunqiinnovation.deepsound.MEDIA_NEXT")) + return "{\"result\": \"已切换到下一首\"}" + } +} + +/** + * 打开录音机处理器 + */ +private class OpenRecorderHandler(private val context: Context?) : FunctionHandler { + override suspend fun handle(arguments: Map): String { + if (context == null) { + return "{\"result\": \"上下文未初始化\"}" + } + + return try { + // 尝试打开录音机应用 + val recorderIntent = Intent(Intent.ACTION_MAIN).apply { + addCategory(Intent.CATEGORY_APP_MUSIC) + flags = Intent.FLAG_ACTIVITY_NEW_TASK + } + + // 或者尝试使用录音Intent + val recordIntent = Intent("android.provider.MediaStore.RECORD_SOUND") + recordIntent.flags = Intent.FLAG_ACTIVITY_NEW_TASK + + when { + context.packageManager?.queryIntentActivities(recordIntent, 0)?.isNotEmpty() == true -> { + context.startActivity(recordIntent) + "{\"result\": \"已打开录音机\"}" + } + context.packageManager?.queryIntentActivities(recorderIntent, 0)?.isNotEmpty() == true -> { + context.startActivity(recorderIntent) + "{\"result\": \"已打开音频应用\"}" + } + else -> { + "{\"result\": \"无法打开录音机应用\"}" + } + } + } catch (e: Exception) { + "{\"result\": \"打开录音机失败:${e.message}\"}" + } + } +} \ No newline at end of file diff --git a/local_plugins/chat_api/android/src/main/kotlin/com/yunqiinnovation/chat_api/Utils.kt b/local_plugins/chat_api/android/src/main/kotlin/com/yunqiinnovation/chat_api/Utils.kt new file mode 100644 index 000000000..0d1b90153 --- /dev/null +++ b/local_plugins/chat_api/android/src/main/kotlin/com/yunqiinnovation/chat_api/Utils.kt @@ -0,0 +1,57 @@ +package com.yunqiinnovation.chat_api + +import io.ktor.client.* +import io.ktor.client.engine.okhttp.* +import io.ktor.client.plugins.sse.* +import okhttp3.OkHttpClient +import okhttp3.logging.HttpLoggingInterceptor +import java.security.SecureRandom +import java.security.cert.X509Certificate +import java.util.concurrent.TimeUnit +import javax.net.ssl.SSLContext +import javax.net.ssl.TrustManager +import javax.net.ssl.X509TrustManager + +/** + * 创建一个信任所有SSL证书的HttpClient + * 仅用于开发环境,生产环境应该使用正确的证书验证 + * MCP SDK 0.5.0应该自动处理SSE相关功能 + */ +fun createSslTrustAllClient(): HttpClient { + // 创建一个信任所有证书的TrustManager + val trustAllCerts = arrayOf(object : X509TrustManager { + override fun checkClientTrusted(chain: Array?, authType: String?) {} + override fun checkServerTrusted(chain: Array?, authType: String?) {} + override fun getAcceptedIssuers(): Array = arrayOf() + }) + + // 创建SSL上下文并初始化它 + val sslContext = SSLContext.getInstance("TLS") + sslContext.init(null, trustAllCerts, SecureRandom()) + + // 创建OkHttpClient并配置信任所有证书 + val okHttpClient = OkHttpClient.Builder() + .sslSocketFactory(sslContext.socketFactory, trustAllCerts[0] as X509TrustManager) + .hostnameVerifier { _, _ -> true } + .connectTimeout(30, TimeUnit.SECONDS) + .readTimeout(30, TimeUnit.SECONDS) + .build() + + // 创建使用OkHttp引擎的HttpClient + return HttpClient(OkHttp) { + engine { + preconfigured = okHttpClient + } + + // 安装SSE插件 - 这是关键! + install(SSE) + } +} + +/** + * 创建一个专门用于MCP连接的HttpClient,包含SSE支持 + * 注意:这个版本尝试不安装SSE插件,让MCP SDK自行处理 + */ +fun createMcpHttpClient(): HttpClient { + return createSslTrustAllClient() +} \ No newline at end of file From efab2b53a4ae21b5689baedfbf51198a3b696eb3 Mon Sep 17 00:00:00 2001 From: wolfplus Date: Sat, 7 Jun 2025 14:53:09 +0100 Subject: [PATCH 2/3] add --- .../yunqiinnovation/deepsound/MainActivity.kt | 26 +- .../agent_service/AgentService.kt | 30 +- .../azure_speech/AzureAsrHelper.kt | 187 ++++---- .../azure_speech/AzureTtsHelper.kt | 79 +++- .../chat_api/ChatApiService.kt | 404 ++++++------------ .../chat_api/CustomSseClientTransport.kt | 14 +- .../com/yunqiinnovation/chat_api/MCPClient.kt | 30 +- .../yunqiinnovation/chat_api/MCPSubClient.kt | 36 +- 8 files changed, 346 insertions(+), 460 deletions(-) diff --git a/android/app/src/main/kotlin/com/yunqiinnovation/deepsound/MainActivity.kt b/android/app/src/main/kotlin/com/yunqiinnovation/deepsound/MainActivity.kt index e7b34ac65..68e3f58b8 100644 --- a/android/app/src/main/kotlin/com/yunqiinnovation/deepsound/MainActivity.kt +++ b/android/app/src/main/kotlin/com/yunqiinnovation/deepsound/MainActivity.kt @@ -19,7 +19,9 @@ class MainActivity: FlutterActivity() { override fun onCreate(savedInstanceState: Bundle?) { super.onCreate(savedInstanceState) - + + // 设置日志级别,减少系统级日志 + suppressSystemLogs() } override fun configureFlutterEngine(flutterEngine: FlutterEngine) { @@ -32,5 +34,27 @@ class MainActivity: FlutterActivity() { override fun onDestroy() { super.onDestroy() } + + /** + * 抑制系统级组件的日志输出 + */ + private fun suppressSystemLogs() { + val systemTags = arrayOf( + "MediaCodec", + "CCodec", + "CCodecBufferChannel", + "OpusManager", + "MediaCodecList", + "MediaPlayer", + "AudioTrack", + "AudioManager", + "AudioFlinger", + "AudioService" + ) + + systemTags.forEach { tag -> + System.setProperty("log.tag.$tag", "WARN") + } + } } diff --git a/local_plugins/agent_service/android/src/main/kotlin/com/yunqiinnovation/agent_service/AgentService.kt b/local_plugins/agent_service/android/src/main/kotlin/com/yunqiinnovation/agent_service/AgentService.kt index d560c4684..90e318861 100644 --- a/local_plugins/agent_service/android/src/main/kotlin/com/yunqiinnovation/agent_service/AgentService.kt +++ b/local_plugins/agent_service/android/src/main/kotlin/com/yunqiinnovation/agent_service/AgentService.kt @@ -152,15 +152,13 @@ object AgentService : CoroutineScope { config["openaiApiKey"]?.toString() ?: "", config["openaiBaseUrl"]?.toString() ?: "", config["openaiModel"]?.toString() ?: "", - "" - //config["mcpServer"]?.toString() ?: "" + config["mcpServer"]?.toString() ?: "" ) // 加载最近的聊天记录 loadChatHistory() isInitialized = true - Log.d(TAG, "代理服务初始化成功") return true } catch (e: Exception) { Log.e(TAG, "初始化失败: ${e.message}") @@ -195,8 +193,6 @@ object AgentService : CoroutineScope { job.cancel() clearListeners() isInitialized = false - - Log.d(TAG, "代理服务资源已释放") } catch (e: Exception) { Log.e(TAG, "释放资源异常: ${e.message}") } @@ -230,9 +226,7 @@ object AgentService : CoroutineScope { language = ttsLanguage ) - if (success) { - Log.d(TAG, "Azure TTS服务初始化成功") - } else { + if (!success) { Log.e(TAG, "Azure TTS服务初始化失败") } @@ -268,13 +262,10 @@ object AgentService : CoroutineScope { } else -> { // 处理其他类型的事件 - Log.d(TAG, "TTS事件: ${event.type}") } } } }) - - Log.d(TAG, "TTS服务配置完成") } catch (e: Exception) { Log.e(TAG, "初始化TTS引擎失败: ${e.message}") } @@ -372,7 +363,6 @@ object AgentService : CoroutineScope { } else { AzureAsrHelper.AudioSourceType.MICROPHONE } - Log.d(TAG, "开始语音识别, 音频源类型: $audioSourceType") azureAsrHelper?.startContinuousRecognition(object : AzureAsrHelper.ContinuousRecognizeCallback { override fun onRecognizing(recognizing: String, detectedLanguage: String) { if (recognizing.isNotEmpty()) { @@ -477,7 +467,6 @@ object AgentService : CoroutineScope { BleService.closeCodec() isRecognitionActive = false stopIdleCheck() - Log.d(TAG, "语音识别已停止") } catch (e: Exception) { Log.e(TAG, "停止语音识别异常: ${e.message}") isRecognitionActive = false @@ -518,8 +507,6 @@ object AgentService : CoroutineScope { // 通知ChatAPI服务终止当前流式请求 chatApiService.cancelCurrentStream() - // 记录日志 - Log.d(TAG, "AI流输出已停止") } catch (e: Exception) { Log.e(TAG, "停止AI流输出异常", e) // 确保状态被重置,即使发生异常 @@ -586,8 +573,6 @@ object AgentService : CoroutineScope { text: String = "", speakResponse: Boolean = false ) { - Log.d(TAG, "处理图片输入: ${if (text.isEmpty()) "无附加文本" else "附带文本: $text"}") - // 创建带图片的用户消息并处理 val userMessage = createUserMessageWithImage(text, imageBase64) // 图片描述用于存储 @@ -687,7 +672,6 @@ object AgentService : CoroutineScope { } override fun onComplete() { - Log.d(TAG, "AI完整回复: $responseBuilder") // 视情况决定是否朗读回复 if (speakResponse) { ttsService?.flushStream() @@ -751,7 +735,6 @@ object AgentService : CoroutineScope { "function_call" to functionCall.toString(), "result" to functionCallResult.toString(), )) - Log.d(TAG, "mcp调用结果: $functionCallResult") val (metestr, broadcast)= autoHandleFcunCallResult(functionCallResult); aiMetadata = metestr; nobroadcast = broadcast @@ -803,8 +786,6 @@ object AgentService : CoroutineScope { addToHistoryMessages(createAssistantMessage(content)) } } - - Log.d(TAG, "已加载${recentMessages.size}条历史记录") } catch (e: Exception) { Log.e(TAG, "加载聊天历史失败: ${e.message}") } @@ -906,7 +887,6 @@ object AgentService : CoroutineScope { historyMessages.remove(0) } } - Log.d(TAG, "聊天历史已清除") } else { Log.e(TAG, "清除聊天历史失败") } @@ -1099,7 +1079,6 @@ object AgentService : CoroutineScope { } } } - Log.d(TAG, "开始播放音频资源") } catch (e: Exception) { Log.e(TAG, "播放音频资源异常: ${e.message}", e) release() @@ -1134,7 +1113,6 @@ object AgentService : CoroutineScope { // Log.d(TAG, "mcp调用结果: $metaStr") if (meta.has("card_music")) { //音乐卡片 val cardMusic = meta.getJSONObject("card_music") - Log.d(TAG, "检查到音乐卡片: $cardMusic") val id = cardMusic.optString("id", "") val url = cardMusic.optString("url", "") val name = cardMusic.optString("name", "") @@ -1157,7 +1135,6 @@ object AgentService : CoroutineScope { playlist.add(mapOf("id" to id,"url" to url, "title" to name, "artist" to sgener,"coverUrl" to image)) } if (playlist.isNotEmpty()) { - Log.w(TAG, "自动播放音乐列表 ${playlist}") broadcast = true; processMusicPlayList(playlist) } else { @@ -1190,7 +1167,6 @@ object AgentService : CoroutineScope { */ fun processMusicPlay(song: Map){ // 在其他 Service、BroadcastReceiver 或 Application 中调用 - Log.i(TAG, "启动音乐服务 播放音乐 $song") MusicServiceStarter.startServiceWithCommand(context, command = "play", song = song) } @@ -1206,7 +1182,6 @@ object AgentService : CoroutineScope { fun processMusicPlayList(songs:List>){ // 如果 playlist 中有有效的歌曲,开始播放 if (songs.isNotEmpty()) { - Log.i(TAG, "启动音乐服务 播放音乐列表 ${songs}") MusicServiceStarter.startServiceWithPlaylist(context, command = "playlist", songs = songs) } else { Log.i(TAG, "音乐列表为空,未启动播放服务") @@ -1224,7 +1199,6 @@ object AgentService : CoroutineScope { fun processNavigation(start:String,end:String){ // 如果 playlist 中有有效的歌曲,开始播放 if (!start.isNullOrEmpty() && !end.isNullOrEmpty()) { - Log.i(TAG, "启动导航服务 ${start} ${end}") NavigationServiceHelper.startNavigation(context,"start", startpos = start, endpos = end) } else { Log.i(TAG, "启动导航服务失败") diff --git a/local_plugins/azure_speech/android/src/main/kotlin/com/yunqiinnovation/azure_speech/AzureAsrHelper.kt b/local_plugins/azure_speech/android/src/main/kotlin/com/yunqiinnovation/azure_speech/AzureAsrHelper.kt index 31a09c506..52a86ce10 100644 --- a/local_plugins/azure_speech/android/src/main/kotlin/com/yunqiinnovation/azure_speech/AzureAsrHelper.kt +++ b/local_plugins/azure_speech/android/src/main/kotlin/com/yunqiinnovation/azure_speech/AzureAsrHelper.kt @@ -123,16 +123,26 @@ class AzureAsrHelper(private val context: Context) { speechRecognitionLanguage = currentLanguage } - // 设置静音超时时间(毫秒) - setProperty("SpeechServiceConnection_EndSilenceTimeoutMs", "800") - setProperty("Speech_SegmentationSilenceTimeoutMs", "800") + // 优化:大幅减少静音超时时间(从800ms到300ms) + setProperty("SpeechServiceConnection_EndSilenceTimeoutMs", "300") + setProperty("Speech_SegmentationSilenceTimeoutMs", "300") + + // 优化:添加低延迟连接配置 + setProperty("SpeechServiceConnection_InitialSilenceTimeoutMs", "200") + setProperty("SpeechServiceConnection_RecoMode", "INTERACTIVE") + setProperty("Speech_PushStreamFormat", "PCM") // 设置分段策略为时间模式 setProperty("Speech_SegmentationStrategy", "Time") } // 创建识别器 - return setupRecognizer() + val setupSuccess = setupRecognizer() + if (setupSuccess) { + // 优化:初始化完成后进行预热 + warmupRecognizer() + } + return setupSuccess } catch (e: Exception) { Log.e(tag, "初始化失败: ${e.message}") return false @@ -216,6 +226,36 @@ class AzureAsrHelper(private val context: Context) { } } + /** + * 预热识别器(减少首次识别延迟) + */ + private fun warmupRecognizer() { + try { + // 创建极短的音频数据进行预热 + val warmupData = ByteArray(320) // 10ms 16kHz 单声道 + // 模拟静音数据 + warmupData.fill(0) + + // 如果是外部音频流,进行预热 + if (audioSourceType == AudioSourceType.EXTERNAL) { + externalAudioStream?.pushAudio(warmupData) + } + Log.d(tag, "ASR预热完成") + } catch (e: Exception) { + // 预热失败不影响正常使用 + Log.d(tag, "ASR预热失败: ${e.message}") + } + } + + /** + * 强制设置单一语言(优化语言检测逻辑) + */ + fun setForceLanguage(language: String) { + this.currentLanguage = language + this.isAutoDetectLanguage = false + speechConfig?.speechRecognitionLanguage = language + } + /** * 设置外部音频流 - 使用拉流方式 */ @@ -337,21 +377,30 @@ class AzureAsrHelper(private val context: Context) { * 设置事件监听器 */ private fun setupEventListeners(callback: ContinuousRecognizeCallback) { - // 识别中事件 + // 优化:识别中事件 - 添加文本长度检查 recognizer?.recognizing?.addEventListener( EventHandler { _, event -> - val detectedLanguage = AutoDetectSourceLanguageResult.fromResult(event.result)?.language ?: "" - // 直接在当前线程调用回调 - callback.onRecognizing(event.result.text, detectedLanguage) + // 优化:只处理非空结果 + if (event.result.text.isNotEmpty()) { + val detectedLanguage = if (isAutoDetectLanguage) { + AutoDetectSourceLanguageResult.fromResult(event.result)?.language ?: "" + } else { + currentLanguage + } + callback.onRecognizing(event.result.text, detectedLanguage) + } } ) - // 识别完成事件 + // 优化:识别完成事件 - 优化语言检测 recognizer?.recognized?.addEventListener( EventHandler { _, event -> - if (event.result.reason == ResultReason.RecognizedSpeech) { - val detectedLanguage = AutoDetectSourceLanguageResult.fromResult(event.result)?.language ?: "" - // 直接在当前线程调用回调 + if (event.result.reason == ResultReason.RecognizedSpeech && event.result.text.isNotEmpty()) { + val detectedLanguage = if (isAutoDetectLanguage) { + AutoDetectSourceLanguageResult.fromResult(event.result)?.language ?: supportedLanguages[0] + } else { + currentLanguage + } callback.onResult(event.result.text, detectedLanguage) } } @@ -538,19 +587,14 @@ class AzureAsrHelper(private val context: Context) { private inner class MicrophoneStream : PullAudioInputStreamCallback() { private var audioRecord: AudioRecord? = null private var echoCanceler: AcousticEchoCanceler? = null - private var currentAudioFile: File? = null - private var fos: FileOutputStream? = null - - // 用于存储音频数据的缓冲区 - private val dataBuffer = mutableListOf() - private var totalBytesWritten = 0 // 音频配置 private val sampleRate = 16000 private val channelConfig = AudioFormat.CHANNEL_IN_MONO private val audioFormat = AudioFormat.ENCODING_PCM_16BIT + // 优化:减少缓冲区大小以降低延迟 private val bufferSize = AudioRecord.getMinBufferSize( sampleRate, channelConfig, audioFormat - ).let { if (it < 0) 4096 else it * 2 } + ).let { if (it < 0) 2048 else it } init { initMicrophone() @@ -561,20 +605,6 @@ class AzureAsrHelper(private val context: Context) { */ private fun initMicrophone() { try { - // 创建新的音频文件 - val timestamp = System.currentTimeMillis() - val filePath = context.getExternalFilesDir(null)?.absolutePath + "/recorded_audio_$timestamp.wav" - currentAudioFile = File(filePath) - currentAudioFile?.createNewFile() - - // 仅写入初始文件头 - fos = FileOutputStream(currentAudioFile).apply { - write(generateWavHeader(0)) // 初始0长度 - close() - } - - // 追加音频数据到文件 - fos = FileOutputStream(currentAudioFile, true) // 创建录音对象 if (android.os.Build.VERSION.SDK_INT >= android.os.Build.VERSION_CODES.M) { val format = AudioFormat.Builder() @@ -619,14 +649,20 @@ class AzureAsrHelper(private val context: Context) { } /** - * 应用音频效果 + * 优化:条件性应用音频效果 */ private fun applyAudioEffects() { val sessionId = audioRecord?.audioSessionId ?: return - // 回音消除 - echoCanceler = AcousticEchoCanceler.create(sessionId) - echoCanceler?.enabled = true; + // 优化:只在真正需要时启用回音消除 + val audioManager = context.getSystemService(Context.AUDIO_SERVICE) as AudioManager + if (audioManager.isSpeakerphoneOn) { + echoCanceler = AcousticEchoCanceler.create(sessionId) + echoCanceler?.enabled = true + Log.d(tag, "启用回音消除") + } else { + Log.d(tag, "跳过回音消除(使用耳机)") + } } /** @@ -638,6 +674,9 @@ class AzureAsrHelper(private val context: Context) { * Azure SDK 调用此方法获取音频数据 * 这是拉流模式的核心方法,由SDK调用以获取音频数据 */ + /** + * 优化:简化音频数据读取(移除文件保存) + */ override fun read(buffer: ByteArray): Int { try { val result = audioRecord?.read(buffer, 0, buffer.size) ?: -1 @@ -645,80 +684,18 @@ class AzureAsrHelper(private val context: Context) { Log.e(tag, "读取音频数据失败: $result") return 0 } - // 将音频数据保存成wav格式的音频文件 - saveAudioDataToWav(buffer, result) return result } catch (e: Exception) { Log.e(tag, "读取音频异常: ${e.message}") return 0 } } - /** - * 保存音频数据到 WAV 文件 - */ -private fun saveAudioDataToWav(buffer: ByteArray, length: Int) { - // 复制数据到新数组 - val copy = ByteArray(length) - System.arraycopy(buffer, 0, copy, 0, length) - - // 添加到缓冲区 - dataBuffer.add(copy) - totalBytesWritten += length - - // 写入文件 - try { - fos?.write(copy, 0, length) - } catch (e: Exception) { - Log.e(tag, "写入音频数据失败: ${e.message}") - } -} - // 更新文件头生成(修正RIFF长度计算) - private fun generateWavHeader(dataLength: Int): ByteArray { - val totalLength = 36 + dataLength // RIFF块总长度 = 头部36字节 + 音频数据 - val byteRate = sampleRate * 2 * 1 // 采样率 * 字节/样本 * 通道数 - - return byteArrayOf( - 'R'.code.toByte(), 'I'.code.toByte(), 'F'.code.toByte(), 'F'.code.toByte(), - (totalLength and 0xFF).toByte(), ((totalLength shr 8) and 0xFF).toByte(), - ((totalLength shr 16) and 0xFF).toByte(), ((totalLength shr 24) and 0xFF).toByte(), - 'W'.code.toByte(), 'A'.code.toByte(), 'V'.code.toByte(), 'E'.code.toByte(), - 'f'.code.toByte(), 'm'.code.toByte(), 't'.code.toByte(), ' '.code.toByte(), - 16, 0, 0, 0, // PCM头长度 - 1, 0, // PCM格式 - 1, 0, // 单声道 - (sampleRate and 0xFF).toByte(), ((sampleRate shr 8) and 0xFF).toByte(), 0, 0, // 采样率 - (byteRate and 0xFF).toByte(), ((byteRate shr 8) and 0xFF).toByte(), 0, 0, // 字节率 - 2, 0, // 块对齐 (通道数 * 样本位数/8) - 16, 0, // 样本位数 - 'd'.code.toByte(), 'a'.code.toByte(), 't'.code.toByte(), 'a'.code.toByte(), - (dataLength and 0xFF).toByte(), ((dataLength shr 8) and 0xFF).toByte(), - ((dataLength shr 16) and 0xFF).toByte(), ((dataLength shr 24) and 0xFF).toByte() - ) - } /** - * 关闭音频资源 + * 优化:简化资源关闭 */ override fun close() { releaseAudioResources() - try { - // 关闭文件流 - fos?.close() - - // 更新WAV文件头 - currentAudioFile?.let { file -> - RandomAccessFile(file, "rw").use { raf -> - raf.seek(0) - raf.write(generateWavHeader(totalBytesWritten)) - } - Log.d(tag, "音频文件保存完成: ${file.absolutePath}") - } - } catch (e: Exception) { - Log.e(tag, "更新WAV文件头失败: ${e.message}") - } finally { - fos = null - currentAudioFile = null - } } /** @@ -763,14 +740,20 @@ private fun saveAudioDataToWav(buffer: ByteArray, length: Int) { } /** - * SDK调用:从队列中拉取数据 + * 修复:外部音频流恢复阻塞读取 * @param buffer SDK提供的缓冲区 * @return 读取的字节数,0表示流结束 */ override fun read(buffer: ByteArray): Int { try { - // 阻塞等待下一块数据 + // 修复:恢复阻塞等待,确保外部音频数据完整性 val chunk = queue.take() + + // 检查是否是结束标志(空数组) + if (chunk.isEmpty()) { + return 0 + } + val toCopy = minOf(chunk.size, buffer.size) System.arraycopy(chunk, 0, buffer, 0, toCopy) return toCopy diff --git a/local_plugins/azure_speech/android/src/main/kotlin/com/yunqiinnovation/azure_speech/AzureTtsHelper.kt b/local_plugins/azure_speech/android/src/main/kotlin/com/yunqiinnovation/azure_speech/AzureTtsHelper.kt index 5f973bb44..1e5251de5 100644 --- a/local_plugins/azure_speech/android/src/main/kotlin/com/yunqiinnovation/azure_speech/AzureTtsHelper.kt +++ b/local_plugins/azure_speech/android/src/main/kotlin/com/yunqiinnovation/azure_speech/AzureTtsHelper.kt @@ -70,7 +70,14 @@ class AzureTtsHelper(private val context: Context) : ITtsService { // 创建语音配置 speechConfig = SpeechConfig.fromSubscription(ttsAppToken, ttsResource) - speechConfig?.setSpeechSynthesisOutputFormat(SpeechSynthesisOutputFormat.Riff24Khz16BitMonoPcm) + + // 优化:使用更低延迟的音频格式(从24KHz降到16KHz) + speechConfig?.setSpeechSynthesisOutputFormat(SpeechSynthesisOutputFormat.Riff16Khz16BitMonoPcm) + + // 优化:设置低延迟连接属性 + speechConfig?.setProperty("SpeechServiceConnection_InitialSilenceTimeoutMs", "300") + speechConfig?.setProperty("SpeechServiceConnection_EndSilenceTimeoutMs", "300") + speechConfig?.setSpeechSynthesisVoiceName(currentVoice) // 创建音频配置 @@ -91,6 +98,9 @@ class AzureTtsHelper(private val context: Context) : ITtsService { isInitialized = true FileLogger.d(TAG, "TTS引擎初始化成功: 区域=$ttsResource") + // 优化:初始化完成后立即预热 + warmupSynthesizer() + return true } catch (e: Exception) { FileLogger.e(TAG, "TTS引擎初始化失败: ${e.message}") @@ -243,6 +253,30 @@ class AzureTtsHelper(private val context: Context) : ITtsService { )) } } + + /** + * 预热TTS引擎(减少首次合成延迟) + */ + private fun warmupSynthesizer() { + try { + // 使用极短文本进行预热 + val warmupSsml = """ + + + + . + + + + """.trimIndent() + + // 静音预热(音量设为0) + synthesizer?.SpeakSsmlAsync(warmupSsml) + } catch (e: Exception) { + // 预热失败不影响正常使用 + FileLogger.d(TAG, "TTS预热失败: ${e.message}") + } + } //清洗播报语音内容 fun cleanTextForTTS(text: String): String { return text @@ -320,6 +354,30 @@ class AzureTtsHelper(private val context: Context) : ITtsService { """.trimIndent() } + /** + * 生成优化的SSML(减少复杂度提升速度) + */ + private fun generateOptimizedSsml(rawText: String): String { + // 优化:简化文本预处理,减少正则表达式使用 + val processedText = rawText + .replace("&", "&") + .replace("<", "<") + .replace(">", ">") + .replace(Regex("[😀-🟿]+"), "") // 简化表情符号移除 + .trim() + + // 优化:简化SSML结构,减少嵌套层级 + return """ + + + + $processedText + + + + """.trimIndent() + } + /** * 单次播放文本(非流式) */ @@ -334,8 +392,8 @@ class AzureTtsHelper(private val context: Context) : ITtsService { } try { - // 生成SSML并播放 - val ssml = generateSsml(text) + // 优化:使用简化的SSML生成 + val ssml = generateOptimizedSsml(text) isSpeaking = true @@ -373,21 +431,22 @@ class AzureTtsHelper(private val context: Context) : ITtsService { // 添加新文本到缓冲区 streamBuffer.append(text) - // 增加500ms防抖逻辑 + // 优化:大幅减少防抖时间从600ms到150ms val currentTime = System.currentTimeMillis() - if (currentTime - lastSpeakTime < 600) { + if (currentTime - lastSpeakTime < 150) { return true } lastSpeakTime = currentTime val currentText = streamBuffer.toString() //cleanTextForTTS(streamBuffer.toString()) - // 定义标点符号列表 - val punctuationMarks = listOf('.', '。', '!', '!', '?', '?', ';', ';', ',', ',', ':', ':', '\n') - // 查找最后一个标点符号的位置 + // 优化:使用字符集合替代列表,提升查找效率 + val punctuationSet = setOf('.', '。', '!', '!', '?', '?', ';', ';', ',', ',', ':', ':', '\n') + + // 优化:从后往前查找最后一个标点符号 var lastPunctuationIndex = -1 - for (i in currentText.indices.reversed()) { - if (currentText[i] in punctuationMarks) { + for (i in currentText.length - 1 downTo 0) { + if (currentText[i] in punctuationSet) { lastPunctuationIndex = i break } diff --git a/local_plugins/chat_api/android/src/main/kotlin/com/yunqiinnovation/chat_api/ChatApiService.kt b/local_plugins/chat_api/android/src/main/kotlin/com/yunqiinnovation/chat_api/ChatApiService.kt index d54ce8298..190e2a1f0 100644 --- a/local_plugins/chat_api/android/src/main/kotlin/com/yunqiinnovation/chat_api/ChatApiService.kt +++ b/local_plugins/chat_api/android/src/main/kotlin/com/yunqiinnovation/chat_api/ChatApiService.kt @@ -68,7 +68,13 @@ private data class ToolCallInfo( var name: String = "", var arguments: String = "" ) { - fun isValid(): Boolean = id.isNotEmpty() && name.isNotEmpty() + fun isValid(): Boolean { + return try { + id.isNotEmpty() && name.isNotEmpty() + } catch (e: Exception) { + false + } + } } /** @@ -108,6 +114,8 @@ class ChatApiService(private val context: android.content.Context? = null) : Cor private var toolCalls: MutableMap = mutableMapOf() private var isCanceled = false + + // JSON处理 private val gson = Gson() @@ -124,13 +132,6 @@ class ChatApiService(private val context: android.content.Context? = null) : Cor * 初始化ChatAPI服务 */ fun initialize(apiKey: String, baseUrl: String, model: String, mcpServer: String): Boolean { - Log.d("ChatApiService", "=== ChatApiService 初始化开始 ===") - Log.d("ChatApiService", "参数检查:") - Log.d("ChatApiService", " - apiKey: ${if (apiKey.isNotEmpty()) "已提供(${apiKey.length}字符)" else "未提供"}") - Log.d("ChatApiService", " - baseUrl: '$baseUrl'") - Log.d("ChatApiService", " - model: '$model'") - Log.d("ChatApiService", " - mcpServer: '$mcpServer'") - this.apiKey = apiKey if (baseUrl.isNotEmpty()) { this.baseUrl = baseUrl @@ -138,7 +139,6 @@ class ChatApiService(private val context: android.content.Context? = null) : Cor if (model.isNotEmpty()) { this.model = model } - Log.d("ChatApiService", "原始 baseUrl: $baseUrl") // 处理 baseUrl:移除末尾的 /chat/completions(如果存在) // 因为 openai-kotlin 会自动拼接 /chat/completions @@ -153,24 +153,16 @@ class ChatApiService(private val context: android.content.Context? = null) : Cor processedBaseUrl += "/" } - Log.d("ChatApiService", "处理后 baseUrl: $processedBaseUrl") - return try { - Log.d("ChatApiService", "开始创建OpenAI配置...") // 创建OpenAI配置 val config = OpenAIConfig( token = apiKey, timeout = Timeout(socket = 60.seconds), host = OpenAIHost(baseUrl = processedBaseUrl) ) - Log.d("ChatApiService", "OpenAI配置创建成功") - Log.d("ChatApiService", "正在创建OpenAI客户端...") openAI = OpenAI(config) - Log.d("ChatApiService", "OpenAI客户端创建成功") - // 初始化MCP客户端 - Log.d("ChatApiService", "准备初始化MCP客户端...") // 保存MCP配置以便后续使用 mcpConfigJson = mcpServer // 异步初始化MCP客户端 @@ -179,7 +171,6 @@ class ChatApiService(private val context: android.content.Context? = null) : Cor } isInitialized = apiKey.isNotEmpty() - Log.d("ChatApiService", "ChatApiService初始化完成,isInitialized: $isInitialized") true } catch (e: Exception) { Log.e("ChatApiService", "ChatApiService初始化失败: ${e.message}", e) @@ -231,72 +222,8 @@ class ChatApiService(private val context: android.content.Context? = null) : Cor // 转换消息格式 val chatMessages = convertToChatMessages(messages) - // 获取MCP工具列表 - val tools = mutableListOf() - mcpClient?.getToolMaps()?.forEach { toolMap -> - try { - Log.d("ChatApiService", "处理工具映射: $toolMap") - - val type = toolMap["type"] as? String - if (type == "function") { - @Suppress("UNCHECKED_CAST") - val functionMap = toolMap["function"] as? Map - if (functionMap == null) { - Log.e("ChatApiService", "工具function映射为null") - return@forEach - } - - val name = functionMap["name"] as? String - if (name == null) { - Log.e("ChatApiService", "工具name为null") - return@forEach - } - - val description = functionMap["description"] as? String ?: "" - val parametersMap = functionMap["parameters"] as? Map - - Log.d("ChatApiService", "工具信息 - 名称: $name, 描述: $description") - Log.d("ChatApiService", "参数映射: $parametersMap") - - if (parametersMap == null) { - Log.e("ChatApiService", "工具 $name 的parameters为null") - return@forEach - } - - // 验证parametersMap的基本结构 - if (!parametersMap.containsKey("type")) { - Log.e("ChatApiService", "工具 $name 的parameters缺少type字段") - return@forEach - } - - val parametersJson = gson.toJson(parametersMap) - Log.d("ChatApiService", "参数JSON: $parametersJson") - - try { - val parameters = com.aallam.openai.api.core.Parameters.fromJsonString(parametersJson) - Log.d("ChatApiService", "成功创建Parameters对象") - - tools.add(Tool.function( - name = name, - description = description, - parameters = parameters - )) - - Log.d("ChatApiService", "成功添加工具: $name") - - } catch (e: Exception) { - Log.e("ChatApiService", "创建Parameters对象失败 for 工具 $name: ${e.message}", e) - Log.e("ChatApiService", "失败的参数JSON: $parametersJson") - // 跳过这个工具,继续处理其他工具 - } - } else { - Log.d("ChatApiService", "跳过非function类型的工具: $type") - } - } catch (e: Exception) { - Log.e("ChatApiService", "处理工具映射时出错: ${e.message}", e) - Log.e("ChatApiService", "出错的工具映射: $toolMap") - } - } + // 直接获取工具列表 + val tools = getOpenAiTools() // 构建请求 val chatCompletionRequest = ChatCompletionRequest( @@ -344,12 +271,6 @@ class ChatApiService(private val context: android.content.Context? = null) : Cor * 发送消息(流式输出) */ fun sendMessageStream(messages: List>) { - Log.d("ChatApiService", "=== sendMessageStream 开始 ===") - Log.d("ChatApiService", "isInitialized: $isInitialized") - Log.d("ChatApiService", "apiKey isEmpty: ${apiKey.isEmpty()}") - Log.d("ChatApiService", "openAI is null: ${openAI == null}") - Log.d("ChatApiService", "messages size: ${messages.size}") - if (!isInitialized || apiKey.isEmpty() || openAI == null) { Log.e("ChatApiService", "ChatAPI服务未初始化,无法发送消息") streamCallback?.onError(ChatApiException("ChatAPI服务未初始化")) @@ -363,140 +284,28 @@ class ChatApiService(private val context: android.content.Context? = null) : Cor // 检查是否包含图片,决定使用哪个模型 val currentModel = if (containsImage(messages)) visionModel else model - Log.d("ChatApiService", "选择的模型: $currentModel (默认模型: $model, 视觉模型: $visionModel)") currentStreamJob = launch { try { - Log.e("ChatApiService", "=== 开始流式请求 ===") - Log.e("ChatApiService", "消息数量: ${messages.size}") - Log.e("ChatApiService", "当前模型: $currentModel") - Log.e("ChatApiService", "MCP客户端状态: ${mcpClient != null}") - // 转换消息格式 - Log.e("ChatApiService", "开始转换消息格式...") val chatMessages = try { convertToChatMessages(messages) } catch (e: Exception) { - Log.e("ChatApiService", "ERROR: 转换消息格式失败: ${e.message}", e) + Log.e("ChatApiService", "转换消息格式失败: ${e.message}", e) throw e } - Log.e("ChatApiService", "转换后的聊天消息数量: ${chatMessages.size}") - // 验证转换后的消息 - chatMessages.forEachIndexed { index, message -> - Log.e("ChatApiService", "消息 $index: role=${message.role}, content类型=${message.content?.javaClass?.simpleName}") - } - - // 获取MCP工具列表 - Log.e("ChatApiService", "开始获取MCP工具列表...") - val tools = mutableListOf() - - val toolMaps = mcpClient?.getToolMaps() - Log.e("ChatApiService", "获取到的工具映射数量: ${toolMaps?.size ?: 0}") - - toolMaps?.forEach { toolMap -> - try { - Log.e("ChatApiService", "处理工具映射: $toolMap") - - val type = toolMap["type"] as? String - if (type == "function") { - @Suppress("UNCHECKED_CAST") - val functionMap = toolMap["function"] as? Map - if (functionMap == null) { - Log.e("ChatApiService", "ERROR: 工具function映射为null") - return@forEach - } - - val name = functionMap["name"] as? String - if (name == null) { - Log.e("ChatApiService", "ERROR: 工具name为null") - return@forEach - } - - val description = functionMap["description"] as? String ?: "" - val parametersMap = functionMap["parameters"] as? Map - - Log.e("ChatApiService", "工具信息 - 名称: $name, 描述: $description") - Log.e("ChatApiService", "参数映射: $parametersMap") - - if (parametersMap == null) { - Log.e("ChatApiService", "ERROR: 工具 $name 的parameters为null") - return@forEach - } - - // 验证parametersMap的基本结构 - if (!parametersMap.containsKey("type")) { - Log.e("ChatApiService", "ERROR: 工具 $name 的parameters缺少type字段") - return@forEach - } - - val parametersJson = gson.toJson(parametersMap) - Log.e("ChatApiService", "参数JSON: $parametersJson") - - try { - Log.e("ChatApiService", "正在调用 Parameters.fromJsonString...") - Log.e("ChatApiService", "参数JSON长度: ${parametersJson.length}") - Log.e("ChatApiService", "参数JSON内容预览: ${parametersJson.take(200)}...") - - if (parametersJson.isBlank()) { - Log.e("ChatApiService", "ERROR: parametersJson为空或空白") - return@forEach - } - - val parameters = com.aallam.openai.api.core.Parameters.fromJsonString(parametersJson) - Log.e("ChatApiService", "SUCCESS: 成功创建Parameters对象") - - Log.e("ChatApiService", "正在创建Tool.function...") - Log.e("ChatApiService", "工具参数 - name: '$name' (isEmpty: ${name.isEmpty()})") - Log.e("ChatApiService", "工具参数 - description: '$description' (isEmpty: ${description.isEmpty()})") - Log.e("ChatApiService", "工具参数 - parameters: ${parameters != null}") - - if (name.isEmpty()) { - Log.e("ChatApiService", "ERROR: 工具名称为空,跳过") - return@forEach - } - - val tool = Tool.function( - name = name, - description = description, - parameters = parameters - ) - tools.add(tool) - - Log.e("ChatApiService", "SUCCESS: 成功添加工具: $name") - - } catch (e: Exception) { - Log.e("ChatApiService", "ERROR: 创建Parameters对象失败 for 工具 $name: ${e.message}", e) - Log.e("ChatApiService", "失败的参数JSON: $parametersJson") - Log.e("ChatApiService", "异常堆栈: ${e.stackTraceToString()}") - // 跳过这个工具,继续处理其他工具 - } - } else { - Log.e("ChatApiService", "跳过非function类型的工具: $type") - } - } catch (e: Exception) { - Log.e("ChatApiService", "ERROR: 处理工具映射时出错: ${e.message}", e) - Log.e("ChatApiService", "出错的工具映射: $toolMap") - } - } - - Log.e("ChatApiService", "最终工具列表大小: ${tools.size}") + // 直接获取工具列表 + val tools = getOpenAiTools() // 构建请求 - Log.e("ChatApiService", "正在构建ChatCompletionRequest...") - Log.e("ChatApiService", "构建参数检查:") - Log.e("ChatApiService", " - currentModel: '$currentModel' (isEmpty: ${currentModel.isEmpty()})") - Log.e("ChatApiService", " - chatMessages size: ${chatMessages.size}") - Log.e("ChatApiService", " - tools size: ${tools.size}") - Log.e("ChatApiService", " - openAI对象: ${openAI != null}") - if (currentModel.isEmpty()) { - Log.e("ChatApiService", "ERROR: 模型名称为空") + Log.e("ChatApiService", "模型名称为空") throw IllegalArgumentException("模型名称不能为空") } if (chatMessages.isEmpty()) { - Log.e("ChatApiService", "ERROR: 消息列表为空") + Log.e("ChatApiService", "消息列表为空") throw IllegalArgumentException("消息列表不能为空") } @@ -508,53 +317,41 @@ class ChatApiService(private val context: android.content.Context? = null) : Cor temperature = 0.7, tools = if (tools.isNotEmpty()) tools else null ) - Log.e("ChatApiService", "SUCCESS: ChatCompletionRequest创建成功") if (openAI == null) { - Log.e("ChatApiService", "ERROR: openAI对象为null") + Log.e("ChatApiService", "openAI对象为null") throw IllegalStateException("OpenAI客户端未初始化") } - Log.e("ChatApiService", "正在调用 openAI.chatCompletions...") val flow = openAI!!.chatCompletions(chatCompletionRequest) - Log.e("ChatApiService", "SUCCESS: 获得chatsFlow,开始收集...") flow } catch (e: Exception) { - Log.e("ChatApiService", "ERROR: 创建ChatCompletionRequest或调用chatCompletions失败: ${e.message}", e) + Log.e("ChatApiService", "创建ChatCompletionRequest或调用chatCompletions失败: ${e.message}", e) throw e } chatsFlow.collect { result -> if (isCanceled) return@collect - Log.e("ChatApiService", "收到流式响应,结果类型: ${result.javaClass.simpleName}") - Log.e("ChatApiService", "result.choices大小: ${result.choices?.size ?: 0}") - val choice = result.choices?.firstOrNull() if (choice == null) { - Log.e("ChatApiService", "ERROR: choice为null") + Log.e("ChatApiService", "choice为null") return@collect } - Log.e("ChatApiService", "choice存在,delta类型: ${choice.delta?.javaClass?.simpleName ?: "null"}") - val delta = choice.delta if (delta == null) { - Log.e("ChatApiService", "ERROR: delta为null") + Log.e("ChatApiService", "delta为null") return@collect } - Log.e("ChatApiService", "delta存在,content: ${delta.content}, toolCalls: ${delta.toolCalls?.size ?: 0}") - // 处理普通文本内容 delta.content?.let { content -> - Log.e("ChatApiService", "收到文本内容: $content") streamCallback?.onToken(content) } // 收集工具调用信息 delta.toolCalls?.forEach { toolCall -> - Log.e("ChatApiService", "处理工具调用: index=${toolCall.index}, id=${toolCall.id}") val index = toolCall.index // 创建或获取现有的工具调用信息 @@ -564,7 +361,6 @@ class ChatApiService(private val context: android.content.Context? = null) : Cor try { toolCall.id?.let { id -> toolCallInfo.id = id.toString() - Log.e("ChatApiService", "更新工具调用ID: $id") } } catch (e: Exception) { Log.e("ChatApiService", "处理工具调用ID异常: ${e.message}") @@ -573,13 +369,20 @@ class ChatApiService(private val context: android.content.Context? = null) : Cor // 安全处理函数信息 try { toolCall.function?.let { function -> - function.name?.let { name -> - toolCallInfo.name = name - Log.e("ChatApiService", "更新工具名称: $name") + try { + function.name?.let { name -> + toolCallInfo.name = name + } + } catch (e: Exception) { + Log.e("ChatApiService", "处理工具调用函数名称异常: ${e.message}") } - function.arguments?.let { args -> - toolCallInfo.arguments += args - Log.e("ChatApiService", "添加参数片段: $args") + + try { + function.arguments?.let { args -> + toolCallInfo.arguments += args + } + } catch (e: Exception) { + Log.e("ChatApiService", "处理工具调用参数异常: ${e.message}") } } } catch (e: Exception) { @@ -589,7 +392,7 @@ class ChatApiService(private val context: android.content.Context? = null) : Cor } if (!isCanceled) { - // 处理工具调用或完成 + // 检查是否有工具调用需要处理 val hasToolCalls = processToolCalls() if (!hasToolCalls) { streamCallback?.onComplete() @@ -608,25 +411,28 @@ class ChatApiService(private val context: android.content.Context? = null) : Cor * 处理工具调用 */ private suspend fun processToolCalls(): Boolean { - Log.d("ChatApiService", "开始处理工具调用,工具调用数量: ${toolCalls.size}") + val firstToolCall = try { + toolCalls.values.firstOrNull { it.isValid() } + } catch (e: Exception) { + Log.e("ChatApiService", "查找有效工具调用异常: ${e.message}") + null + } - val firstToolCall = toolCalls.values.firstOrNull { it.isValid() } if (firstToolCall == null) { - Log.d("ChatApiService", "没有有效的工具调用") return false } - Log.d("ChatApiService", "第一个有效的工具调用: name=${firstToolCall.name}, id=${firstToolCall.id}") - Log.d("ChatApiService", "工具调用参数: ${firstToolCall.arguments}") - // 创建函数调用字典 - val functionCall = mapOf( - "name" to firstToolCall.name, - "arguments" to firstToolCall.arguments, - "id" to firstToolCall.id - ) - - Log.d("ChatApiService", "创建的函数调用字典: $functionCall") + val functionCall = try { + mapOf( + "name" to (firstToolCall.name.takeIf { it.isNotEmpty() } ?: ""), + "arguments" to (firstToolCall.arguments.takeIf { it.isNotEmpty() } ?: "{}"), + "id" to (firstToolCall.id.takeIf { it.isNotEmpty() } ?: "") + ) + } catch (e: Exception) { + Log.e("ChatApiService", "创建函数调用字典异常: ${e.message}") + return false + } // 通知上层工具调用事件 streamCallback?.onFunctionCall(convertMapToJsonObject(functionCall)) @@ -639,36 +445,25 @@ class ChatApiService(private val context: android.content.Context? = null) : Cor val functionName = firstToolCall.name val argumentsJson = firstToolCall.arguments - Log.d("ChatApiService", "准备调用MCP工具: $functionName") - Log.d("ChatApiService", "参数JSON: $argumentsJson") - val result = if (_mcpClient?.hasToolWithName(functionName) == true) { - Log.d("ChatApiService", "MCP客户端中找到工具: $functionName") - // 解析参数 val arguments = _mcpClient?.parseJsonArguments(argumentsJson) ?: emptyMap() - Log.d("ChatApiService", "解析后的参数: $arguments") // 调用MCP工具 - Log.d("ChatApiService", "开始调用MCP工具...") val toolResult = _mcpClient?.callTool(functionName, arguments) - Log.d("ChatApiService", "MCP工具调用结果: $toolResult") // 处理结果 if (toolResult != null) { if (toolResult["isError"] == true) { - Log.w("ChatApiService", "MCP工具调用返回错误") // 处理错误情况 val content = toolResult["content"] as? List<*> val firstContent = content?.firstOrNull() as? Map<*, *> val errorText = firstContent?.get("text") as? String ?: "Tool execution failed" mapOf("context" to errorText) } else if (toolResult.containsKey("context")) { - Log.d("ChatApiService", "本地函数结果") // 本地函数结果 toolResult } else { - Log.d("ChatApiService", "MCP工具结果") // MCP工具结果 val content = toolResult["content"] as? List<*> val firstContent = content?.firstOrNull() as? Map<*, *> @@ -680,13 +475,10 @@ class ChatApiService(private val context: android.content.Context? = null) : Cor mapOf("context" to "Tool call failed") } } else { - Log.w("ChatApiService", "MCP客户端中未找到工具: $functionName") // 工具不存在 mapOf("context" to "Tool not found: $functionName") } - Log.d("ChatApiService", "最终处理结果: $result") - if (!isCanceled) { // 处理结果 streamCallback?.onFunctionCallResult( @@ -808,6 +600,8 @@ class ChatApiService(private val context: android.content.Context? = null) : Cor handler = handler ) + + true } catch (e: Exception) { Log.e("ChatApiService", "Failed to register function: $name", e) @@ -826,8 +620,8 @@ class ChatApiService(private val context: android.content.Context? = null) : Cor // 直接使用类的CoroutineScope启动协程 launch { try { - val result = _mcpClient?.connectToSSE(serverUrl) ?: false - Log.d("ChatApiService", "MCP客户端初始化${if (result) "成功" else "失败"}") + _mcpClient?.connectToSSE(serverUrl) + } catch (e: Exception) { Log.e("ChatApiService", "MCP客户端初始化失败: ${e.message}", e) } @@ -843,6 +637,88 @@ class ChatApiService(private val context: android.content.Context? = null) : Cor return _mcpClient?.isConnected() ?: false } + /** + * 直接从MCPClient获取OpenAI工具格式 + * 将MCP工具映射转换为OpenAI工具格式 + */ + private fun getOpenAiTools(): List { + try { + val tools = mutableListOf() + val toolMaps = mcpClient?.getToolMaps() ?: return emptyList() + + toolMaps.forEach { toolMap -> + try { + val type = toolMap["type"] as? String + if (type == "function") { + @Suppress("UNCHECKED_CAST") + val functionMap = toolMap["function"] as? Map + if (functionMap == null) { + Log.e("ChatApiService", "工具function映射为null") + return@forEach + } + + val name = functionMap["name"] as? String + if (name == null) { + Log.e("ChatApiService", "工具name为null") + return@forEach + } + + val description = functionMap["description"] as? String ?: "" + val parametersMap = functionMap["parameters"] as? Map + + if (parametersMap == null) { + Log.e("ChatApiService", "工具 $name 的parameters为null") + return@forEach + } + + // 验证parametersMap的基本结构 + if (!parametersMap.containsKey("type")) { + Log.e("ChatApiService", "工具 $name 的parameters缺少type字段") + return@forEach + } + + val parametersJson = gson.toJson(parametersMap) + + try { + if (parametersJson.isBlank()) { + Log.e("ChatApiService", "parametersJson为空或空白") + return@forEach + } + + val parameters = com.aallam.openai.api.core.Parameters.fromJsonString(parametersJson) + + if (name.isEmpty()) { + Log.e("ChatApiService", "工具名称为空,跳过") + return@forEach + } + + val tool = Tool.function( + name = name, + description = description, + parameters = parameters + ) + tools.add(tool) + + } catch (e: Exception) { + Log.e("ChatApiService", "创建Parameters对象失败 for 工具 $name: ${e.message}", e) + // 跳过这个工具,继续处理其他工具 + } + } + } catch (e: Exception) { + Log.e("ChatApiService", "处理工具映射时出错: ${e.message}", e) + } + } + + return tools + + } catch (e: Exception) { + Log.e("ChatApiService", "获取OpenAI工具格式失败: ${e.message}", e) + return emptyList() + } + } + + + /** * 关闭MCP客户端 */ @@ -864,7 +740,7 @@ class ChatApiService(private val context: android.content.Context? = null) : Cor return "MCP客户端未初始化" } - try { + return try { // 获取函数名称 val name = functionCall.getString("name") @@ -876,10 +752,10 @@ class ChatApiService(private val context: android.content.Context? = null) : Cor val result = _mcpClient?.callTool(name, arguments) // 返回工具调用结果 - return result?.get("context") as? String ?: "工具调用失败" + result?.get("context") as? String ?: "工具调用失败" } catch (e: Exception) { Log.e("ChatApiService", "处理MCP工具调用失败: ${e.message}", e) - return "处理MCP工具调用失败: ${e.message}" + "处理MCP工具调用失败: ${e.message}" } } @@ -890,6 +766,8 @@ class ChatApiService(private val context: android.content.Context? = null) : Cor return _mcpClient?.getToolMaps() ?: emptyList() } + + /** * 处理图片 */ diff --git a/local_plugins/chat_api/android/src/main/kotlin/com/yunqiinnovation/chat_api/CustomSseClientTransport.kt b/local_plugins/chat_api/android/src/main/kotlin/com/yunqiinnovation/chat_api/CustomSseClientTransport.kt index 945dde9f3..2e6fa877e 100644 --- a/local_plugins/chat_api/android/src/main/kotlin/com/yunqiinnovation/chat_api/CustomSseClientTransport.kt +++ b/local_plugins/chat_api/android/src/main/kotlin/com/yunqiinnovation/chat_api/CustomSseClientTransport.kt @@ -105,13 +105,12 @@ class CustomSseClientTransport( } "open" -> { - Log.d(TAG, "SSE连接已打开") + // SSE连接已打开 } "endpoint" -> { try { val eventData = event.data ?: "" - Log.d(TAG, "收到endpoint事件: $eventData") // 构建完整的端点URL val fullEndpoint = if (eventData.contains(hostPart)) { @@ -135,7 +134,6 @@ class CustomSseClientTransport( fullEndpoint } - Log.d(TAG, "最终消息端点: $endpointWithParams") endpoint.complete(endpointWithParams) } catch (e: Exception) { Log.e(TAG, "处理endpoint事件失败: ${e.message}", e) @@ -182,11 +180,6 @@ class CustomSseClientTransport( hostPart = urlInfo.first pathPart = urlInfo.second queryParams = urlInfo.third - - Log.d(TAG, "原始URL: $urlString") - Log.d(TAG, "主机部分: $hostPart") - Log.d(TAG, "路径部分: $pathPart") - Log.d(TAG, "查询参数: $queryParams") } // 创建SSE会话 - 直接使用原始URL @@ -202,8 +195,6 @@ class CustomSseClientTransport( "$hostPart$pathPart" } - Log.d(TAG, "SSE连接URL: $sseConnectUrl") - client.sseSession( urlString = sseConnectUrl, reconnectionTime = reconnectionTime, @@ -219,7 +210,6 @@ class CustomSseClientTransport( // 等待endpoint就绪 endpoint.await() - Log.d(TAG, "传输层启动完成,消息端点已就绪") } /** @@ -234,7 +224,6 @@ class CustomSseClientTransport( try { val messageEndpoint = endpoint.getCompleted() - Log.d(TAG, "发送消息到: $messageEndpoint") val jsonString = json.encodeToString(message) @@ -266,6 +255,5 @@ class CustomSseClientTransport( session.cancel() _onClose() job?.cancelAndJoin() - Log.d(TAG, "传输层已关闭") } } \ No newline at end of file diff --git a/local_plugins/chat_api/android/src/main/kotlin/com/yunqiinnovation/chat_api/MCPClient.kt b/local_plugins/chat_api/android/src/main/kotlin/com/yunqiinnovation/chat_api/MCPClient.kt index 9aadf3d43..7231f4999 100644 --- a/local_plugins/chat_api/android/src/main/kotlin/com/yunqiinnovation/chat_api/MCPClient.kt +++ b/local_plugins/chat_api/android/src/main/kotlin/com/yunqiinnovation/chat_api/MCPClient.kt @@ -98,14 +98,12 @@ class MCPClient(private val context: Context? = null) : AutoCloseable { if (subClient.connect()) { subClients[serverId] = subClient connectedCount++ - Log.d(TAG, "Connected to MCP server: $serverId") } else { Log.w(TAG, "Failed to connect to MCP server: $serverId") } } isConnectedFlag = connectedCount > 0 - Log.d(TAG, "Connected to $connectedCount MCP servers") connectedCount > 0 } catch (e: Exception) { @@ -121,7 +119,6 @@ class MCPClient(private val context: Context? = null) : AutoCloseable { subClients.forEach { (serverId, client) -> try { client.close() - Log.d(TAG, "已关闭子客户端 [$serverId]") } catch (e: Exception) { Log.e(TAG, "关闭子客户端 [$serverId] 失败: ${e.message}") } @@ -177,7 +174,6 @@ class MCPClient(private val context: Context? = null) : AutoCloseable { ) localFunctionDefs[name] = functionDef - Log.d(TAG, "Registered local function: $name") return true } catch (e: Exception) { Log.e(TAG, "注册本地函数失败: ${e.message}", e) @@ -192,7 +188,6 @@ class MCPClient(private val context: Context? = null) : AutoCloseable { val removed = localFunctions.remove(name) != null if (removed) { localFunctionDefs.remove(name) - Log.d(TAG, "Unregistered local function: $name") } return removed } @@ -284,10 +279,29 @@ class MCPClient(private val context: Context? = null) : AutoCloseable { */ fun parseJsonArguments(json: String): Map { return try { - val jsonObject = JSONObject(json) + // 处理空字符串或空白字符串 + val trimmedJson = json.trim() + if (trimmedJson.isEmpty()) { + return emptyMap() + } + + // 如果不是以{开头,尝试包装为{} + val jsonToUse = if (!trimmedJson.startsWith("{")) { + if (trimmedJson.contains("=") || trimmedJson.contains(":")) { + // 简单的键值对,包装成JSON对象 + "{$trimmedJson}" + } else { + // 空参数,返回空对象 + "{}" + } + } else { + trimmedJson + } + + val jsonObject = JSONObject(jsonToUse) convertJsonObjectToMap(jsonObject) } catch (e: Exception) { - Log.w(TAG, "Failed to parse JSON arguments", e) + Log.w(TAG, "Failed to parse JSON arguments: '$json'", e) emptyMap() } } @@ -304,7 +318,6 @@ class MCPClient(private val context: Context? = null) : AutoCloseable { */ suspend fun disconnectAll() { closeAllConnections() - Log.d(TAG, "Disconnected all MCP clients") } /** @@ -315,7 +328,6 @@ class MCPClient(private val context: Context? = null) : AutoCloseable { closeAllConnections() localFunctions.clear() localFunctionDefs.clear() - Log.d(TAG, "已关闭MCP客户端") } } diff --git a/local_plugins/chat_api/android/src/main/kotlin/com/yunqiinnovation/chat_api/MCPSubClient.kt b/local_plugins/chat_api/android/src/main/kotlin/com/yunqiinnovation/chat_api/MCPSubClient.kt index eb3a7d4f4..60eacea0a 100644 --- a/local_plugins/chat_api/android/src/main/kotlin/com/yunqiinnovation/chat_api/MCPSubClient.kt +++ b/local_plugins/chat_api/android/src/main/kotlin/com/yunqiinnovation/chat_api/MCPSubClient.kt @@ -41,8 +41,6 @@ class MCPSubClient( if (isConnected) return true return try { - Log.d(TAG, "[$serverId] 开始连接到MCP服务器: $serverUrl") - // 创建MCP客户端实例 val client = Client( clientInfo = Implementation( @@ -55,7 +53,6 @@ class MCPSubClient( val transport = when { serverUrl.startsWith("http://") || serverUrl.startsWith("https://") -> { // SSE传输 - 使用自定义的CustomSseClientTransport - Log.d(TAG, "[$serverId] 创建自定义SSE传输,URL: $serverUrl") val mcpHttpClient = httpClient ?: createMcpHttpClient() CustomSseClientTransport( client = mcpHttpClient, @@ -77,12 +74,6 @@ class MCPSubClient( if (toolsResult != null) { availableTools.clear() availableTools.addAll(toolsResult.tools) - Log.d(TAG, "[$serverId] 获取到 ${availableTools.size} 个工具") - - // 打印工具信息 - availableTools.forEach { tool -> - Log.d(TAG, "[$serverId] 工具: ${tool.name} - ${tool.description}") - } } } catch (e: Exception) { Log.w(TAG, "[$serverId] 获取工具列表失败: ${e.message}") @@ -91,7 +82,6 @@ class MCPSubClient( mcpClient = client isConnected = true - Log.d(TAG, "[$serverId] MCP连接成功") true } catch (e: Exception) { @@ -111,13 +101,7 @@ class MCPSubClient( * 获取工具映射列表 */ fun getToolMaps(): List> { - Log.d(TAG, "[$serverId] 开始获取工具映射,工具数量: ${availableTools.size}") - val toolMaps = availableTools.map { tool -> - Log.d(TAG, "[$serverId] 处理工具: ${tool.name}") - Log.d(TAG, "[$serverId] 工具描述: ${tool.description}") - Log.d(TAG, "[$serverId] 输入Schema: ${tool.inputSchema}") - val parametersMap = tool.inputSchema?.let { inputSchema -> convertInputSchemaToMap(inputSchema) } ?: mapOf( @@ -126,8 +110,6 @@ class MCPSubClient( "required" to emptyList() ) - Log.d(TAG, "[$serverId] 转换后的参数映射: $parametersMap") - val toolMap = mapOf( "type" to "function", "function" to mapOf( @@ -137,11 +119,9 @@ class MCPSubClient( ) ) - Log.d(TAG, "[$serverId] 最终工具映射: $toolMap") toolMap } - Log.d(TAG, "[$serverId] 完成工具映射生成,返回 ${toolMaps.size} 个工具") return toolMaps } @@ -152,8 +132,6 @@ class MCPSubClient( val client = mcpClient ?: return null return try { - Log.d(TAG, "[$serverId] 调用工具: $name, 参数: $arguments") - // 创建工具调用请求 - 将Map转换为JsonObject val argumentsJson = kotlinx.serialization.json.buildJsonObject { arguments.forEach { (key, value) -> @@ -175,8 +153,6 @@ class MCPSubClient( val result = client.callTool(request) result?.let { callResult -> - Log.d(TAG, "[$serverId] 工具调用结果: ${callResult.content.size} 个内容项") - // 将结果转换为统一格式 val contentList = callResult.content.map { contentItem -> // 根据不同的内容类型处理 @@ -208,16 +184,12 @@ class MCPSubClient( * 将Tool.Input转换为Map格式,供OpenAI使用 */ private fun convertInputSchemaToMap(inputSchema: Tool.Input): Map { - Log.d(TAG, "开始转换Input Schema: $inputSchema") - val properties = mutableMapOf() val required = mutableListOf() // 处理properties inputSchema.properties?.let { propsJsonObject -> - Log.d(TAG, "处理properties: $propsJsonObject") for ((key, value) in propsJsonObject) { - Log.d(TAG, "处理属性: $key = $value (${value::class.java.simpleName})") when (value) { is JsonPrimitive -> { if (value.isString) { @@ -235,13 +207,12 @@ class MCPSubClient( } } } - } ?: Log.d(TAG, "properties为null") + } // 处理required inputSchema.required?.let { requiredList -> - Log.d(TAG, "处理required: $requiredList") required.addAll(requiredList) - } ?: Log.d(TAG, "required为null") + } val result = mapOf( "type" to "object", @@ -249,7 +220,6 @@ class MCPSubClient( "required" to required ) - Log.d(TAG, "转换后的Schema Map: $result") return result } @@ -296,7 +266,6 @@ class MCPSubClient( if (toolsResult != null) { availableTools.clear() availableTools.addAll(toolsResult.tools) - Log.d(TAG, "[$serverId] 刷新工具列表成功,共 ${availableTools.size} 个工具") true } else { false @@ -318,7 +287,6 @@ class MCPSubClient( mcpClient = null isConnected = false availableTools.clear() - Log.d(TAG, "[$serverId] MCP连接已关闭") } catch (e: Exception) { Log.e(TAG, "[$serverId] 关闭MCP连接时出错: ${e.message}", e) } From 60f01f7b30f25b9cd563151d12ec1d2f8e659ea2 Mon Sep 17 00:00:00 2001 From: wolfplus Date: Sat, 7 Jun 2025 17:51:46 +0100 Subject: [PATCH 3/3] add --- lib/data/services/asr_service.dart | 6 ------ .../controllers/translation_controller.dart | 10 +++++----- .../yunqiinnovation/azure_speech/AzureAsrHelper.kt | 11 ----------- 3 files changed, 5 insertions(+), 22 deletions(-) diff --git a/lib/data/services/asr_service.dart b/lib/data/services/asr_service.dart index 9cc9a2f0e..75c2457c6 100644 --- a/lib/data/services/asr_service.dart +++ b/lib/data/services/asr_service.dart @@ -29,12 +29,6 @@ abstract class AsrService { /// 检查连续识别是否活跃 bool isContinuousRecognitionActive(); - /// 开启录音 - Future enableRecord(); - - /// 停止录音 - Future stopRecord(); - /// 释放资源 Future dispose(); } diff --git a/lib/modules/translation/controllers/translation_controller.dart b/lib/modules/translation/controllers/translation_controller.dart index 52503548f..47bfbd1b6 100644 --- a/lib/modules/translation/controllers/translation_controller.dart +++ b/lib/modules/translation/controllers/translation_controller.dart @@ -196,11 +196,11 @@ class TranslationController extends GetxController { // 切换录音功能 Future toggleisRecord() async { isRecording.toggle(); - if (isRecording.value) { - await _asrService.enableRecord(); - } else { - await _asrService.stopRecord(); - } + // if (isRecording.value) { + // await _asrService.enableRecord(); + // } else { + // await _asrService.stopRecord(); + // } } // 开始语音识别 diff --git a/local_plugins/azure_speech/android/src/main/kotlin/com/yunqiinnovation/azure_speech/AzureAsrHelper.kt b/local_plugins/azure_speech/android/src/main/kotlin/com/yunqiinnovation/azure_speech/AzureAsrHelper.kt index a5db177bc..48c10d5bc 100644 --- a/local_plugins/azure_speech/android/src/main/kotlin/com/yunqiinnovation/azure_speech/AzureAsrHelper.kt +++ b/local_plugins/azure_speech/android/src/main/kotlin/com/yunqiinnovation/azure_speech/AzureAsrHelper.kt @@ -128,7 +128,6 @@ class AzureAsrHelper(private val context: Context) { speechRecognitionLanguage = currentLanguage } - // 优化:大幅减少静音超时时间(从800ms到300ms) setProperty("SpeechServiceConnection_EndSilenceTimeoutMs", "300") setProperty("Speech_SegmentationSilenceTimeoutMs", "300") @@ -256,14 +255,6 @@ class AzureAsrHelper(private val context: Context) { } } - /** - * 强制设置单一语言(优化语言检测逻辑) - */ - fun setForceLanguage(language: String) { - this.currentLanguage = language - this.isAutoDetectLanguage = false - speechConfig?.speechRecognitionLanguage = language - } /** * 设置外部音频流 - 使用拉流方式 @@ -796,7 +787,6 @@ class AzureAsrHelper(private val context: Context) { */ override fun read(buffer: ByteArray): Int { try { - // 修复:恢复阻塞等待,确保外部音频数据完整性 val chunk = queue.take() // 检查是否是结束标志(空数组) @@ -806,7 +796,6 @@ class AzureAsrHelper(private val context: Context) { val toCopy = minOf(chunk.size, buffer.size) System.arraycopy(chunk, 0, buffer, 0, toCopy) - Log.i(tag, "外部音频:") return toCopy } catch (e: InterruptedException) { Thread.currentThread().interrupt()