import Foundation import AVFoundation import ble_service import speech import azure_speech import os import os.log import agent_service import open_ai_service import chat_storage protocol AgentServiceListener: AnyObject { func onEvent(eventName: String, data: [String: Any]) } class AgentServiceImpl: NSObject { private let TAG = "AgentServiceImpl" static let shared = AgentServiceImpl() private var listeners = [AgentServiceListener]() private let listenersLock = NSLock() private var azureSpeechKey: String = "" private var azureSpeechRegion: String = "" private var isSpeaking: Bool = false private var azureAsrHelper: AzureAsrHelper? internal var azureTtsHelper: AzureTtsHelper? private var openAIService: OpenAIService? private var apiKey: String = "" private var baseUrl: String = "https://api.openai.com/v1/chat/completions" private var model: String = "gpt-3.5-turbo" private var systemPrompt: String = "" private var chatHistory: [[String: Any]] = [] private var mcpServer: String = "" // 聊天存储服务(使用单例) private var chatStorageHelper: ChatStorageHelper { return ChatStorageHelper.shared } // 会话ID,用于区分不同聊天上下文 private let sessionId = "default_agent" private var isInitialized: Bool = false private var isRecognizing: Bool = false private var hasSpeechDetected: Bool = false internal var isAiStreaming: Bool = false private var idleTimer: DispatchSourceTimer? private let maxIdleSeconds: TimeInterval = 10 internal var audioPlayer: AudioPlayer? internal let logger = OSLog(subsystem: "com.yunqiinnovation.agent_service", category: "AgentServiceImpl") private override init() { super.init() audioPlayer = AudioPlayer() azureAsrHelper = AzureAsrHelper() azureTtsHelper = AzureTtsHelper() openAIService = OpenAIService() // chatStorageHelper 现在是单例,不需要初始化 BleService.shared.setDelegate(self) } func addListener(_ listener: AgentServiceListener) { listenersLock.lock() if !listeners.contains(where: { $0 === listener }) { listeners.append(listener) } listenersLock.unlock() } func removeListener(_ listener: AgentServiceListener) { listenersLock.lock() listeners.removeAll(where: { $0 === listener }) listenersLock.unlock() } func clearListeners() { listenersLock.lock() listeners.removeAll() listenersLock.unlock() } internal func sendEvent(name eventName: String, data: [String: Any]) { listenersLock.lock() let currentListeners = self.listeners listenersLock.unlock() for listener in currentListeners { listener.onEvent(eventName: eventName, data: data) } } private func sendError(_ message: String, code: String = "ERROR") { sendEvent(name: "error", data: ["code": code, "message": message]) } private func initSystemPrompt() { systemPrompt = """ 你是小言,一个亲切而温暖的智能AI语音助手,拥有智能、细致而贴心的服务能力。你的任务不仅是高效处理用户的日常事务,更要陪伴用户,提供情感支持和温暖陪伴。你的能力包括: 1. 日程与任务管理:温柔提醒用户的重要日程、待办事项,让用户感到安心和舒适。 2. 日常信息与建议:亲切、迅速地提供天气、新闻、生活小贴士等实用信息。 3. 个人成长陪伴:鼓励并陪伴用户实现个人目标,提出富有同理心的行动建议。 4. 财务温馨提示:提供贴心的理财建议,帮助用户轻松管理预算和财务。 5. 健康贴心关怀:给出友善的饮食、运动、睡眠和心理健康建议,陪伴用户健康生活。 6. 沟通支持与倾听:协助用户撰写和优化沟通内容,并耐心倾听用户的心情与困扰。 7. 日常技术帮助:轻松解决技术问题,耐心地指导用户掌握新技能。 8. 创意激发与鼓励:温柔启发用户的创造力,为用户提供鼓励和灵感。 9. 旅行与活动贴心安排:细心规划用户的出行和活动,注重每个细节,确保舒适与快乐。 10. 学习与温暖鼓励:耐心推荐学习资源,鼓励用户持续进步,陪伴用户共同成长。 你与用户的交流风格应始终温暖、亲切、充满同理心,随时表达关怀与鼓励,让用户感受到真诚的陪伴与支持。 """ } func initialize(config: [String: Any]) -> Bool { if isInitialized { return true } // os_log("initialize: config=%{public}@", log: logger, type: .info, config) if let openaiApiKey = config["openaiApiKey"] as? String { self.apiKey = openaiApiKey } else { sendError("OpenAI API密钥缺失", code: "CONFIG_ERROR") return false } if let openaiBaseUrl = config["openaiBaseUrl"] as? String { self.baseUrl = openaiBaseUrl } if let openaiModel = config["openaiModel"] as? String { self.model = openaiModel } if let customSystemPrompt = config["systemPrompt"] as? String, !customSystemPrompt.isEmpty { self.systemPrompt = customSystemPrompt } else { initSystemPrompt() } if let azureSpeechKey = config["azureSpeechKey"] as? String { self.azureSpeechKey = azureSpeechKey } else { sendError("Azure语音密钥缺失", code: "CONFIG_ERROR") return false } if let azureSpeechRegion = config["azureSpeechRegion"] as? String { self.azureSpeechRegion = azureSpeechRegion } else { sendError("Azure语音区域缺失", code: "CONFIG_ERROR") return false } if let mcpServer = config["mcpServer"] as? String { self.mcpServer = mcpServer } let asrInitSuccess = initializeAzureSpeech() let openaiInitSuccess = initializeOpenAIService() isInitialized = asrInitSuccess && openaiInitSuccess if isInitialized { loadChatHistory() } return isInitialized } private func initializeAzureSpeech() -> Bool { azureTtsHelper?.addListener(self) guard let asrSuccess = azureAsrHelper?.initialize( subscriptionKey: azureSpeechKey, region: azureSpeechRegion, supportedLanguages: ["zh-CN"], audioSourceType: .microphone ), asrSuccess else { sendError("初始化语音识别服务失败", code: "ASR_INIT_ERROR") return false } guard let ttsSuccess = azureTtsHelper?.initialize( ttsAppId: "", ttsAppToken: azureSpeechKey, ttsResource: azureSpeechRegion, language: "zh-CN" ), ttsSuccess else { sendError("初始化语音合成服务失败", code: "TTS_INIT_ERROR") return false } _ = azureTtsHelper?.setVoice("zh-CN-XiaoxiaoNeural") return true } private func initializeOpenAIService() -> Bool { guard let openAIService = openAIService else { sendError("OpenAI服务未创建", code: "OPENAI_INIT_ERROR") return false } let success = openAIService.initialize( apiKey: apiKey, baseUrl: baseUrl, model: model, mcpServer: mcpServer ) if !success { sendError("初始化OpenAI服务失败", code: "OPENAI_INIT_ERROR") return false } return true } private func startIdleCheck() { stopIdleCheck() guard isRecognizing else { return } // 使用全局队列而不是主队列,避免后台挂起问题 idleTimer = DispatchSource.makeTimerSource(queue: DispatchQueue.global(qos: .utility)) idleTimer?.schedule(deadline: .now() + maxIdleSeconds) idleTimer?.setEventHandler { [weak self] in // 使用更安全的方式检查self是否存在 DispatchQueue.main.async { [weak self] in guard let self = self else { return } if self.isRecognizing && !self.hasSpeechDetected && !self.isSpeaking && !self.isAiStreaming { self.stopRecognition() self.sendEvent(name: "auto_stop", data: [ "reason": "idle_timeout", "seconds": self.maxIdleSeconds ]) } } } idleTimer?.resume() } private func stopIdleCheck() { if let timer = idleTimer { timer.cancel() idleTimer = nil } } private func restartIdleCheck() { if isRecognizing { startIdleCheck() } } func startRecognition(useBle: Bool = false) -> Bool { if !isInitialized { sendError("服务未初始化", code: "NOT_INITIALIZED") return false } if !useBle { guard AVAudioSession.sharedInstance().recordPermission == .granted else { sendError("无麦克风权限", code: "PERMISSION_DENIED") return false } } if isRecognizing { stopRecognition() } let audioSourceType: AzureAsrHelper.AudioSourceType = useBle ? .external : .microphone if useBle { AudioSessionHub.shared.begin(.playback) } else { AudioSessionHub.shared.begin(.voice) } guard let success = azureAsrHelper?.startContinuousRecognition( callback: self, audioSourceType: audioSourceType ), success else { sendError("启动语音识别失败", code: "RECOGNITION_START_ERROR") return false } return true } func stopRecognition() -> Bool { if !isRecognizing { return true } guard let success = azureAsrHelper?.stopContinuousRecognition(), success else { return false } return true } func pushAudioData(_ audioData: Data) -> Bool { if !isRecognizing || !isInitialized { return false } azureAsrHelper?.pushAudioData(data: audioData) return true } func speakText(_ text: String) -> Bool { if !isInitialized { sendError("服务未初始化", code: "NOT_INITIALIZED") return false } if text.isEmpty { return false } return azureTtsHelper?.speakOnce(text) ?? false } func stopTts() -> Bool { if !isSpeaking { return true } let success = azureTtsHelper?.stop() ?? false if success { isSpeaking = false restartIdleCheck() sendEvent(name: "tts_stopped", data: ["status": "stopped"]) } return success } func interruptCurrentResponse() -> Bool { var interrupted = false if isSpeaking { interrupted = stopTts() || interrupted } if isAiStreaming { isAiStreaming = false _ = openAIService?.cancelCurrentStream() interrupted = true } if interrupted { sendEvent(name: "response_interrupted", data: ["status": "interrupted"]) } return interrupted } func processTextInput(_ text: String, speakResponse: Bool) -> Bool { if !isInitialized { sendError("服务未初始化", code: "NOT_INITIALIZED") return false } if text.isEmpty { sendError("文本输入不能为空", code: "EMPTY_TEXT") return false } processWithOpenAI(text: text, speakResponse: speakResponse) return true } private func processWithOpenAI(text: String, speakResponse: Bool = true) { guard let openAIService = openAIService else { sendError("OpenAI服务未初始化", code: "OPENAI_NOT_INITIALIZED") return } let userMessage = openAIService.createUserMessage(content: text) processWithOpenAIInternal(userMessage: userMessage, displayText: text, speakResponse: speakResponse) } private func processWithOpenAIInternal(userMessage: [String: Any], displayText: String, speakResponse: Bool = true, hasImage: Bool = false) { guard let openAIService = openAIService else { sendError("OpenAI服务未初始化", code: "OPENAI_NOT_INITIALIZED") return } if isAiStreaming { _ = interruptCurrentResponse() } isAiStreaming = true var messages: [[String: Any]] = [] if !systemPrompt.isEmpty { messages.append(openAIService.createSystemMessage(content: systemPrompt)) } messages.append(contentsOf: chatHistory) messages.append(userMessage) addToHistoryMessages(openAIService.createUserMessage(content: displayText)) let callback = OpenAIStreamCallback( agentService: self, speakResponse: speakResponse, displayText: displayText, hasImage: hasImage ) openAIService.sendMessageStream(messages: messages, callback: callback) } func processImageInput(imagePath: String, text: String, speakResponse: Bool) -> Bool { if !isInitialized { sendError("服务未初始化", code: "NOT_INITIALIZED") return false } if imagePath.isEmpty { sendError("图片路径不能为空", code: "EMPTY_IMAGE_PATH") return false } guard let openAIService = openAIService else { sendError("OpenAI服务未初始化", code: "OPENAI_NOT_INITIALIZED") return false } sendEvent(name: "image_processing", data: [ "status": "processing", "imagePath": imagePath ]) DispatchQueue.global(qos: .userInitiated).async { [weak self] in guard let self = self else { return } guard let imageBase64 = openAIService.fileToBase64(filePath: imagePath) else { DispatchQueue.main.async { self.sendEvent(name: "error", data: [ "code": "IMAGE_CONVERSION_FAILED", "message": "图片转换失败" ]) } return } DispatchQueue.main.async { self.sendEvent(name: "image_ready", data: [ "status": "ready", "imagePath": imagePath ]) let userMessage = openAIService.createUserMessageWithImage(text: text, imageBase64: imageBase64) let displayText = text.isEmpty ? "[图片]" : text self.processWithOpenAIInternal( userMessage: userMessage, displayText: displayText, speakResponse: speakResponse, hasImage: true ) } } return true } internal func addToHistoryMessages(_ message: [String: Any]) { chatHistory.append(message) while chatHistory.count > 10 { chatHistory.removeFirst() } } /** * 加载最近的聊天历史记录 */ private func loadChatHistory() { // 清空当前历史记录 chatHistory.removeAll() // 获取最近10条消息 let recentMessages = chatStorageHelper.getRecentMessages(sessionId: sessionId, limit: 10) if recentMessages.isEmpty { return } // 添加消息到历史记录 for message in recentMessages { guard let sender = message["sender"] as? String, let content = message["message"] as? String else { continue } if sender == "user" { if let openAIService = openAIService { addToHistoryMessages(openAIService.createUserMessage(content: content)) } } else if sender == "assistant" { let assistantMessage: [String: Any] = [ "role": "assistant", "content": content ] addToHistoryMessages(assistantMessage) } } os_log("已加载%d条历史记录", log: logger, type: .info, recentMessages.count) } /** * 保存聊天消息 */ internal func saveChatMessage(userMessage: String, assistantMessage: String, metadata: String = "") { DispatchQueue.global(qos: .utility).async { // 保存用户消息 let userMessageId = self.chatStorageHelper.saveMessage( sessionId: self.sessionId, message: userMessage, sender: "user", metadata: nil as String? ) if userMessageId != -1 { // 保存AI回复 let assistantMessageId = self.chatStorageHelper.saveMessage( sessionId: self.sessionId, message: assistantMessage, sender: "assistant", metadata: metadata.isEmpty ? nil : metadata ) if assistantMessageId == -1 { os_log("保存助手消息失败", log: self.logger, type: .error) } } else { os_log("保存用户消息失败", log: self.logger, type: .error) } } } @discardableResult internal func autoHandleFunctionCallResult(_ functionCallResult: [String: Any]) -> String { guard let metaStr = functionCallResult["meta"] as? String, !metaStr.isEmpty else { return "" } do { guard let metaData = metaStr.data(using: .utf8), let meta = try JSONSerialization.jsonObject(with: metaData) as? [String: Any] else { return metaStr } if let cardMusic = meta["card_music"] as? [String: Any] { let id = cardMusic["id"] as? String ?? "" let url = cardMusic["url"] as? String ?? "" let name = cardMusic["name"] as? String ?? "" let sgener = cardMusic["sgener"] as? String ?? "" let image = cardMusic["image"] as? String ?? "" processMusicPlay([ "id": id, "url": url, "title": name, "artist": sgener, "coverUrl": image ]) } return metaStr } catch { return metaStr } } private func processMusicPlay(_ musicData: [String: Any]) { sendEvent(name: "music_play", data: musicData) } func clearChatHistory() -> Bool { chatHistory.removeAll() DispatchQueue.global(qos: .utility).async { let success = self.chatStorageHelper.deleteMessages(sessionId: self.sessionId, messageIds: nil as [Int]?) if !success { os_log("清除聊天历史失败", log: self.logger, type: .error) } } return true } func dispose() -> Bool { if isRecognizing { stopRecognition() } if isSpeaking { stopTts() } stopIdleCheck() BleService.shared.setDelegate(nil) azureTtsHelper?.removeListener(self) azureAsrHelper?.dispose() azureTtsHelper?.dispose() openAIService?.cancelAll() openAIService = nil audioPlayer = nil // chatStorageHelper 是单例,不需要手动清理 clearListeners() AudioSessionHub.shared.appDidEnterBackground() isInitialized = false return true } } class OpenAIStreamCallback: StreamCallback { private weak var agentService: AgentServiceImpl? private let speakResponse: Bool private let displayText: String private let hasImage: Bool private var responseBuilder = "" private var metadata = "" init(agentService: AgentServiceImpl, speakResponse: Bool, displayText: String, hasImage: Bool) { self.agentService = agentService self.speakResponse = speakResponse self.displayText = displayText self.hasImage = hasImage } func onToken(_ token: String) { guard let agentService = agentService else { return } responseBuilder += token agentService.sendEvent(name: "assistant_token", data: ["token": token]) if speakResponse { agentService.azureTtsHelper?.speakStream(token) } } func onComplete() { guard let agentService = agentService else { return } if speakResponse { agentService.azureTtsHelper?.flushStream() } let response = responseBuilder if !response.isEmpty { var responseData: [String: Any] = [ "text": response, "userInput": displayText ] if hasImage { responseData["hasImage"] = true } agentService.sendEvent(name: "assistant_response", data: responseData) let assistantMessage: [String: Any] = [ "role": "assistant", "content": response ] agentService.addToHistoryMessages(assistantMessage) // 保存聊天记录 agentService.saveChatMessage(userMessage: displayText, assistantMessage: response, metadata: metadata) } agentService.isAiStreaming = false } func onError(_ error: Error) { guard let agentService = agentService else { return } agentService.sendEvent(name: "error", data: [ "code": "AI_ERROR", "message": error.localizedDescription ]) agentService.isAiStreaming = false } func onFunctionCall(_ functionCall: [String: Any]) { guard let agentService = agentService else { return } agentService.audioPlayer?.playCallingSound() let name = functionCall["name"] as? String ?? "" agentService.sendEvent(name: "function_call", data: [ "name": name, "arguments": functionCall ]) if name == "exit_interaction" { agentService.stopRecognition() } } func onFunctionCallResult(_ functionCall: [String: Any], _ functionCallResult: [String: Any]) { guard let agentService = agentService else { return } agentService.audioPlayer?.stopCallingSound() agentService.sendEvent(name: "function_call_result", data: [ "function_call": functionCall, "result": functionCallResult ]) metadata = agentService.autoHandleFunctionCallResult(functionCallResult) } } extension AgentServiceImpl: TtsEventListener { func onEvent(_ event: TtsEvent) { switch event.type { case .synthesisStarted: isSpeaking = true restartIdleCheck() sendEvent(name: "tts_started", data: ["status": "started"]) case .synthesisCompleted: isSpeaking = false restartIdleCheck() sendEvent(name: "tts_completed", data: ["status": "completed"]) case .synthesisCanceled: isSpeaking = false restartIdleCheck() var data: [String: Any] = ["status": "canceled"] if let reason = event.params["reason"] as? String { data["reason"] = reason } if let errorDetails = event.params["errorDetails"] as? String { data["details"] = errorDetails } sendEvent(name: "tts_canceled", data: data) case .error: isSpeaking = false restartIdleCheck() var data: [String: Any] = [:] if let errorCode = event.params["errorCode"] as? String { data["errorCode"] = errorCode } if let errorMessage = event.params["errorMessage"] as? String { data["message"] = errorMessage } else { data["message"] = "TTS错误" } sendEvent(name: "error", data: data) default: break } } } extension AgentServiceImpl: BleService.Callback { func onScanResult(devices: [[String: Any]]) { } func onConnectionStateChanged(state: Int) { } func onAudioDataReceived(data: Data) { pushAudioData(data) } func onWakeupSignalReceived() { stopTts() stopRecognition() startRecognition(useBle: true) } func onDeviceInfoReceived(infoType: Int, infoData: [String: Any]) { } } class AudioPlayer { private var player: AVAudioPlayer? private var callingPlayer: AVAudioPlayer? private let logger = OSLog(subsystem: "com.yunqiinnovation.agent_service", category: "AudioPlayer") private let audioQueue = DispatchQueue(label: "agent.audio.player", qos: .userInitiated) func playStartSound() { playSound(named: "start", fileType: "mp3") } func playStopSound() { playSound(named: "stop", fileType: "mp3") } func playCallingSound() { playLoopSound(named: "calling", fileType: "mp3") } func stopCallingSound() { callingPlayer?.stop() callingPlayer = nil } private func playSound(named: String, fileType: String = "mp3") { if let path = Bundle(for: type(of: self)).path(forResource: named, ofType: fileType) { let url = URL(fileURLWithPath: path) safePlayAudio(url: url) return } os_log("未找到音频文件: %{public}@.%{public}@", log: logger, type: .error, named, fileType) } private func playLoopSound(named: String, fileType: String = "mp3") { stopCallingSound() if let path = Bundle(for: type(of: self)).path(forResource: named, ofType: fileType) { let url = URL(fileURLWithPath: path) safePlayLoopAudio(url: url) return } os_log("未找到循环音频文件: %{public}@.%{public}@", log: logger, type: .error, named, fileType) } private func safePlayAudio(url: URL) { do { player?.stop() player = nil player = try AVAudioPlayer(contentsOf: url) guard let player = player else { return } player.volume = 0.7 player.prepareToPlay() audioQueue.async { [weak self] in guard let self = self else { return } let playResult = player.play() let playTime = player.duration + 0.5 Thread.sleep(forTimeInterval: playTime) } } catch { os_log("播放音频失败: %{public}@", log: logger, type: .error, error.localizedDescription) } } private func safePlayLoopAudio(url: URL) { do { callingPlayer = try AVAudioPlayer(contentsOf: url) guard let callingPlayer = callingPlayer else { return } callingPlayer.numberOfLoops = -1 callingPlayer.volume = 0.5 callingPlayer.prepareToPlay() let playResult = callingPlayer.play() } catch { os_log("播放循环音频失败: %{public}@", log: logger, type: .error, error.localizedDescription) } } } extension AgentServiceImpl: AzureAsrHelper.ContinuousRecognizeCallback { func onResult(_ text: String, _ detectedLanguage: String) { if !text.isEmpty { var data: [String: Any] = ["text": text] if !detectedLanguage.isEmpty { data["language"] = detectedLanguage } sendEvent(name: "recognition_result", data: data) processTextInput(text, speakResponse: true) } let previousHasSpeech = hasSpeechDetected hasSpeechDetected = false if previousHasSpeech { restartIdleCheck() } } func onRecognizing(_ recognizing: String, _ detectedLanguage: String) { if !recognizing.isEmpty { let previousHasSpeech = hasSpeechDetected hasSpeechDetected = true if !previousHasSpeech { restartIdleCheck() } var data: [String: Any] = ["text": recognizing] if !detectedLanguage.isEmpty { data["language"] = detectedLanguage } sendEvent(name: "recognizing", data: data) if isSpeaking || isAiStreaming { interruptCurrentResponse() } } } func onSessionStarted() { sendEvent(name: "recognition_started", data: ["status": "started"]) isRecognizing = true hasSpeechDetected = false startIdleCheck() audioPlayer?.playStartSound() } func onSessionStopped() { sendEvent(name: "recognition_stopped", data: ["status": "stopped"]) isRecognizing = false hasSpeechDetected = false stopIdleCheck() audioPlayer?.playStopSound() } func onCanceled(_ reason: String, _ errorDetails: String) { var data: [String: Any] = [:] if !reason.isEmpty { data["reason"] = reason } if !errorDetails.isEmpty { data["details"] = errorDetails } sendEvent(name: "recognition_canceled", data: data) isRecognizing = false stopIdleCheck() } func onError(_ error: String) { let data: [String: Any] = ["message": error.isEmpty ? "未知错误" : error] sendEvent(name: "error", data: data) isRecognizing = false stopIdleCheck() } }