import Foundation import MicrosoftCognitiveServicesSpeech import open_ai_service import chat_storage /// 代理服务 - 处理语音识别、TTS和AI对话相关逻辑 /// /// 负责集成Azure语音服务、OpenAI服务和本地存储服务, /// 提供语音识别、语音合成、AI对话等功能 class AgentService { // 常量定义 private let TAG = "AgentService" // 上下文和监听器 private var listeners = NSHashTable.weakObjects() // 配置参数 private var azureSpeechKey: String = "" private var azureSpeechRegion: String = "" private var openaiApiKey: String = "" private var openaiBaseUrl: String = "" private var openaiModel: String = "gpt-3.5-turbo" private var mcpServer: String = "" // Azure服务 private var azureAsrHelper: AzureAsrHelper? private var azureTtsHelper: AzureTtsHelper? // OpenAI服务 private var openAIService: OpenAIServiceBridge? // 聊天存储服务 private var chatStorageHelper: ChatStorageHelperBridge? // 会话ID,用于区分不同聊天上下文 private let sessionId = "default_agent" // 历史聊天消息缓存 private var historyMessages = [[String: Any]]() // 系统提示词 private var systemPrompt = "" // 状态 private var isInitialized = false private var isRecognitionActive = false private var isTtsSpeaking = false private var hasSpeechDetected = false private var isAiStreaming = false // AI流生成相关 private var currentAiTask: DispatchWorkItem? // 空闲检测相关 private var idleCheckTimer: Timer? private let maxIdleSeconds: TimeInterval = 10 // 最大空闲秒数 /// 初始化 init() { // 创建OpenAI服务桥接器 openAIService = OpenAIServiceBridge() // 创建聊天存储桥接器 chatStorageHelper = ChatStorageHelperBridge() } /** * 初始化系统提示词 */ private func initSystemPrompt() { systemPrompt = """ 你是一个友好、专业的语音助手,名叫"小语"。你的目标是通过对话为用户提供帮助、解答问题和完成任务。 遵循以下指导原则: 1. 保持简短精炼的回答,因为用户是通过语音与你交流 2. 优先使用中文回复,除非用户明确要求使用其他语言 3. 当用户问题不明确时,礼貌地请求更多信息 4. 避免过长的列表,尽量将信息分成小段 5. 不要使用需要视觉展示的元素(如表格、图表或代码块) 6. 记住用户之前的对话内容,保持对话连贯 7. 如果用户发送了图片,请根据图片内容和文字要求回答问题 你不仅可以回答知识性问题,还可以帮助用户设置提醒、提供建议,或进行轻松愉快的对话。 无论遇到什么问题,都要尽力以温暖、贴心的语气提供最佳帮助。 """ } /** * 初始化 * * - Parameter config: 配置参数,包含所需的所有API密钥和设置 * - Returns: 是否初始化成功 */ func initialize(config: [String: Any]) -> Bool { if isInitialized { return true } do { // 从配置中获取参数 if let azureKey = config["azureSpeechKey"] as? String { azureSpeechKey = azureKey } if let azureRegion = config["azureSpeechRegion"] as? String { azureSpeechRegion = azureRegion } if let openaiKey = config["openaiApiKey"] as? String { openaiApiKey = openaiKey } if let baseUrl = config["openaiBaseUrl"] as? String { openaiBaseUrl = baseUrl } if let model = config["openaiModel"] as? String { openaiModel = model } if let server = config["mcpServer"] as? String { mcpServer = server } // 自定义系统提示词 if let customSystemPrompt = config["systemPrompt"] as? String, !customSystemPrompt.isEmpty { systemPrompt = customSystemPrompt } else { // 使用默认系统提示词 initSystemPrompt() } // 检查必要参数 if azureSpeechKey.isEmpty || azureSpeechRegion.isEmpty || openaiApiKey.isEmpty { print("\(TAG): 初始化失败:关键配置参数缺失") return false } // 初始化OpenAI服务 openAIService?.initialize(apiKey: openaiApiKey, baseUrl: openaiBaseUrl, model: openaiModel, mcpServer: mcpServer) // 初始化Azure ASR azureAsrHelper = AzureAsrHelper(self) let asrInitResult = azureAsrHelper?.initialize( subscriptionKey: azureSpeechKey, region: azureSpeechRegion, audioSourceType: .microphone ) ?? false if !asrInitResult { print("\(TAG): Azure ASR初始化失败") return false } // 初始化Azure TTS azureTtsHelper = AzureTtsHelper(self) let ttsInitResult = azureTtsHelper?.initialize( subscriptionKey: azureSpeechKey, region: azureSpeechRegion, callback: self ) ?? false if !ttsInitResult { print("\(TAG): Azure TTS初始化失败") return false } // 加载最近的聊天记录 loadChatHistory() isInitialized = true print("\(TAG): 代理服务初始化成功") return true } catch { print("\(TAG): 初始化失败:\(error.localizedDescription)") return false } } /** * 添加事件监听器 * * - Parameter listener: 要添加的监听器 */ func addListener(_ listener: AgentServiceListener) { listeners.add(listener as AnyObject) } /** * 设置事件监听器(替换现有监听器) * * - Parameter listener: 要设置的监听器 */ func setListener(_ listener: AgentServiceListener) { listeners.removeAllObjects() listeners.add(listener as AnyObject) } /** * 移除事件监听器 * * - Parameter listener: 要移除的监听器 */ func removeListener(_ listener: AgentServiceListener) { listeners.remove(listener as AnyObject) } /** * 移除所有事件监听器 */ func clearListeners() { listeners.removeAllObjects() } /** * 启动空闲检测 */ private func startIdleCheck() { stopIdleCheck() // 先停止现有的检查 if !isRecognitionActive { return } // 创建定时器 idleCheckTimer = Timer.scheduledTimer(withTimeInterval: maxIdleSeconds, repeats: false) { [weak self] _ in guard let self = self else { return } // 如果状态仍然是空闲的,则停止识别 if self.isRecognitionActive && !self.hasSpeechDetected && !self.isTtsSpeaking && !self.isAiStreaming { self.stopRecognition() self.sendEvent("auto_stop", ["reason": "idle_timeout", "seconds": self.maxIdleSeconds]) } } } /** * 停止空闲检测 */ private func stopIdleCheck() { idleCheckTimer?.invalidate() idleCheckTimer = nil } /** * 重启空闲检测 * 当状态发生变化时调用 */ private func restartIdleCheck() { if isRecognitionActive { startIdleCheck() } } /** * 开始语音识别 * * - Returns: 是否成功开始识别 */ func startRecognition() -> Bool { if !isInitialized { print("\(TAG): 服务未初始化") return false } if isRecognitionActive { return true } isRecognitionActive = true hasSpeechDetected = false do { // 创建连续识别回调 class ContinuousRecognizeCallbackImpl: AzureAsrHelper.ContinuousRecognizeCallback { private weak var service: AgentService? init(_ service: AgentService) { self.service = service } func onRecognizing(recognizing: String, detectedLanguage: String) { guard let service = service else { return } if !recognizing.isEmpty { // 检测到语音,更新状态 let previousHasSpeech = service.hasSpeechDetected service.hasSpeechDetected = true // 状态发生变化时重启空闲检测 if !previousHasSpeech { service.restartIdleCheck() } service.sendEvent("recognizing", [ "text": recognizing, "language": detectedLanguage ]) // 如果TTS正在播放或AI正在生成,则触发打断 if service.isTtsSpeaking || service.isAiStreaming { service.interruptCurrentResponse() } } } func onResult(text: String, detectedLanguage: String) { guard let service = service else { return } if !text.isEmpty { service.sendEvent("recognition_result", [ "text": text, "language": detectedLanguage ]) service.processWithOpenAI(text: text) } // 重置状态,继续识别 let previousHasSpeech = service.hasSpeechDetected service.hasSpeechDetected = false // 状态发生变化时重启空闲检测 if previousHasSpeech { service.restartIdleCheck() } } func onSessionStarted() { guard let service = service else { return } service.sendEvent("recognition_started", ["status": "started"]) // 启动空闲检测 service.startIdleCheck() } func onSessionStopped() { guard let service = service else { return } service.isRecognitionActive = false service.stopIdleCheck() service.sendEvent("recognition_stopped", ["status": "stopped"]) } func onCanceled(reason: String, errorDetails: String) { guard let service = service else { return } service.isRecognitionActive = false service.stopIdleCheck() service.sendEvent("recognition_canceled", [ "reason": reason, "details": errorDetails ]) } func onError(error: String) { guard let service = service else { return } service.isRecognitionActive = false service.stopIdleCheck() print("\(service.TAG): 语音识别出错: \(error)") service.sendEvent("error", [ "code": "RECOGNITION_ERROR", "message": error ]) } } let callback = ContinuousRecognizeCallbackImpl(self) let result = azureAsrHelper?.startContinuousRecognition(callback) ?? false if !result { isRecognitionActive = false print("\(TAG): 启动语音识别失败") sendEvent("error", [ "code": "RECOGNITION_START_ERROR", "message": "启动语音识别失败" ]) } return result } catch { isRecognitionActive = false print("\(TAG): 启动语音识别失败: \(error.localizedDescription)") sendEvent("error", [ "code": "RECOGNITION_START_ERROR", "message": error.localizedDescription ]) return false } } /** * 停止语音识别 */ func stopRecognition() { if !isRecognitionActive { return } print("\(TAG): 停止语音识别") // 停止识别 let _ = azureAsrHelper?.stopContinuousRecognition() isRecognitionActive = false stopIdleCheck() print("\(TAG): 语音识别已停止") } /** * 打断当前响应 * 停止TTS播放和AI流输出 */ func interruptCurrentResponse() { if isAiStreaming || isTtsSpeaking { // 停止TTS播放 stopTts() // 停止AI流输出 stopAiStream() // 发送打断事件 sendEvent("response_interrupted", ["status": "interrupted"]) } } /** * 停止AI流输出 */ private func stopAiStream() { if isAiStreaming { // 取消当前AI生成任务 currentAiTask?.cancel() currentAiTask = nil // 通知OpenAI服务终止当前流式请求 openAIService?.cancelCurrentStream() // 更新状态 isAiStreaming = false // 记录日志 print("\(TAG): AI流输出已停止") } } /** * 处理文本输入 * 作为语音输入的补充,直接处理文本并通过事件返回结果 * * - Parameters: * - text: 用户输入文本 * - speakResponse: 是否朗读回复,默认为false * - Returns: 是否成功开始处理 */ func processTextInput(text: String, speakResponse: Bool = false) -> Bool { if !isInitialized { print("\(TAG): 服务未初始化") sendEvent("error", ["code": "NOT_INITIALIZED", "message": "服务未初始化"]) return false } if text.isEmpty { print("\(TAG): 文本输入不能为空") sendEvent("error", ["code": "EMPTY_TEXT", "message": "文本输入不能为空"]) return false } // 使用OpenAI处理文本 processWithOpenAI(text: text, speakResponse: speakResponse) return true } /** * 使用OpenAI处理语音识别结果 * * - Parameter text: 用户输入文本 */ private func processWithOpenAI(text: String) { processWithOpenAI(text: text, speakResponse: true) } /** * 使用OpenAI处理文本消息 * * - Parameters: * - text: 用户输入文本 * - speakResponse: 是否使用TTS朗读回复 */ private func processWithOpenAI(text: String, speakResponse: Bool = true) { print("\(TAG): 用户问题: \(text)") // 创建用户文本消息并处理 guard let userMessage = openAIService?.createUserMessage(text: text) else { print("\(TAG): 创建用户消息失败") return } processWithOpenAIInternal(userMessage: userMessage, displayText: text, speakResponse: speakResponse) } /** * 使用OpenAI处理图片 * * - Parameters: * - imagePath: 图片文件路径 * - text: 可选的文本描述或问题 * - speakResponse: 是否朗读回复 * - Returns: 是否成功开始处理 */ func processImageInput(imagePath: String, text: String = "", speakResponse: Bool = false) -> Bool { if !isInitialized { print("\(TAG): 服务未初始化") sendEvent("error", ["code": "NOT_INITIALIZED", "message": "服务未初始化"]) return false } if imagePath.isEmpty { print("\(TAG): 图片路径不能为空") sendEvent("error", ["code": "EMPTY_IMAGE_PATH", "message": "图片路径不能为空"]) return false } // 通知开始处理图片 sendEvent("image_processing", [ "status": "processing", "imagePath": imagePath ]) // 异步处理图片 DispatchQueue.global(qos: .userInitiated).async { [weak self] in guard let self = self else { return } // 将图片转换为Base64格式 guard let imageBase64 = self.openAIService?.fileToBase64(filePath: imagePath) else { DispatchQueue.main.async { print("\(self.TAG): 图片转换失败: \(imagePath)") self.sendEvent("error", [ "code": "IMAGE_CONVERSION_FAILED", "message": "图片转换失败" ]) } return } // 通知图片准备完成 DispatchQueue.main.async { self.sendEvent("image_ready", [ "status": "ready", "imagePath": imagePath ]) // 处理包含图片的消息 self.processImageWithOpenAI(imageBase64: imageBase64, text: text, speakResponse: speakResponse) } } return true } /** * 使用OpenAI处理图片 * * - Parameters: * - imageBase64: Base64编码的图片数据 * - text: 可选的文本描述或问题 * - speakResponse: 是否朗读回复 */ private func processImageWithOpenAI(imageBase64: String, text: String = "", speakResponse: Bool = false) { print("\(TAG): 处理图片输入: \(text.isEmpty ? "无附加文本" : "附带文本: \(text)")") // 创建带图片的用户消息并处理 guard let userMessage = openAIService?.createUserMessageWithImage(text: text, imageBase64: imageBase64) else { print("\(TAG): 创建带图片的用户消息失败") return } // 图片描述用于存储 let displayText = text.isEmpty ? "[图片]" : text processWithOpenAIInternal(userMessage: userMessage, displayText: displayText, speakResponse: speakResponse, hasImage: true) } /** * 内部方法:通用的OpenAI处理逻辑 * * - Parameters: * - userMessage: 用户消息(可以是文本或图片格式) * - displayText: 用于显示和存储的文本 * - speakResponse: 是否朗读回复 * - hasImage: 是否包含图片 */ private func processWithOpenAIInternal(userMessage: [String: Any], displayText: String, speakResponse: Bool = true, hasImage: Bool = false) { // 如果有正在进行的AI流式输出,先停止它 stopAiStream() // 创建AI任务 let workItem = DispatchWorkItem { [weak self] in guard let self = self else { return } // 设置状态为正在流式输出 self.isAiStreaming = true // 使用历史记录作为上下文发送到OpenAI var responseBuilder = "" // 添加系统提示词到历史消息的副本中 var messagesWithSystemPrompt: [[String: Any]] = [] // 先添加系统提示词 if !self.systemPrompt.isEmpty { if let systemMessage = self.openAIService?.createSystemMessage(text: self.systemPrompt) { messagesWithSystemPrompt.append(systemMessage) } } // 再添加历史消息 messagesWithSystemPrompt.append(contentsOf: self.historyMessages) messagesWithSystemPrompt.append(userMessage) // 将用户消息添加到历史记录 self.addToHistoryMessages(userMessage) // 流式回调 class StreamCallbackBridge: NSObject, OpenAIStreamCallback { private weak var service: AgentService? private var responseBuilder: String private let speakResponse: Bool private let displayText: String private let hasImage: Bool init(_ service: AgentService, responseBuilder: String = "", speakResponse: Bool, displayText: String, hasImage: Bool) { self.service = service self.responseBuilder = responseBuilder self.speakResponse = speakResponse self.displayText = displayText self.hasImage = hasImage super.init() } func onToken(token: String) { guard let service = service else { return } responseBuilder.append(token) if speakResponse { _ = service.azureTtsHelper?.speakStream(token) } // 发送流式回复token service.sendEvent("assistant_token", ["token": token]) } func onComplete() { guard let service = service else { return } // 视情况决定是否朗读回复 if speakResponse { _ = service.azureTtsHelper?.flushStream() } if !responseBuilder.isEmpty { // 发送完整回复,包含是否有图片的标记 var responseData: [String: Any] = [ "text": responseBuilder, "userInput": displayText ] if hasImage { responseData["hasImage"] = true } service.sendEvent("assistant_response", responseData) // 添加AI回复到历史记录 if let assistantMessage = service.openAIService?.createAssistantMessage(text: responseBuilder) { service.addToHistoryMessages(assistantMessage) } // 保存聊天记录 service.saveChatMessage(userMessage: displayText, assistantMessage: responseBuilder) } // 标记AI流式输出已完成 service.isAiStreaming = false service.currentAiTask = nil } func onError(error: Error) { guard let service = service else { return } print("\(service.TAG): AI处理出错: \(error.localizedDescription)") service.sendEvent("error", [ "code": "AI_ERROR", "message": error.localizedDescription ]) // 标记AI流式输出已完成 service.isAiStreaming = false service.currentAiTask = nil } func onFunctionCall(call: [String: Any]) { guard let service = service else { return } if let name = call["name"] as? String { service.sendEvent("function_call", [ "name": name, "arguments": call ]) if name == "exit_interaction" { service.stopRecognition() } } } func onFunctionCallResult(functionCall: [String: Any], functionCallResult: [String: Any]) { guard let service = service else { return } service.sendEvent("function_call_result", [ "function_call": functionCall, "result": functionCallResult, ]) } } let callback = StreamCallbackBridge(self, responseBuilder: responseBuilder, speakResponse: speakResponse, displayText: displayText, hasImage: hasImage) self.openAIService?.sendMessageStream(messages: messagesWithSystemPrompt, callback: callback) } // 保存任务引用并在全局队列中执行 currentAiTask = workItem DispatchQueue.global(qos: .userInitiated).async(execute: workItem) } /** * 加载最近的聊天历史记录 */ private func loadChatHistory() { // 清空当前历史记录 historyMessages.removeAll() // 使用ChatStorageHelper直接获取最近消息 guard let recentMessages = chatStorageHelper?.getRecentMessages(sessionId: sessionId, count: 10) else { print("\(TAG): 没有找到历史记录") return } // 将消息添加到历史记录 for message in recentMessages { if let sender = message["sender"] as? String, let content = message["message"] as? String { if sender == "user" { if let userMessage = openAIService?.createUserMessage(text: content) { addToHistoryMessages(userMessage) } } else if sender == "assistant" { if let assistantMessage = openAIService?.createAssistantMessage(text: content) { addToHistoryMessages(assistantMessage) } } } } print("\(TAG): 已加载\(recentMessages.count)条历史记录") } /** * 添加消息到历史记录,保持最近10条 */ private func addToHistoryMessages(_ message: [String: Any]) { // 添加新消息 historyMessages.append(message) // 如果超过10条,删除最早的消息 while historyMessages.count > 10 { historyMessages.removeFirst() } } /** * TTS播放函数 * * - Parameter text: 要播放的文本 */ func speakText(text: String) { if text.isEmpty { return } // 更新状态 isTtsSpeaking = true restartIdleCheck() // 状态变化,重启检测 // 播放文本 _ = azureTtsHelper?.speakText(text) } /** * 停止TTS播放 */ func stopTts() { if isTtsSpeaking { _ = azureTtsHelper?.stopSpeaking() isTtsSpeaking = false restartIdleCheck() // 状态变化,重启检测 sendEvent("tts_stopped", ["status": "stopped"]) } } /** * 保存聊天消息 */ private func saveChatMessage(userMessage: String, assistantMessage: String) { DispatchQueue.global(qos: .background).async { [weak self] in guard let self = self else { return } // 保存用户消息 let userMessageId = self.chatStorageHelper?.saveMessage( sessionId: self.sessionId, message: userMessage, sender: "user" ) ?? -1 if userMessageId != -1 { // 保存AI回复 let assistantMessageId = self.chatStorageHelper?.saveMessage( sessionId: self.sessionId, message: assistantMessage, sender: "assistant" ) ?? -1 if assistantMessageId == -1 { print("\(self.TAG): 保存助手消息失败") } } else { print("\(self.TAG): 保存用户消息失败") } } } /** * 清除聊天历史 */ func clearChatHistory(completion: @escaping (Bool) -> Void) { DispatchQueue.global(qos: .background).async { [weak self] in guard let self = self else { DispatchQueue.main.async { completion(false) } return } // 清除指定会话的所有消息 let success = self.chatStorageHelper?.deleteMessages(sessionId: self.sessionId) ?? false if success { // 清空内存中的历史记录 self.historyMessages.removeAll() print("\(self.TAG): 聊天历史已清除") } else { print("\(self.TAG): 清除聊天历史失败") } DispatchQueue.main.async { completion(success) } } } /** * 发送事件 */ private func sendEvent(_ eventName: String, _ data: [String: Any]) { // 向所有监听器发送事件 let allListeners = listeners.allObjects for case let listener as AgentServiceListener in allListeners { listener.onEvent(eventName: eventName, data: data) } } /** * 释放资源 */ func dispose() { stopRecognition() stopTts() stopAiStream() stopIdleCheck() azureAsrHelper?.dispose() azureTtsHelper?.dispose() // 清除所有监听器 clearListeners() // 重置单例状态,以便下次使用时可以重新初始化 isInitialized = false print("\(TAG): 代理服务资源已释放") } } // MARK: - AzureTtsHelper.TtsCallback extension AgentService: AzureTtsHelper.TtsCallback { func onSynthesisStarted() { isTtsSpeaking = true // 状态变化,重置空闲检测 restartIdleCheck() sendEvent("tts_started", ["status": "started"]) } func onSynthesizing() { // 可以在此添加TTS合成中的处理逻辑 } func onSynthesisCompleted() { isTtsSpeaking = false // 状态变化,重启空闲检测 restartIdleCheck() sendEvent("tts_completed", ["status": "completed"]) } func onSynthesisCanceled() { isTtsSpeaking = false // 状态变化,重启空闲检测 restartIdleCheck() sendEvent("tts_canceled", ["status": "canceled"]) } } // MARK: - 外部插件桥接器 /// 桥接OpenAIService插件 class OpenAIServiceBridge { private let service = OpenAIService() /// 初始化服务 func initialize(apiKey: String, baseUrl: String, model: String, mcpServer: String) { service.initialize(apiKey: apiKey, baseUrl: baseUrl, model: model, mcpServer: mcpServer) } /// 创建系统消息 func createSystemMessage(text: String) -> [String: Any] { return service.createSystemMessage(text) } /// 创建用户消息 func createUserMessage(text: String) -> [String: Any] { return service.createUserMessage(text) } /// 创建助手消息 func createAssistantMessage(text: String) -> [String: Any] { return service.createAssistantMessage(text) } /// 创建带图片的用户消息 func createUserMessageWithImage(text: String, imageBase64: String) -> [String: Any]? { return service.createUserMessageWithImage(text: text, imageBase64: imageBase64) } /// 文件转Base64 func fileToBase64(filePath: String) -> String? { return service.fileToBase64(filePath: filePath) } /// 发送流式消息 func sendMessageStream(messages: [[String: Any]], callback: OpenAIStreamCallback) { service.sendMessageStream(messages: messages, callback: callback) } /// 取消当前流 func cancelCurrentStream() { service.cancelCurrentStream() } } /// OpenAI流式回调协议 @objc protocol OpenAIStreamCallback { func onToken(token: String) func onComplete() func onError(error: Error) func onFunctionCall(call: [String: Any]) func onFunctionCallResult(functionCall: [String: Any], functionCallResult: [String: Any]) } /// 桥接ChatStorageHelper插件 class ChatStorageHelperBridge { private let helper = ChatStorageHelper() /// 保存消息 func saveMessage(sessionId: String, message: String, sender: String) -> Int64 { return helper.saveMessage(sessionId: sessionId, message: message, sender: sender) } /// 获取最近消息 func getRecentMessages(sessionId: String, count: Int) -> [[String: Any]]? { return helper.getRecentMessages(sessionId: sessionId, count: count) } /// 删除消息 func deleteMessages(sessionId: String) -> Bool { return helper.deleteMessages(sessionId: sessionId) } }