import 'dart:async'; import 'dart:io'; import 'dart:convert'; import '../../data/models/appconfig.dart'; import '../../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? _eventSubscription; final StreamController _tokenStreamController = StreamController.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?)); // _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 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; final apiKey = AppConfig.env('OPENAI_API_KEY') ?? ''; final baseUrl = AppConfig.env('OPENAI_BASE_URL') ?? ''; final model = AppConfig.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; // 读取.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 sendMessage({ required List> messages, String? systemPrompt, String? userProperties, }) async { try { // 在方法内直接转换 final convertedMessages = messages.map((m) => Map.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> messages, String? systemPrompt, //提示词 String? userProperties, //用户属性 }) async* { try { // 在方法内直接转换 final convertedMessages = messages.map((m) => Map.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 registerFunction( String name, String description, Map 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'); } } }