Browse Source

优化android mcp 重连和心跳代码

newdev_shunjiawei
liwei1dao 1 year ago
parent
commit
8a737d8224
  1. 2
      lib/data/models/appconfig_model.dart
  2. 2
      lib/data/models/appconfig_model.g.dart
  3. 131
      local_plugins/agent_service/android/src/main/kotlin/com/yunqiinnovation/agent_service/AgentService.kt
  4. 2
      local_plugins/chat_api/android/src/main/kotlin/com/yunqiinnovation/chat_api/CustomSseClientTransport.kt
  5. 5
      local_plugins/chat_api/android/src/main/kotlin/com/yunqiinnovation/chat_api/MCPClient.kt
  6. 129
      local_plugins/chat_api/android/src/main/kotlin/com/yunqiinnovation/chat_api/MCPSubClient.kt

2
lib/data/models/appconfig_model.dart

@ -52,7 +52,7 @@ class DBAgent {
class DBMCPServer { class DBMCPServer {
final String servername; final String servername;
final String url; final String url;
final String? tools; final String tools;
DBMCPServer({ DBMCPServer({
required this.servername, required this.servername,
required this.url, required this.url,

2
lib/data/models/appconfig_model.g.dart

@ -57,7 +57,7 @@ Map<String, dynamic> _$DBAgentToJson(DBAgent instance) => <String, dynamic>{
DBMCPServer _$DBMCPServerFromJson(Map<String, dynamic> json) => DBMCPServer( DBMCPServer _$DBMCPServerFromJson(Map<String, dynamic> json) => DBMCPServer(
servername: json['servername'] as String, servername: json['servername'] as String,
url: json['url'] as String, url: json['url'] as String,
tools: json['tools'] as String?, tools: json['tools'] as String? ?? '',
); );
Map<String, dynamic> _$DBMCPServerToJson(DBMCPServer instance) => Map<String, dynamic> _$DBMCPServerToJson(DBMCPServer instance) =>

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

@ -673,7 +673,7 @@ object AgentService : CoroutineScope {
override fun onComplete() { override fun onComplete() {
// 视情况决定是否朗读回复 // 视情况决定是否朗读回复
if (speakResponse) { if (speakResponse && !nobroadcast) {
ttsService?.flushStream() ttsService?.flushStream()
} }
val response = responseBuilder.toString() val response = responseBuilder.toString()
@ -731,13 +731,16 @@ object AgentService : CoroutineScope {
override fun onFunctionCallResult(functionCall: JSONObject, functionCallResult: JSONObject) { override fun onFunctionCallResult(functionCall: JSONObject, functionCallResult: JSONObject) {
audioPlayer?.stopAudio() audioPlayer?.stopAudio()
val name = functionCall.get("name") as String;
val (metestr, broadcast)= autoHandleFcunCallResult(name,functionCallResult);
aiMetadata = metestr;
nobroadcast = broadcast
sendEvent("function_call_result", mapOf( sendEvent("function_call_result", mapOf(
"function_call" to functionCall.toString(), "function_call" to functionCall.toString(),
"result" to functionCallResult.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<String, Boolean>{ fun autoHandleFcunCallResult(toolname:String,functionCallResult: JSONObject): Pair<String, Boolean>{
val metaStr = functionCallResult.optString("meta") println("协议工具返回数据 检查是否存在卡片! functionCallResult: ${functionCallResult}")
var broadcast = false; try {
if (metaStr.isNotEmpty()) { val contentStr = functionCallResult.optString("context")
val meta = JSONObject(metaStr) println("协议工具返回数据 检查是否存在卡片! context: ${contentStr}")
// Log.d(TAG, "mcp调用结果: $metaStr") val content = JSONObject(contentStr)
if (meta.has("card_music")) { //音乐卡片 val textStr = content.optString("text")
val cardMusic = meta.getJSONObject("card_music") println("协议工具返回数据 检查是否存在卡片! textStr: ${textStr}")
val id = cardMusic.optString("id", "") val metadata: MutableMap<String, Any> = mutableMapOf()
val url = cardMusic.optString("url", "") var broadcast = false;
val name = cardMusic.optString("name", "") println("协议工具返回数据 检查是否存在卡片!")
val sgener = cardMusic.optString("sgener", "") if (textStr.isNotEmpty()) {
val image = cardMusic.optString("image", "") val meta = JSONObject(textStr)
processMusicPlay(mapOf("id" to id,"url" to url, "title" to name, "artist" to sgener,"coverUrl" to image)) metadata[toolname] = meta // 现在可以赋值
broadcast = true; val metaStr = JSONObject(metadata as Map<*, *>).toString()
}else if (meta.has("card_musiclist")) { //音乐列表 // Log.d(TAG, "mcp调用结果: $metaStr")
val cardMusiclist = meta.getJSONObject("card_musiclist") if (meta.has("card_music")) { //音乐卡片
val musics = cardMusiclist.optJSONArray("musics") val cardMusic = meta.getJSONObject("card_music")
if (musics != null) { val id = cardMusic.optString("id", "")
val playlist = mutableListOf<Map<String, String>>() val url = cardMusic.optString("url", "")
for (i in 0 until musics.length()) { val name = cardMusic.optString("name", "")
val item = musics.optJSONObject(i) ?: continue val sgener = cardMusic.optString("sgener", "")
val id = item.optString("id", "") val image = cardMusic.optString("image", "")
val url = item.optString("url", "") processMusicPlay(
val name = item.optString("name", "") mapOf(
val sgener = item.optString("sgener", "") "id" to id,
val image = item.optString("image", "") "url" to url,
playlist.add(mapOf("id" to id,"url" to url, "title" to name, "artist" to sgener,"coverUrl" to image)) "title" to name,
} "artist" to sgener,
if (playlist.isNotEmpty()) { "coverUrl" to image
broadcast = true; )
processMusicPlayList(playlist) )
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<Map<String, String>>()
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 { } else {
Log.w(TAG, "card_musiclist 中没有有效的音乐条目") Log.w(TAG, "card_musiclist.musics 不是一个数组")
} }
} else { } else if (meta.has("card_navigation")) { //导航服务
Log.w(TAG, "card_musiclist.musics 不是一个数组") 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")) { //导航服务 return Pair(metaStr, broadcast);
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); } catch (e: Exception) {
// 捕获其他类型的异常(可选)
println("为解析到卡片数据!")
} }
return Pair("", false); return Pair("", false);
} }
@ -1167,6 +1201,7 @@ object AgentService : CoroutineScope {
*/ */
fun processMusicPlay(song: Map<String,String>){ fun processMusicPlay(song: Map<String,String>){
// 在其他 Service、BroadcastReceiver 或 Application 中调用 // 在其他 Service、BroadcastReceiver 或 Application 中调用
Log.i(TAG, "播放音乐: ${song}")
MusicServiceStarter.startServiceWithCommand(context, command = "play", song = song) MusicServiceStarter.startServiceWithCommand(context, command = "play", song = song)
} }
@ -1184,7 +1219,7 @@ object AgentService : CoroutineScope {
if (songs.isNotEmpty()) { if (songs.isNotEmpty()) {
MusicServiceStarter.startServiceWithPlaylist(context, command = "playlist", songs = songs) MusicServiceStarter.startServiceWithPlaylist(context, command = "playlist", songs = songs)
} else { } else {
Log.i(TAG, "音乐列表为空,未启动播放服务") Log.e(TAG, "音乐列表为空,未启动播放服务")
} }
} }

2
local_plugins/chat_api/android/src/main/kotlin/com/yunqiinnovation/chat_api/CustomSseClientTransport.kt

@ -234,7 +234,7 @@ class CustomSseClientTransport(
if (!response.status.isSuccess()) { if (!response.status.isSuccess()) {
val text = response.bodyAsText() 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") error("Error POSTing to endpoint (HTTP ${response.status}): $text")
} }
} catch (e: Exception) { } catch (e: Exception) {

5
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 serverConfig = mcpServers.optJSONObject(serverId) ?: continue
val url = serverConfig.optString("url", "") val url = serverConfig.optString("url", "")
val tools = serverConfig.optString("tools", "")
if (url.isEmpty()) continue if (url.isEmpty()) continue
val subClient = MCPSubClient(serverId, url,tools, sharedHttpClient)
val subClient = MCPSubClient(serverId, url, sharedHttpClient)
if (subClient.connect()) { if (subClient.connect()) {
subClients[serverId] = subClient subClients[serverId] = subClient
@ -369,4 +369,5 @@ class MCPClient(private val context: Context? = null) : AutoCloseable {
return list return list
} }
} }

129
local_plugins/chat_api/android/src/main/kotlin/com/yunqiinnovation/chat_api/MCPSubClient.kt

@ -1,6 +1,7 @@
package com.yunqiinnovation.chat_api package com.yunqiinnovation.chat_api
import android.util.Log import android.util.Log
import com.google.gson.GsonBuilder
import io.ktor.client.* import io.ktor.client.*
import io.modelcontextprotocol.kotlin.sdk.* import io.modelcontextprotocol.kotlin.sdk.*
import io.modelcontextprotocol.kotlin.sdk.client.* import io.modelcontextprotocol.kotlin.sdk.client.*
@ -21,13 +22,23 @@ import kotlin.collections.mutableListOf
class MCPSubClient( class MCPSubClient(
private val serverId: String, private val serverId: String,
private val serverUrl: String, private val serverUrl: String,
private val filterTools: String,
private val httpClient: HttpClient? = null private val httpClient: HttpClient? = null
) : AutoCloseable { ) : AutoCloseable {
companion object { companion object {
private const val TAG = "MCPSubClient" 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 scope = CoroutineScope(Dispatchers.IO + SupervisorJob())
private val connectionMutex = Mutex() private val connectionMutex = Mutex()
private var mcpClient: Client? = null private var mcpClient: Client? = null
@ -66,14 +77,18 @@ class MCPSubClient(
} }
// 连接到服务器 // 连接到服务器
client.connect(transport) client?.connect(transport)
// 获取可用工具列表 // 获取可用工具列表
try { try {
val toolsResult = client.listTools() val toolsResult = client?.listTools()
if (toolsResult != null) { if (toolsResult != null) {
availableTools.clear() 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) { } catch (e: Exception) {
Log.w(TAG, "[$serverId] 获取工具列表失败: ${e.message}") Log.w(TAG, "[$serverId] 获取工具列表失败: ${e.message}")
@ -82,6 +97,11 @@ class MCPSubClient(
mcpClient = client mcpClient = client
isConnected = true isConnected = true
retryCount = 0
currentReconnectDelay = initialReconnectDelay
// 启动心跳检测
startHeartbeat()
true true
} catch (e: Exception) { } catch (e: Exception) {
@ -129,8 +149,15 @@ class MCPSubClient(
* 调用MCP工具 * 调用MCP工具
*/ */
suspend fun callTool(name: String, arguments: Map<String, Any>): Map<String, Any>? { suspend fun callTool(name: String, arguments: Map<String, Any>): Map<String, Any>? {
val client = mcpClient ?: return null // 先检查连接状态
if (!checkConnection()) {
Log.e(TAG, "无法连接到服务器,工具调用失败")
return null
}
if (!containsTool(name)) {
Log.w(TAG, "此客户端不包含工具: $name")
return null
}
return try { return try {
// 创建工具调用请求 - 将Map转换为JsonObject // 创建工具调用请求 - 将Map转换为JsonObject
val argumentsJson = kotlinx.serialization.json.buildJsonObject { 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 -> result?.let { callResult ->
// 将结果转换为统一格式 // 将结果转换为统一格式
@ -158,7 +185,7 @@ class MCPSubClient(
// 根据不同的内容类型处理 // 根据不同的内容类型处理
mapOf( mapOf(
"type" to "text", "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使用 * 将Tool.Input转换为Map格式,供OpenAI使用
*/ */
@ -275,7 +316,77 @@ class MCPSubClient(
false 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),停止重连")
// 可以在这里添加通知或回调,告知上层连接彻底失败
}
}
/** /**
* 关闭连接 * 关闭连接
*/ */

Loading…
Cancel
Save