diff --git a/lib/modules/agent/controllers/agent_controller.dart b/lib/modules/agent/controllers/agent_controller.dart index cc98107c4..7cdebb046 100644 --- a/lib/modules/agent/controllers/agent_controller.dart +++ b/lib/modules/agent/controllers/agent_controller.dart @@ -350,7 +350,7 @@ class AgentController extends GetxController { final token = event.data['token'] ?? ''; final responseId = event.data['responseId'] ?? ''; - Logger.i(TAG, 'AI回复Token: $token, responseId: $responseId'); + // Logger.i(TAG, 'AI回复Token: $token, responseId: $responseId'); if (token.isNotEmpty) { // 如果是新的回复或者响应ID改变,创建新消息 if (_isNewAssistantResponse || diff --git a/local_plugins/agent_service/android/src/main/kotlin/com/yunqiinnovation/agent_service/AgentService.kt b/local_plugins/agent_service/android/src/main/kotlin/com/yunqiinnovation/agent_service/AgentService.kt index ea0a1b19a..e18de276f 100644 --- a/local_plugins/agent_service/android/src/main/kotlin/com/yunqiinnovation/agent_service/AgentService.kt +++ b/local_plugins/agent_service/android/src/main/kotlin/com/yunqiinnovation/agent_service/AgentService.kt @@ -16,6 +16,9 @@ import java.util.Collections import kotlin.coroutines.CoroutineContext import com.yunqiinnovation.ble_service.BleService import android.media.MediaPlayer +import java.util.concurrent.atomic.AtomicBoolean +import kotlinx.coroutines.sync.Mutex +import kotlinx.coroutines.sync.withLock import com.deep_voice.speech.tts.TtsEvent import com.deep_voice.speech.tts.TtsEventListener import com.deep_voice.speech.tts.TtsEventType @@ -43,7 +46,7 @@ object AgentService : CoroutineScope { // 协程相关 private val job = SupervisorJob() override val coroutineContext: CoroutineContext - get() = Dispatchers.Main + job + get() = Dispatchers.IO + job // 上下文和监听器 private lateinit var context: Context @@ -73,16 +76,24 @@ object AgentService : CoroutineScope { // 音频播放器 private var audioPlayer: AudioPlayer? = null // 初始化音频播放器 - // 状态 - private var isInitialized = false - var isRecognitionActive = false - private set - var isTtsSpeaking = false - private set - var hasSpeechDetected = false - private set - var isAiStreaming = false - private set + // 状态 - 使用原子类型确保线程安全 + private val _isInitialized = AtomicBoolean(false) + val isInitialized: Boolean get() = _isInitialized.get() + + private val _isRecognitionActive = AtomicBoolean(false) + val isRecognitionActive: Boolean get() = _isRecognitionActive.get() + + private val _isTtsSpeaking = AtomicBoolean(false) + val isTtsSpeaking: Boolean get() = _isTtsSpeaking.get() + + private val _hasSpeechDetected = AtomicBoolean(false) + val hasSpeechDetected: Boolean get() = _hasSpeechDetected.get() + + private val _isAiStreaming = AtomicBoolean(false) + val isAiStreaming: Boolean get() = _isAiStreaming.get() + + // 用于保护复杂状态操作的互斥锁 + private val stateMutex = Mutex() // AI流生成相关 private var currentAiJob: Job? = null @@ -175,7 +186,7 @@ object AgentService : CoroutineScope { // 加载最近的聊天记录 loadChatHistory() - isInitialized = true + _isInitialized.set(true) return true } catch (e: Exception) { Log.e(TAG, "初始化失败: ${e.message}") @@ -201,6 +212,7 @@ object AgentService : CoroutineScope { // 释放ChatAPI服务 if (::chatApiService.isInitialized) { chatApiService.cancelCurrentStream() + chatApiService.dispose() } // 释放TTS服务 @@ -209,7 +221,7 @@ object AgentService : CoroutineScope { job.cancel() clearListeners() - isInitialized = false + _isInitialized.set(false) } catch (e: Exception) { Log.e(TAG, "释放资源异常: ${e.message}") } @@ -252,17 +264,17 @@ object AgentService : CoroutineScope { override fun onEvent(event: TtsEvent) { when (event.type) { TtsEventType.SYNTHESIS_STARTED -> { - isTtsSpeaking = true + _isTtsSpeaking.set(true) restartIdleCheck() sendEvent("tts_started", mapOf("status" to "started")) } TtsEventType.SYNTHESIS_COMPLETED -> { - isTtsSpeaking = false + _isTtsSpeaking.set(false) restartIdleCheck() sendEvent("tts_completed", mapOf("status" to "completed")) } TtsEventType.SYNTHESIS_CANCELED -> { - isTtsSpeaking = false + _isTtsSpeaking.set(false) restartIdleCheck() sendEvent("tts_canceled", mapOf("status" to "canceled")) } @@ -275,7 +287,7 @@ object AgentService : CoroutineScope { sendEvent("playback_completed", mapOf("status" to "playback_completed")) } TtsEventType.ERROR -> { - isTtsSpeaking = false + _isTtsSpeaking.set(false) restartIdleCheck() val params = event.params val code = params["errorCode"] as? String ?: "UNKNOWN_ERROR" @@ -379,8 +391,8 @@ object AgentService : CoroutineScope { return false } - isRecognitionActive = true - hasSpeechDetected = false + _isRecognitionActive.set(true) + _hasSpeechDetected.set(false) try { val audioSourceType = if (isExternalActive) { @@ -393,7 +405,7 @@ object AgentService : CoroutineScope { if (recognizing.isNotEmpty()) { // 检测到语音,更新状态 val previousHasSpeech = hasSpeechDetected - hasSpeechDetected = true + _hasSpeechDetected.set(true) if (!previousHasSpeech) { restartIdleCheck() @@ -408,21 +420,21 @@ object AgentService : CoroutineScope { if (isTtsSpeaking || isAiStreaming) { val currentTime = System.currentTimeMillis() - // // 防抖处理:避免过于频繁的打断 - // if (currentTime - lastInterruptTime < INTERRUPT_DEBOUNCE_MS) { - // return - // } + // 防抖处理:避免过于频繁的打断 + if (currentTime - lastInterruptTime < INTERRUPT_DEBOUNCE_MS) { + return + } - // // 简单过滤:太短的内容可能是噪音 - // if (recognizing.trim().length < 2) { - // return - // } + // 简单过滤:太短的内容可能是噪音 + if (recognizing.trim().length < 2) { + return + } - // // 过滤纯语气词 - // val trimmedText = recognizing.trim().lowercase() - // if (FILLER_WORDS.contains(trimmedText)) { - // return - // } + // 过滤纯语气词 + val trimmedText = recognizing.trim().lowercase() + if (FILLER_WORDS.contains(trimmedText)) { + return + } // 执行打断 lastInterruptTime = currentTime @@ -443,7 +455,7 @@ object AgentService : CoroutineScope { // 重置状态,继续识别 val previousHasSpeech = hasSpeechDetected - hasSpeechDetected = false + _hasSpeechDetected.set(false) if (previousHasSpeech) { restartIdleCheck() @@ -460,13 +472,13 @@ object AgentService : CoroutineScope { override fun onSessionStopped() { sendEvent("recognition_stopped", mapOf("status" to "stopped")) - isRecognitionActive = false + _isRecognitionActive.set(false) stopIdleCheck() audioPlayer?.playAudio(R.raw.stop) } override fun onCanceled(reason: String, errorDetails: String) { - isRecognitionActive = false + _isRecognitionActive.set(false) stopIdleCheck() BleService.closeCodec() sendEvent("recognition_canceled", mapOf( @@ -476,7 +488,7 @@ object AgentService : CoroutineScope { } override fun onError(error: String) { - isRecognitionActive = false + _isRecognitionActive.set(false) stopIdleCheck() BleService.closeCodec() Log.e(TAG, "语音识别错误: $error") @@ -488,7 +500,7 @@ object AgentService : CoroutineScope { }, audioSourceType) return true } catch (e: Exception) { - isRecognitionActive = false + _isRecognitionActive.set(false) Log.e(TAG, "启动语音识别失败: ${e.message}") sendEvent("error", mapOf( "code" to "RECOGNITION_START_ERROR", @@ -510,11 +522,11 @@ object AgentService : CoroutineScope { try { azureAsrHelper?.stopContinuousRecognition() BleService.closeCodec() - isRecognitionActive = false + _isRecognitionActive.set(false) stopIdleCheck() } catch (e: Exception) { Log.e(TAG, "停止语音识别异常: ${e.message}") - isRecognitionActive = false + _isRecognitionActive.set(false) stopIdleCheck() } } @@ -549,7 +561,7 @@ object AgentService : CoroutineScope { if (isAiStreaming) { try { // 先更新状态,避免回调时的状态不一致 - isAiStreaming = false + _isAiStreaming.set(false) // 取消当前AI生成任务 currentAiJob?.cancel() @@ -561,7 +573,7 @@ object AgentService : CoroutineScope { } catch (e: Exception) { Log.e(TAG, "停止AI流输出异常", e) // 确保状态被重置,即使发生异常 - isAiStreaming = false + _isAiStreaming.set(false) currentAiJob = null } } @@ -655,7 +667,7 @@ object AgentService : CoroutineScope { currentAiJob = launch { try { // 设置状态为正在流式输出 - isAiStreaming = true + _isAiStreaming.set(true) // 使用历史记录作为上下文发送到OpenAI val responseBuilder = StringBuilder() @@ -753,7 +765,7 @@ object AgentService : CoroutineScope { // 保存聊天记录 saveChatMessage(displayText, response,aiMetadata,userMetadata.toString()) // 标记AI流式输出已完成 - isAiStreaming = false + _isAiStreaming.set(false) currentAiJob = null } @@ -765,7 +777,7 @@ object AgentService : CoroutineScope { )) // 标记AI流式输出已完成 - isAiStreaming = false + _isAiStreaming.set(false) currentAiJob = null } @@ -809,7 +821,7 @@ object AgentService : CoroutineScope { )) // 确保状态被重置 - isAiStreaming = false + _isAiStreaming.set(false) currentAiJob = null } } @@ -873,7 +885,7 @@ object AgentService : CoroutineScope { if (text.isEmpty()) return // 更新状态 - isTtsSpeaking = true + _isTtsSpeaking.set(true) restartIdleCheck() // 状态变化,重启检测 // 直接调用TTS,无需协程包装 @@ -887,7 +899,7 @@ object AgentService : CoroutineScope { if (isTtsSpeaking) { ttsService?.stop() - isTtsSpeaking = false + _isTtsSpeaking.set(false) restartIdleCheck() // 状态变化,重启检测 sendEvent("tts_stopped", mapOf("status" to "stopped")) } @@ -964,9 +976,8 @@ object AgentService : CoroutineScope { * 发送事件 */ private fun sendEvent(eventName: String, data: Map) { - // 使用协程确保在主线程上执行 - launch { - // 我们已在主线程上下文中启动协程,无需再切换线程 + // 切换到主线程执行监听器回调,避免 UI 更新问题 + launch(Dispatchers.Main) { listeners.forEach { listener -> try { listener.onEvent(eventName, data) diff --git a/local_plugins/bytedance_speech/android/src/main/kotlin/com/deep_voice/bytedance_speech/BytedanceAudioPlayer.kt b/local_plugins/bytedance_speech/android/src/main/kotlin/com/deep_voice/bytedance_speech/BytedanceAudioPlayer.kt index ac7cf0afd..b3cffcbd1 100644 --- a/local_plugins/bytedance_speech/android/src/main/kotlin/com/deep_voice/bytedance_speech/BytedanceAudioPlayer.kt +++ b/local_plugins/bytedance_speech/android/src/main/kotlin/com/deep_voice/bytedance_speech/BytedanceAudioPlayer.kt @@ -24,6 +24,7 @@ class BytedanceAudioPlayer(private val context: Context) : AudioDataListener { private const val SAMPLE_RATE = 24000 private const val CHANNEL_CONFIG = AudioFormat.CHANNEL_OUT_MONO private const val AUDIO_FORMAT = AudioFormat.ENCODING_PCM_16BIT + private const val VOLUME_REDUCTION = 0.7f // 降低音量以减少回音 } private var audioTrack: AudioTrack? = null @@ -206,6 +207,7 @@ class BytedanceAudioPlayer(private val context: Context) : AudioDataListener { } } + /** * 初始化 AudioTrack */ @@ -213,13 +215,18 @@ class BytedanceAudioPlayer(private val context: Context) : AudioDataListener { audioTrack?.release() val minBufferSize = AudioTrack.getMinBufferSize(SAMPLE_RATE, CHANNEL_CONFIG, AUDIO_FORMAT) - val bufferSize = minBufferSize * 2 + val bufferSize = minBufferSize * 2 // 使用较小的缓冲区以降低延迟 audioTrack = if (Build.VERSION.SDK_INT >= Build.VERSION_CODES.M) { - AudioTrack.Builder() + val builder = AudioTrack.Builder() .setAudioAttributes( AudioAttributes.Builder() - .setUsage(AudioAttributes.USAGE_MEDIA) + .setUsage( + // 扬声器模式下默认使用语音通信模式以启用回声消除 + if (audioOutputDevice == AudioOutputDevice.SPEAKER) + AudioAttributes.USAGE_VOICE_COMMUNICATION + else AudioAttributes.USAGE_MEDIA + ) .setContentType(AudioAttributes.CONTENT_TYPE_SPEECH) .build() ) @@ -232,7 +239,13 @@ class BytedanceAudioPlayer(private val context: Context) : AudioDataListener { ) .setBufferSizeInBytes(bufferSize) .setTransferMode(AudioTrack.MODE_STREAM) - .build() + + // API 26+ 设置低延迟模式 + if (Build.VERSION.SDK_INT >= Build.VERSION_CODES.O) { + builder.setPerformanceMode(AudioTrack.PERFORMANCE_MODE_LOW_LATENCY) + } + + builder.build() } else { @Suppress("DEPRECATION") AudioTrack( @@ -247,6 +260,9 @@ class BytedanceAudioPlayer(private val context: Context) : AudioDataListener { audioTrack?.setPlaybackPositionUpdateListener(playbackListener) + // 设置音量以减少回音 + audioTrack?.setVolume(VOLUME_REDUCTION) + if (Build.VERSION.SDK_INT >= Build.VERSION_CODES.S) { setPreferredDeviceForTrack() } else { @@ -268,13 +284,17 @@ class BytedanceAudioPlayer(private val context: Context) : AudioDataListener { AudioOutputDevice.DEFAULT -> { track.setPreferredDevice(null) manager.isSpeakerphoneOn = false - manager.mode = AudioManager.MODE_NORMAL + // 默认模式下使用通信模式以获得更好的音频处理 + manager.mode = AudioManager.MODE_IN_COMMUNICATION } AudioOutputDevice.SPEAKER -> { val speaker = manager.getDevices(AudioManager.GET_DEVICES_OUTPUTS) .firstOrNull { it.type == AudioDeviceInfo.TYPE_BUILTIN_SPEAKER } speaker?.let { track.setPreferredDevice(it) } + // 扬声器模式下始终启用通信模式以获得回声消除 + manager.mode = AudioManager.MODE_IN_COMMUNICATION + manager.isSpeakerphoneOn = true } AudioOutputDevice.HEADPHONES -> { @@ -287,6 +307,9 @@ class BytedanceAudioPlayer(private val context: Context) : AudioDataListener { } headphones?.let { track.setPreferredDevice(it) } ?: track.setPreferredDevice(null) + // 耳机模式下可以使用普通模式 + manager.mode = AudioManager.MODE_NORMAL + manager.isSpeakerphoneOn = false } } } @@ -300,16 +323,19 @@ class BytedanceAudioPlayer(private val context: Context) : AudioDataListener { audioManager?.let { manager -> when (audioOutputDevice) { AudioOutputDevice.DEFAULT -> { - manager.mode = AudioManager.MODE_NORMAL + // 默认模式下使用通信模式以获得更好的音频处理 + manager.mode = AudioManager.MODE_IN_COMMUNICATION manager.isSpeakerphoneOn = false } AudioOutputDevice.SPEAKER -> { - manager.mode = AudioManager.MODE_NORMAL + // 扬声器模式下始终启用通信模式以获得回声消除 + manager.mode = AudioManager.MODE_IN_COMMUNICATION manager.isSpeakerphoneOn = true } AudioOutputDevice.HEADPHONES -> { + // 耳机模式下可以使用普通模式 manager.mode = AudioManager.MODE_NORMAL manager.isSpeakerphoneOn = false } diff --git a/local_plugins/chat_api/android/src/main/kotlin/com/yunqiinnovation/chat_api/ChatApiService.kt b/local_plugins/chat_api/android/src/main/kotlin/com/yunqiinnovation/chat_api/ChatApiService.kt index 55485623b..ee7fae90a 100644 --- a/local_plugins/chat_api/android/src/main/kotlin/com/yunqiinnovation/chat_api/ChatApiService.kt +++ b/local_plugins/chat_api/android/src/main/kotlin/com/yunqiinnovation/chat_api/ChatApiService.kt @@ -111,6 +111,10 @@ class ChatApiService(private val context: android.content.Context? = null) : Cor override val coroutineContext: CoroutineContext = Dispatchers.IO + SupervisorJob() + companion object { + private const val TAG = "ChatApiService" + } + // MARK: - 属性 private var baseUrl = "https://api.openai.com/v1/" private var apiKey = "" @@ -218,8 +222,8 @@ class ChatApiService(private val context: android.content.Context? = null) : Cor mcpConfigJson = mcpServer // 异步初始化MCP客户端 launch { - // initializeMcpClient(mcpServer) - initializeMcpClient("{}") + initializeMcpClient(mcpServer) + // initializeMcpClient("{}") } isInitialized = apiKey.isNotEmpty() @@ -514,8 +518,8 @@ class ChatApiService(private val context: android.content.Context? = null) : Cor ) // 通知上层工具调用事件 currSessionCallback?.onFunctionCall(convertMapToJsonObject(functionCall)) - // 在后台队列处理工具调用 - launch { + // 在当前协程作用域内处理工具调用,使用async确保生命周期管理 + val toolCallDeferred = async { try { if (sessionid == currSessionId) { // 通过MCP客户端处理工具调用 @@ -587,6 +591,15 @@ class ChatApiService(private val context: android.content.Context? = null) : Cor } } + // 等待工具调用完成,确保生命周期管理 + try { + toolCallDeferred.await() + } catch (e: CancellationException) { + // 协程被取消,确保子任务也被取消 + toolCallDeferred.cancel() + throw e + } + return true } @@ -780,7 +793,7 @@ class ChatApiService(private val context: android.content.Context? = null) : Cor description = description, parameters = parameters ) - Log.d("ChatApiService", "AI携带工具: $name, 参数定义: $parameters") + // Log.d("ChatApiService", "AI携带工具: $name, 参数定义: $parameters") tools.add(tool) } catch (e: Exception) { @@ -813,6 +826,30 @@ class ChatApiService(private val context: android.content.Context? = null) : Cor } } + /** + * 释放所有资源 + * 统一的资源管理方法,确保所有协程被正确取消 + */ + fun dispose() { + runBlocking { + // 取消当前流式任务 + currentStreamJob?.cancelAndJoin() + currentStreamJob = null + + // 关闭MCP客户端 + closeMcpClient() + + // 取消所有子协程 + coroutineContext[Job]?.cancelChildren() + + // 清理其他资源 + currSessionCallback = null + currSessionId = "" + currentMessages = emptyList() + toolCalls.clear() + } + } + /** * 处理MCP工具调用 * diff --git a/local_plugins/chat_api/android/src/main/kotlin/com/yunqiinnovation/chat_api/MCPSubClient.kt b/local_plugins/chat_api/android/src/main/kotlin/com/yunqiinnovation/chat_api/MCPSubClient.kt index 3fe6b926a..77393a44a 100644 --- a/local_plugins/chat_api/android/src/main/kotlin/com/yunqiinnovation/chat_api/MCPSubClient.kt +++ b/local_plugins/chat_api/android/src/main/kotlin/com/yunqiinnovation/chat_api/MCPSubClient.kt @@ -129,7 +129,7 @@ class MCPSubClient( "properties" to emptyMap(), "required" to emptyList() ) - Log.w(TAG, "解析工具数据: ${tool.name} ${parametersMap}") + // Log.w(TAG, "解析工具数据: ${tool.name} ${parametersMap}") val toolMap = mapOf( "type" to "function", "function" to mapOf( @@ -175,7 +175,7 @@ class MCPSubClient( name = name, arguments = argumentsJson ) - + Log.d(TAG, "[$serverId] 调用工具 '$name',参数: ${argumentsJson.toString()}") // 调用工具 val result = mcpClient?.callTool(request)