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.
247 lines
7.7 KiB
247 lines
7.7 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';
|
|
|
|
/// OpenAI服务适配器 - 连接AiService接口与OpenAIService插件
|
|
class OpenAIServiceAdapter implements AiService {
|
|
final OpenAIService _openAIService = OpenAIService();
|
|
StreamSubscription<OpenAIEvent>? _eventSubscription;
|
|
final StreamController<String> _tokenStreamController = StreamController<String>.broadcast();
|
|
bool _isProcessingStream = false;
|
|
|
|
/// 构造函数
|
|
OpenAIServiceAdapter() {
|
|
printInfo(info: '创建OpenAIServiceAdapter实例');
|
|
_setupEventListener();
|
|
}
|
|
|
|
/// 设置事件监听器
|
|
void _setupEventListener() {
|
|
try {
|
|
// 首先确保访问eventStream以初始化底层事件通道
|
|
_openAIService.eventStream;
|
|
|
|
// 设置事件处理
|
|
_eventSubscription = _openAIService.eventStream.listen(
|
|
(event) {
|
|
if (!_isProcessingStream) return;
|
|
|
|
try {
|
|
switch (event.type) {
|
|
case OpenAIEventType.token:
|
|
if (event.content is String) {
|
|
_tokenStreamController.add(event.content as String);
|
|
} else {
|
|
printInfo(info: '收到非字符串类型的token: ${event.content}');
|
|
}
|
|
break;
|
|
case OpenAIEventType.complete:
|
|
_isProcessingStream = false;
|
|
break;
|
|
case OpenAIEventType.error:
|
|
if (event.content is String) {
|
|
_tokenStreamController.addError(event.content as String);
|
|
} else {
|
|
_tokenStreamController.addError('未知错误: ${event.content}');
|
|
}
|
|
_isProcessingStream = false;
|
|
break;
|
|
case OpenAIEventType.functionCall:
|
|
try {
|
|
if (event.content is Map) {
|
|
_tokenStreamController.addError('收到函数调用,该流仅支持文本响应');
|
|
} else {
|
|
_tokenStreamController.addError('收到未知格式的函数调用');
|
|
printError(info: '函数调用格式错误: ${event.content}');
|
|
}
|
|
} catch (e) {
|
|
printError(info: '处理函数调用事件出错: $e');
|
|
_tokenStreamController.addError('处理函数调用失败: $e');
|
|
}
|
|
_isProcessingStream = false;
|
|
break;
|
|
}
|
|
} catch (e) {
|
|
printError(info: '处理事件出错: $e');
|
|
_tokenStreamController.addError('处理事件失败: $e');
|
|
_isProcessingStream = false;
|
|
}
|
|
},
|
|
onError: (error) {
|
|
printError(info: '事件流错误: $error');
|
|
_tokenStreamController.addError('事件流错误: $error');
|
|
_isProcessingStream = false;
|
|
},
|
|
onDone: () {
|
|
printInfo(info: '事件流已关闭');
|
|
_isProcessingStream = false;
|
|
},
|
|
);
|
|
} 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>();
|
|
|
|
// 标记开始处理流
|
|
_isProcessingStream = true;
|
|
|
|
// 添加从广播流到本地流的订阅
|
|
final subscription = _tokenStreamController.stream.listen(
|
|
(token) => localController.add(token),
|
|
onError: (error) {
|
|
printError(info: '令牌流错误: $error');
|
|
localController.addError(error);
|
|
localController.close();
|
|
},
|
|
onDone: () {
|
|
printInfo(info: '令牌流已完成');
|
|
localController.close();
|
|
}
|
|
);
|
|
|
|
// 当本地控制器关闭时,取消订阅
|
|
localController.onCancel = () {
|
|
subscription.cancel();
|
|
};
|
|
|
|
// 启动流式消息请求
|
|
bool started = false;
|
|
try {
|
|
started = await _openAIService.sendMessageStream(
|
|
messages: convertedMessages,
|
|
);
|
|
} catch (e) {
|
|
printError(info: '启动消息流失败: $e');
|
|
localController.addError('启动消息流失败: $e');
|
|
localController.close();
|
|
_isProcessingStream = false;
|
|
throw '启动消息流失败: $e';
|
|
}
|
|
|
|
if (!started) {
|
|
printError(info: '无法启动消息流');
|
|
localController.addError('无法启动消息流');
|
|
localController.close();
|
|
_isProcessingStream = false;
|
|
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 {
|
|
_isProcessingStream = false;
|
|
_eventSubscription?.cancel();
|
|
_tokenStreamController.close();
|
|
printInfo(info: 'OpenAIServiceAdapter资源已释放');
|
|
} catch (e) {
|
|
printError(info: '释放资源时出错: $e');
|
|
}
|
|
}
|
|
|
|
}
|