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.
300 lines
9.7 KiB
300 lines
9.7 KiB
import 'dart:async';
|
|
import 'dart:io';
|
|
import 'dart:convert';
|
|
import 'package:deep_voice/data/models/appconfig_model.dart';
|
|
import 'package:flutter_dotenv/flutter_dotenv.dart';
|
|
import 'package:get_storage/get_storage.dart';
|
|
import 'package:open_ai_service/open_ai_service.dart';
|
|
import 'package:get/get.dart';
|
|
import 'ai_service.dart';
|
|
import 'package:flutter/services.dart';
|
|
import 'package:path_provider/path_provider.dart';
|
|
|
|
/// 定义流事件类型,用于区分不同类型的事件
|
|
enum StreamEventType { token, functionCall, complete, error }
|
|
|
|
/// 流事件包装类
|
|
class StreamEvent {
|
|
final StreamEventType type;
|
|
final String? content;
|
|
final Map? meta;
|
|
|
|
StreamEvent(this.type, {this.content, this.meta});
|
|
}
|
|
|
|
/// OpenAI服务适配器 - 连接AiService接口与OpenAIService插件
|
|
class OpenAIServiceAdapter implements AiService {
|
|
final OpenAIService _openAIService = OpenAIService();
|
|
StreamSubscription<OpenAIEvent>? _eventSubscription;
|
|
final StreamController<StreamEvent> _tokenStreamController =
|
|
StreamController<StreamEvent>.broadcast();
|
|
final GetStorage _storage = GetStorage();
|
|
|
|
/// 构造函数
|
|
OpenAIServiceAdapter() {
|
|
printInfo(info: '创建OpenAIServiceAdapter实例');
|
|
_setupEventListener();
|
|
}
|
|
|
|
/// 设置事件监听器
|
|
void _setupEventListener() {
|
|
try {
|
|
// 首先确保访问eventStream以初始化底层事件通道
|
|
_openAIService.eventStream;
|
|
|
|
// 设置事件处理
|
|
_eventSubscription = _openAIService.eventStream.listen(
|
|
(event) {
|
|
try {
|
|
switch (event.type) {
|
|
case OpenAIEventType.token:
|
|
if (event.content is String) {
|
|
_tokenStreamController.add(StreamEvent(StreamEventType.token,
|
|
content: event.content as String));
|
|
} else {
|
|
printInfo(info: '收到非字符串类型的token: ${event.content}');
|
|
}
|
|
break;
|
|
case OpenAIEventType.complete:
|
|
_tokenStreamController
|
|
.add(StreamEvent(StreamEventType.complete));
|
|
break;
|
|
case OpenAIEventType.error:
|
|
_tokenStreamController.add(StreamEvent(StreamEventType.error,
|
|
content: '未知错误: ${event.content}'));
|
|
break;
|
|
case OpenAIEventType.functionCall:
|
|
printInfo(
|
|
info:
|
|
'接收到 functionCall content: ${event.content} meta: ${event.meta}');
|
|
_tokenStreamController.add(StreamEvent(
|
|
StreamEventType.functionCall,
|
|
meta: event.meta as Map<String, dynamic>?));
|
|
// _tokenStreamController.add(StreamEvent(
|
|
// StreamEventType.error,
|
|
// content: '收到函数调用,该流仅支持文本响应'
|
|
// ));
|
|
break;
|
|
}
|
|
} catch (e) {
|
|
printError(info: '处理事件出错: $e');
|
|
_tokenStreamController
|
|
.add(StreamEvent(StreamEventType.error, content: '处理事件失败: $e'));
|
|
}
|
|
},
|
|
onError: (error) {
|
|
printError(info: '事件流错误: $error');
|
|
_tokenStreamController.add(
|
|
StreamEvent(StreamEventType.error, content: '事件流错误: $error'));
|
|
},
|
|
onDone: () {
|
|
printInfo(info: '事件流已关闭');
|
|
},
|
|
);
|
|
} catch (e) {
|
|
printError(info: '设置事件监听器失败: $e');
|
|
}
|
|
}
|
|
|
|
/// 初始化OpenAI服务
|
|
Future<bool> initialize() async {
|
|
try {
|
|
// 从.env文件中读取配置
|
|
// final apiKey = dotenv.env['OPENAI_API_KEY'] ?? '';
|
|
// final baseUrl = dotenv.env['OPENAI_BASE_URL'] ?? '';
|
|
// final model = dotenv.env['OPENAI_MODEL'] ?? '';
|
|
final _env = _storage.read("ENV") as Map<String, String>;
|
|
final apiKey = _env['OPENAI_API_KEY'] ?? '';
|
|
final baseUrl = _env['OPENAI_BASE_URL'] ?? '';
|
|
final model = _env['OPENAI_MODEL'] ?? '';
|
|
|
|
printInfo(info: '从.env读取OpenAI配置');
|
|
printInfo(info: '基础URL: $baseUrl');
|
|
printInfo(info: '模型名称: $model');
|
|
|
|
if (apiKey.isEmpty) {
|
|
printError(info: '错误: OpenAI API密钥未配置,请在.env文件中设置OPENAI_API_KEY');
|
|
return false;
|
|
}
|
|
final _mcps = _storage.read("MCPS") as Map<String, DBMCPServer>;
|
|
// 读取.mcp.json文件
|
|
String mcpConfig = jsonEncode(_mcps);
|
|
// try {
|
|
// // 从Flutter资源包中加载.mcp.json
|
|
// mcpConfig = await rootBundle.loadString('.mcp.json');
|
|
// printInfo(info: '成功从资源包加载.mcp.json配置文件');
|
|
|
|
// // 验证JSON格式
|
|
// final jsonData = jsonDecode(mcpConfig);
|
|
// if (jsonData is Map && jsonData.containsKey('mcpServers')) {
|
|
// printInfo(info: '解析到有效的mcpServers配置');
|
|
// } else {
|
|
// printInfo(info: '.mcp.json内容格式不正确,期望包含mcpServers字段');
|
|
// }
|
|
// } catch (e) {
|
|
// printError(info: '加载.mcp.json资源文件时出错: $e');
|
|
// }
|
|
|
|
// 初始化OpenAI服务
|
|
final result = await _openAIService.initialize(
|
|
apiKey: apiKey,
|
|
baseUrl: baseUrl,
|
|
model: model,
|
|
mcpServer: mcpConfig,
|
|
);
|
|
|
|
if (result) {
|
|
printInfo(info: 'OpenAI服务初始化成功');
|
|
} else {
|
|
printError(info: 'OpenAI服务初始化失败');
|
|
}
|
|
|
|
return result;
|
|
} catch (e) {
|
|
printError(info: 'OpenAI服务初始化异常: $e');
|
|
return false;
|
|
}
|
|
}
|
|
|
|
/// 发送消息并获取回复
|
|
@override
|
|
Future<String> sendMessage({
|
|
required List<Map<String, String>> messages,
|
|
String? systemPrompt,
|
|
String? userProperties,
|
|
}) async {
|
|
try {
|
|
// 在方法内直接转换
|
|
final convertedMessages =
|
|
messages.map((m) => Map<String, dynamic>.from(m)).toList();
|
|
|
|
// 发送消息并获取回复
|
|
final response = await _openAIService.sendMessage(
|
|
messages: convertedMessages,
|
|
systemPrompt: systemPrompt,
|
|
userProperties: userProperties,
|
|
);
|
|
|
|
return response;
|
|
} catch (e) {
|
|
printError(info: 'OpenAI发送消息失败: $e');
|
|
throw '发送消息失败: $e';
|
|
}
|
|
}
|
|
|
|
/// 发送消息并获取流式回复
|
|
@override
|
|
Stream sendMessageStream({
|
|
required List<Map<String, String>> messages,
|
|
String? systemPrompt, //提示词
|
|
String? userProperties, //用户属性
|
|
}) async* {
|
|
try {
|
|
// 在方法内直接转换
|
|
final convertedMessages =
|
|
messages.map((m) => Map<String, dynamic>.from(m)).toList();
|
|
|
|
// 创建用于接收token的控制器
|
|
final localController = StreamController();
|
|
|
|
// 添加从广播流到本地流的订阅
|
|
final subscription = _tokenStreamController.stream.listen((streamEvent) {
|
|
switch (streamEvent.type) {
|
|
case StreamEventType.token:
|
|
if (streamEvent.content != null && !localController.isClosed) {
|
|
localController.add(streamEvent.content!);
|
|
}
|
|
break;
|
|
case StreamEventType.functionCall:
|
|
if (streamEvent.meta != null && !localController.isClosed) {
|
|
localController.add(streamEvent.meta!);
|
|
}
|
|
break;
|
|
case StreamEventType.complete:
|
|
if (!localController.isClosed) {
|
|
localController.close();
|
|
}
|
|
break;
|
|
case StreamEventType.error:
|
|
if (!localController.isClosed) {
|
|
printError(info: '令牌流错误: ${streamEvent.content}');
|
|
localController.addError(streamEvent.content ?? '未知错误');
|
|
localController.close();
|
|
}
|
|
break;
|
|
}
|
|
});
|
|
|
|
// 当本地控制器关闭时,取消订阅
|
|
localController.onCancel = () {
|
|
subscription.cancel();
|
|
};
|
|
|
|
// 启动流式消息请求
|
|
bool started = false;
|
|
try {
|
|
started = await _openAIService.sendMessageStream(
|
|
messages: convertedMessages,
|
|
systemPrompt: systemPrompt,
|
|
userProperties: userProperties,
|
|
);
|
|
} catch (e) {
|
|
printError(info: '启动消息流失败: $e');
|
|
if (!localController.isClosed) {
|
|
localController.addError('启动消息流失败: $e');
|
|
localController.close();
|
|
}
|
|
throw '启动消息流失败: $e';
|
|
}
|
|
|
|
if (!started) {
|
|
printError(info: '无法启动消息流');
|
|
if (!localController.isClosed) {
|
|
localController.addError('无法启动消息流');
|
|
localController.close();
|
|
}
|
|
throw '无法启动消息流';
|
|
}
|
|
|
|
// 通过yield*将controller的流转发
|
|
yield* localController.stream;
|
|
} catch (e) {
|
|
printError(info: 'OpenAI流式消息处理失败: $e');
|
|
throw '流式消息处理失败: $e';
|
|
}
|
|
}
|
|
|
|
/// 注册函数
|
|
Future<bool> registerFunction(
|
|
String name, String description, Map<String, dynamic> parameters) async {
|
|
try {
|
|
final result = await _openAIService.registerFunction(
|
|
name: name,
|
|
description: description,
|
|
parameters: parameters,
|
|
);
|
|
|
|
if (result) {
|
|
printInfo(info: '函数 "$name" 注册成功');
|
|
} else {
|
|
printError(info: '函数 "$name" 注册失败');
|
|
}
|
|
|
|
return result;
|
|
} catch (e) {
|
|
printError(info: '注册函数失败: $e');
|
|
return false;
|
|
}
|
|
}
|
|
|
|
/// 释放资源
|
|
void dispose() {
|
|
try {
|
|
_eventSubscription?.cancel();
|
|
_tokenStreamController.close();
|
|
printInfo(info: 'OpenAIServiceAdapter资源已释放');
|
|
} catch (e) {
|
|
printError(info: '释放资源时出错: $e');
|
|
}
|
|
}
|
|
}
|
|
|