18 changed files with 2236 additions and 453 deletions
@ -0,0 +1,514 @@ |
|||
import Foundation |
|||
import UIKit |
|||
import AVFoundation |
|||
|
|||
/** |
|||
* 语音交互服务,iOS原生实现 |
|||
* |
|||
* 负责: |
|||
* 1) 监听蓝牙耳机按键 |
|||
* 2) 处理语音识别 |
|||
* 3) 与火山AI服务交互 |
|||
* 4) 文本转语音播放 |
|||
*/ |
|||
@available(iOS 13.0, *) |
|||
class VoiceInteractionService: NSObject { |
|||
// 常量 |
|||
private let TAG = "VoiceInteractionService" |
|||
private let RECOGNITION_TIMEOUT: TimeInterval = 8.0 |
|||
|
|||
// 单例 |
|||
static let shared = VoiceInteractionService() |
|||
|
|||
// 服务状态 |
|||
private var isRunning = false |
|||
private var isActive = false |
|||
private var isRecognitionActive = false |
|||
private var isTimeoutPaused = false |
|||
private var hasSpeechDetected = false |
|||
private var isTtsSpeaking = false |
|||
|
|||
// 事件处理器 |
|||
private weak var eventHandler: VoiceInteractionEventHandler? |
|||
|
|||
// 按键处理 |
|||
private var lastKeyEventTime: TimeInterval = 0 |
|||
private var keyEventCount = 0 |
|||
|
|||
// 活动时间 |
|||
private var lastActivityTime: TimeInterval = 0 |
|||
|
|||
// 当前用户输入 |
|||
private var currentUserInput = "" |
|||
|
|||
// 服务组件 |
|||
private var azureAsrHelper: AzureAsrHelper? |
|||
private var azureTtsHelper: AzureTtsHelper? |
|||
private var volcanoAIService: VolcanoAIService? |
|||
|
|||
|
|||
private var silenceTimer: Timer? |
|||
|
|||
// 系统提示词 |
|||
private let systemPrompt = """ |
|||
你是一个智能语音助手,能够简洁明了地回答用户的问题。 |
|||
请保持回答简短、准确,避免过长的解释。 |
|||
如果用户的问题不清楚,请礼貌地请求澄清。 |
|||
不要使用复杂的术语,除非用户明确要求。 |
|||
用户用语音和你交互。 |
|||
""" |
|||
|
|||
// 初始化方法 |
|||
private override init() { |
|||
|
|||
|
|||
super.init() |
|||
NSLog("%@: VoiceInteractionService 初始化中", TAG) |
|||
} |
|||
|
|||
/** |
|||
* 启动服务 |
|||
* |
|||
* @param apiKey Azure语音服务API密钥 |
|||
* @param region Azure语音服务区域 |
|||
* @param volcanoKey 火山AI服务API密钥 |
|||
* @return 启动是否成功 |
|||
*/ |
|||
func start(azureKey: String, azureRegion: String, volcanoKey: String) -> Bool { |
|||
guard !isRunning else { |
|||
NSLog("%@: 服务已经在运行", TAG) |
|||
return true |
|||
} |
|||
|
|||
NSLog("%@: 启动服务中", TAG) |
|||
|
|||
// 重置状态 |
|||
resetState() |
|||
|
|||
// 初始化各个服务组件 |
|||
initServices(azureKey: azureKey, azureRegion: azureRegion, volcanoKey: volcanoKey) |
|||
|
|||
// 播放静音音频,确保音频会话活跃 |
|||
// startSilencePlayback() |
|||
|
|||
// 启动monitorService定时器 |
|||
silenceTimer = Timer.scheduledTimer(timeInterval: 1.0, target: self, selector: #selector(monitorService), userInfo: nil, repeats: true) |
|||
NSLog("%@: monitorService定时器已启动", TAG) |
|||
|
|||
// 标记服务为运行状态 |
|||
isRunning = true |
|||
|
|||
NSLog("%@: 服务启动完成,等待蓝牙按键事件", TAG) |
|||
return true |
|||
} |
|||
|
|||
/** |
|||
* 停止服务 |
|||
*/ |
|||
func stop() { |
|||
guard isRunning else { |
|||
NSLog("%@: 服务未运行", TAG) |
|||
return |
|||
} |
|||
|
|||
NSLog("%@: 停止服务中", TAG) |
|||
|
|||
// 停止monitorService定时器 |
|||
silenceTimer?.invalidate() |
|||
silenceTimer = nil |
|||
|
|||
// 停止语音识别 |
|||
if isRecognitionActive { |
|||
stopVoiceRecognition() |
|||
} |
|||
|
|||
// 停止TTS |
|||
stopCurrentTTS() |
|||
|
|||
// 释放服务实例 |
|||
azureAsrHelper?.dispose() |
|||
azureTtsHelper?.dispose() |
|||
|
|||
// 更新状态 |
|||
resetState() |
|||
isRunning = false |
|||
|
|||
NSLog("%@: 服务已停止", TAG) |
|||
} |
|||
|
|||
/** |
|||
* 获取服务运行状态 |
|||
*/ |
|||
func isServiceRunning() -> Bool { |
|||
return isRunning |
|||
} |
|||
|
|||
/** |
|||
* 重置状态 |
|||
*/ |
|||
private func resetState() { |
|||
isActive = false |
|||
isRecognitionActive = false |
|||
isTimeoutPaused = false |
|||
hasSpeechDetected = false |
|||
isTtsSpeaking = false |
|||
} |
|||
|
|||
/** |
|||
* 初始化服务 |
|||
*/ |
|||
private func initServices(azureKey: String, azureRegion: String, volcanoKey: String) { |
|||
NSLog("%@: 初始化服务组件", TAG) |
|||
|
|||
// 初始化Azure ASR |
|||
azureAsrHelper = AzureAsrHelper { [weak self] eventName, eventData in |
|||
// 处理ASR事件 |
|||
self?.handleAsrEvent(eventName: eventName, eventData: eventData) |
|||
} |
|||
|
|||
// 初始化Azure TTS |
|||
azureTtsHelper = AzureTtsHelper { [weak self] eventName, eventData in |
|||
// 处理TTS事件 |
|||
self?.handleTtsEvent(eventName: eventName, eventData: eventData) |
|||
} |
|||
|
|||
// 初始化火山AI服务 |
|||
volcanoAIService = VolcanoAIService() |
|||
|
|||
// 配置Azure ASR |
|||
if !azureKey.isEmpty && !azureRegion.isEmpty { |
|||
azureAsrHelper?.initialize(speechSubscriptionKey: azureKey, serviceRegion: azureRegion, supportedLanguages: ["zh-CN"]) |
|||
NSLog("%@: Azure ASR初始化完成", TAG) |
|||
} else { |
|||
NSLog("%@: Azure配置信息不完整,无法初始化Azure ASR", TAG) |
|||
} |
|||
|
|||
// 配置Azure TTS |
|||
if !azureKey.isEmpty && !azureRegion.isEmpty { |
|||
azureTtsHelper?.initialize(speechSubscriptionKey: azureKey, serviceRegion: azureRegion, language: "zh-CN") |
|||
NSLog("%@: Azure TTS初始化完成", TAG) |
|||
} else { |
|||
NSLog("%@: Azure配置信息不完整,无法初始化Azure TTS", TAG) |
|||
} |
|||
|
|||
// 配置火山AI |
|||
if !volcanoKey.isEmpty { |
|||
volcanoAIService?.initialize(apiKey: volcanoKey) |
|||
NSLog("%@: 火山AI服务初始化完成", TAG) |
|||
} else { |
|||
NSLog("%@: 火山AI配置信息不完整,无法初始化火山AI服务", TAG) |
|||
} |
|||
} |
|||
|
|||
/** |
|||
* 监控服务状态 |
|||
*/ |
|||
@objc private func monitorService() { |
|||
// 确保服务保持活跃状态 |
|||
if !isActive { |
|||
isActive = true |
|||
} |
|||
|
|||
// 检查语音识别状态 |
|||
if isRecognitionActive { |
|||
let currentTime = Date().timeIntervalSince1970 |
|||
let elapsedTime = currentTime - lastActivityTime |
|||
|
|||
// NSLog("%@: 已过去时间: %@, 是否检测到语音: %@, 是否TTS播放: %@", TAG, String(elapsedTime), String(hasSpeechDetected), String(isTtsSpeaking), String(elapsedTime)) |
|||
|
|||
// 如果超过超时时间没有检测到语音,且不在TTS播放中,暂停语音识别 |
|||
if !hasSpeechDetected && !isTtsSpeaking && elapsedTime >= RECOGNITION_TIMEOUT { |
|||
// NSLog("%@: 超过%@秒未检测到语音,停止识别", TAG, String(RECOGNITION_TIMEOUT)) |
|||
isTimeoutPaused = true |
|||
playNotification("没有听到您说话,已暂停对话。双击耳机按钮可重新开始。") |
|||
|
|||
// 停止语音识别并确保资源完全释放 |
|||
stopVoiceRecognition() |
|||
|
|||
// 清理识别状态 |
|||
isRecognitionActive = false |
|||
hasSpeechDetected = false |
|||
} |
|||
} |
|||
|
|||
} |
|||
|
|||
/** |
|||
* 处理ASR事件 |
|||
*/ |
|||
private func handleAsrEvent(eventName: String, eventData: [String: Any]) { |
|||
// NSLog("%@: 处理ASR事件: %@", TAG, eventName) |
|||
switch eventName { |
|||
case "recognizing": |
|||
if let text = eventData["text"] as? String, !text.isEmpty { |
|||
hasSpeechDetected = true |
|||
stopCurrentTTS() |
|||
updateLastActivityTime() |
|||
} |
|||
|
|||
case "result": |
|||
if let text = eventData["text"] as? String, !text.isEmpty { |
|||
updateLastActivityTime() |
|||
|
|||
processWithVolcanoAI(text) |
|||
} |
|||
|
|||
// 重置状态,继续识别 |
|||
hasSpeechDetected = false |
|||
|
|||
case "sessionStarted": |
|||
updateLastActivityTime() |
|||
|
|||
case "sessionStopped": |
|||
isRecognitionActive = false |
|||
|
|||
case "canceled": |
|||
isRecognitionActive = false |
|||
|
|||
case "error": |
|||
isRecognitionActive = false |
|||
playNotification("语音识别出错") |
|||
|
|||
default: |
|||
break |
|||
} |
|||
} |
|||
|
|||
/** |
|||
* 处理TTS事件 |
|||
*/ |
|||
private func handleTtsEvent(eventName: String, eventData: [String: Any]) { |
|||
// NSLog("%@: 处理TTS事件: %@", TAG, eventName) |
|||
switch eventName { |
|||
case "started": |
|||
isTtsSpeaking = true |
|||
|
|||
case "completed": |
|||
isTtsSpeaking = false |
|||
updateLastActivityTime() |
|||
|
|||
case "canceled": |
|||
isTtsSpeaking = false |
|||
|
|||
case "error": |
|||
isTtsSpeaking = false |
|||
NSLog("%@: TTS错误: %@", TAG, eventData["error"] as? String ?? "未知错误") |
|||
|
|||
default: |
|||
break |
|||
} |
|||
} |
|||
|
|||
/** |
|||
* 处理双击事件 |
|||
*/ |
|||
func handleMediaButtonAction() { |
|||
NSLog("%@: 处理媒体按钮事件", TAG) |
|||
|
|||
// 屏蔽短时间重复响应的问题 |
|||
let currentTime = Date().timeIntervalSince1970 |
|||
let timeDiff = currentTime - lastKeyEventTime |
|||
if timeDiff < 0.5 { |
|||
NSLog("%@: 短时间重复响应,忽略", TAG) |
|||
return |
|||
} |
|||
lastKeyEventTime = currentTime |
|||
|
|||
// 停止当前TTS播放 |
|||
stopCurrentTTS() |
|||
|
|||
// 播放提示音 |
|||
playPrompt("我在!") |
|||
|
|||
// 重置超时暂停标志 |
|||
isTimeoutPaused = false |
|||
|
|||
// 启动或重置语音识别 |
|||
if !isRecognitionActive { |
|||
NSLog("%@: 语音识别未激活,开始启动", TAG) |
|||
startVoiceRecognition() |
|||
} else { |
|||
NSLog("%@: 语音识别已激活,更新活动时间", TAG) |
|||
updateLastActivityTime() |
|||
hasSpeechDetected = false |
|||
} |
|||
} |
|||
|
|||
/** |
|||
* 开始语音识别 |
|||
*/ |
|||
private func startVoiceRecognition() { |
|||
if isRecognitionActive { return } |
|||
|
|||
// 通知Flutter语音识别已启动 |
|||
notifyVoiceRecognitionStarted() |
|||
|
|||
isActive = true |
|||
isRecognitionActive = true |
|||
hasSpeechDetected = false |
|||
updateLastActivityTime() |
|||
|
|||
if let asrHelper = azureAsrHelper { |
|||
let success = asrHelper.startContinuousRecognition() |
|||
if !success { |
|||
isRecognitionActive = false |
|||
NSLog("%@: 启动语音识别失败", TAG) |
|||
playNotification("启动语音识别失败") |
|||
} |
|||
} else { |
|||
isRecognitionActive = false |
|||
NSLog("%@: Azure ASR服务未初始化", TAG) |
|||
playNotification("语音识别服务未初始化") |
|||
} |
|||
} |
|||
|
|||
/** |
|||
* 停止语音识别 |
|||
*/ |
|||
private func stopVoiceRecognition() { |
|||
if !isRecognitionActive { return } |
|||
|
|||
NSLog("%@: 停止语音识别", TAG) |
|||
|
|||
if let asrHelper = azureAsrHelper { |
|||
let success = asrHelper.stopContinuousRecognition() |
|||
if !success { |
|||
NSLog("%@: 停止语音识别失败", TAG) |
|||
} |
|||
} |
|||
|
|||
isRecognitionActive = false |
|||
hasSpeechDetected = false |
|||
} |
|||
|
|||
/** |
|||
* 使用火山AI处理语音识别结果 |
|||
*/ |
|||
private func processWithVolcanoAI(_ text: String) { |
|||
// 保存当前用户输入 |
|||
currentUserInput = text |
|||
|
|||
// 在后台线程处理 |
|||
DispatchQueue.global(qos: .userInitiated).async { [weak self] in |
|||
guard let self = self, let volcanoAIService = self.volcanoAIService else { |
|||
return |
|||
} |
|||
|
|||
do { |
|||
// 创建消息 |
|||
let message = volcanoAIService.createUserMessage(content: text) |
|||
let messages = [message] |
|||
|
|||
// 发送请求 |
|||
let response = try volcanoAIService.sendMessage(messages: messages, systemPrompt: self.systemPrompt) |
|||
|
|||
// 播放AI回复 |
|||
self.speakAIResponse(response) |
|||
|
|||
// 同步聊天记录到Flutter端 |
|||
self.notifyChatHistoryUpdated(userMessage: text, assistantMessage: response) |
|||
} catch { |
|||
NSLog("%@: AI处理出错: %@", self.TAG, error.localizedDescription) |
|||
self.playNotification("AI处理出错") |
|||
} |
|||
} |
|||
} |
|||
|
|||
/** |
|||
* 播放AI回复 |
|||
*/ |
|||
private func speakAIResponse(_ text: String) { |
|||
isTtsSpeaking = true |
|||
azureTtsHelper?.speakText(text: text) |
|||
} |
|||
|
|||
/** |
|||
* 播放提示音 |
|||
*/ |
|||
private func playPrompt(_ message: String) { |
|||
isTtsSpeaking = true |
|||
azureTtsHelper?.speakText(text: message) |
|||
} |
|||
|
|||
/** |
|||
* 播放通知提示音 |
|||
*/ |
|||
private func playNotification(_ message: String) { |
|||
isTtsSpeaking = true |
|||
azureTtsHelper?.speakText(text: message) |
|||
} |
|||
|
|||
/** |
|||
* 停止当前TTS播放 |
|||
*/ |
|||
private func stopCurrentTTS() { |
|||
if isTtsSpeaking { |
|||
azureTtsHelper?.stopSpeaking() |
|||
isTtsSpeaking = false |
|||
} |
|||
} |
|||
|
|||
/** |
|||
* 更新最后活动时间 |
|||
*/ |
|||
private func updateLastActivityTime() { |
|||
lastActivityTime = Date().timeIntervalSince1970 |
|||
} |
|||
|
|||
/** |
|||
* 暂停语音交互 |
|||
*/ |
|||
func pauseVoiceInteraction() { |
|||
NSLog("%@: 暂停语音交互", TAG) |
|||
|
|||
// 停止当前语音播放 |
|||
stopCurrentTTS() |
|||
|
|||
// 停止语音识别 |
|||
if isRecognitionActive { |
|||
stopVoiceRecognition() |
|||
} |
|||
|
|||
// 设置状态为暂停,但服务保持运行 |
|||
isTimeoutPaused = true |
|||
} |
|||
|
|||
/** |
|||
* 通知语音识别已开始 |
|||
*/ |
|||
private func notifyVoiceRecognitionStarted() { |
|||
// 发送识别开始事件 |
|||
sendEvent(type: "recognition_started", data: [:]) |
|||
} |
|||
|
|||
/** |
|||
* 通知聊天历史更新 |
|||
*/ |
|||
private func notifyChatHistoryUpdated(userMessage: String, assistantMessage: String) { |
|||
// 发送聊天历史更新事件 |
|||
sendEvent(type: "chat_history_updated", data: [ |
|||
"agentId": "personal_assistant", |
|||
"userMessage": userMessage, |
|||
"assistantMessage": assistantMessage |
|||
]) |
|||
} |
|||
|
|||
/** |
|||
* 设置事件处理器 |
|||
*/ |
|||
func setEventHandler(_ handler: VoiceInteractionEventHandler) { |
|||
self.eventHandler = handler |
|||
} |
|||
|
|||
/** |
|||
* 发送事件到Flutter端 |
|||
*/ |
|||
private func sendEvent(type: String, data: [String: Any] = [:]) { |
|||
var eventData = data |
|||
eventData["type"] = type |
|||
eventData["timestamp"] = Int(Date().timeIntervalSince1970 * 1000) |
|||
|
|||
// 使用事件处理器发送事件 |
|||
eventHandler?.sendEvent(eventData) |
|||
} |
|||
} |
|||
@ -0,0 +1,391 @@ |
|||
import Foundation |
|||
|
|||
/** |
|||
* 火山AI服务的iOS原生实现 |
|||
* |
|||
* 参考Android端的VolcanoAIService实现,提供同步和异步的API调用方式 |
|||
*/ |
|||
class VolcanoAIService { |
|||
private let TAG = "VolcanoAIService" |
|||
private let baseUrl = "https://ark.cn-beijing.volces.com/api/v3" |
|||
private let chatEndpoint = "/chat/completions" |
|||
private let session: URLSession |
|||
|
|||
private var apiKey: String = "" |
|||
private var isInitialized = false |
|||
|
|||
init() { |
|||
let config = URLSessionConfiguration.default |
|||
config.timeoutIntervalForRequest = 30.0 |
|||
config.timeoutIntervalForResource = 30.0 |
|||
self.session = URLSession(configuration: config) |
|||
} |
|||
|
|||
/** |
|||
* 初始化火山AI服务 |
|||
* |
|||
* @param apiKey 火山AI API密钥 |
|||
* @return 初始化是否成功 |
|||
*/ |
|||
func initialize(apiKey: String) -> Bool { |
|||
self.apiKey = apiKey |
|||
isInitialized = !apiKey.isEmpty |
|||
|
|||
if !isInitialized { |
|||
NSLog("%@: 初始化失败:API key 不能为空", TAG) |
|||
} else { |
|||
NSLog("%@: 火山AI服务初始化成功", TAG) |
|||
} |
|||
|
|||
return isInitialized |
|||
} |
|||
|
|||
/** |
|||
* 生成个性化问候语 |
|||
* |
|||
* @param agentName 代理名称 |
|||
* @param systemPrompt 系统提示词 |
|||
* @param callback 回调函数,返回生成的问候语 |
|||
*/ |
|||
func generateGreeting(agentName: String, systemPrompt: String, callback: @escaping (String?, Error?) -> Void) { |
|||
let messages: [[String: Any]] = [ |
|||
["role": "system", "content": systemPrompt], |
|||
["role": "user", "content": "请用一句简短的话向我打个招呼,要符合你的身份和性格特点,不要超过18个字。"] |
|||
] |
|||
|
|||
let messagesData = try? JSONSerialization.data(withJSONObject: messages, options: []) |
|||
let messagesArray = try? JSONSerialization.jsonObject(with: messagesData!, options: []) as? [[String: Any]] |
|||
|
|||
var result = "" |
|||
|
|||
sendMessageStream(messages: messagesArray!, systemPrompt: systemPrompt, streamCallback: StreamCallback( |
|||
onToken: { token in |
|||
result.append(token) |
|||
}, |
|||
onComplete: { |
|||
callback(result, nil) |
|||
}, |
|||
onError: { error in |
|||
callback(nil, error) |
|||
} |
|||
)) |
|||
} |
|||
|
|||
/** |
|||
* 发送消息(非流式输出) |
|||
* |
|||
* @param messages 消息列表 |
|||
* @param systemPrompt 系统提示词 |
|||
* @return 返回AI的回复 |
|||
* @throws VolcanoAIError 如果API调用失败 |
|||
*/ |
|||
func sendMessage(messages: [[String: Any]], systemPrompt: String) throws -> String { |
|||
// 检查是否已初始化 |
|||
if !isInitialized || apiKey.isEmpty { |
|||
throw VolcanoAIError.serviceNotInitialized |
|||
} |
|||
|
|||
var fullMessages: [[String: Any]] = [ |
|||
["role": "system", "content": systemPrompt] |
|||
] |
|||
|
|||
fullMessages.append(contentsOf: messages) |
|||
|
|||
let requestBody: [String: Any] = [ |
|||
"model": "doubao-1-5-lite-32k-250115", |
|||
"messages": fullMessages, |
|||
"temperature": 0.7, |
|||
"max_tokens": 2000, |
|||
"stream": false |
|||
] |
|||
|
|||
guard let url = URL(string: "\(baseUrl)\(chatEndpoint)") else { |
|||
throw VolcanoAIError.invalidURL |
|||
} |
|||
|
|||
var request = URLRequest(url: url) |
|||
request.httpMethod = "POST" |
|||
request.addValue("application/json", forHTTPHeaderField: "Content-Type") |
|||
request.addValue("Bearer \(apiKey)", forHTTPHeaderField: "Authorization") |
|||
|
|||
do { |
|||
request.httpBody = try JSONSerialization.data(withJSONObject: requestBody, options: []) |
|||
} catch { |
|||
throw VolcanoAIError.invalidRequestBody |
|||
} |
|||
|
|||
let semaphore = DispatchSemaphore(value: 0) |
|||
var responseData: Data? |
|||
var responseError: Error? |
|||
|
|||
let task = session.dataTask(with: request) { data, response, error in |
|||
if let error = error { |
|||
responseError = VolcanoAIError.networkError(error.localizedDescription) |
|||
semaphore.signal() |
|||
return |
|||
} |
|||
|
|||
guard let httpResponse = response as? HTTPURLResponse else { |
|||
responseError = VolcanoAIError.invalidResponse |
|||
semaphore.signal() |
|||
return |
|||
} |
|||
|
|||
if !(200...299).contains(httpResponse.statusCode) { |
|||
var errorMessage = "Unknown error occurred" |
|||
if let data = data, let json = try? JSONSerialization.jsonObject(with: data) as? [String: Any], |
|||
let error = json["error"] as? [String: Any], |
|||
let message = error["message"] as? String { |
|||
errorMessage = message |
|||
} |
|||
responseError = VolcanoAIError.apiError(errorMessage) |
|||
semaphore.signal() |
|||
return |
|||
} |
|||
|
|||
responseData = data |
|||
semaphore.signal() |
|||
} |
|||
|
|||
task.resume() |
|||
_ = semaphore.wait(timeout: .distantFuture) |
|||
|
|||
if let error = responseError { |
|||
throw error |
|||
} |
|||
|
|||
guard let data = responseData else { |
|||
throw VolcanoAIError.emptyResponse |
|||
} |
|||
|
|||
do { |
|||
guard let json = try JSONSerialization.jsonObject(with: data) as? [String: Any], |
|||
let choices = json["choices"] as? [[String: Any]], |
|||
let firstChoice = choices.first, |
|||
let message = firstChoice["message"] as? [String: Any], |
|||
let content = message["content"] as? String else { |
|||
throw VolcanoAIError.invalidResponseFormat |
|||
} |
|||
|
|||
return content |
|||
} catch { |
|||
throw VolcanoAIError.invalidResponseFormat |
|||
} |
|||
} |
|||
|
|||
/** |
|||
* 发送消息(流式输出) |
|||
* |
|||
* @param messages 消息列表 |
|||
* @param systemPrompt 系统提示词 |
|||
* @param streamCallback 回调函数,用于接收流式输出的结果 |
|||
*/ |
|||
func sendMessageStream(messages: [[String: Any]], systemPrompt: String, streamCallback: StreamCallback) { |
|||
// 检查是否已初始化 |
|||
if !isInitialized || apiKey.isEmpty { |
|||
streamCallback.onError(VolcanoAIError.serviceNotInitialized) |
|||
return |
|||
} |
|||
|
|||
var fullMessages: [[String: Any]] = [ |
|||
["role": "system", "content": systemPrompt] |
|||
] |
|||
|
|||
fullMessages.append(contentsOf: messages) |
|||
|
|||
let requestBody: [String: Any] = [ |
|||
"model": "doubao-1-5-lite-32k-250115", |
|||
"messages": fullMessages, |
|||
"temperature": 0.7, |
|||
"max_tokens": 2000, |
|||
"stream": true |
|||
] |
|||
|
|||
guard let url = URL(string: "\(baseUrl)\(chatEndpoint)") else { |
|||
streamCallback.onError(VolcanoAIError.invalidURL) |
|||
return |
|||
} |
|||
|
|||
var request = URLRequest(url: url) |
|||
request.httpMethod = "POST" |
|||
request.addValue("application/json", forHTTPHeaderField: "Content-Type") |
|||
request.addValue("Bearer \(apiKey)", forHTTPHeaderField: "Authorization") |
|||
request.addValue("text/event-stream", forHTTPHeaderField: "Accept") |
|||
|
|||
do { |
|||
request.httpBody = try JSONSerialization.data(withJSONObject: requestBody, options: []) |
|||
} catch { |
|||
streamCallback.onError(VolcanoAIError.invalidRequestBody) |
|||
return |
|||
} |
|||
|
|||
let task = session.dataTask(with: request) { data, response, error in |
|||
if let error = error { |
|||
streamCallback.onError(VolcanoAIError.networkError(error.localizedDescription)) |
|||
return |
|||
} |
|||
|
|||
guard let httpResponse = response as? HTTPURLResponse else { |
|||
streamCallback.onError(VolcanoAIError.invalidResponse) |
|||
return |
|||
} |
|||
|
|||
if !(200...299).contains(httpResponse.statusCode) { |
|||
var errorMessage = "Unknown error occurred" |
|||
if let data = data, let json = try? JSONSerialization.jsonObject(with: data) as? [String: Any], |
|||
let error = json["error"] as? [String: Any], |
|||
let message = error["message"] as? String { |
|||
errorMessage = message |
|||
} |
|||
streamCallback.onError(VolcanoAIError.apiError(errorMessage)) |
|||
return |
|||
} |
|||
|
|||
guard let data = data else { |
|||
streamCallback.onError(VolcanoAIError.emptyResponse) |
|||
return |
|||
} |
|||
|
|||
// 处理SSE流数据 |
|||
let responseString = String(data: data, encoding: .utf8) ?? "" |
|||
let lines = responseString.components(separatedBy: "\n") |
|||
|
|||
for line in lines { |
|||
if line.isEmpty { continue } |
|||
|
|||
if line.hasPrefix("data: ") { |
|||
let data = String(line.dropFirst(6)) |
|||
if data == "[DONE]" { |
|||
streamCallback.onComplete() |
|||
break |
|||
} |
|||
|
|||
do { |
|||
if let jsonData = data.data(using: .utf8), |
|||
let json = try JSONSerialization.jsonObject(with: jsonData) as? [String: Any], |
|||
let choices = json["choices"] as? [[String: Any]], |
|||
let firstChoice = choices.first, |
|||
let delta = firstChoice["delta"] as? [String: Any], |
|||
let content = delta["content"] as? String { |
|||
streamCallback.onToken(content) |
|||
} |
|||
} catch { |
|||
// 忽略无效的JSON数据 |
|||
continue |
|||
} |
|||
} |
|||
} |
|||
} |
|||
|
|||
task.resume() |
|||
} |
|||
|
|||
/** |
|||
* 同步方式发送消息(流式输出) |
|||
* |
|||
* 注意:此方法会阻塞当前线程,请在后台线程中调用 |
|||
* |
|||
* @param messages 消息列表 |
|||
* @param systemPrompt 系统提示词 |
|||
* @return 返回完整的AI回复 |
|||
* @throws VolcanoAIError 如果API调用失败 |
|||
*/ |
|||
func sendMessageStreamSync(messages: [[String: Any]], systemPrompt: String) throws -> String { |
|||
var result = "" |
|||
let semaphore = DispatchSemaphore(value: 0) |
|||
var responseError: Error? |
|||
|
|||
sendMessageStream(messages: messages, systemPrompt: systemPrompt, streamCallback: StreamCallback( |
|||
onToken: { token in |
|||
result.append(token) |
|||
}, |
|||
onComplete: { |
|||
semaphore.signal() |
|||
}, |
|||
onError: { error in |
|||
responseError = error |
|||
semaphore.signal() |
|||
} |
|||
)) |
|||
|
|||
// 等待流式输出完成或出错 |
|||
_ = semaphore.wait(timeout: .now() + 60) |
|||
|
|||
if let error = responseError { |
|||
throw error |
|||
} |
|||
|
|||
return result |
|||
} |
|||
|
|||
/** |
|||
* 创建用户消息 |
|||
*/ |
|||
func createUserMessage(content: String) -> [String: Any] { |
|||
return ["role": "user", "content": content] |
|||
} |
|||
|
|||
/** |
|||
* 创建系统消息 |
|||
*/ |
|||
func createSystemMessage(content: String) -> [String: Any] { |
|||
return ["role": "system", "content": content] |
|||
} |
|||
|
|||
/** |
|||
* 创建助手消息 |
|||
*/ |
|||
func createAssistantMessage(content: String) -> [String: Any] { |
|||
return ["role": "assistant", "content": content] |
|||
} |
|||
|
|||
/** |
|||
* 流式输出回调类 |
|||
*/ |
|||
class StreamCallback { |
|||
let onToken: (String) -> Void |
|||
let onComplete: () -> Void |
|||
let onError: (Error) -> Void |
|||
|
|||
init(onToken: @escaping (String) -> Void, onComplete: @escaping () -> Void, onError: @escaping (Error) -> Void) { |
|||
self.onToken = onToken |
|||
self.onComplete = onComplete |
|||
self.onError = onError |
|||
} |
|||
} |
|||
} |
|||
|
|||
/** |
|||
* 火山AI错误枚举 |
|||
*/ |
|||
enum VolcanoAIError: Error { |
|||
case serviceNotInitialized |
|||
case invalidURL |
|||
case invalidRequestBody |
|||
case networkError(String) |
|||
case invalidResponse |
|||
case emptyResponse |
|||
case invalidResponseFormat |
|||
case apiError(String) |
|||
|
|||
var localizedDescription: String { |
|||
switch self { |
|||
case .serviceNotInitialized: |
|||
return "火山AI服务未初始化或API key为空,请先调用initialize方法" |
|||
case .invalidURL: |
|||
return "无效的URL" |
|||
case .invalidRequestBody: |
|||
return "无效的请求体" |
|||
case .networkError(let message): |
|||
return "网络错误: \(message)" |
|||
case .invalidResponse: |
|||
return "无效的响应" |
|||
case .emptyResponse: |
|||
return "空响应" |
|||
case .invalidResponseFormat: |
|||
return "无效的响应格式" |
|||
case .apiError(let message): |
|||
return "API错误: \(message)" |
|||
} |
|||
} |
|||
} |
|||
@ -1,33 +1,275 @@ |
|||
import 'dart:async'; |
|||
import 'package:flutter/services.dart'; |
|||
import 'package:get/get.dart'; |
|||
import 'package:flutter_dotenv/flutter_dotenv.dart'; |
|||
import '../models/events/voice_interaction_event.dart'; |
|||
import '../../core/utils/logger.dart'; |
|||
import '../../modules/chat/models/message_model.dart'; |
|||
import 'chat_history_service.dart'; |
|||
|
|||
/// 语音交互服务抽象基类 |
|||
/// 语音交互服务接口 |
|||
/// |
|||
/// 定义了iOS和Android平台共用的语音交互服务接口 |
|||
abstract class VoiceInteractionService { |
|||
/// 事件流,用于发布语音交互相关事件 |
|||
Stream<VoiceInteractionEvent> get eventStream; |
|||
/// 管理与平台原生语音交互服务的通信,提供统一的接口供应用使用 |
|||
class VoiceInteractionService extends GetxService { |
|||
static VoiceInteractionService get to => Get.find<VoiceInteractionService>(); |
|||
|
|||
// 方法通道 |
|||
static const MethodChannel _channel = MethodChannel('com.deep_voice.voice_interaction'); |
|||
|
|||
// 事件通道 |
|||
static const EventChannel _eventChannel = EventChannel('com.deep_voice.voice_interaction_events'); |
|||
|
|||
// 服务状态 |
|||
final _isServiceRunning = false.obs; |
|||
bool get isServiceRunning => _isServiceRunning.value; |
|||
|
|||
// 流控制器 |
|||
final _eventStreamController = StreamController<VoiceInteractionEvent>.broadcast(); |
|||
|
|||
// 事件流 |
|||
Stream<VoiceInteractionEvent> get eventStream => _eventStreamController.stream; |
|||
|
|||
// 事件通道订阅 |
|||
StreamSubscription? _eventSubscription; |
|||
|
|||
// 标记是否初始化 |
|||
bool _isInitialized = false; |
|||
|
|||
// 配置信息 |
|||
late String _azureSpeechKey; |
|||
late String _azureSpeechRegion; |
|||
late String _volcanoAiKey; |
|||
|
|||
// 聊天历史服务 |
|||
late final ChatHistoryService _chatHistoryService; |
|||
|
|||
|
|||
/// 设置事件通道 |
|||
void _setupEventChannel() { |
|||
_eventSubscription = _eventChannel |
|||
.receiveBroadcastStream() |
|||
.listen(_handleVoiceInteractionEvent, onError: (error) { |
|||
Logger.error('语音交互事件通道错误: $error'); |
|||
}); |
|||
} |
|||
|
|||
/// 从环境变量加载配置 |
|||
void _loadConfig() { |
|||
_azureSpeechKey = dotenv.env['AZURE_SPEECH_KEY'] ?? ''; |
|||
_azureSpeechRegion = dotenv.env['AZURE_SPEECH_REGION'] ?? ''; |
|||
_volcanoAiKey = dotenv.env['VOLCANO_AI_API_KEY'] ?? ''; |
|||
|
|||
if (_azureSpeechKey.isEmpty || _azureSpeechRegion.isEmpty) { |
|||
Logger.warning('未找到 Azure 语音服务配置。请在 .env 文件中设置 AZURE_SPEECH_KEY 和 AZURE_SPEECH_REGION'); |
|||
} |
|||
|
|||
if (_volcanoAiKey.isEmpty) { |
|||
Logger.warning('未找到火山 AI API 密钥。请在 .env 文件中设置 VOLCANO_AI_API_KEY'); |
|||
} |
|||
} |
|||
|
|||
/// 处理来自原生层的事件 |
|||
void _handleVoiceInteractionEvent(dynamic event) { |
|||
if (event is! Map) return; |
|||
|
|||
final eventMap = event as Map<dynamic, dynamic>; |
|||
final String eventType = eventMap['type'] as String? ?? ''; |
|||
// final int timestamp = eventMap['timestamp'] as int? ?? DateTime.now().millisecondsSinceEpoch; |
|||
|
|||
Logger.info('收到语音交互事件: $eventType'); |
|||
|
|||
switch (eventType) { |
|||
case 'recognition_started': |
|||
// 语音识别开始事件 |
|||
final recognitionEvent = RecognitionStartedEvent( |
|||
timestamp: DateTime.now().millisecondsSinceEpoch, |
|||
); |
|||
_eventStreamController.add(recognitionEvent); |
|||
break; |
|||
|
|||
case 'chat_history_updated': |
|||
// 聊天历史更新事件 |
|||
final String agentId = eventMap['agentId'] as String? ?? ''; |
|||
final String userMessage = eventMap['userMessage'] as String? ?? ''; |
|||
final String assistantMessage = eventMap['assistantMessage'] as String? ?? ''; |
|||
|
|||
final chatHistoryEvent = ChatHistoryEvent( |
|||
agentId: agentId, |
|||
userMessage: userMessage, |
|||
assistantMessage: assistantMessage, |
|||
timestamp: DateTime.now().millisecondsSinceEpoch, |
|||
); |
|||
_eventStreamController.add(chatHistoryEvent); |
|||
|
|||
// 保存聊天历史到ChatHistoryService |
|||
_saveChatHistory(agentId, userMessage, assistantMessage, DateTime.now().millisecondsSinceEpoch); |
|||
break; |
|||
} |
|||
} |
|||
|
|||
/// 保存聊天历史 |
|||
void _saveChatHistory(String agentId, String userMessage, String assistantMessage, int timestamp) { |
|||
try { |
|||
// 检查参数有效性 |
|||
if (userMessage.isEmpty) { |
|||
return; |
|||
} |
|||
|
|||
// 创建用户消息和助手消息 |
|||
final userMsg = Message( |
|||
role: 'user', |
|||
content: userMessage, |
|||
timestamp: DateTime.fromMillisecondsSinceEpoch(timestamp), |
|||
); |
|||
|
|||
final assistantMsg = Message( |
|||
role: 'assistant', |
|||
content: assistantMessage, |
|||
timestamp: DateTime.fromMillisecondsSinceEpoch(timestamp + 1), // 确保助手消息时间戳晚于用户消息 |
|||
); |
|||
|
|||
// 加载现有历史记录 |
|||
final existingMessages = _chatHistoryService.loadHistory(agentId); |
|||
|
|||
// 添加新消息 |
|||
existingMessages.addAll([userMsg, assistantMsg]); |
|||
|
|||
// 保存更新后的历史记录 |
|||
_chatHistoryService.saveHistory(agentId, existingMessages); |
|||
|
|||
Logger.info('已保存聊天历史: agentId=$agentId'); |
|||
} catch (e) { |
|||
Logger.error('保存聊天历史失败: $e'); |
|||
} |
|||
} |
|||
|
|||
/// 初始化服务 |
|||
/// |
|||
/// 返回服务实例自身,以支持链式调用 |
|||
Future<VoiceInteractionService> initialize(); |
|||
Future<bool> initialize() async { |
|||
if (_isInitialized) return true; |
|||
|
|||
try { |
|||
Logger.info('正在初始化语音交互服务...'); |
|||
|
|||
// 获取聊天历史服务 |
|||
try { |
|||
_chatHistoryService = Get.find<ChatHistoryService>(); |
|||
} catch (e) { |
|||
Logger.warning('获取ChatHistoryService失败,将创建新实例'); |
|||
_chatHistoryService = Get.put(ChatHistoryService()); |
|||
} |
|||
|
|||
// 设置事件通道 |
|||
_setupEventChannel(); |
|||
|
|||
// 加载配置 |
|||
_loadConfig(); |
|||
|
|||
// 检查服务是否已在运行 |
|||
final bool running = await checkServiceStatus(); |
|||
|
|||
// 如果服务未运行,启动服务 |
|||
if (!running) { |
|||
// 启动语音交互服务 |
|||
final success = await startService(); |
|||
if (!success) { |
|||
Logger.error('语音交互服务启动失败'); |
|||
return false; |
|||
} |
|||
} |
|||
|
|||
_isInitialized = true; |
|||
Logger.info('语音交互服务初始化完成'); |
|||
return true; |
|||
} catch (e) { |
|||
Logger.error('语音交互服务初始化失败: $e'); |
|||
return false; |
|||
} |
|||
} |
|||
|
|||
/// 启动语音交互服务 |
|||
Future<bool> startService() async { |
|||
try { |
|||
Logger.info('启动语音交互服务...'); |
|||
|
|||
// 使用已加载的配置信息 |
|||
final result = await _channel.invokeMethod<bool>('startService', { |
|||
'azureSpeechKey': _azureSpeechKey, |
|||
'azureSpeechRegion': _azureSpeechRegion, |
|||
'volcanoAiKey': _volcanoAiKey, |
|||
}) ?? false; |
|||
|
|||
if (result) { |
|||
_isServiceRunning.value = true; |
|||
Logger.info('语音交互服务已启动'); |
|||
} else { |
|||
Logger.error('启动语音交互服务失败'); |
|||
} |
|||
|
|||
return result; |
|||
} catch (e) { |
|||
Logger.error('启动语音交互服务时发生错误: $e'); |
|||
return false; |
|||
} |
|||
} |
|||
|
|||
/// 停止语音交互服务 |
|||
Future<bool> stopService() async { |
|||
try { |
|||
Logger.info('停止语音交互服务...'); |
|||
|
|||
final result = await _channel.invokeMethod<bool>('stopService') ?? false; |
|||
|
|||
if (result) { |
|||
_isServiceRunning.value = false; |
|||
Logger.info('语音交互服务已停止'); |
|||
} else { |
|||
Logger.error('停止语音交互服务失败'); |
|||
} |
|||
|
|||
return result; |
|||
} catch (e) { |
|||
Logger.error('停止语音交互服务时发生错误: $e'); |
|||
return false; |
|||
} |
|||
} |
|||
|
|||
/// 检查服务是否运行 |
|||
Future<bool> checkServiceStatus() async { |
|||
try { |
|||
final bool result = await _channel.invokeMethod<bool>('isServiceRunning') ?? false; |
|||
_isServiceRunning.value = result; |
|||
return result; |
|||
} catch (e) { |
|||
Logger.error('检查服务状态时发生错误: $e'); |
|||
return false; |
|||
} |
|||
} |
|||
|
|||
/// 暂停语音交互 |
|||
/// |
|||
/// 停止正在进行的语音识别和TTS播放 |
|||
void pauseVoiceInteraction(); |
|||
Future<bool> pauseVoiceInteraction() async { |
|||
try { |
|||
Logger.info('暂停语音交互...'); |
|||
|
|||
/// 唤醒语音交互 |
|||
void onWakeup(); |
|||
final result = await _channel.invokeMethod<bool>('pauseVoiceInteraction') ?? false; |
|||
|
|||
/// 是否处于活跃状态 |
|||
bool isActive(); |
|||
if (result) { |
|||
Logger.info('语音交互已暂停'); |
|||
} else { |
|||
Logger.error('暂停语音交互失败'); |
|||
} |
|||
|
|||
/// 获取当前对话历史 |
|||
List<Map<String, String>> getMessageHistory(); |
|||
return result; |
|||
} catch (e) { |
|||
Logger.error('暂停语音交互时发生错误: $e'); |
|||
return false; |
|||
} |
|||
} |
|||
|
|||
/// 获取当前服务ID |
|||
String getServiceId(); |
|||
/// 资源释放 |
|||
@override |
|||
void onClose() { |
|||
_eventSubscription?.cancel(); |
|||
_eventStreamController.close(); |
|||
super.onClose(); |
|||
} |
|||
} |
|||
Loading…
Reference in new issue