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.
283 lines
8.7 KiB
283 lines
8.7 KiB
import 'dart:async';
|
|
import 'dart:io';
|
|
import 'dart:convert';
|
|
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';
|
|
import 'package:flutter/services.dart';
|
|
import 'package:path_provider/path_provider.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;
|
|
}
|
|
|
|
// 读取.mcp.json文件
|
|
String mcpConfig = '';
|
|
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,
|
|
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');
|
|
}
|
|
}
|
|
}
|