import 'dart:async'; import 'dart:convert'; import 'dart:io'; import 'package:flutter_riverpod/flutter_riverpod.dart'; import 'auth_provider.dart'; import 'conversation_history_provider.dart'; import 'data_providers.dart'; import '../utils/sse_handler.dart'; enum MessageType { text, dataConfirm, agentWelcome, taskCard } class ChatMessage { final String id; final String role; String content; final DateTime createdAt; MessageType type; Map? metadata; bool confirmed; ChatMessage({ required this.id, required this.role, required this.content, required this.createdAt, this.type = MessageType.text, this.metadata, this.confirmed = false, }); bool get isUser => role == 'user'; bool get isReadOnly => metadata?['readOnly'] == true; } enum ActiveAgent { default_, consultation, health, diet, medication, report, exercise, } const _keepConversationId = Object(); class ChatState { final ActiveAgent activeAgent; final List messages; final String? conversationId; final bool isStreaming; final String? thinkingText; final bool isViewingHistory; const ChatState({ this.activeAgent = ActiveAgent.default_, this.messages = const [], this.conversationId, this.isStreaming = false, this.thinkingText, this.isViewingHistory = false, }); ChatState copyWith({ ActiveAgent? activeAgent, List? messages, Object? conversationId = _keepConversationId, bool? isStreaming, String? thinkingText, bool? isViewingHistory, }) => ChatState( activeAgent: activeAgent ?? this.activeAgent, messages: messages ?? this.messages, conversationId: identical(conversationId, _keepConversationId) ? this.conversationId : conversationId as String?, isStreaming: isStreaming ?? this.isStreaming, thinkingText: thinkingText ?? this.thinkingText, isViewingHistory: isViewingHistory ?? this.isViewingHistory, ); } final chatProvider = NotifierProvider( ChatNotifier.new, ); class ChatNotifier extends Notifier { StreamSubscription>? _subscription; Completer? _streamDone; ActiveAgent? _lastTriggeredAgent; Timer? _agentTapLockTimer; bool _loadingConversation = false; /// 重置整个会话:取消正在进行的 SSE,清空消息和会话 ID。 /// 历史记录页一键清空 / 删除当前会话时调用。 Future resetSession() async { await _cancelActiveStream(); _cancelPendingAgentWelcome(); _lastTriggeredAgent = null; state = const ChatState(); } /// 不可变消息操作方法(供 chat_messages_view 新版代码调用) Future confirmMessage(String id) async { final msgs = state.messages.toList(); final i = msgs.indexWhere((m) => m.id == id); if (i < 0) return '确认卡片不存在'; if (msgs[i].isReadOnly) return '历史记录中的录入卡片仅供查看'; final rawIds = msgs[i].metadata?['confirmationIds']; final confirmationIds = rawIds is List ? rawIds.map((e) => e.toString()).where((e) => e.isNotEmpty).toList() : []; if (confirmationIds.isEmpty) { return '确认信息已失效,请重新发送需要录入的内容'; } try { final api = ref.read(apiClientProvider); for (final confirmationId in confirmationIds.toList()) { final response = await api.post( '/api/ai/confirm-write/$confirmationId', ); final body = response.data; if (body is! Map || body['code'] != 0) { return body is Map ? body['message']?.toString() ?? '录入失败' : '录入失败'; } confirmationIds.remove(confirmationId); msgs[i].metadata?['confirmationIds'] = confirmationIds.toList(); state = state.copyWith(messages: msgs); } } catch (e) { return '录入失败,请检查网络后重试'; } msgs[i].confirmed = true; state = state.copyWith(messages: msgs); ref.invalidate(medicationListProvider); ref.invalidate(medicationReminderProvider); ref.invalidate(latestHealthProvider); ref.invalidate(currentExercisePlanProvider); return null; } @override ChatState build() { ref.onDispose(() { _subscription?.cancel(); _agentTapLockTimer?.cancel(); _subscription = null; if (_streamDone != null && !_streamDone!.isCompleted) { _streamDone!.complete(); } _streamDone = null; }); Future.microtask(() { insertTaskCard(); }); return const ChatState(); } void insertTaskCard() { if (state.messages.any((m) => m.type == MessageType.taskCard)) return; state = state.copyWith( messages: [ ChatMessage( id: 'task_card', role: 'assistant', content: '', createdAt: DateTime.now(), type: MessageType.taskCard, ), ...state.messages, ], ); } Future loadConversation(String convId) async { if (state.isStreaming) return '小脉正在回复,请稍后再切换对话'; if (_loadingConversation) return '正在加载其他对话,请稍候'; _loadingConversation = true; await _cancelActiveStream(); _cancelPendingAgentWelcome(); try { final api = ref.read(apiClientProvider); final res = await api.get('/api/ai/conversations/$convId'); final rawMessages = (res.data['data'] as List?) ?? []; if (rawMessages.isEmpty) return '该对话已不存在,请刷新记录'; final messages = rawMessages.map((m) { final map = m as Map; final role = map['role']?.toString().toLowerCase() == 'user' ? 'user' : 'assistant'; final metadata = _parseMetadata(map['metadataJson']) ?? {}; metadata['readOnly'] = true; return ChatMessage( id: map['id']?.toString() ?? '', role: role, content: map['content']?.toString() ?? '', createdAt: DateTime.tryParse(map['createdAt']?.toString() ?? '') ?? DateTime.now(), type: _messageTypeFromMetadata(metadata), metadata: metadata, confirmed: metadata['confirmationIds'] is! List, ); }).toList(); state = ChatState( messages: messages, conversationId: convId, activeAgent: ActiveAgent.default_, isViewingHistory: true, ); _lastTriggeredAgent = null; return null; } catch (_) { return '会话加载失败,请稍后重试'; } finally { _loadingConversation = false; } } /// 点击胶囊:用户标签和欢迎卡片立即出现,不走 AI。 /// 300ms 点击锁只防止误触,不延迟界面反馈。 /// 重复点击同一胶囊不重复弹卡片 void triggerAgent(ActiveAgent agent, String label) { if (_agentTapLockTimer != null || _lastTriggeredAgent == agent) return; _resumeConversationFromHistory(); _lastTriggeredAgent = agent; final now = DateTime.now(); final userMsg = ChatMessage( id: 'agent_trigger_${now.microsecondsSinceEpoch}', role: 'user', content: label, createdAt: now, ); final welcomeMsg = ChatMessage( id: 'welcome_${agent.name}_${now.microsecondsSinceEpoch}', role: 'assistant', content: '', createdAt: now, type: MessageType.agentWelcome, metadata: {'agent': agent.name}, ); state = state.copyWith( messages: [...state.messages, userMsg, welcomeMsg], activeAgent: agent, ); _agentTapLockTimer = Timer(const Duration(milliseconds: 300), () { _agentTapLockTimer = null; }); } void _cancelPendingAgentWelcome() { _agentTapLockTimer?.cancel(); _agentTapLockTimer = null; } Future sendImage(String imagePath, String text) async { if (state.isStreaming) return; final file = File(imagePath); if (!await file.exists()) return; _lastTriggeredAgent = null; _cancelPendingAgentWelcome(); _resumeConversationFromHistory(); // 先显示用户消息(本地显示图片路径) final userMsg = ChatMessage( id: '${DateTime.now().millisecondsSinceEpoch}', role: 'user', content: text.isNotEmpty ? text : '[图片]', createdAt: DateTime.now(), metadata: {'localImagePath': imagePath}, ); state = state.copyWith( messages: [...state.messages, userMsg], isStreaming: true, ); // 异步上传图片 String? uploadedUrl; Object? uploadError; try { final api = ref.read(apiClientProvider); uploadedUrl = await api.uploadFile('/api/files/upload', file); } catch (e) { uploadError = e; } // 更新消息元数据(保留本地路径 + 添加远程URL) final updatedMsgs = state.messages.toList(); final idx = updatedMsgs.indexWhere((m) => m.id == userMsg.id); if (idx >= 0) { final meta = {'localImagePath': imagePath}; if (uploadedUrl != null) meta['imageUrl'] = uploadedUrl; updatedMsgs[idx] = ChatMessage( id: userMsg.id, role: 'user', content: userMsg.content, createdAt: userMsg.createdAt, metadata: meta, ); state = state.copyWith(messages: updatedMsgs); } if (uploadedUrl == null) { final errorMsg = ChatMessage( id: '${DateTime.now().millisecondsSinceEpoch}_upload_error', role: 'assistant', content: uploadError == null ? '图片上传失败,请稍后重试。' : '图片上传失败,请检查文件大小或网络后重试。', createdAt: DateTime.now(), ); state = state.copyWith(messages: [...state.messages, errorMsg]); state = state.copyWith(isStreaming: false); return; } // 把图片 URL 透传给后端,后端会调 VLM 识图并把描述拼到 LLM 上下文 final userText = text.isNotEmpty ? text : '请帮我看看这张图片'; await _sendToAI(userText, imageUrl: uploadedUrl); } /// 发送 PDF 附件 + 文字(PDF 解析在后端做)。 Future sendPdf(String pdfPath, String fileName, String text) async { if (state.isStreaming) return; final file = File(pdfPath); if (!await file.exists()) return; _lastTriggeredAgent = null; _cancelPendingAgentWelcome(); _resumeConversationFromHistory(); final userMsg = ChatMessage( id: '${DateTime.now().millisecondsSinceEpoch}', role: 'user', content: text.isNotEmpty ? text : '请帮我看看这份 PDF', createdAt: DateTime.now(), metadata: {'pdfFileName': fileName}, ); state = state.copyWith( messages: [...state.messages, userMsg], isStreaming: true, ); String? uploadedUrl; try { final api = ref.read(apiClientProvider); uploadedUrl = await api.uploadFile('/api/files/upload', file); } catch (_) { // ignore,下方统一处理 } // 更新消息附带的远程 URL if (uploadedUrl != null) { final updatedMsgs = state.messages.toList(); final idx = updatedMsgs.indexWhere((m) => m.id == userMsg.id); if (idx >= 0) { updatedMsgs[idx] = ChatMessage( id: userMsg.id, role: 'user', content: userMsg.content, createdAt: userMsg.createdAt, metadata: {'pdfFileName': fileName, 'pdfUrl': uploadedUrl}, ); state = state.copyWith(messages: updatedMsgs); } } else { final errorMsg = ChatMessage( id: '${DateTime.now().millisecondsSinceEpoch}_upload_error', role: 'assistant', content: 'PDF 上传失败,请检查文件大小或网络后重试。', createdAt: DateTime.now(), ); state = state.copyWith(messages: [...state.messages, errorMsg]); state = state.copyWith(isStreaming: false); return; } await _sendToAI(userMsg.content, pdfUrl: uploadedUrl); } Future sendMessage(String text) async { if (text.trim().isEmpty || state.isStreaming) return; _lastTriggeredAgent = null; _cancelPendingAgentWelcome(); _resumeConversationFromHistory(); final userMsg = ChatMessage( id: '${DateTime.now().millisecondsSinceEpoch}', role: 'user', content: text, createdAt: DateTime.now(), ); state = state.copyWith( messages: [...state.messages, userMsg], isStreaming: true, ); await _sendToAI(text); } Future _sendToAI( String text, { String? imageUrl, String? pdfUrl, }) async { final aiMsg = ChatMessage( id: '${DateTime.now().millisecondsSinceEpoch}_ai', role: 'assistant', content: '', createdAt: DateTime.now(), ); // 立即加入空 AI 消息,让思考动画有载体 state = state.copyWith( messages: [...state.messages, aiMsg], isStreaming: true, ); try { final token = await ref.read(apiClientProvider).accessToken; if (token == null) { _addError(aiMsg, '未登录,请重新登录'); return; } // 始终用 unified 智能体,AI 自动判断意图分配工具 final stream = SseHandler.connect( agentType: 'unified', message: text, conversationId: state.conversationId, imageUrl: imageUrl, pdfUrl: pdfUrl, token: token, ); await _cancelActiveStream(); final done = Completer(); _streamDone = done; _subscription = stream.listen( (event) => _processEvent(event, aiMsg), onError: (_) { _addError(aiMsg, '网络异常,请稍后重试'); if (!done.isCompleted) done.complete(); }, onDone: () { if (!done.isCompleted) done.complete(); }, cancelOnError: true, ); await done.future; if (_streamDone == done) { _streamDone = null; _subscription = null; } if (state.isStreaming) { _done(aiMsg); } } catch (e) { _addError(aiMsg, '网络异常,请稍后重试'); } } void _resumeConversationFromHistory() { if (!state.isViewingHistory) return; state = state.copyWith(isViewingHistory: false); } Future _cancelActiveStream() async { await _subscription?.cancel(); _subscription = null; if (_streamDone != null && !_streamDone!.isCompleted) { _streamDone!.complete(); } _streamDone = null; } void _addError(ChatMessage aiMsg, String errorText) { aiMsg.content = errorText; aiMsg.type = MessageType.text; final u = state.messages.toList(); final i = u.indexWhere((x) => x.id == aiMsg.id); if (i >= 0) { u[i] = aiMsg; } else { u.add(aiMsg); } state = state.copyWith(messages: u, isStreaming: false, thinkingText: null); } void _processEvent(Map j, ChatMessage aiMsg) { final a = j['action'] as String?; switch (a) { case 'conversation_id': state = state.copyWith( conversationId: j['data']?.toString(), isViewingHistory: false, ); case 'answer': final messageType = j['type'] as String? ?? 'text'; aiMsg.type = _parseMessageType(messageType); if (j['metadata'] is Map) { aiMsg.metadata = Map.from(j['metadata']); } aiMsg.content += (j['data'] as String?) ?? ''; state = state.copyWith(thinkingText: null); _update(aiMsg); case 'notice': state = state.copyWith(thinkingText: j['message'] as String?); case 'tool_result': final tool = j['tool'] as String? ?? ''; if (tool == 'record_health_data') { ref.invalidate(latestHealthProvider); } case 'status': _done(aiMsg); case 'error': _done(aiMsg); } } MessageType _parseMessageType(String type) { switch (type) { case 'data_confirm': case 'medication_confirm': return MessageType.dataConfirm; case 'agent_welcome': return MessageType.agentWelcome; default: return MessageType.text; } } Map? _parseMetadata(dynamic raw) { if (raw == null) return null; if (raw is Map) return Map.from(raw); if (raw is! String || raw.trim().isEmpty) return null; try { final decoded = jsonDecode(raw); return decoded is Map ? Map.from(decoded) : null; } catch (_) { return null; } } MessageType _messageTypeFromMetadata(Map? metadata) { if (metadata == null) return MessageType.text; final type = metadata['messageType']?.toString(); if (type != null && type.isNotEmpty) return _parseMessageType(type); if (metadata['confirmationIds'] is List) return MessageType.dataConfirm; return MessageType.text; } void _update(ChatMessage m) { final u = state.messages.toList(); final i = u.indexWhere((x) => x.id == m.id); if (i >= 0) { u[i] = m; } else { u.add(m); // 空内容也加入,让思考气泡有载体 } state = state.copyWith(messages: u); } void _done(ChatMessage m) { final u = state.messages.toList(); final content = m.content.trim(); if (content.isEmpty && m.type == MessageType.text) { m.content = '暂时没有收到回复,请稍后再试。'; } final i = u.indexWhere((x) => x.id == m.id); if (i >= 0) { u[i] = m; } else { u.add(m); } state = state.copyWith(messages: u, isStreaming: false, thinkingText: null); ref.invalidate(conversationHistoryProvider); } }