From 78a762477a96c2365d1a78b09be9f9c7f005d691 Mon Sep 17 00:00:00 2001 From: liwei1dao Date: Thu, 12 Jun 2025 17:03:35 +0800 Subject: [PATCH] =?UTF-8?q?=E4=BC=98=E5=8C=96android=20mcp=20=E9=87=8D?= =?UTF-8?q?=E8=BF=9E=E5=92=8C=E5=BF=83=E8=B7=B3=E4=BB=A3=E7=A0=81?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- lib/data/models/appconfig_model.dart | 2 +- lib/data/models/appconfig_model.g.dart | 2 +- .../agent_service/AgentService.kt | 131 +++++++++++------- .../chat_api/CustomSseClientTransport.kt | 2 +- .../com/yunqiinnovation/chat_api/MCPClient.kt | 5 +- .../yunqiinnovation/chat_api/MCPSubClient.kt | 129 +++++++++++++++-- 6 files changed, 209 insertions(+), 62 deletions(-) diff --git a/lib/data/models/appconfig_model.dart b/lib/data/models/appconfig_model.dart index b95b5fae2..602abf4dd 100644 --- a/lib/data/models/appconfig_model.dart +++ b/lib/data/models/appconfig_model.dart @@ -52,7 +52,7 @@ class DBAgent { class DBMCPServer { final String servername; final String url; - final String? tools; + final String tools; DBMCPServer({ required this.servername, required this.url, diff --git a/lib/data/models/appconfig_model.g.dart b/lib/data/models/appconfig_model.g.dart index 6d47874e6..c4ac4cc90 100644 --- a/lib/data/models/appconfig_model.g.dart +++ b/lib/data/models/appconfig_model.g.dart @@ -57,7 +57,7 @@ Map _$DBAgentToJson(DBAgent instance) => { DBMCPServer _$DBMCPServerFromJson(Map json) => DBMCPServer( servername: json['servername'] as String, url: json['url'] as String, - tools: json['tools'] as String?, + tools: json['tools'] as String? ?? '', ); Map _$DBMCPServerToJson(DBMCPServer instance) => 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 90e318861..c3854b568 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 @@ -673,7 +673,7 @@ object AgentService : CoroutineScope { override fun onComplete() { // 视情况决定是否朗读回复 - if (speakResponse) { + if (speakResponse && !nobroadcast) { ttsService?.flushStream() } val response = responseBuilder.toString() @@ -731,13 +731,16 @@ object AgentService : CoroutineScope { override fun onFunctionCallResult(functionCall: JSONObject, functionCallResult: JSONObject) { audioPlayer?.stopAudio() + val name = functionCall.get("name") as String; + val (metestr, broadcast)= autoHandleFcunCallResult(name,functionCallResult); + aiMetadata = metestr; + nobroadcast = broadcast sendEvent("function_call_result", mapOf( "function_call" to functionCall.toString(), "result" to functionCallResult.toString(), + "meta" to metestr, )) - val (metestr, broadcast)= autoHandleFcunCallResult(functionCallResult); - aiMetadata = metestr; - nobroadcast = broadcast + } } ) @@ -1102,55 +1105,86 @@ object AgentService : CoroutineScope { } } } - /* + + /* * 自动播放音乐 * */ - fun autoHandleFcunCallResult(functionCallResult: JSONObject): Pair{ - val metaStr = functionCallResult.optString("meta") - var broadcast = false; - if (metaStr.isNotEmpty()) { - val meta = JSONObject(metaStr) - // Log.d(TAG, "mcp调用结果: $metaStr") - if (meta.has("card_music")) { //音乐卡片 - val cardMusic = meta.getJSONObject("card_music") - val id = cardMusic.optString("id", "") - val url = cardMusic.optString("url", "") - val name = cardMusic.optString("name", "") - val sgener = cardMusic.optString("sgener", "") - val image = cardMusic.optString("image", "") - processMusicPlay(mapOf("id" to id,"url" to url, "title" to name, "artist" to sgener,"coverUrl" to image)) - broadcast = true; - }else if (meta.has("card_musiclist")) { //音乐列表 - val cardMusiclist = meta.getJSONObject("card_musiclist") - val musics = cardMusiclist.optJSONArray("musics") - if (musics != null) { - val playlist = mutableListOf>() - for (i in 0 until musics.length()) { - val item = musics.optJSONObject(i) ?: continue - val id = item.optString("id", "") - val url = item.optString("url", "") - val name = item.optString("name", "") - val sgener = item.optString("sgener", "") - val image = item.optString("image", "") - playlist.add(mapOf("id" to id,"url" to url, "title" to name, "artist" to sgener,"coverUrl" to image)) - } - if (playlist.isNotEmpty()) { - broadcast = true; - processMusicPlayList(playlist) + fun autoHandleFcunCallResult(toolname:String,functionCallResult: JSONObject): Pair{ + println("协议工具返回数据 检查是否存在卡片! functionCallResult: ${functionCallResult}") + try { + val contentStr = functionCallResult.optString("context") + println("协议工具返回数据 检查是否存在卡片! context: ${contentStr}") + val content = JSONObject(contentStr) + val textStr = content.optString("text") + println("协议工具返回数据 检查是否存在卡片! textStr: ${textStr}") + val metadata: MutableMap = mutableMapOf() + var broadcast = false; + println("协议工具返回数据 检查是否存在卡片!") + if (textStr.isNotEmpty()) { + val meta = JSONObject(textStr) + metadata[toolname] = meta // 现在可以赋值 + val metaStr = JSONObject(metadata as Map<*, *>).toString() + // Log.d(TAG, "mcp调用结果: $metaStr") + if (meta.has("card_music")) { //音乐卡片 + val cardMusic = meta.getJSONObject("card_music") + val id = cardMusic.optString("id", "") + val url = cardMusic.optString("url", "") + val name = cardMusic.optString("name", "") + val sgener = cardMusic.optString("sgener", "") + val image = cardMusic.optString("image", "") + processMusicPlay( + mapOf( + "id" to id, + "url" to url, + "title" to name, + "artist" to sgener, + "coverUrl" to image + ) + ) + broadcast = true; + } else if (meta.has("card_musiclist")) { //音乐列表 + val cardMusiclist = meta.getJSONObject("card_musiclist") + val musics = cardMusiclist.optJSONArray("musics") + if (musics != null) { + val playlist = mutableListOf>() + for (i in 0 until musics.length()) { + val item = musics.optJSONObject(i) ?: continue + val id = item.optString("id", "") + val url = item.optString("url", "") + val name = item.optString("name", "") + val sgener = item.optString("sgener", "") + val image = item.optString("image", "") + playlist.add( + mapOf( + "id" to id, + "url" to url, + "title" to name, + "artist" to sgener, + "coverUrl" to image + ) + ) + } + if (playlist.isNotEmpty()) { + broadcast = true; + processMusicPlayList(playlist) + } else { + Log.w(TAG, "card_musiclist 中没有有效的音乐条目") + } } else { - Log.w(TAG, "card_musiclist 中没有有效的音乐条目") + Log.w(TAG, "card_musiclist.musics 不是一个数组") } - } else { - Log.w(TAG, "card_musiclist.musics 不是一个数组") + } else if (meta.has("card_navigation")) { //导航服务 + val cardNavigation = meta.getJSONObject("card_navigation") + val start = cardNavigation.optString("start", "") + val end = cardNavigation.optString("end", "") + broadcast = true; + processNavigation(start, end); } - }else if (meta.has("card_navigation")) { //导航服务 - val cardNavigation = meta.getJSONObject("card_navigation") - val start = cardNavigation.optString("start", "") - val end = cardNavigation.optString("end", "") - broadcast = true; - processNavigation(start,end); + return Pair(metaStr, broadcast); } - return Pair(metaStr, broadcast); + } catch (e: Exception) { + // 捕获其他类型的异常(可选) + println("为解析到卡片数据!") } return Pair("", false); } @@ -1167,6 +1201,7 @@ object AgentService : CoroutineScope { */ fun processMusicPlay(song: Map){ // 在其他 Service、BroadcastReceiver 或 Application 中调用 + Log.i(TAG, "播放音乐: ${song}") MusicServiceStarter.startServiceWithCommand(context, command = "play", song = song) } @@ -1184,7 +1219,7 @@ object AgentService : CoroutineScope { if (songs.isNotEmpty()) { MusicServiceStarter.startServiceWithPlaylist(context, command = "playlist", songs = songs) } else { - Log.i(TAG, "音乐列表为空,未启动播放服务") + Log.e(TAG, "音乐列表为空,未启动播放服务") } } 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 2e6fa877e..14777d02a 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 @@ -234,7 +234,7 @@ class CustomSseClientTransport( if (!response.status.isSuccess()) { val text = response.bodyAsText() - Log.e(TAG, "发送消息失败: HTTP ${response.status}, $text") + Log.e(TAG, "发送消息失败:URL:${urlString} HTTP ${response.status}, $text") error("Error POSTing to endpoint (HTTP ${response.status}): $text") } } catch (e: Exception) { 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 7231f4999..23c538115 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 @@ -91,9 +91,9 @@ class MCPClient(private val context: Context? = null) : AutoCloseable { val serverConfig = mcpServers.optJSONObject(serverId) ?: continue val url = serverConfig.optString("url", "") + val tools = serverConfig.optString("tools", "") if (url.isEmpty()) continue - - val subClient = MCPSubClient(serverId, url, sharedHttpClient) + val subClient = MCPSubClient(serverId, url,tools, sharedHttpClient) if (subClient.connect()) { subClients[serverId] = subClient @@ -369,4 +369,5 @@ class MCPClient(private val context: Context? = null) : AutoCloseable { 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 index 60eacea0a..3290c2297 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 @@ -1,6 +1,7 @@ package com.yunqiinnovation.chat_api import android.util.Log +import com.google.gson.GsonBuilder import io.ktor.client.* import io.modelcontextprotocol.kotlin.sdk.* import io.modelcontextprotocol.kotlin.sdk.client.* @@ -21,13 +22,23 @@ import kotlin.collections.mutableListOf class MCPSubClient( private val serverId: String, private val serverUrl: String, + private val filterTools: String, private val httpClient: HttpClient? = null ) : AutoCloseable { companion object { private const val TAG = "MCPSubClient" } - + // 心跳检测 + private var heartbeatJob: Job? = null + private val heartbeatInterval = 30000L // 30秒 + // 重连配置 + private val maxRetryAttempts = 3 + private val initialReconnectDelay = 1000L // 初始重连延迟1秒 + private val maxReconnectDelay = 30000L // 最大重连延迟30秒 + private var currentReconnectDelay = initialReconnectDelay + private var retryCount = 0 + private val scope = CoroutineScope(Dispatchers.IO + SupervisorJob()) private val connectionMutex = Mutex() private var mcpClient: Client? = null @@ -66,14 +77,18 @@ class MCPSubClient( } // 连接到服务器 - client.connect(transport) + client?.connect(transport) // 获取可用工具列表 try { - val toolsResult = client.listTools() + val toolsResult = client?.listTools() if (toolsResult != null) { availableTools.clear() - availableTools.addAll(toolsResult.tools) + val filtered = toolsResult.tools.filter { tool -> + filterTools.isEmpty() || filterTools.contains(tool.name) + } + Log.w(TAG, "[$serverId] [${filterTools}] 获取工具列表: ${filtered} 原始列表:${toolsResult.tools}") + availableTools.addAll(filtered) } } catch (e: Exception) { Log.w(TAG, "[$serverId] 获取工具列表失败: ${e.message}") @@ -82,6 +97,11 @@ class MCPSubClient( mcpClient = client isConnected = true + + retryCount = 0 + currentReconnectDelay = initialReconnectDelay + // 启动心跳检测 + startHeartbeat() true } catch (e: Exception) { @@ -129,8 +149,15 @@ class MCPSubClient( * 调用MCP工具 */ suspend fun callTool(name: String, arguments: Map): Map? { - val client = mcpClient ?: return null - + // 先检查连接状态 + if (!checkConnection()) { + Log.e(TAG, "无法连接到服务器,工具调用失败") + return null + } + if (!containsTool(name)) { + Log.w(TAG, "此客户端不包含工具: $name") + return null + } return try { // 创建工具调用请求 - 将Map转换为JsonObject val argumentsJson = kotlinx.serialization.json.buildJsonObject { @@ -150,7 +177,7 @@ class MCPSubClient( ) // 调用工具 - val result = client.callTool(request) + val result = mcpClient?.callTool(request) result?.let { callResult -> // 将结果转换为统一格式 @@ -158,7 +185,7 @@ class MCPSubClient( // 根据不同的内容类型处理 mapOf( "type" to "text", - "text" to (contentItem.toString()) + "text" to (contentItem.toJsonString()) ) } @@ -180,6 +207,20 @@ class MCPSubClient( } } + // 外部扩展函数 + fun PromptMessageContent.toJsonString(indent: Int = 0): String { + val gson = GsonBuilder() + .setPrettyPrinting() + .create() + + return if (indent > 0) { + gson.toJson(this) + } else { + // 移除空格和换行符 + gson.toJson(this).replace("\\s+".toRegex(), "") + } + } + /** * 将Tool.Input转换为Map格式,供OpenAI使用 */ @@ -275,7 +316,77 @@ class MCPSubClient( false } } - + + /** + * 检查连接状态并自动重连 + */ + suspend fun checkConnection(): Boolean { + if (!isConnected) { + Log.d(TAG, "当前未连接,尝试重新连接...") + return connect() + } + + try { + // 简单ping测试连接状态 + mcpClient?.ping() + return true + } catch (e: Exception) { + Log.e(TAG, "连接检查失败: ${e.message}") + isConnected = false + return false + } + } + /** + * 停止心跳检测 + */ + private fun stopHeartbeat() { + heartbeatJob?.cancel() + heartbeatJob = null + } + /** + * 启动心跳检测 + */ + private fun startHeartbeat() { + heartbeatJob?.cancel() + heartbeatJob = CoroutineScope(Dispatchers.IO).launch { + while (isActive && isConnected) { + try { + delay(heartbeatInterval) + + // 发送心跳请求 + val startTime = System.currentTimeMillis() + val response = mcpClient?.ping() + val latency = System.currentTimeMillis() - startTime + + Log.d(TAG, "心跳检测成功,延迟: ${latency}ms") + } catch (e: Exception) { + Log.e(TAG, "心跳检测失败: ${e.message}") + handleConnectionError(e) + break + } + } + } + } + /** + * 处理连接错误并尝试重连 + */ + private suspend fun handleConnectionError(e: Exception) { + isConnected = false + stopHeartbeat() + + if (retryCount < maxRetryAttempts) { + retryCount++ + currentReconnectDelay = minOf(currentReconnectDelay * 2, maxReconnectDelay) + + Log.w(TAG, "连接失败,将在 ${currentReconnectDelay}ms 后尝试重连 (尝试 $retryCount/$maxRetryAttempts)") + + delay(currentReconnectDelay) + connect() + } else { + Log.e(TAG, "已达到最大重试次数($maxRetryAttempts),停止重连") + // 可以在这里添加通知或回调,告知上层连接彻底失败 + } + } /** * 关闭连接 */