Browse Source

优化mcp代码

newdev_shunjiawei
liwei1dao 1 year ago
parent
commit
56fd5f10df
  1. 57
      local_plugins/chat_api/android/src/main/kotlin/com/yunqiinnovation/chat_api/CustomSseClientTransport.kt
  2. 39
      local_plugins/chat_api/android/src/main/kotlin/com/yunqiinnovation/chat_api/MCPClient.kt
  3. 175
      local_plugins/chat_api/android/src/main/kotlin/com/yunqiinnovation/chat_api/MCPSubClient.kt
  4. 10
      local_plugins/location_service/android/src/main/kotlin/com/yunqiinnovation/location_service/LocationService.kt

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

@ -22,7 +22,8 @@ import java.util.concurrent.atomic.AtomicBoolean
*/
class CustomSseClientTransport(
private val client: HttpClient,
private val urlString: String?,
private val serviceidString: String?,
public val urlString: String?,
private val reconnectionTime: Duration? = null,
private val requestBuilder: HttpRequestBuilder.() -> Unit = {},
private val onConnectionLost: (() -> Unit)? = null
@ -88,7 +89,7 @@ class CustomSseClientTransport(
Triple(hostUrl, path, params)
} catch (e: Exception) {
Log.e(TAG, "解析URL失败: $url, ${e.message}")
Log.e(TAG, "$serviceidString 解析URL失败: $url, ${e.message}")
Triple(url, "", emptyMap())
}
}
@ -102,9 +103,9 @@ class CustomSseClientTransport(
session.incoming.collect { event ->
when (event.event) {
"error" -> {
Log.e(TAG, "SSE错误: ${event.data}")
Log.e(TAG, "$serviceidString SSE错误: ${event.data}")
isConnected.set(false)
val exception = Exception("SSE Error: ${event.data}")
val exception = Exception("$serviceidString SSE Error: ${event.data}")
_onError(exception)
onConnectionLost?.invoke()
throw exception
@ -112,7 +113,7 @@ class CustomSseClientTransport(
"open" -> {
// SSE连接已打开
Log.d(TAG, "SSE连接已打开")
Log.d(TAG, "$serviceidString SSE连接已打开")
isConnected.set(true)
}
"ping" -> {
@ -146,7 +147,7 @@ class CustomSseClientTransport(
endpoint.complete(endpointWithParams)
} catch (e: Exception) {
Log.e(TAG, "处理endpoint事件失败: ${e.message}", e)
Log.e(TAG, "$serviceidString 处理endpoint事件失败: ${e.message}", e)
_onError(e)
close()
error(e)
@ -161,22 +162,22 @@ class CustomSseClientTransport(
val message = json.decodeFromString<JSONRPCMessage>(data)
_onMessage(message)
} catch (e: Exception) {
Log.e(TAG, "解析JSON-RPC消息失败: ${e.message}", e)
Log.e(TAG, "$serviceidString 解析JSON-RPC消息失败: ${e.message}", e)
_onError(e)
}
}
} catch (e: Exception) {
Log.e(TAG, "处理事件失败: ${e.message}", e)
Log.e(TAG, "$serviceidString 处理事件失败: ${e.message}", e)
_onError(e)
}
}
}
}
} catch (e: CancellationException) {
Log.d(TAG, "SSE事件收集被取消")
Log.d(TAG, "$serviceidString SSE事件收集被取消")
throw e
} catch (e: Exception) {
Log.e(TAG, "SSE连接异常断开: ${e.message}", e)
Log.e(TAG, "$serviceidString SSE连接异常断开: ${e.message}", e)
isConnected.set(false)
_onError(e)
onConnectionLost?.invoke()
@ -199,13 +200,13 @@ class CustomSseClientTransport(
// 检查session是否仍然活跃
if (session.coroutineContext[Job]?.isCancelled == true) {
Log.w(TAG, "检测到SSE会话已取消")
Log.w(TAG, "$serviceidString 检测到SSE会话已取消")
isConnected.set(false)
onConnectionLost?.invoke()
break
}
} catch (e: Exception) {
Log.e(TAG, "连接监控异常: ${e.message}", e)
Log.e(TAG, "$serviceidString 连接监控异常: ${e.message}", e)
isConnected.set(false)
onConnectionLost?.invoke()
break
@ -219,7 +220,7 @@ class CustomSseClientTransport(
*/
override suspend fun start() {
if (!initialized.compareAndSet(false, true)) {
Log.e(TAG, "传输层已经启动,不能重复启动")
Log.e(TAG, "$serviceidString 传输层已经启动,不能重复启动")
error("CustomSseClientTransport already started!")
}
@ -283,12 +284,30 @@ class CustomSseClientTransport(
setBody(jsonString)
}
if (!response.status.isSuccess()) {
val text = response.bodyAsText()
// Log.e(TAG, "发送消息失败:URL:${urlString} HTTP ${response.status}, $text")
// 优化处理百度API的202状态码
when {
response.status.isSuccess() -> {
// 2xx状态码都视为成功
Log.d(TAG, "$serviceidString 消息发送成功: HTTP ${response.status}")
}
response.status == HttpStatusCode.Accepted -> {
// 202 Accepted - 百度API异步处理中,这是正常状态
Log.d(TAG, "$serviceidString 消息已被接受,正在异步处理: HTTP ${response.status}")
}
else -> {
// 其他非成功状态码才记录为错误
val text = response.bodyAsText()
Log.w(TAG, "$serviceidString 发送消息收到非成功状态码: URL:${messageEndpoint} HTTP ${response.status}, $text")
// 根据具体状态码决定是否抛出异常
if (response.status.value >= 400) {
// 4xx和5xx错误才抛出异常
throw Exception("HTTP ${response.status}: $text")
}
}
}
} catch (e: Exception) {
Log.e(TAG, "发送消息异常: ${e.message}", e)
Log.e(TAG, "$serviceidString 发送消息异常: ${e.message}", e)
_onError(e)
throw e
}
@ -306,7 +325,7 @@ class CustomSseClientTransport(
*/
override suspend fun close() {
if (!initialized.get()) {
Log.e(TAG, "关闭失败: 传输层未初始化")
Log.e(TAG, "$serviceidString 关闭失败: 传输层未初始化")
error("CustomSseClientTransport is not initialized!")
}
@ -317,6 +336,6 @@ class CustomSseClientTransport(
job?.cancelAndJoin()
connectionMonitorJob?.cancelAndJoin()
Log.d(TAG, "CustomSseClientTransport已关闭")
Log.d(TAG, "$serviceidString CustomSseClientTransport已关闭")
}
}

39
local_plugins/chat_api/android/src/main/kotlin/com/yunqiinnovation/chat_api/MCPClient.kt

@ -67,45 +67,54 @@ class MCPClient(private val context: Context? = null) : AutoCloseable {
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", "")
val tools = serverConfig.optString("tools", "")
if (url.isEmpty()) continue
val subClient = MCPSubClient(serverId, url,tools, sharedHttpClient)
if (subClient.connect()) {
subClients[serverId] = subClient
connectedCount++
} else {
Log.w(TAG, "Failed to connect to MCP server: $serverId")
Log.d(TAG, "开始连接MCP服务器: $serverId")
val subClient = MCPSubClient(serverId, url, tools, sharedHttpClient)
// 使用协程并发连接,但每个服务器都会进行重试
try {
if (subClient.connect()) {
subClients[serverId] = subClient
connectedCount++
Log.d(TAG, "MCP服务器连接成功: $serverId")
} else {
Log.w(TAG, "MCP服务器连接失败: $serverId")
}
} catch (e: Exception) {
Log.e(TAG, "MCP服务器连接异常: $serverId, ${e.message}", e)
}
}
isConnectedFlag = connectedCount > 0
Log.d(TAG, "MCP连接完成,成功连接 $connectedCount 个服务器")
connectedCount > 0
} catch (e: Exception) {
Log.e(TAG, "Failed to connect to SSE", e)
false

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

@ -33,7 +33,7 @@ class MCPSubClient(
private var heartbeatJob: Job? = null
private val heartbeatInterval = 30000L // 30秒
// 重连配置
private val maxRetryAttempts = 3
private val maxRetryAttempts = 10
private val initialReconnectDelay = 1000L // 初始重连延迟1秒
private val maxReconnectDelay = 30000L // 最大重连延迟30秒
private var currentReconnectDelay = initialReconnectDelay
@ -45,95 +45,122 @@ class MCPSubClient(
private var isConnected = false
private var availableTools = mutableListOf<Tool>()
private var transport: CustomSseClientTransport? = null
/**
* 连接到MCP服务器
*/
suspend fun connect(): Boolean = connectionMutex.withLock {
// Log.e(TAG, "liwei---------Mcp [$serverId] 连接 1")
if (isConnected) return true
Log.e(TAG, "[$serverId] 开始连接mcp服务器: $serverUrl")
return try {
// 创建MCP客户端实例
val client = Client(
clientInfo = Implementation(
name = "deep-voice-chat-api",
version = "1.0.0"
// 重试连接逻辑
for (attempt in 0 until maxRetryAttempts) {
try {
Log.d(TAG, "[$serverId] 连接尝试 ${attempt + 1}/$maxRetryAttempts")
// 创建MCP客户端实例
val client = Client(
clientInfo = Implementation(
name = "deep-voice-chat-api",
version = "1.0.0"
)
)
)
// Log.e(TAG, "liwei---------Mcp [$serverId] 连接 2")
// 根据URL类型选择传输方式
val newTransport = when {
serverUrl.startsWith("http://") || serverUrl.startsWith("https://") -> {
// SSE传输 - 使用自定义的CustomSseClientTransport
val mcpHttpClient = httpClient ?: createMcpHttpClient()
CustomSseClientTransport(
client = mcpHttpClient,
urlString = serverUrl,
onConnectionLost = {
// 连接断开回调
Log.w(TAG, "[$serverId] 检测到连接断开")
scope.launch {
handleConnectionLost()
// 根据URL类型选择传输方式
val newTransport = when {
serverUrl.startsWith("http://") || serverUrl.startsWith("https://") -> {
// SSE传输 - 使用自定义的CustomSseClientTransport
val mcpHttpClient = httpClient ?: createMcpHttpClient()
CustomSseClientTransport(
client = mcpHttpClient,
serviceidString = serverId,
urlString = serverUrl,
onConnectionLost = {
// 连接断开回调
Log.w(TAG, "[$serverId] 检测到连接断开")
scope.launch {
handleConnectionLost()
}
}
}
)
}
else -> {
Log.e(TAG, "[$serverId] 不支持的服务器URL格式: $serverUrl")
return false
)
}
else -> {
Log.e(TAG, "[$serverId] 不支持的服务器URL格式: $serverUrl")
return false
}
}
}
// Log.e(TAG, "liwei---------Mcp [$serverId] 连接 3")
transport = newTransport
try {
// 连接到服务器
withTimeout(10000) { // 10秒超时
client.connect(newTransport)
transport = newTransport
// 连接到服务器 - 增加超时时间
try {
Log.d(TAG, "[$serverId] 尝试建立连接 ${transport?.urlString}")
withTimeout(30000) { // 30秒超时
client.connect(newTransport)
}
Log.d(TAG, "[$serverId] 连接建立成功")
} catch (e: TimeoutCancellationException) {
Log.w(TAG, "[$serverId] 连接超时 (尝试 ${attempt + 1}/$maxRetryAttempts)")
if (attempt < maxRetryAttempts - 1) {
delay(currentReconnectDelay)
currentReconnectDelay = (currentReconnectDelay * 2).coerceAtMost(maxReconnectDelay)
continue // 继续下一次重试
} else {
Log.e(TAG, "[$serverId] 所有连接尝试都超时,连接失败")
return false
}
} catch (e: Exception) {
Log.w(TAG, "[$serverId] 连接异常 (尝试 ${attempt + 1}/$maxRetryAttempts): ${e.message}")
if (attempt < maxRetryAttempts - 1) {
delay(currentReconnectDelay)
currentReconnectDelay = (currentReconnectDelay * 2).coerceAtMost(maxReconnectDelay)
continue // 继续下一次重试
} else {
Log.e(TAG, "[$serverId] 所有连接尝试都失败: ${e.message}", e)
return false
}
}
// Log.e(TAG, "liwei---------Mcp [$serverId] 连接 4")
} catch (e: TimeoutCancellationException) {
Log.e(TAG, "liwei---------Mcp [$serverId] 连接超时")
return false
} catch (e: Exception) {
Log.e(TAG, "liwei---------Mcp [$serverId] 连接异常: ${e.message}", e)
return false
}
// 获取可用工具列表
try {
val toolsResult = client.listTools()
// Log.e(TAG, "liwei---------Mcp [$serverId] 连接 5")
if (toolsResult != null) {
availableTools.clear()
val filtered = toolsResult.tools.filter { tool ->
filterTools.isEmpty() || filterTools.contains(tool.name)
// 获取可用工具列表
try {
val toolsResult = client.listTools()
if (toolsResult != null) {
availableTools.clear()
val allToolNames = toolsResult.tools.map { it.name }
Log.d(TAG, "$serverId:所有工具名称列表: $allToolNames")
val filtered = toolsResult.tools.filter { tool ->
filterTools.isEmpty() || filterTools.contains(tool.name)
}
val filteredToolNames = filtered.map { it.name }
Log.d(TAG, "$serverId: 过滤后的工具: $filteredToolNames")
availableTools.addAll(filtered)
}
Log.w(TAG, "[$serverId] [${filterTools}] 获取工具列表: ${filtered} 原始列表:${toolsResult.tools}")
availableTools.addAll(filtered)
// Log.e(TAG, "liwei---------Mcp [$serverId] 连接 6")
} catch (e: Exception) {
Log.w(TAG, "[$serverId] 获取工具列表失败: ${e.message}")
// 即使获取工具失败,连接也可能是成功的
}
mcpClient = client
isConnected = true
retryCount = 0
currentReconnectDelay = initialReconnectDelay
Log.e(TAG, "[$serverId] 连接mcp服务器成功: $serverUrl")
return true
} catch (e: Exception) {
Log.w(TAG, "[$serverId] 获取工具列表失败: ${e.message}")
// 即使获取工具失败,连接也可能是成功的
Log.w(TAG, "[$serverId] 连接尝试 ${attempt + 1} 失败: ${e.message}")
if (attempt < maxRetryAttempts - 1) {
delay(currentReconnectDelay)
currentReconnectDelay = (currentReconnectDelay * 2).coerceAtMost(maxReconnectDelay)
}
}
// Log.e(TAG, "liwei---------Mcp [$serverId] 连接 7")
mcpClient = client
isConnected = true
retryCount = 0
currentReconnectDelay = initialReconnectDelay
// Log.e(TAG, "liwei---------Mcp [$serverId] 连接 8")
// 启动心跳检测
// startHeartbeat()
Log.e(TAG, "[$serverId] 连接mcp服务器成功: $serverUrl")
true
} catch (e: Exception) {
Log.e(TAG, "[$serverId] MCP连接失败: ${e.message}", e)
false
}
Log.e(TAG, "[$serverId] 所有重试尝试都失败,MCP连接失败")
return false
}
/**
* 检查是否包含指定工具
*/

10
local_plugins/location_service/android/src/main/kotlin/com/yunqiinnovation/location_service/LocationService.kt

@ -185,10 +185,10 @@ object LocationService {
Log.d(TAG, "定位结果: 来源=${location.provider}, 精度=${location.accuracy}米, 坐标=(${location.latitude}, ${location.longitude})")
if (location.errorCode == AMapLocation.LOCATION_SUCCESS) {
val accuracy = location.accuracy
Log.d(TAG, "GPS定位成功: 精度=${accuracy}米")
// Log.d(TAG, "GPS定位成功: 精度=${accuracy}米")
// 只接受高精度GPS定位
if (accuracy <= HIGH_ACCURACY_THRESHOLD) {
// if (accuracy <= HIGH_ACCURACY_THRESHOLD) {
Log.d(TAG, "接受高精度GPS定位: ${accuracy}米")
// 停止超时计时器
@ -222,9 +222,9 @@ object LocationService {
}
// 清空服务回调列表(单次定位)
serviceCallbacks.clear()
} else {
Log.w(TAG, "GPS精度不足(${accuracy}米),继续等待更高精度定位")
}
// } else {
// Log.w(TAG, "GPS精度不足(${accuracy}米),继续等待更高精度定位")
// }
} else {
Log.e(TAG, "GPS定位失败: ${location.errorCode} - ${location.errorInfo}")
}

Loading…
Cancel
Save