You can not select more than 25 topics Topics must start with a letter or number, can include dashes ('-') and can be up to 35 characters long.

878 lines
26 KiB

import Foundation
import AVFoundation
import ble_service
import speech
import azure_speech
import os
import os.log
import agent_service
import open_ai_service
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 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()
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 = """
你是小言,一个具备专业能力的智能助手,需根据用户场景灵活切换回答模式,确保服务精准高效。
"""
}
func initialize(config: [String: Any]) -> Bool {
if isInitialized { return true }
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
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.main)
idleTimer?.schedule(deadline: .now() + maxIdleSeconds)
idleTimer?.setEventHandler { [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() {
guard let timer = idleTimer else { return }
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()
}
}
@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()
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
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 = ""
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.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
])
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()
}
}