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.

259 lines
7.9 KiB

import 'dart:async';
import 'package:flutter_dotenv/flutter_dotenv.dart';
import 'package:open_ai_service/open_ai_service.dart';
import 'package:get/get.dart';
import 'ai_service.dart';
/// 定义流事件类型,用于区分不同类型的事件
enum StreamEventType { token, complete, error }
/// 流事件包装类
class StreamEvent {
final StreamEventType type;
final String? content;
StreamEvent(this.type, {this.content});
}
/// OpenAI服务适配器 - 连接AiService接口与OpenAIService插件
class OpenAIServiceAdapter implements AiService {
final OpenAIService _openAIService = OpenAIService();
StreamSubscription<OpenAIEvent>? _eventSubscription;
final StreamController<StreamEvent> _tokenStreamController = StreamController<StreamEvent>.broadcast();
/// 构造函数
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:
_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'] ?? '';
printInfo(info: '从.env读取OpenAI配置');
printInfo(info: '基础URL: $baseUrl');
printInfo(info: '模型名称: $model');
if (apiKey.isEmpty) {
printError(info: '错误: OpenAI API密钥未配置,请在.env文件中设置OPENAI_API_KEY');
return false;
}
// 初始化OpenAI服务
final result = await _openAIService.initialize(
apiKey: apiKey,
baseUrl: baseUrl,
model: model,
);
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,
required String systemPrompt,
}) async {
try {
// 在方法内直接转换
final convertedMessages = messages.map((m) =>
Map<String, dynamic>.from(m)).toList();
// 发送消息并获取回复
final response = await _openAIService.sendMessage(
messages: convertedMessages,
);
return response;
} catch (e) {
printError(info: 'OpenAI发送消息失败: $e');
throw '发送消息失败: $e';
}
}
/// 发送消息并获取流式回复
@override
Stream<String> sendMessageStream({
required List<Map<String, String>> messages,
required String systemPrompt,
}) async* {
try {
// 在方法内直接转换
final convertedMessages = messages.map((m) =>
Map<String, dynamic>.from(m)).toList();
// 创建用于接收token的控制器
final localController = StreamController<String>();
// 添加从广播流到本地流的订阅
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.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,
);
} 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');
}
}
}