Browse Source

Merge branch 'dev' of https://github.com/deepcloud2048/deep_voice into dev

newdev_shunjiawei
lxm 1 year ago
parent
commit
baa8b699cb
  1. 26
      android/app/src/main/kotlin/com/yunqiinnovation/deepsound/MainActivity.kt
  2. 6
      lib/data/services/asr_service.dart
  3. 10
      lib/modules/translation/controllers/translation_controller.dart
  4. 27
      local_plugins/agent_service/android/src/main/kotlin/com/yunqiinnovation/agent_service/AgentService.kt
  5. 113
      local_plugins/azure_speech/android/src/main/kotlin/com/yunqiinnovation/azure_speech/AzureAsrHelper.kt
  6. 79
      local_plugins/azure_speech/android/src/main/kotlin/com/yunqiinnovation/azure_speech/AzureTtsHelper.kt
  7. 70
      local_plugins/chat_api/android/README.md
  8. 12
      local_plugins/chat_api/android/build.gradle.kts
  9. 291
      local_plugins/chat_api/android/src/main/kotlin/com/yunqiinnovation/chat_api/ChatApiPlugin.kt
  10. 434
      local_plugins/chat_api/android/src/main/kotlin/com/yunqiinnovation/chat_api/ChatApiService.kt
  11. 259
      local_plugins/chat_api/android/src/main/kotlin/com/yunqiinnovation/chat_api/CustomSseClientTransport.kt
  12. 372
      local_plugins/chat_api/android/src/main/kotlin/com/yunqiinnovation/chat_api/MCPClient.kt
  13. 297
      local_plugins/chat_api/android/src/main/kotlin/com/yunqiinnovation/chat_api/MCPSubClient.kt
  14. 471
      local_plugins/chat_api/android/src/main/kotlin/com/yunqiinnovation/chat_api/SystemFunctionHandler.kt
  15. 57
      local_plugins/chat_api/android/src/main/kotlin/com/yunqiinnovation/chat_api/Utils.kt

26
android/app/src/main/kotlin/com/yunqiinnovation/deepsound/MainActivity.kt

@ -19,7 +19,9 @@ class MainActivity: FlutterActivity() {
override fun onCreate(savedInstanceState: Bundle?) { override fun onCreate(savedInstanceState: Bundle?) {
super.onCreate(savedInstanceState) super.onCreate(savedInstanceState)
// 设置日志级别,减少系统级日志
suppressSystemLogs()
} }
override fun configureFlutterEngine(flutterEngine: FlutterEngine) { override fun configureFlutterEngine(flutterEngine: FlutterEngine) {
@ -32,5 +34,27 @@ class MainActivity: FlutterActivity() {
override fun onDestroy() { override fun onDestroy() {
super.onDestroy() super.onDestroy()
} }
/**
* 抑制系统级组件的日志输出
*/
private fun suppressSystemLogs() {
val systemTags = arrayOf(
"MediaCodec",
"CCodec",
"CCodecBufferChannel",
"OpusManager",
"MediaCodecList",
"MediaPlayer",
"AudioTrack",
"AudioManager",
"AudioFlinger",
"AudioService"
)
systemTags.forEach { tag ->
System.setProperty("log.tag.$tag", "WARN")
}
}
} }

6
lib/data/services/asr_service.dart

@ -29,12 +29,6 @@ abstract class AsrService {
/// 检查连续识别是否活跃 /// 检查连续识别是否活跃
bool isContinuousRecognitionActive(); bool isContinuousRecognitionActive();
/// 开启录音
Future<bool> enableRecord();
/// 停止录音
Future<bool> stopRecord();
/// 释放资源 /// 释放资源
Future<void> dispose(); Future<void> dispose();
} }

10
lib/modules/translation/controllers/translation_controller.dart

@ -196,11 +196,11 @@ class TranslationController extends GetxController {
// 切换录音功能 // 切换录音功能
Future<void> toggleisRecord() async { Future<void> toggleisRecord() async {
isRecording.toggle(); isRecording.toggle();
if (isRecording.value) { // if (isRecording.value) {
await _asrService.enableRecord(); // await _asrService.enableRecord();
} else { // } else {
await _asrService.stopRecord(); // await _asrService.stopRecord();
} // }
} }
// 开始语音识别 // 开始语音识别

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

@ -159,7 +159,6 @@ object AgentService : CoroutineScope {
loadChatHistory() loadChatHistory()
isInitialized = true isInitialized = true
Log.d(TAG, "代理服务初始化成功")
return true return true
} catch (e: Exception) { } catch (e: Exception) {
Log.e(TAG, "初始化失败: ${e.message}") Log.e(TAG, "初始化失败: ${e.message}")
@ -194,8 +193,6 @@ object AgentService : CoroutineScope {
job.cancel() job.cancel()
clearListeners() clearListeners()
isInitialized = false isInitialized = false
Log.d(TAG, "代理服务资源已释放")
} catch (e: Exception) { } catch (e: Exception) {
Log.e(TAG, "释放资源异常: ${e.message}") Log.e(TAG, "释放资源异常: ${e.message}")
} }
@ -229,9 +226,7 @@ object AgentService : CoroutineScope {
language = ttsLanguage language = ttsLanguage
) )
if (success) { if (!success) {
Log.d(TAG, "Azure TTS服务初始化成功")
} else {
Log.e(TAG, "Azure TTS服务初始化失败") Log.e(TAG, "Azure TTS服务初始化失败")
} }
@ -267,13 +262,10 @@ object AgentService : CoroutineScope {
} }
else -> { else -> {
// 处理其他类型的事件 // 处理其他类型的事件
Log.d(TAG, "TTS事件: ${event.type}")
} }
} }
} }
}) })
Log.d(TAG, "TTS服务配置完成")
} catch (e: Exception) { } catch (e: Exception) {
Log.e(TAG, "初始化TTS引擎失败: ${e.message}") Log.e(TAG, "初始化TTS引擎失败: ${e.message}")
} }
@ -371,7 +363,6 @@ object AgentService : CoroutineScope {
} else { } else {
AzureAsrHelper.AudioSourceType.MICROPHONE AzureAsrHelper.AudioSourceType.MICROPHONE
} }
Log.d(TAG, "开始语音识别, 音频源类型: $audioSourceType")
azureAsrHelper?.startContinuousRecognition(object : AzureAsrHelper.ContinuousRecognizeCallback { azureAsrHelper?.startContinuousRecognition(object : AzureAsrHelper.ContinuousRecognizeCallback {
override fun onRecognizing(recognizing: String, detectedLanguage: String) { override fun onRecognizing(recognizing: String, detectedLanguage: String) {
if (recognizing.isNotEmpty()) { if (recognizing.isNotEmpty()) {
@ -476,7 +467,6 @@ object AgentService : CoroutineScope {
BleService.closeCodec() BleService.closeCodec()
isRecognitionActive = false isRecognitionActive = false
stopIdleCheck() stopIdleCheck()
Log.d(TAG, "语音识别已停止")
} catch (e: Exception) { } catch (e: Exception) {
Log.e(TAG, "停止语音识别异常: ${e.message}") Log.e(TAG, "停止语音识别异常: ${e.message}")
isRecognitionActive = false isRecognitionActive = false
@ -517,8 +507,6 @@ object AgentService : CoroutineScope {
// 通知ChatAPI服务终止当前流式请求 // 通知ChatAPI服务终止当前流式请求
chatApiService.cancelCurrentStream() chatApiService.cancelCurrentStream()
// 记录日志
Log.d(TAG, "AI流输出已停止")
} catch (e: Exception) { } catch (e: Exception) {
Log.e(TAG, "停止AI流输出异常", e) Log.e(TAG, "停止AI流输出异常", e)
// 确保状态被重置,即使发生异常 // 确保状态被重置,即使发生异常
@ -585,8 +573,6 @@ object AgentService : CoroutineScope {
text: String = "", text: String = "",
speakResponse: Boolean = false speakResponse: Boolean = false
) { ) {
Log.d(TAG, "处理图片输入: ${if (text.isEmpty()) "无附加文本" else "附带文本: $text"}")
// 创建带图片的用户消息并处理 // 创建带图片的用户消息并处理
val userMessage = createUserMessageWithImage(text, imageBase64) val userMessage = createUserMessageWithImage(text, imageBase64)
// 图片描述用于存储 // 图片描述用于存储
@ -686,7 +672,6 @@ object AgentService : CoroutineScope {
} }
override fun onComplete() { override fun onComplete() {
Log.d(TAG, "AI完整回复: $responseBuilder")
// 视情况决定是否朗读回复 // 视情况决定是否朗读回复
if (speakResponse) { if (speakResponse) {
ttsService?.flushStream() ttsService?.flushStream()
@ -750,7 +735,6 @@ object AgentService : CoroutineScope {
"function_call" to functionCall.toString(), "function_call" to functionCall.toString(),
"result" to functionCallResult.toString(), "result" to functionCallResult.toString(),
)) ))
Log.d(TAG, "mcp调用结果: $functionCallResult")
val (metestr, broadcast)= autoHandleFcunCallResult(functionCallResult); val (metestr, broadcast)= autoHandleFcunCallResult(functionCallResult);
aiMetadata = metestr; aiMetadata = metestr;
nobroadcast = broadcast nobroadcast = broadcast
@ -802,8 +786,6 @@ object AgentService : CoroutineScope {
addToHistoryMessages(createAssistantMessage(content)) addToHistoryMessages(createAssistantMessage(content))
} }
} }
Log.d(TAG, "已加载${recentMessages.size}条历史记录")
} catch (e: Exception) { } catch (e: Exception) {
Log.e(TAG, "加载聊天历史失败: ${e.message}") Log.e(TAG, "加载聊天历史失败: ${e.message}")
} }
@ -905,7 +887,6 @@ object AgentService : CoroutineScope {
historyMessages.remove(0) historyMessages.remove(0)
} }
} }
Log.d(TAG, "聊天历史已清除")
} else { } else {
Log.e(TAG, "清除聊天历史失败") Log.e(TAG, "清除聊天历史失败")
} }
@ -1098,7 +1079,6 @@ object AgentService : CoroutineScope {
} }
} }
} }
Log.d(TAG, "开始播放音频资源")
} catch (e: Exception) { } catch (e: Exception) {
Log.e(TAG, "播放音频资源异常: ${e.message}", e) Log.e(TAG, "播放音频资源异常: ${e.message}", e)
release() release()
@ -1133,7 +1113,6 @@ object AgentService : CoroutineScope {
// Log.d(TAG, "mcp调用结果: $metaStr") // Log.d(TAG, "mcp调用结果: $metaStr")
if (meta.has("card_music")) { //音乐卡片 if (meta.has("card_music")) { //音乐卡片
val cardMusic = meta.getJSONObject("card_music") val cardMusic = meta.getJSONObject("card_music")
Log.d(TAG, "检查到音乐卡片: $cardMusic")
val id = cardMusic.optString("id", "") val id = cardMusic.optString("id", "")
val url = cardMusic.optString("url", "") val url = cardMusic.optString("url", "")
val name = cardMusic.optString("name", "") val name = cardMusic.optString("name", "")
@ -1156,7 +1135,6 @@ object AgentService : CoroutineScope {
playlist.add(mapOf("id" to id,"url" to url, "title" to name, "artist" to sgener,"coverUrl" to image)) playlist.add(mapOf("id" to id,"url" to url, "title" to name, "artist" to sgener,"coverUrl" to image))
} }
if (playlist.isNotEmpty()) { if (playlist.isNotEmpty()) {
Log.w(TAG, "自动播放音乐列表 ${playlist}")
broadcast = true; broadcast = true;
processMusicPlayList(playlist) processMusicPlayList(playlist)
} else { } else {
@ -1189,7 +1167,6 @@ 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)
} }
@ -1205,7 +1182,6 @@ object AgentService : CoroutineScope {
fun processMusicPlayList(songs:List<Map<String,String>>){ fun processMusicPlayList(songs:List<Map<String,String>>){
// 如果 playlist 中有有效的歌曲,开始播放 // 如果 playlist 中有有效的歌曲,开始播放
if (songs.isNotEmpty()) { if (songs.isNotEmpty()) {
Log.i(TAG, "启动音乐服务 播放音乐列表 ${songs}")
MusicServiceStarter.startServiceWithPlaylist(context, command = "playlist", songs = songs) MusicServiceStarter.startServiceWithPlaylist(context, command = "playlist", songs = songs)
} else { } else {
Log.i(TAG, "音乐列表为空,未启动播放服务") Log.i(TAG, "音乐列表为空,未启动播放服务")
@ -1223,7 +1199,6 @@ object AgentService : CoroutineScope {
fun processNavigation(start:String,end:String){ fun processNavigation(start:String,end:String){
// 如果 playlist 中有有效的歌曲,开始播放 // 如果 playlist 中有有效的歌曲,开始播放
if (!start.isNullOrEmpty() && !end.isNullOrEmpty()) { if (!start.isNullOrEmpty() && !end.isNullOrEmpty()) {
Log.i(TAG, "启动导航服务 ${start} ${end}")
NavigationServiceHelper.startNavigation(context,"start", startpos = start, endpos = end) NavigationServiceHelper.startNavigation(context,"start", startpos = start, endpos = end)
} else { } else {
Log.i(TAG, "启动导航服务失败") Log.i(TAG, "启动导航服务失败")

113
local_plugins/azure_speech/android/src/main/kotlin/com/yunqiinnovation/azure_speech/AzureAsrHelper.kt

@ -127,18 +127,27 @@ class AzureAsrHelper(private val context: Context) {
// 设置指定的识别语言 // 设置指定的识别语言
speechRecognitionLanguage = currentLanguage speechRecognitionLanguage = currentLanguage
} }
// 设置静音超时时间(毫秒) setProperty("SpeechServiceConnection_EndSilenceTimeoutMs", "300")
setProperty("SpeechServiceConnection_EndSilenceTimeoutMs", "800") setProperty("Speech_SegmentationSilenceTimeoutMs", "300")
setProperty("Speech_SegmentationSilenceTimeoutMs", "800")
// 优化:添加低延迟连接配置
setProperty("SpeechServiceConnection_InitialSilenceTimeoutMs", "200")
setProperty("SpeechServiceConnection_RecoMode", "INTERACTIVE")
setProperty("Speech_PushStreamFormat", "PCM")
// 设置分段策略为时间模式 // 设置分段策略为时间模式
setProperty("Speech_SegmentationStrategy", "Time") setProperty("Speech_SegmentationStrategy", "Time")
} }
// 录音文件类 // 录音文件类
recordfile = RecordingFile() recordfile = RecordingFile()
// 创建识别器 // 创建识别器
return setupRecognizer() val setupSuccess = setupRecognizer()
if (setupSuccess) {
// 优化:初始化完成后进行预热
warmupRecognizer()
}
return setupSuccess
} catch (e: Exception) { } catch (e: Exception) {
Log.e(tag, "初始化失败: ${e.message}") Log.e(tag, "初始化失败: ${e.message}")
return false return false
@ -224,7 +233,29 @@ class AzureAsrHelper(private val context: Context) {
microphoneStream = null microphoneStream = null
} }
} }
/**
* 预热识别器(减少首次识别延迟)
*/
private fun warmupRecognizer() {
try {
// 创建极短的音频数据进行预热
val warmupData = ByteArray(320) // 10ms 16kHz 单声道
// 模拟静音数据
warmupData.fill(0)
// 如果是外部音频流,进行预热
if (audioSourceType == AudioSourceType.EXTERNAL) {
externalAudioStream?.pushAudio(warmupData)
}
Log.d(tag, "ASR预热完成")
} catch (e: Exception) {
// 预热失败不影响正常使用
Log.d(tag, "ASR预热失败: ${e.message}")
}
}
/** /**
* 设置外部音频流 - 使用拉流方式 * 设置外部音频流 - 使用拉流方式
*/ */
@ -357,23 +388,30 @@ class AzureAsrHelper(private val context: Context) {
* 设置事件监听器 * 设置事件监听器
*/ */
private fun setupEventListeners(callback: ContinuousRecognizeCallback) { private fun setupEventListeners(callback: ContinuousRecognizeCallback) {
// 识别中事件 // 优化:识别中事件 - 添加文本长度检查
recognizer?.recognizing?.addEventListener( recognizer?.recognizing?.addEventListener(
EventHandler<SpeechRecognitionEventArgs> { _, event -> EventHandler<SpeechRecognitionEventArgs> { _, event ->
val detectedLanguage = // 优化:只处理非空结果
AutoDetectSourceLanguageResult.fromResult(event.result)?.language ?: "" if (event.result.text.isNotEmpty()) {
// 直接在当前线程调用回调 val detectedLanguage = if (isAutoDetectLanguage) {
callback.onRecognizing(event.result.text, detectedLanguage) AutoDetectSourceLanguageResult.fromResult(event.result)?.language ?: ""
} else {
currentLanguage
}
callback.onRecognizing(event.result.text, detectedLanguage)
}
} }
) )
// 识别完成事件 // 优化:识别完成事件 - 优化语言检测
recognizer?.recognized?.addEventListener( recognizer?.recognized?.addEventListener(
EventHandler<SpeechRecognitionEventArgs> { _, event -> EventHandler<SpeechRecognitionEventArgs> { _, event ->
if (event.result.reason == ResultReason.RecognizedSpeech) { if (event.result.reason == ResultReason.RecognizedSpeech && event.result.text.isNotEmpty()) {
val detectedLanguage = val detectedLanguage = if (isAutoDetectLanguage) {
AutoDetectSourceLanguageResult.fromResult(event.result)?.language ?: "" AutoDetectSourceLanguageResult.fromResult(event.result)?.language ?: supportedLanguages[0]
// 直接在当前线程调用回调 } else {
currentLanguage
}
callback.onResult(event.result.text, detectedLanguage) callback.onResult(event.result.text, detectedLanguage)
} }
} }
@ -587,15 +625,15 @@ class AzureAsrHelper(private val context: Context) {
private inner class MicrophoneStream : PullAudioInputStreamCallback() { private inner class MicrophoneStream : PullAudioInputStreamCallback() {
private var audioRecord: AudioRecord? = null private var audioRecord: AudioRecord? = null
private var echoCanceler: AcousticEchoCanceler? = null private var echoCanceler: AcousticEchoCanceler? = null
// 音频配置 // 音频配置
private val sampleRate = 16000 private val sampleRate = 16000
private val channelConfig = AudioFormat.CHANNEL_IN_MONO private val channelConfig = AudioFormat.CHANNEL_IN_MONO
private val audioFormat = AudioFormat.ENCODING_PCM_16BIT private val audioFormat = AudioFormat.ENCODING_PCM_16BIT
// 优化:减少缓冲区大小以降低延迟
private val bufferSize = AudioRecord.getMinBufferSize( private val bufferSize = AudioRecord.getMinBufferSize(
sampleRate, channelConfig, audioFormat sampleRate, channelConfig, audioFormat
).let { if (it < 0) 4096 else it * 2 } ).let { if (it < 0) 2048 else it }
init { init {
initMicrophone() initMicrophone()
} }
@ -649,14 +687,20 @@ class AzureAsrHelper(private val context: Context) {
} }
/** /**
* 应用音频效果 * 优化:条件性应用音频效果
*/ */
private fun applyAudioEffects() { private fun applyAudioEffects() {
val sessionId = audioRecord?.audioSessionId ?: return val sessionId = audioRecord?.audioSessionId ?: return
// 回音消除 // 优化:只在真正需要时启用回音消除
echoCanceler = AcousticEchoCanceler.create(sessionId) val audioManager = context.getSystemService(Context.AUDIO_SERVICE) as AudioManager
echoCanceler?.enabled = true; if (audioManager.isSpeakerphoneOn) {
echoCanceler = AcousticEchoCanceler.create(sessionId)
echoCanceler?.enabled = true
Log.d(tag, "启用回音消除")
} else {
Log.d(tag, "跳过回音消除(使用耳机)")
}
} }
/** /**
@ -669,6 +713,9 @@ class AzureAsrHelper(private val context: Context) {
* Azure SDK 调用此方法获取音频数据 * Azure SDK 调用此方法获取音频数据
* 这是拉流模式的核心方法,由SDK调用以获取音频数据 * 这是拉流模式的核心方法,由SDK调用以获取音频数据
*/ */
/**
* 优化:简化音频数据读取(移除文件保存)
*/
override fun read(buffer: ByteArray): Int { override fun read(buffer: ByteArray): Int {
try { try {
val result = audioRecord?.read(buffer, 0, buffer.size) ?: -1 val result = audioRecord?.read(buffer, 0, buffer.size) ?: -1
@ -676,9 +723,6 @@ class AzureAsrHelper(private val context: Context) {
Log.e(tag, "读取音频数据失败: $result") Log.e(tag, "读取音频数据失败: $result")
return 0 return 0
} }
Log.i(tag, "麦克风:")
// 将音频数据保存成wav格式的音频文件
recordfile?.saveAudioDataToWav(buffer)
return result return result
} catch (e: Exception) { } catch (e: Exception) {
Log.e(tag, "读取音频异常: ${e.message}") Log.e(tag, "读取音频异常: ${e.message}")
@ -686,13 +730,11 @@ class AzureAsrHelper(private val context: Context) {
} }
} }
/** /**
* 关闭音频资源 * 优化:简化资源关闭
*/ */
override fun close() { override fun close() {
releaseAudioResources() releaseAudioResources()
} }
/** /**
@ -739,18 +781,21 @@ class AzureAsrHelper(private val context: Context) {
} }
/** /**
* SDK调用:从队列中拉取数据 * 修复:外部音频流恢复阻塞读取
* @param buffer SDK提供的缓冲区 * @param buffer SDK提供的缓冲区
* @return 读取的字节数,0表示流结束 * @return 读取的字节数,0表示流结束
*/ */
override fun read(buffer: ByteArray): Int { override fun read(buffer: ByteArray): Int {
try { try {
// 阻塞等待下一块数据
val chunk = queue.take() val chunk = queue.take()
// 检查是否是结束标志(空数组)
if (chunk.isEmpty()) {
return 0
}
val toCopy = minOf(chunk.size, buffer.size) val toCopy = minOf(chunk.size, buffer.size)
System.arraycopy(chunk, 0, buffer, 0, toCopy) System.arraycopy(chunk, 0, buffer, 0, toCopy)
Log.i(tag, "外部音频:")
return toCopy return toCopy
} catch (e: InterruptedException) { } catch (e: InterruptedException) {
Thread.currentThread().interrupt() Thread.currentThread().interrupt()

79
local_plugins/azure_speech/android/src/main/kotlin/com/yunqiinnovation/azure_speech/AzureTtsHelper.kt

@ -70,7 +70,14 @@ class AzureTtsHelper(private val context: Context) : ITtsService {
// 创建语音配置 // 创建语音配置
speechConfig = SpeechConfig.fromSubscription(ttsAppToken, ttsResource) speechConfig = SpeechConfig.fromSubscription(ttsAppToken, ttsResource)
speechConfig?.setSpeechSynthesisOutputFormat(SpeechSynthesisOutputFormat.Riff24Khz16BitMonoPcm)
// 优化:使用更低延迟的音频格式(从24KHz降到16KHz)
speechConfig?.setSpeechSynthesisOutputFormat(SpeechSynthesisOutputFormat.Riff16Khz16BitMonoPcm)
// 优化:设置低延迟连接属性
speechConfig?.setProperty("SpeechServiceConnection_InitialSilenceTimeoutMs", "300")
speechConfig?.setProperty("SpeechServiceConnection_EndSilenceTimeoutMs", "300")
speechConfig?.setSpeechSynthesisVoiceName(currentVoice) speechConfig?.setSpeechSynthesisVoiceName(currentVoice)
// 创建音频配置 // 创建音频配置
@ -91,6 +98,9 @@ class AzureTtsHelper(private val context: Context) : ITtsService {
isInitialized = true isInitialized = true
FileLogger.d(TAG, "TTS引擎初始化成功: 区域=$ttsResource") FileLogger.d(TAG, "TTS引擎初始化成功: 区域=$ttsResource")
// 优化:初始化完成后立即预热
warmupSynthesizer()
return true return true
} catch (e: Exception) { } catch (e: Exception) {
FileLogger.e(TAG, "TTS引擎初始化失败: ${e.message}") FileLogger.e(TAG, "TTS引擎初始化失败: ${e.message}")
@ -243,6 +253,30 @@ class AzureTtsHelper(private val context: Context) : ITtsService {
)) ))
} }
} }
/**
* 预热TTS引擎(减少首次合成延迟)
*/
private fun warmupSynthesizer() {
try {
// 使用极短文本进行预热
val warmupSsml = """
<speak version="1.0" xmlns="http://www.w3.org/2001/10/synthesis" xml:lang="zh-CN">
<voice name="$currentVoice">
<prosody rate="$currentRate" pitch="$currentPitch" volume="0%">
.
</prosody>
</voice>
</speak>
""".trimIndent()
// 静音预热(音量设为0)
synthesizer?.SpeakSsmlAsync(warmupSsml)
} catch (e: Exception) {
// 预热失败不影响正常使用
FileLogger.d(TAG, "TTS预热失败: ${e.message}")
}
}
//清洗播报语音内容 //清洗播报语音内容
fun cleanTextForTTS(text: String): String { fun cleanTextForTTS(text: String): String {
return text return text
@ -320,6 +354,30 @@ class AzureTtsHelper(private val context: Context) : ITtsService {
""".trimIndent() """.trimIndent()
} }
/**
* 生成优化的SSML(减少复杂度提升速度)
*/
private fun generateOptimizedSsml(rawText: String): String {
// 优化:简化文本预处理,减少正则表达式使用
val processedText = rawText
.replace("&", "&amp;")
.replace("<", "&lt;")
.replace(">", "&gt;")
.replace(Regex("[😀-🟿]+"), "") // 简化表情符号移除
.trim()
// 优化:简化SSML结构,减少嵌套层级
return """
<speak version="1.0" xmlns="http://www.w3.org/2001/10/synthesis" xml:lang="zh-CN">
<voice name="$currentVoice">
<prosody rate="$currentRate" pitch="$currentPitch" volume="$currentVolume">
$processedText
</prosody>
</voice>
</speak>
""".trimIndent()
}
/** /**
* 单次播放文本(非流式) * 单次播放文本(非流式)
*/ */
@ -334,8 +392,8 @@ class AzureTtsHelper(private val context: Context) : ITtsService {
} }
try { try {
// 生成SSML并播放 // 优化:使用简化的SSML生成
val ssml = generateSsml(text) val ssml = generateOptimizedSsml(text)
isSpeaking = true isSpeaking = true
@ -373,21 +431,22 @@ class AzureTtsHelper(private val context: Context) : ITtsService {
// 添加新文本到缓冲区 // 添加新文本到缓冲区
streamBuffer.append(text) streamBuffer.append(text)
// 增加500ms防抖逻辑 // 优化:大幅减少防抖时间从600ms到150ms
val currentTime = System.currentTimeMillis() val currentTime = System.currentTimeMillis()
if (currentTime - lastSpeakTime < 600) { if (currentTime - lastSpeakTime < 150) {
return true return true
} }
lastSpeakTime = currentTime lastSpeakTime = currentTime
val currentText = streamBuffer.toString() //cleanTextForTTS(streamBuffer.toString()) val currentText = streamBuffer.toString() //cleanTextForTTS(streamBuffer.toString())
// 定义标点符号列表
val punctuationMarks = listOf('.', '。', '!', '!', '?', '?', ';', ';', ',', ',', ':', ':', '\n')
// 查找最后一个标点符号的位置 // 优化:使用字符集合替代列表,提升查找效率
val punctuationSet = setOf('.', '。', '!', '!', '?', '?', ';', ';', ',', ',', ':', ':', '\n')
// 优化:从后往前查找最后一个标点符号
var lastPunctuationIndex = -1 var lastPunctuationIndex = -1
for (i in currentText.indices.reversed()) { for (i in currentText.length - 1 downTo 0) {
if (currentText[i] in punctuationMarks) { if (currentText[i] in punctuationSet) {
lastPunctuationIndex = i lastPunctuationIndex = i
break break
} }

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 // OpenAI Kotlin 4.0.1
implementation("com.aallam.openai:openai-client: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处理 // JSON处理
implementation("org.jetbrains.kotlinx:kotlinx-serialization-json:1.6.0") 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 package com.yunqiinnovation.chat_api
import android.os.Handler import androidx.annotation.NonNull
import android.os.Looper
import io.flutter.embedding.engine.plugins.FlutterPlugin import io.flutter.embedding.engine.plugins.FlutterPlugin
import io.flutter.plugin.common.EventChannel import io.flutter.plugin.common.EventChannel
import io.flutter.plugin.common.MethodCall 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.MethodCallHandler
import io.flutter.plugin.common.MethodChannel.Result import io.flutter.plugin.common.MethodChannel.Result
import kotlinx.coroutines.* import kotlinx.coroutines.*
import android.util.Log
/** /**
* ChatApiPlugin * ChatApiPlugin
@ -17,215 +17,144 @@ import kotlinx.coroutines.*
* 与 iOS 版本接口完全一致 * 与 iOS 版本接口完全一致
*/ */
class ChatApiPlugin : FlutterPlugin, MethodCallHandler, EventChannel.StreamHandler { class ChatApiPlugin : FlutterPlugin, MethodCallHandler, EventChannel.StreamHandler {
private lateinit var methodChannel: MethodChannel private lateinit var channel: MethodChannel
private lateinit var eventChannel: EventChannel private lateinit var eventChannel: EventChannel
private var eventSink: EventChannel.EventSink? = null private var eventSink: EventChannel.EventSink? = null
private val chatApiService = ChatApiService() private var chatApiService: ChatApiService? = null
private val mainHandler = Handler(Looper.getMainLooper())
// 协程作用域 // 插件协程作用域
private val pluginScope = CoroutineScope(Dispatchers.Main + SupervisorJob()) private val pluginScope = CoroutineScope(Dispatchers.Main + SupervisorJob())
override fun onAttachedToEngine(flutterPluginBinding: FlutterPlugin.FlutterPluginBinding) { override fun onAttachedToEngine(@NonNull flutterPluginBinding: FlutterPlugin.FlutterPluginBinding) {
methodChannel = MethodChannel(flutterPluginBinding.binaryMessenger, "com.yunqiinnovation.chat_api/methods") channel = MethodChannel(flutterPluginBinding.binaryMessenger, "chat_api")
methodChannel.setMethodCallHandler(this) channel.setMethodCallHandler(this)
eventChannel = EventChannel(flutterPluginBinding.binaryMessenger, "com.yunqiinnovation.chat_api/events") eventChannel = EventChannel(flutterPluginBinding.binaryMessenger, "com.yunqiinnovation.chat_api/events")
eventChannel.setStreamHandler(this) 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) { override fun onToken(token: String) {
sendEvent("token", token) channel.invokeMethod("onToken", token)
} }
override fun onComplete() { override fun onComplete() {
sendEvent("complete", null) channel.invokeMethod("onComplete", null)
} }
override fun onError(error: Exception) { override fun onError(error: Exception) {
sendEvent("error", error.message ?: "未知错误") channel.invokeMethod("onError", error.message)
} }
override fun onFunctionCall(functionCall: org.json.JSONObject) { override fun onFunctionCall(functionCall: org.json.JSONObject) {
try { channel.invokeMethod("onFunctionCall", functionCall.toString())
val jsonString = functionCall.toString()
sendEvent("functionCall", jsonString)
} catch (e: Exception) {
sendEvent("error", "Failed to serialize function call: ${e.message}")
}
} }
override fun onFunctionCallResult(functionCall: org.json.JSONObject, functionCallResult: org.json.JSONObject) { override fun onFunctionCallResult(functionCall: org.json.JSONObject, functionCallResult: org.json.JSONObject) {
try { channel.invokeMethod("onFunctionCallResult", mapOf(
val meta = org.json.JSONObject().apply { "functionCall" to functionCall.toString(),
put("functionCall", functionCall) "functionCallResult" to functionCallResult.toString()
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}")
}
} }
}) })
} }
override fun onDetachedFromEngine(binding: FlutterPlugin.FlutterPluginBinding) { override fun onDetachedFromEngine(@NonNull binding: FlutterPlugin.FlutterPluginBinding) {
methodChannel.setMethodCallHandler(null) channel.setMethodCallHandler(null)
eventChannel.setStreamHandler(null) eventChannel.setStreamHandler(null)
// 清理资源
pluginScope.cancel() 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) { when (call.method) {
"initialize" -> handleInitialize(call, result) "initialize" -> {
"createUserMessage" -> handleCreateUserMessage(call, result) val apiKey = call.argument<String>("apiKey") ?: ""
"createAssistantMessage" -> handleCreateAssistantMessage(call, result) val baseUrl = call.argument<String>("baseUrl") ?: ""
"sendMessage" -> handleSendMessage(call, result) val model = call.argument<String>("model") ?: ""
"sendMessageStream" -> handleSendMessageStream(call, result) val mcpServer = call.argument<String>("mcpServer") ?: ""
"sendFunctionCallResult" -> handleSendFunctionCallResult(call, result) val success = chatApiService?.initialize(apiKey, baseUrl, model, mcpServer) ?: false
"cancelCurrentStream" -> handleCancelCurrentStream(result) result.success(success)
"registerFunction" -> handleRegisterFunction(call, result) }
"initializeMcpClient" -> handleInitializeMcpClient(call, result) "chatCompletionStream" -> {
"isMcpInitialized" -> handleIsMcpInitialized(result) val messages = call.argument<List<Map<String, Any>>>("messages") ?: emptyList()
"closeMcpClient" -> handleCloseMcpClient(result) val tool = call.argument<Boolean>("tool") ?: false
"handleMcpToolCall" -> handleMcpToolCall(call, result) pluginScope.launch {
else -> result.notImplemented() chatApiService?.chatCompletionStream(messages, tool)
} }
} result.success(true)
}
// MARK: - Method Handlers "cancelChatStream" -> {
chatApiService?.cancelChatStream()
private fun handleInitialize(call: MethodCall, result: Result) { result.success(true)
val apiKey = call.argument<String>("apiKey") }
if (apiKey.isNullOrEmpty()) { "processImage" -> {
result.error("INVALID_ARGS", "缺少必要参数", null) val imagePath = call.argument<String>("imagePath") ?: ""
return val prompt = call.argument<String>("prompt") ?: ""
} val maxWidth = call.argument<Double>("maxWidth") ?: 2048.0
val detail = call.argument<String>("detail") ?: "auto"
val baseUrl = call.argument<String>("baseUrl") ?: "" val base64Image = chatApiService?.processImage(imagePath, prompt, maxWidth, detail)
val model = call.argument<String>("model") ?: "" result.success(base64Image)
val mcpServer = call.argument<String>("mcpServer") ?: "" }
"initializeMcpClient" -> {
val success = chatApiService.initialize(apiKey, baseUrl, model, mcpServer) val serverUrl = call.argument<String>("serverUrl") ?: ""
result.success(success) val success = chatApiService?.initializeMcpClient(serverUrl) ?: false
} result.success(success)
}
private fun handleCreateUserMessage(call: MethodCall, result: Result) { "isMcpInitialized" -> {
val content = call.argument<String>("content") val initialized = chatApiService?.isMcpInitialized() ?: false
if (content.isNullOrEmpty()) { result.success(initialized)
result.error("INVALID_ARGS", "缺少必要参数", null) }
return "closeMcpClient" -> {
} chatApiService?.closeMcpClient()
result.success(true)
val message = chatApiService.createUserMessage(content) }
result.success(message) "getToolMaps" -> {
} val mcpClient = chatApiService?.mcpClient
if (mcpClient != null) {
private fun handleCreateAssistantMessage(call: MethodCall, result: Result) { val toolMaps = mcpClient.getToolMaps()
val content = call.argument<String>("content") result.success(toolMaps)
if (content.isNullOrEmpty()) { } else {
result.error("INVALID_ARGS", "缺少必要参数", null) result.success(emptyList<Map<String, Any>>())
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)
} }
} catch (e: Exception) { }
mainHandler.post { "hasToolWithName" -> {
result.error("SEND_ERROR", e.message, null) 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 // MARK: - Event Stream Handler
@ -238,15 +167,13 @@ class ChatApiPlugin : FlutterPlugin, MethodCallHandler, EventChannel.StreamHandl
} }
private fun sendEvent(type: String, content: Any?, meta: String? = null) { private fun sendEvent(type: String, content: Any?, meta: String? = null) {
mainHandler.post { eventSink?.let { sink ->
eventSink?.let { sink -> val event = mutableMapOf<String, Any>(
val event = mutableMapOf<String, Any>( "type" to type
"type" to type )
) content?.let { event["content"] = it }
content?.let { event["content"] = it } meta?.let { event["meta"] = it }
meta?.let { event["meta"] = it } sink.success(event)
sink.success(event)
}
} }
} }
} }

434
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.math.sqrt
import kotlin.time.Duration.Companion.seconds import kotlin.time.Duration.Companion.seconds
import android.util.Log import android.util.Log
import org.json.JSONObject
import android.os.Handler
import android.os.Looper
/** /**
* ChatAPI服务异常 * ChatAPI服务异常
@ -65,7 +68,13 @@ private data class ToolCallInfo(
var name: String = "", var name: String = "",
var arguments: String = "" var arguments: String = ""
) { ) {
fun isValid(): Boolean = id.isNotEmpty() && name.isNotEmpty() fun isValid(): Boolean {
return try {
id.isNotEmpty() && name.isNotEmpty()
} catch (e: Exception) {
false
}
}
} }
/** /**
@ -83,9 +92,21 @@ class ChatApiService(private val context: android.content.Context? = null) : Cor
private var visionModel = "gpt-4-vision-preview" private var visionModel = "gpt-4-vision-preview"
private var isInitialized = false private var isInitialized = false
// Handler for main thread
private val mainHandler = Handler(Looper.getMainLooper())
// OpenAI 客户端 // OpenAI 客户端
private var openAI: OpenAI? = null 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 currentStreamJob: Job? = null
private var streamCallback: StreamCallback? = null private var streamCallback: StreamCallback? = null
@ -93,6 +114,8 @@ class ChatApiService(private val context: android.content.Context? = null) : Cor
private var toolCalls: MutableMap<Int, ToolCallInfo> = mutableMapOf() private var toolCalls: MutableMap<Int, ToolCallInfo> = mutableMapOf()
private var isCanceled = false private var isCanceled = false
// JSON处理 // JSON处理
private val gson = Gson() private val gson = Gson()
@ -116,7 +139,6 @@ class ChatApiService(private val context: android.content.Context? = null) : Cor
if (model.isNotEmpty()) { if (model.isNotEmpty()) {
this.model = model this.model = model
} }
Log.d("ChatApiService", "原始 baseUrl: $baseUrl")
// 处理 baseUrl:移除末尾的 /chat/completions(如果存在) // 处理 baseUrl:移除末尾的 /chat/completions(如果存在)
// 因为 openai-kotlin 会自动拼接 /chat/completions // 因为 openai-kotlin 会自动拼接 /chat/completions
@ -131,8 +153,6 @@ class ChatApiService(private val context: android.content.Context? = null) : Cor
processedBaseUrl += "/" processedBaseUrl += "/"
} }
Log.d("ChatApiService", "处理后 baseUrl: $processedBaseUrl")
return try { return try {
// 创建OpenAI配置 // 创建OpenAI配置
val config = OpenAIConfig( val config = OpenAIConfig(
@ -143,15 +163,17 @@ class ChatApiService(private val context: android.content.Context? = null) : Cor
openAI = OpenAI(config) openAI = OpenAI(config)
// 初始化MCP客户端 (暂时留空,但保持接口一致) // 保存MCP配置以便后续使用
if (mcpServer.isNotEmpty()) { mcpConfigJson = mcpServer
// MCP功能暂时留空,但记录服务器地址以便后续实现 // 异步初始化MCP客户端
// initializeMcpClient(mcpServer) launch {
initializeMcpClient(mcpServer)
} }
isInitialized = apiKey.isNotEmpty() isInitialized = apiKey.isNotEmpty()
true true
} catch (e: Exception) { } catch (e: Exception) {
Log.e("ChatApiService", "ChatApiService初始化失败: ${e.message}", e)
false false
} }
} }
@ -200,12 +222,16 @@ class ChatApiService(private val context: android.content.Context? = null) : Cor
// 转换消息格式 // 转换消息格式
val chatMessages = convertToChatMessages(messages) val chatMessages = convertToChatMessages(messages)
// 直接获取工具列表
val tools = getOpenAiTools()
// 构建请求 // 构建请求
val chatCompletionRequest = ChatCompletionRequest( val chatCompletionRequest = ChatCompletionRequest(
model = ModelId(currentModel), model = ModelId(currentModel),
messages = chatMessages, messages = chatMessages,
maxTokens = 2000, maxTokens = 2000,
temperature = 0.7 temperature = 0.7,
tools = if (tools.isNotEmpty()) tools else null
) )
val result = openAI!!.chatCompletion(chatCompletionRequest) val result = openAI!!.chatCompletion(chatCompletionRequest)
@ -246,6 +272,7 @@ class ChatApiService(private val context: android.content.Context? = null) : Cor
*/ */
fun sendMessageStream(messages: List<Map<String, Any>>) { fun sendMessageStream(messages: List<Map<String, Any>>) {
if (!isInitialized || apiKey.isEmpty() || openAI == null) { if (!isInitialized || apiKey.isEmpty() || openAI == null) {
Log.e("ChatApiService", "ChatAPI服务未初始化,无法发送消息")
streamCallback?.onError(ChatApiException("ChatAPI服务未初始化")) streamCallback?.onError(ChatApiException("ChatAPI服务未初始化"))
return return
} }
@ -261,23 +288,62 @@ class ChatApiService(private val context: android.content.Context? = null) : Cor
currentStreamJob = launch { currentStreamJob = launch {
try { try {
// 转换消息格式 // 转换消息格式
val chatMessages = convertToChatMessages(messages) val chatMessages = try {
convertToChatMessages(messages)
} catch (e: Exception) {
Log.e("ChatApiService", "转换消息格式失败: ${e.message}", e)
throw e
}
// 直接获取工具列表
val tools = getOpenAiTools()
// 构建请求 // 构建请求
val chatCompletionRequest = ChatCompletionRequest( if (currentModel.isEmpty()) {
model = ModelId(currentModel), Log.e("ChatApiService", "模型名称为空")
messages = chatMessages, throw IllegalArgumentException("模型名称不能为空")
maxTokens = 2000, }
temperature = 0.7
) if (chatMessages.isEmpty()) {
Log.e("ChatApiService", "消息列表为空")
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
)
if (openAI == null) {
Log.e("ChatApiService", "openAI对象为null")
throw IllegalStateException("OpenAI客户端未初始化")
}
val flow = openAI!!.chatCompletions(chatCompletionRequest)
flow
} catch (e: Exception) {
Log.e("ChatApiService", "创建ChatCompletionRequest或调用chatCompletions失败: ${e.message}", e)
throw e
}
chatsFlow.collect { result -> chatsFlow.collect { result ->
if (isCanceled) return@collect if (isCanceled) return@collect
val choice = result.choices.firstOrNull() ?: return@collect val choice = result.choices?.firstOrNull()
val delta = choice.delta ?: return@collect if (choice == null) {
Log.e("ChatApiService", "choice为null")
return@collect
}
val delta = choice.delta
if (delta == null) {
Log.e("ChatApiService", "delta为null")
return@collect
}
// 处理普通文本内容 // 处理普通文本内容
delta.content?.let { content -> delta.content?.let { content ->
@ -291,19 +357,42 @@ class ChatApiService(private val context: android.content.Context? = null) : Cor
// 创建或获取现有的工具调用信息 // 创建或获取现有的工具调用信息
val toolCallInfo = toolCalls.getOrPut(index) { ToolCallInfo() } val toolCallInfo = toolCalls.getOrPut(index) { ToolCallInfo() }
// 更新ID // 安全处理工具调用ID
toolCall.id?.let { toolCallInfo.id = it.toString() } try {
toolCall.id?.let { id ->
toolCallInfo.id = id.toString()
}
} catch (e: Exception) {
Log.e("ChatApiService", "处理工具调用ID异常: ${e.message}")
}
// 更新函数信息 // 安全处理函数信息
toolCall.function?.let { function -> try {
function.name?.let { toolCallInfo.name = it } toolCall.function?.let { function ->
function.arguments?.let { toolCallInfo.arguments += it } try {
function.name?.let { name ->
toolCallInfo.name = name
}
} catch (e: Exception) {
Log.e("ChatApiService", "处理工具调用函数名称异常: ${e.message}")
}
try {
function.arguments?.let { args ->
toolCallInfo.arguments += args
}
} catch (e: Exception) {
Log.e("ChatApiService", "处理工具调用参数异常: ${e.message}")
}
}
} catch (e: Exception) {
Log.e("ChatApiService", "处理工具调用函数信息异常: ${e.message}")
} }
} }
} }
if (!isCanceled) { if (!isCanceled) {
// 处理工具调用或完成 // 检查是否有工具调用需要处理
val hasToolCalls = processToolCalls() val hasToolCalls = processToolCalls()
if (!hasToolCalls) { if (!hasToolCalls) {
streamCallback?.onComplete() streamCallback?.onComplete()
@ -322,14 +411,28 @@ class ChatApiService(private val context: android.content.Context? = null) : Cor
* 处理工具调用 * 处理工具调用
*/ */
private suspend fun processToolCalls(): Boolean { private suspend fun processToolCalls(): Boolean {
val firstToolCall = toolCalls.values.firstOrNull { it.isValid() } ?: return false val firstToolCall = try {
toolCalls.values.firstOrNull { it.isValid() }
} catch (e: Exception) {
Log.e("ChatApiService", "查找有效工具调用异常: ${e.message}")
null
}
if (firstToolCall == null) {
return false
}
// 创建函数调用字典 // 创建函数调用字典
val functionCall = mapOf( val functionCall = try {
"name" to firstToolCall.name, mapOf(
"arguments" to firstToolCall.arguments, "name" to (firstToolCall.name.takeIf { it.isNotEmpty() } ?: ""),
"id" to firstToolCall.id "arguments" to (firstToolCall.arguments.takeIf { it.isNotEmpty() } ?: "{}"),
) "id" to (firstToolCall.id.takeIf { it.isNotEmpty() } ?: "")
)
} catch (e: Exception) {
Log.e("ChatApiService", "创建函数调用字典异常: ${e.message}")
return false
}
// 通知上层工具调用事件 // 通知上层工具调用事件
streamCallback?.onFunctionCall(convertMapToJsonObject(functionCall)) streamCallback?.onFunctionCall(convertMapToJsonObject(functionCall))
@ -338,8 +441,43 @@ class ChatApiService(private val context: android.content.Context? = null) : Cor
launch { launch {
try { try {
if (!isCanceled) { if (!isCanceled) {
// 这里暂时返回占位符结果,实际MCP功能留空 // 通过MCP客户端处理工具调用
val result = mapOf("context" to "MCP功能暂未实现") val functionName = firstToolCall.name
val argumentsJson = firstToolCall.arguments
val result = if (_mcpClient?.hasToolWithName(functionName) == true) {
// 解析参数
val arguments = _mcpClient?.parseJsonArguments(argumentsJson) ?: emptyMap()
// 调用MCP工具
val toolResult = _mcpClient?.callTool(functionName, arguments)
// 处理结果
if (toolResult != null) {
if (toolResult["isError"] == true) {
// 处理错误情况
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")) {
// 本地函数结果
toolResult
} else {
// 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 {
// 工具不存在
mapOf("context" to "Tool not found: $functionName")
}
if (!isCanceled) { if (!isCanceled) {
// 处理结果 // 处理结果
@ -357,6 +495,7 @@ class ChatApiService(private val context: android.content.Context? = null) : Cor
} }
} }
} catch (e: Exception) { } catch (e: Exception) {
Log.e("ChatApiService", "工具调用处理过程中出错: ${e.message}", e)
if (!isCanceled) { if (!isCanceled) {
val errorMessage = "工具调用处理失败: ${e.message}" val errorMessage = "工具调用处理失败: ${e.message}"
sendFunctionCallResultInternal( sendFunctionCallResultInternal(
@ -435,40 +574,221 @@ class ChatApiService(private val context: android.content.Context? = null) : Cor
/** /**
* 注册函数 * 注册函数
*/ */
fun registerFunction(name: String, description: String, parameters: Map<String, Any>): Boolean { fun registerFunction(name: String, description: String, parameters: JSONObject): Boolean {
// 暂时返回true,实际功能留空,但保持与iOS版本接口一致 // 将JSONObject转换为Map
return true 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客户端 * 初始化MCP客户端
*/ */
fun initializeMcpClient(serverUrl: String): Boolean { fun initializeMcpClient(serverUrl: String): Boolean {
// MCP功能暂时留空,但保持与iOS版本接口一致 if (_mcpClient == null) {
return false _mcpClient = MCPClient(context)
}
// 直接使用类的CoroutineScope启动协程
launch {
try {
_mcpClient?.connectToSSE(serverUrl)
} catch (e: Exception) {
Log.e("ChatApiService", "MCP客户端初始化失败: ${e.message}", e)
}
}
return true // 立即返回,实际连接在后台进行
} }
/** /**
* MCP客户端是否已初始化 * MCP客户端是否已初始化
*/ */
fun isMcpInitialized(): Boolean { fun isMcpInitialized(): Boolean {
// MCP功能暂时留空,但保持与iOS版本接口一致 return _mcpClient?.isConnected() ?: false
return false
} }
/**
* 直接从MCPClient获取OpenAI工具格式
* 将MCP工具映射转换为OpenAI工具格式
*/
private fun getOpenAiTools(): List<Tool> {
try {
val tools = mutableListOf<Tool>()
val toolMaps = mcpClient?.getToolMaps() ?: return emptyList()
toolMaps.forEach { toolMap ->
try {
val type = toolMap["type"] as? String
if (type == "function") {
@Suppress("UNCHECKED_CAST")
val functionMap = toolMap["function"] as? Map<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>
if (parametersMap == null) {
Log.e("ChatApiService", "工具 $name 的parameters为null")
return@forEach
}
// 验证parametersMap的基本结构
if (!parametersMap.containsKey("type")) {
Log.e("ChatApiService", "工具 $name 的parameters缺少type字段")
return@forEach
}
val parametersJson = gson.toJson(parametersMap)
try {
if (parametersJson.isBlank()) {
Log.e("ChatApiService", "parametersJson为空或空白")
return@forEach
}
val parameters = com.aallam.openai.api.core.Parameters.fromJsonString(parametersJson)
if (name.isEmpty()) {
Log.e("ChatApiService", "工具名称为空,跳过")
return@forEach
}
val tool = Tool.function(
name = name,
description = description,
parameters = parameters
)
tools.add(tool)
} catch (e: Exception) {
Log.e("ChatApiService", "创建Parameters对象失败 for 工具 $name: ${e.message}", e)
// 跳过这个工具,继续处理其他工具
}
}
} catch (e: Exception) {
Log.e("ChatApiService", "处理工具映射时出错: ${e.message}", e)
}
}
return tools
} catch (e: Exception) {
Log.e("ChatApiService", "获取OpenAI工具格式失败: ${e.message}", e)
return emptyList()
}
}
/** /**
* 关闭MCP客户端 * 关闭MCP客户端
*/ */
fun closeMcpClient() { fun closeMcpClient() {
// MCP功能暂时留空,但保持与iOS版本接口一致 runBlocking {
_mcpClient?.disconnectAll()
_mcpClient = null
}
} }
/** /**
* 处理MCP工具调用 * 处理MCP工具调用
*
* @param functionCall 函数调用JSON对象,必须包含name和arguments字段
* @return 工具调用结果
*/
suspend fun handleMcpToolCall(functionCall: JSONObject): String {
if (_mcpClient == null || !isMcpInitialized()) {
return "MCP客户端未初始化"
}
return try {
// 获取函数名称
val name = functionCall.getString("name")
// 获取参数
val argumentsJson = functionCall.getString("arguments")
val arguments = _mcpClient?.parseJsonArguments(argumentsJson) ?: mapOf()
// 调用工具
val result = _mcpClient?.callTool(name, arguments)
// 返回工具调用结果
result?.get("context") as? String ?: "工具调用失败"
} catch (e: Exception) {
Log.e("ChatApiService", "处理MCP工具调用失败: ${e.message}", e)
"处理MCP工具调用失败: ${e.message}"
}
}
/**
* 获取所有可用的工具定义
*/ */
suspend fun handleMcpToolCall(functionCallJson: String): String { fun getToolDefinitions(): List<Map<String, Any>> {
// MCP功能暂时留空,但保持与iOS版本接口一致 return _mcpClient?.getToolMaps() ?: emptyList()
return "MCP功能暂未实现" }
/**
* 处理图片
*/
fun processImage(imagePath: String, prompt: String, maxWidth: Double, detail: String): String? {
return fileToBase64(imagePath, (maxWidth * 2).toInt()) // 简化处理,使用宽度的两倍作为最大KB数
}
/**
* 取消聊天流
*/
fun cancelChatStream() {
cancelCurrentStream()
}
/**
* 聊天完成流式接口
* 与iOS版本保持一致的接口
*/
fun chatCompletionStream(messages: List<Map<String, Any>>, tool: Boolean = false) {
// 直接调用sendMessageStream,因为该方法已经处理了工具调用
sendMessageStream(messages)
} }
// MARK: - 工具方法 // MARK: - 工具方法
@ -796,4 +1116,28 @@ class ChatApiService(private val context: android.content.Context? = null) : Cor
return jsonObject 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"
}
}
} }

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

@ -0,0 +1,259 @@
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" -> {
// SSE连接已打开
}
"endpoint" -> {
try {
val eventData = event.data ?: ""
// 构建完整的端点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
}
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
}
// 创建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"
}
client.sseSession(
urlString = sseConnectUrl,
reconnectionTime = reconnectionTime,
block = requestBuilder,
)
} ?: client.sseSession(
reconnectionTime = reconnectionTime,
block = requestBuilder,
)
// 收集SSE事件
collectEvents()
// 等待endpoint就绪
endpoint.await()
}
/**
* 发送消息
*/
@OptIn(ExperimentalCoroutinesApi::class)
override suspend fun send(message: JSONRPCMessage) {
if (!endpoint.isCompleted) {
Log.e(TAG, "发送失败: 未连接")
error("Not connected")
}
try {
val messageEndpoint = endpoint.getCompleted()
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()
}
}

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

@ -0,0 +1,372 @@
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++
} else {
Log.w(TAG, "Failed to connect to MCP server: $serverId")
}
}
isConnectedFlag = connectedCount > 0
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()
} 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
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)
}
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 trimmedJson = json.trim()
if (trimmedJson.isEmpty()) {
return emptyMap()
}
// 如果不是以{开头,尝试包装为{}
val jsonToUse = if (!trimmedJson.startsWith("{")) {
if (trimmedJson.contains("=") || trimmedJson.contains(":")) {
// 简单的键值对,包装成JSON对象
"{$trimmedJson}"
} else {
// 空参数,返回空对象
"{}"
}
} else {
trimmedJson
}
val jsonObject = JSONObject(jsonToUse)
convertJsonObjectToMap(jsonObject)
} catch (e: Exception) {
Log.w(TAG, "Failed to parse JSON arguments: '$json'", e)
emptyMap()
}
}
/**
* 是否已连接
*/
fun isConnected(): Boolean {
return isConnectedFlag || localFunctions.isNotEmpty()
}
/**
* 断开所有连接
*/
suspend fun disconnectAll() {
closeAllConnections()
}
/**
* 关闭连接
*/
override fun close() {
runBlocking {
closeAllConnections()
localFunctions.clear()
localFunctionDefs.clear()
}
}
/**
* 将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
}
}

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

@ -0,0 +1,297 @@
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 {
// 创建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
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)
}
} catch (e: Exception) {
Log.w(TAG, "[$serverId] 获取工具列表失败: ${e.message}")
// 即使获取工具失败,连接也可能是成功的
}
mcpClient = client
isConnected = true
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>> {
val toolMaps = availableTools.map { tool ->
val parametersMap = tool.inputSchema?.let { inputSchema ->
convertInputSchemaToMap(inputSchema)
} ?: mapOf(
"type" to "object",
"properties" to emptyMap<String, Any>(),
"required" to emptyList<String>()
)
val toolMap = mapOf(
"type" to "function",
"function" to mapOf(
"name" to tool.name,
"description" to (tool.description ?: ""),
"parameters" to parametersMap
)
)
toolMap
}
return toolMaps
}
/**
* 调用MCP工具
*/
suspend fun callTool(name: String, arguments: Map<String, Any>): Map<String, Any>? {
val client = mcpClient ?: return null
return try {
// 创建工具调用请求 - 将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 ->
// 将结果转换为统一格式
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> {
val properties = mutableMapOf<String, Any>()
val required = mutableListOf<String>()
// 处理properties
inputSchema.properties?.let { propsJsonObject ->
for ((key, value) in propsJsonObject) {
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()
}
}
}
}
// 处理required
inputSchema.required?.let { requiredList ->
required.addAll(requiredList)
}
val result = mapOf(
"type" to "object",
"properties" to properties,
"required" to required
)
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)
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()
} 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