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

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');
}
}
}