Files
AI-Health/health_app/lib/providers/chat_provider.dart
MingNian 13714d9ed8 feat: 应用内通知系统 + 结构化手术史/用药等相关改动
- 新增用户通知 outbox 流水线(EfUserNotificationPipeline)与后台投递 worker
- 通知中心页面及前端通知服务接入
- 健康指标异常、用药/运动提醒等事件统一产出站内通知
- 健康档案结构化手术史、用药提醒扫描、医生/用户端点等配套调整
- AppDbContext 注册通知相关实体
2026-06-21 21:06:29 +08:00

483 lines
14 KiB
Dart
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

import 'dart:async';
import 'dart:io';
import 'package:flutter_riverpod/flutter_riverpod.dart';
import 'auth_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<String, dynamic>? 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';
}
enum ActiveAgent {
default_,
consultation,
health,
diet,
medication,
report,
exercise,
}
class ChatState {
final ActiveAgent activeAgent;
final List<ChatMessage> messages;
final String? conversationId;
final bool isStreaming;
final String? thinkingText;
const ChatState({
this.activeAgent = ActiveAgent.default_,
this.messages = const [],
this.conversationId,
this.isStreaming = false,
this.thinkingText,
});
ChatState copyWith({
ActiveAgent? activeAgent,
List<ChatMessage>? messages,
String? conversationId,
bool? isStreaming,
String? thinkingText,
}) => ChatState(
activeAgent: activeAgent ?? this.activeAgent,
messages: messages ?? this.messages,
conversationId: conversationId ?? this.conversationId,
isStreaming: isStreaming ?? this.isStreaming,
thinkingText: thinkingText ?? this.thinkingText,
);
}
class SelectedAgentNotifier extends Notifier<ActiveAgent?> {
@override
ActiveAgent? build() => null;
void select(ActiveAgent? a) => state = a;
}
final selectedAgentProvider =
NotifierProvider<SelectedAgentNotifier, ActiveAgent?>(
SelectedAgentNotifier.new,
);
final chatProvider = NotifierProvider<ChatNotifier, ChatState>(
ChatNotifier.new,
);
class ChatNotifier extends Notifier<ChatState> {
StreamSubscription<Map<String, dynamic>>? _subscription;
Completer<void>? _streamDone;
ActiveAgent? _lastTriggeredAgent;
void markNeedsRebuild() => state = state.copyWith();
/// 不可变消息操作方法(供 chat_messages_view 新版代码调用)
Future<String?> confirmMessage(String id) async {
final msgs = state.messages.toList();
final i = msgs.indexWhere((m) => m.id == id);
if (i < 0) return '确认卡片不存在';
final rawIds = msgs[i].metadata?['confirmationIds'];
final confirmationIds = rawIds is List
? rawIds.map((e) => e.toString()).where((e) => e.isNotEmpty).toList()
: <String>[];
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();
_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,
],
);
}
void setAgent(ActiveAgent a) {
// 流式回复中忽略胶囊切换,防止状态混乱
if (state.isStreaming) return;
_cancelActiveStream();
state = state.copyWith(activeAgent: a);
ref.read(selectedAgentProvider.notifier).select(a);
}
/// 根据 AI 调用的工具自动切换智能体胶囊
void _switchAgentByTool(String tool) {
ActiveAgent? agent;
switch (tool) {
case 'record_health_data':
case 'query_health_records':
agent = ActiveAgent.health;
break;
case 'estimate_food_text':
agent = ActiveAgent.diet;
break;
case 'manage_medication':
agent = ActiveAgent.medication;
break;
case 'manage_exercise':
agent = ActiveAgent.exercise;
break;
case 'request_doctor':
agent = ActiveAgent.consultation;
break;
case 'analyze_report':
agent = ActiveAgent.report;
break;
}
if (agent != null) {
ref.read(selectedAgentProvider.notifier).select(agent);
state = state.copyWith(activeAgent: agent);
}
}
Future<void> loadConversation(String convId) async {
await _cancelActiveStream();
try {
final api = ref.read(apiClientProvider);
final res = await api.get('/api/ai/conversations/$convId');
final rawMessages = (res.data['data'] as List?) ?? [];
final messages = rawMessages.map((m) {
final map = m as Map<String, dynamic>;
return ChatMessage(
id: map['id']?.toString() ?? '',
role: map['role']?.toString() ?? 'user',
content: map['content']?.toString() ?? '',
createdAt:
DateTime.tryParse(map['createdAt']?.toString() ?? '') ??
DateTime.now(),
type: MessageType.text,
);
}).toList();
state = state.copyWith(
messages: messages,
conversationId: convId,
activeAgent: ActiveAgent.default_,
);
ref.read(selectedAgentProvider.notifier).select(ActiveAgent.default_);
} catch (_) {}
}
void insertAgentWelcome(ActiveAgent agent) {
state = state.copyWith(
messages: [
...state.messages,
ChatMessage(
id: 'welcome_${agent.name}_${DateTime.now().millisecondsSinceEpoch}',
role: 'assistant',
content: '',
createdAt: DateTime.now(),
type: MessageType.agentWelcome,
metadata: {'agent': agent.name},
),
],
);
}
/// 点击胶囊:先出用户标签 → 0.4 秒后出欢迎卡片,不走 AI
/// 重复点击同一胶囊不重复弹卡片
void triggerAgent(ActiveAgent agent, String label) {
if (_lastTriggeredAgent == agent) return;
_lastTriggeredAgent = agent;
final userMsg = ChatMessage(
id: 'agent_trigger_${DateTime.now().millisecondsSinceEpoch}',
role: 'user',
content: label,
createdAt: DateTime.now(),
);
state = state.copyWith(messages: [...state.messages, userMsg]);
Future.delayed(const Duration(milliseconds: 400), () {
insertAgentWelcome(agent);
});
}
Future<void> sendImage(String imagePath, String text) async {
final file = File(imagePath);
if (!await file.exists()) return;
_lastTriggeredAgent = null;
// 先显示用户消息(本地显示图片路径)
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]);
// 异步上传图片
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 = <String, dynamic>{'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]);
return;
}
// 将图片 URL 作为消息内容发送给 AI
final msgWithImage = text.isNotEmpty ? '$text\n[图片已上传]' : '[图片已上传]';
await _sendToAI(msgWithImage);
}
Future<void> sendMessage(String text) async {
if (text.trim().isEmpty || state.isStreaming) return;
_lastTriggeredAgent = null;
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<void> _sendToAI(String text) 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,
token: token,
);
await _cancelActiveStream();
final done = Completer<void>();
_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, '网络异常,请稍后重试');
}
}
Future<void> _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<String, dynamic> j, ChatMessage aiMsg) {
final a = j['action'] as String?;
switch (a) {
case 'conversation_id':
state = state.copyWith(conversationId: j['data']?.toString());
case 'answer':
final messageType = j['type'] as String? ?? 'text';
aiMsg.type = _parseMessageType(messageType);
if (j['metadata'] is Map) {
aiMsg.metadata = Map<String, dynamic>.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? ?? '';
// 根据 AI 调用的工具自动切换智能体胶囊
_switchAgentByTool(tool);
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;
}
}
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);
}
}