wolfplus2048 1 year ago
parent
commit
bcda9422e8
  1. 2
      lib/modules/agent/controllers/agent_controller.dart
  2. 113
      local_plugins/agent_service/android/src/main/kotlin/com/yunqiinnovation/agent_service/AgentService.kt
  3. 40
      local_plugins/bytedance_speech/android/src/main/kotlin/com/deep_voice/bytedance_speech/BytedanceAudioPlayer.kt
  4. 47
      local_plugins/chat_api/android/src/main/kotlin/com/yunqiinnovation/chat_api/ChatApiService.kt
  5. 4
      local_plugins/chat_api/android/src/main/kotlin/com/yunqiinnovation/chat_api/MCPSubClient.kt

2
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 ||

113
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<String, Any>) {
// 使用协程确保在主线程上执行
launch {
// 我们已在主线程上下文中启动协程,无需再切换线程
// 切换到主线程执行监听器回调,避免 UI 更新问题
launch(Dispatchers.Main) {
listeners.forEach { listener ->
try {
listener.onEvent(eventName, data)

40
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
}

47
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工具调用
*

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

@ -129,7 +129,7 @@ class MCPSubClient(
"properties" to emptyMap<String, Any>(),
"required" to emptyList<String>()
)
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)

Loading…
Cancel
Save