wolfplus 1 year ago
parent
commit
c7ab1fe5f3
  1. 3
      local_plugins/agent_service/android/src/main/kotlin/com/yunqiinnovation/agent_service/AgentService.kt
  2. 70
      local_plugins/chat_api/android/README.md
  3. 12
      local_plugins/chat_api/android/build.gradle.kts
  4. 291
      local_plugins/chat_api/android/src/main/kotlin/com/yunqiinnovation/chat_api/ChatApiPlugin.kt
  5. 536
      local_plugins/chat_api/android/src/main/kotlin/com/yunqiinnovation/chat_api/ChatApiService.kt
  6. 271
      local_plugins/chat_api/android/src/main/kotlin/com/yunqiinnovation/chat_api/CustomSseClientTransport.kt
  7. 360
      local_plugins/chat_api/android/src/main/kotlin/com/yunqiinnovation/chat_api/MCPClient.kt
  8. 329
      local_plugins/chat_api/android/src/main/kotlin/com/yunqiinnovation/chat_api/MCPSubClient.kt
  9. 471
      local_plugins/chat_api/android/src/main/kotlin/com/yunqiinnovation/chat_api/SystemFunctionHandler.kt
  10. 57
      local_plugins/chat_api/android/src/main/kotlin/com/yunqiinnovation/chat_api/Utils.kt

3
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() ?: ""
)
// 加载最近的聊天记录

70
local_plugins/chat_api/android/README.md

@ -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. 支持取消正在进行的请求

12
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")

291
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<String>("apiKey")
if (apiKey.isNullOrEmpty()) {
result.error("INVALID_ARGS", "缺少必要参数", null)
return
}
val baseUrl = call.argument<String>("baseUrl") ?: ""
val model = call.argument<String>("model") ?: ""
val mcpServer = call.argument<String>("mcpServer") ?: ""
val success = chatApiService.initialize(apiKey, baseUrl, model, mcpServer)
result.success(success)
}
private fun handleCreateUserMessage(call: MethodCall, result: Result) {
val content = call.argument<String>("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<String>("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<List<Map<String, Any>>>("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<String>("apiKey") ?: ""
val baseUrl = call.argument<String>("baseUrl") ?: ""
val model = call.argument<String>("model") ?: ""
val mcpServer = call.argument<String>("mcpServer") ?: ""
val success = chatApiService?.initialize(apiKey, baseUrl, model, mcpServer) ?: false
result.success(success)
}
"chatCompletionStream" -> {
val messages = call.argument<List<Map<String, Any>>>("messages") ?: emptyList()
val tool = call.argument<Boolean>("tool") ?: false
pluginScope.launch {
chatApiService?.chatCompletionStream(messages, tool)
}
result.success(true)
}
"cancelChatStream" -> {
chatApiService?.cancelChatStream()
result.success(true)
}
"processImage" -> {
val imagePath = call.argument<String>("imagePath") ?: ""
val prompt = call.argument<String>("prompt") ?: ""
val maxWidth = call.argument<Double>("maxWidth") ?: 2048.0
val detail = call.argument<String>("detail") ?: "auto"
val base64Image = chatApiService?.processImage(imagePath, prompt, maxWidth, detail)
result.success(base64Image)
}
"initializeMcpClient" -> {
val serverUrl = call.argument<String>("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<Map<String, Any>>())
}
} catch (e: Exception) {
mainHandler.post {
result.error("SEND_ERROR", e.message, null)
}
"hasToolWithName" -> {
val name = call.argument<String>("name") ?: ""
val mcpClient = chatApiService?.mcpClient
val hasTool = mcpClient?.hasToolWithName(name) ?: false
result.success(hasTool)
}
"getToolType" -> {
val name = call.argument<String>("name") ?: ""
val mcpClient = chatApiService?.mcpClient
val toolType = mcpClient?.getToolType(name)
result.success(toolType?.name)
}
"registerFunction" -> {
val name = call.argument<String>("name") ?: ""
val description = call.argument<String>("description") ?: ""
val parametersJson = call.argument<String>("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<List<Map<String, Any>>>("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<String>("name")
val description = call.argument<String>("description")
@Suppress("UNCHECKED_CAST")
val parameters = call.argument<Map<String, Any>>("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<String>("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<String>("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<String, Any>(
"type" to type
)
content?.let { event["content"] = it }
meta?.let { event["meta"] = it }
sink.success(event)
}
eventSink?.let { sink ->
val event = mutableMapOf<String, Any>(
"type" to type
)
content?.let { event["content"] = it }
meta?.let { event["meta"] = it }
sink.success(event)
}
}
}

536
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<Tool>()
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<String, Any>
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<String, Any>
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<Map<String, Any>>) {
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<Tool>()
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<String, Any>
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<String, Any>
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<String, Any>): 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<String, Any>): 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<Map<String, Any>> {
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<Map<String, Any>>, 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, Any>): 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"
}
}
}

271
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<String>()
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<String, String> = emptyMap()
/**
* 解析URL,分离主机、路径和查询参数
*/
private fun parseUrl(url: String): Triple<String, String, Map<String, String>> {
return try {
val params = mutableMapOf<String, String>()
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<JSONRPCMessage>(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, "传输层已关闭")
}
}

360
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, Any>): String
}
/**
* MCP客户端
* 与 iOS 版本 MCPClient 功能对等
*/
class MCPClient(private val context: Context? = null) : AutoCloseable {
companion object {
private const val TAG = "MCPClient"
}
// 本地函数Map,函数名 -> 处理器
private val localFunctions = mutableMapOf<String, FunctionHandler>()
// 本地函数定义Map,函数名 -> 定义
private val localFunctionDefs = mutableMapOf<String, Map<String, Any>>()
// 子客户端列表,每个连接一个MCP服务器
private val subClients = mutableMapOf<String, MCPSubClient>()
// 是否已连接
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<String, Any> = when (parameters) {
is Map<*, *> -> {
@Suppress("UNCHECKED_CAST")
parameters as Map<String, Any>
}
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<Map<String, Any>> {
val allToolMaps = mutableListOf<Map<String, Any>>()
// 添加本地函数
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<String, Any>): Map<String, Any>? {
// 首先检查本地函数
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<String, Any> {
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<String, Any> {
val map = mutableMapOf<String, Any>()
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<Any> {
val list = mutableListOf<Any>()
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
}
}

329
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<Tool>()
/**
* 连接到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<Map<String, Any>> {
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<String, Any>(),
"required" to emptyList<String>()
)
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<String, Any>): Map<String, Any>? {
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<String, Any> {
Log.d(TAG, "开始转换Input Schema: $inputSchema")
val properties = mutableMapOf<String, Any>()
val required = mutableListOf<String>()
// 处理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<String, Any> {
val map = mutableMapOf<String, Any>()
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()
}
}

471
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<String, Any>(),
"required" to emptyList<String>()
),
handler = ExitInteractionHandler(context)
)
// 注册翻译模式函数
client.registerLocalFunction(
name = "enter_translation_mode",
description = "用户请求进入实时翻译模式时,启动实时翻译功能",
parameters = mapOf(
"type" to "object",
"properties" to emptyMap<String, Any>(),
"required" to emptyList<String>()
),
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<String>()
),
handler = GetCurrentTimeHandler()
)
// 注册获取当前位置函数
client.registerLocalFunction(
name = "get_current_location",
description = "获取当前地理位置",
parameters = mapOf(
"type" to "object",
"properties" to emptyMap<String, Any>(),
"required" to emptyList<String>()
),
handler = GetCurrentLocationHandler(context)
)
// 注册媒体播放功能
client.registerLocalFunction(
name = "media_play",
description = "播放媒体",
parameters = mapOf(
"type" to "object",
"properties" to emptyMap<String, Any>(),
"required" to emptyList<String>()
),
handler = MediaPlayHandler(context)
)
// 注册媒体暂停功能
client.registerLocalFunction(
name = "media_pause",
description = "暂停媒体播放",
parameters = mapOf(
"type" to "object",
"properties" to emptyMap<String, Any>(),
"required" to emptyList<String>()
),
handler = MediaPauseHandler(context)
)
// 注册媒体上一首功能
client.registerLocalFunction(
name = "media_previous",
description = "播放上一首",
parameters = mapOf(
"type" to "object",
"properties" to emptyMap<String, Any>(),
"required" to emptyList<String>()
),
handler = MediaPreviousHandler(context)
)
// 注册媒体下一首功能
client.registerLocalFunction(
name = "media_next",
description = "播放下一首",
parameters = mapOf(
"type" to "object",
"properties" to emptyMap<String, Any>(),
"required" to emptyList<String>()
),
handler = MediaNextHandler(context)
)
// 注册打开录音机功能
client.registerLocalFunction(
name = "open_recorder",
description = "打开系统录音机并开始录音",
parameters = mapOf(
"type" to "object",
"properties" to emptyMap<String, Any>(),
"required" to emptyList<String>()
),
handler = OpenRecorderHandler(context)
)
}
}
/**
* 退出交互处理器
*/
private class ExitInteractionHandler(private val context: Context?) : FunctionHandler {
override suspend fun handle(arguments: Map<String, Any>): 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, Any>): 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, Any>): 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, Any>): 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, Any>): 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, Any>): 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, Any>): 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, Any>): 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, Any>): 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, Any>): 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, Any>): 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, Any>): 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}\"}"
}
}
}

57
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<TrustManager>(object : X509TrustManager {
override fun checkClientTrusted(chain: Array<out X509Certificate>?, authType: String?) {}
override fun checkServerTrusted(chain: Array<out X509Certificate>?, authType: String?) {}
override fun getAcceptedIssuers(): Array<X509Certificate> = 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()
}
Loading…
Cancel
Save