Files
ai-agent-admin/web/apps/web-ele/src/components/ai-chat-panel/composables/useEventHandler.ts
T

851 lines
26 KiB
TypeScript
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.
/**
* 流式事件处理器
*
* 统一处理工作流和 Agent 的流式事件
*/
import type { Ref } from 'vue';
import type {
AiChatPanelEmits,
AiChatPanelProps,
ChatMessage,
DesignPreviewData,
ReasoningStep,
StreamEvent,
} from '../types';
/** 事件处理器配置 */
interface EventHandlerConfig {
props: AiChatPanelProps;
emit: AiChatPanelEmits;
messages: Ref<ChatMessage[]>;
currentSteps: Ref<ReasoningStep[]>;
currentRunId: Ref<string>;
currentAssistantMsgId: Ref<null | string>;
waitingForInput: Ref<boolean>;
waitingConfig: Ref<any>;
running: Ref<boolean>;
conversationId: Ref<null | string>;
// 设计预览回调
onDesignPreview?: (data: DesignPreviewData) => void;
}
/** 生成消息 ID */
const generateId = () =>
`msg-${Date.now()}-${Math.random().toString(36).slice(2, 8)}`;
/**
* 创建事件处理器
*/
export function useEventHandler(config: EventHandlerConfig) {
const {
props,
emit,
messages,
currentSteps,
currentRunId,
// currentAssistantMsgId 保留用于未来扩展
waitingForInput,
waitingConfig,
running,
onDesignPreview,
} = config;
/** 更新助手消息 */
const updateAssistantMessage = (
msgId: string,
updates: Partial<ChatMessage>,
) => {
const index = messages.value.findIndex((m) => m.id === msgId);
if (index !== -1) {
const currentMsg = messages.value[index];
if (currentMsg) {
messages.value[index] = { ...currentMsg, ...updates } as ChatMessage;
}
}
};
/** 获取等待提示文本 */
const getWaitingPrompt = (waitConfig: any): string => {
if (!waitConfig) return '等待输入...';
switch (waitConfig.type) {
case 'choice': {
return waitConfig.question || '请选择:';
}
case 'confirm': {
return `**${waitConfig.title || '确认'}**\n\n${waitConfig.content || '请确认是否继续?'}`;
}
case 'question': {
return waitConfig.question || '请输入:';
}
default: {
return '等待输入...';
}
}
};
/** 处理设计预览事件 */
const handleDesignPreview = (
event: StreamEvent,
msgId: string,
configData: any,
) => {
const previewType = configData.preview_type;
const previewData: DesignPreviewData = {
type: previewType,
title: configData.title || '设计预览',
data: configData.data || {},
nodeId: event.node_id || '',
data_sources: configData.data_sources || configData.table_configs,
schema_fields: configData.schema_fields || configData.form_fields,
layoutOptions: configData.layoutOptions,
themeOptions: configData.themeOptions,
};
// 调用设计预览回调
onDesignPreview?.(previewData);
// 显示提示消息
const waitingContent = `**${configData.title}**\n\n${configData.message || '请在右侧面板中确认或编辑设计'}`;
const initialMsg = messages.value.find((m) => m.id === msgId);
if (initialMsg && !initialMsg.content?.trim()) {
updateAssistantMessage(msgId, {
status: 'completed',
content: waitingContent,
});
} else {
const newMsg: ChatMessage = {
id: generateId(),
role: 'assistant',
content: waitingContent,
timestamp: new Date(),
status: 'completed',
};
messages.value.push(newMsg);
}
};
/** 处理普通对话流交互 */
const handleDialogFlowInteraction = (msgId: string, configData: any) => {
const waitingContent = getWaitingPrompt(configData);
const initialMsg = messages.value.find((m) => m.id === msgId);
if (initialMsg && !initialMsg.content?.trim()) {
updateAssistantMessage(msgId, {
status: 'completed',
content: waitingContent,
interaction: configData,
});
} else {
const newMsg: ChatMessage = {
id: generateId(),
role: 'assistant',
content: waitingContent,
timestamp: new Date(),
status: 'completed',
interaction: configData,
};
messages.value.push(newMsg);
}
};
/** 获取节点显示名称(优先从前端 nodes 中获取) */
const getNodeLabel = (
nodeId?: string,
nodeType?: string,
fallbackLabel?: string,
) => {
if (!nodeId) return fallbackLabel || nodeType || '';
if (props.nodes?.length) {
const node = props.nodes.find((n) => n.id === nodeId);
if (node?.data?.label) return node.data.label as string;
}
return fallbackLabel || nodeType || '';
};
const getItemCount = (value: any): number => {
if (Array.isArray(value)) return value.length;
if (value && typeof value === 'object') return Object.keys(value).length;
return 0;
};
const buildReasoningMeta = (event: StreamEvent): Partial<ReasoningStep> => {
const collaboration = event.collaboration || {};
const communication = event.communication || {};
return {
branch_id: event.branch_id,
branch_label: event.branch_label,
agent_code: collaboration.agent_code || event.agent_code,
agent_name: collaboration.agent_name || event.agent_name,
model:
collaboration.model ||
collaboration.model_name ||
event.model ||
event.model_name,
model_id: collaboration.model_id || event.model_id,
provider_name: collaboration.provider_name || event.provider_name,
provider_type: collaboration.provider_type || event.provider_type,
subflow_name:
event.subflow_name || collaboration.subflow_name || event.event?.subflow_name,
from_subflow: Boolean(
event.from_subflow || collaboration.from_subflow || event.event?.from_subflow,
),
collaboration_role:
collaboration.collaboration_role || collaboration.role || event.collaboration_role,
collaboration_mode: event.collaboration_mode || collaboration.collaboration_mode,
communication,
};
};
const buildMessageMeta = (event: StreamEvent): Partial<ChatMessage> => {
const collaboration = event.collaboration || {};
const modelName =
event.model ||
event.model_name ||
collaboration.model ||
collaboration.model_name ||
collaboration.model_id;
return {
...(modelName ? { model_name: modelName } : {}),
...(event.model_id || collaboration.model_id
? { model_id: event.model_id || collaboration.model_id }
: {}),
...(event.provider_name || collaboration.provider_name
? { provider_name: event.provider_name || collaboration.provider_name }
: {}),
...(collaboration.agent_code
? { agent_code: collaboration.agent_code }
: {}),
...(collaboration.agent_name
? { agent_name: collaboration.agent_name }
: {}),
...(collaboration.collaboration_mode
? { collaboration_mode: collaboration.collaboration_mode }
: {}),
};
};
const getStreamErrorMessage = (event: StreamEvent) =>
event.error_message ||
event.message ||
event.error ||
event.content ||
'执行失败';
const isSuccessEvent = (event: StreamEvent) =>
!event.error_message &&
!event.error &&
(!event.status || event.status === 'success' || event.status === 'completed');
const getEventDisplayName = (event: StreamEvent) => {
const typeMap: Record<string, string> = {
complete: '执行完成',
loop_start: '循环开始',
node_event: '节点事件',
resume: '恢复执行',
start: '开始执行',
};
return (
event.node_label ||
event.event?.message ||
event.event?.content ||
event.event?.type ||
event.message ||
event.content ||
typeMap[event.type] ||
event.type ||
'运行事件'
);
};
const appendObservedStep = (
event: StreamEvent,
msgId: string,
status: 'completed' | 'failed' | 'running' = 'completed',
) => {
currentSteps.value.push({
type: status === 'failed' ? 'error' : 'observation',
content: getEventDisplayName(event),
node_id: event.node_id,
node_type: event.node_type,
...buildReasoningMeta(event),
params: event.config || event.waiting_config || event.event || event.params,
output: event.outputs || event.output || event.result,
status,
timestamp: new Date().toISOString(),
});
updateAssistantMessage(msgId, {
reasoning_steps: [...currentSteps.value],
...buildMessageMeta(event),
});
};
/** 处理流式事件 */
const handleStreamEvent = (event: StreamEvent, msgId: string) => {
switch (event.type) {
// ========== Agent 特有事件 ==========
case 'action': {
currentSteps.value.push({
type: 'action',
content: '',
tool: event.tool,
params: event.params,
timestamp: new Date().toISOString(),
});
updateAssistantMessage(msgId, {
reasoning_steps: [...currentSteps.value],
});
break;
}
case 'annotation_reply':
case 'knowledge_retrieval':
case 'observation':
case 'thought':
case 'tool_call':
case 'tool_result': {
currentSteps.value.push({
type: event.type as ReasoningStep['type'],
content:
event.content ||
event.message ||
event.tool ||
event.type,
tool: event.tool,
params: event.params,
result: event.result,
node_id: event.node_id,
node_type: event.node_type,
...buildReasoningMeta(event),
status: 'completed',
timestamp: new Date().toISOString(),
});
updateAssistantMessage(msgId, {
reasoning_steps: [...currentSteps.value],
...buildMessageMeta(event),
});
break;
}
case 'answer': {
const initialMsg = messages.value.find((m) => m.id === msgId);
if (initialMsg && !initialMsg.content?.trim()) {
updateAssistantMessage(msgId, {
status: 'completed',
content: event.content || '',
...buildMessageMeta(event),
});
} else {
const answerMsg: ChatMessage = {
id: generateId(),
role: 'assistant',
content: event.content || '',
timestamp: new Date(),
status: 'completed',
...buildMessageMeta(event),
};
messages.value.push(answerMsg);
}
break;
}
case 'complete': {
let endOutput = '';
if (
event.outputs &&
typeof event.outputs.output === 'string' &&
event.outputs.output.trim()
) {
endOutput = event.outputs.output;
}
const initialMsg = messages.value.find((m) => m.id === msgId);
const initialMsgHasContent = initialMsg?.content?.trim();
if (endOutput) {
if (initialMsg && !initialMsgHasContent) {
updateAssistantMessage(msgId, {
status: 'completed',
content: endOutput,
elapsed_time: event.elapsed_time,
tokens_used: event.tokens_used || event.total_tokens,
...buildMessageMeta(event),
});
} else {
const endMsg: ChatMessage = {
id: generateId(),
role: 'assistant',
content: endOutput,
timestamp: new Date(),
status: 'completed',
elapsed_time: event.elapsed_time,
tokens_used: event.tokens_used || event.total_tokens,
...buildMessageMeta(event),
};
messages.value.push(endMsg);
}
} else if (!initialMsgHasContent) {
const otherAssistantMessages = messages.value.filter(
(m) =>
m.role === 'assistant' && m.id !== msgId && m.content?.trim(),
);
if (otherAssistantMessages.length > 0) {
const initialMsgIndex = messages.value.findIndex(
(m) => m.id === msgId,
);
if (initialMsgIndex !== -1) {
messages.value.splice(initialMsgIndex, 1);
}
} else {
updateAssistantMessage(msgId, {
status: 'completed',
elapsed_time: event.elapsed_time,
tokens_used: event.tokens_used || event.total_tokens,
content: '工作流执行完成',
...buildMessageMeta(event),
});
}
} else if (initialMsg) {
// 流式输出已有内容,只更新状态和统计信息
updateAssistantMessage(msgId, {
status: 'completed',
elapsed_time: event.elapsed_time,
tokens_used: event.tokens_used || event.total_tokens,
...buildMessageMeta(event),
});
}
running.value = false;
// 触发消息完成事件,用于刷新对话列表(更新标题)
emit('message-complete', config.conversationId.value);
break;
}
case 'error': {
const errorMessage = getStreamErrorMessage(event);
currentSteps.value.push({
type: 'error',
content: errorMessage,
node_id: event.node_id,
node_type: event.node_type,
...buildReasoningMeta(event),
status: 'failed',
timestamp: new Date().toISOString(),
});
updateAssistantMessage(msgId, {
status: 'failed',
error_message: errorMessage,
reasoning_steps: [...currentSteps.value],
...buildMessageMeta(event),
});
running.value = false;
break;
}
case 'llm_chunk': {
// 更新助手消息内容(流式输出)
const initialMsg = messages.value.find((m) => m.id === msgId);
if (initialMsg) {
updateAssistantMessage(msgId, {
content: event.accumulated_content || event.content || '',
status: 'streaming',
...buildMessageMeta(event),
});
} else {
// 如果消息不存在,创建新消息
const streamMsg: ChatMessage = {
id: msgId,
role: 'assistant',
content: event.accumulated_content || event.content || '',
timestamp: new Date(),
status: 'streaming',
...buildMessageMeta(event),
};
messages.value.push(streamMsg);
}
if (props.enableNodeEvents) {
emit('node-streaming', {
node_id: event.node_id,
content: event.content,
accumulated_content: event.accumulated_content,
});
}
break;
}
// ========== 循环事件 ==========
case 'loop_complete': {
const loopEvent = event as any;
const outputVar = loopEvent.output_variable || 'loop_results';
const outputs: Record<string, any> = {
total_iterations: loopEvent.total_iterations,
results_count: loopEvent.results_count,
};
if (loopEvent.loop_results) {
outputs[outputVar] = loopEvent.loop_results;
}
currentSteps.value.push({
type: 'loop_complete',
content: '循环执行完成',
node_id: loopEvent.node_id,
node_type: 'loop',
...buildReasoningMeta(event),
output: outputs,
status: 'completed',
timestamp: new Date().toISOString(),
});
updateAssistantMessage(msgId, {
reasoning_steps: [...currentSteps.value],
});
if (props.enableNodeEvents) {
emit('node-complete', {
node_id: loopEvent.node_id,
node_type: 'loop',
status: 'success',
elapsed_time: 0,
tokens_used: 0,
outputs,
});
}
break;
}
case 'loop_iteration_complete': {
currentSteps.value.push({
type: 'loop_iteration_complete',
content: `循环第 ${event.iteration || '-'} 次完成`,
node_id: event.node_id,
node_type: 'loop',
...buildReasoningMeta(event),
output: event.output,
status: 'completed',
timestamp: new Date().toISOString(),
});
updateAssistantMessage(msgId, {
reasoning_steps: [...currentSteps.value],
});
if (props.enableNodeEvents) {
emit('loop-iteration', {
node_id: event.node_id,
type: 'complete',
iteration: event.iteration,
total: event.total,
output: event.output,
});
}
break;
}
case 'loop_iteration_error': {
currentSteps.value.push({
type: 'loop_iteration_error',
content: `循环第 ${event.iteration || '-'} 次失败:${event.error || '执行失败'}`,
node_id: event.node_id,
node_type: 'loop',
...buildReasoningMeta(event),
status: 'completed',
timestamp: new Date().toISOString(),
});
updateAssistantMessage(msgId, {
reasoning_steps: [...currentSteps.value],
});
if (props.enableNodeEvents) {
emit('loop-iteration', {
node_id: event.node_id,
type: 'error',
iteration: event.iteration,
total: event.total,
error: event.error,
});
}
break;
}
case 'loop_iteration_start': {
currentSteps.value.push({
type: 'loop_iteration_start',
content: `循环第 ${event.iteration || '-'} 次开始`,
node_id: event.node_id,
node_type: 'loop',
...buildReasoningMeta(event),
params: { item: event.item, total: event.total },
status: 'running',
timestamp: new Date().toISOString(),
});
updateAssistantMessage(msgId, {
reasoning_steps: [...currentSteps.value],
});
if (props.enableNodeEvents) {
emit('loop-iteration', {
node_id: event.node_id,
type: 'start',
iteration: event.iteration,
total: event.total,
item: event.item,
});
}
break;
}
case 'node_complete': {
const existingIndex = currentSteps.value.findIndex(
(s) => s.node_id === event.node_id && s.type === 'node_start',
);
const completeLabel = getNodeLabel(
event.node_id,
event.node_type,
event.node_label,
);
const nodeSucceeded = isSuccessEvent(event);
const nodeStatus = nodeSucceeded ? 'completed' : 'failed';
const nodeContent = nodeSucceeded
? completeLabel
: `${completeLabel || event.node_id || '节点'}${getStreamErrorMessage(event)}`;
if (existingIndex === -1) {
currentSteps.value.push({
type: 'node_complete',
content: nodeContent,
node_id: event.node_id,
node_type: event.node_type,
...buildReasoningMeta(event),
output: event.outputs,
status: nodeStatus,
timestamp: new Date().toISOString(),
});
} else {
currentSteps.value[existingIndex] = {
...currentSteps.value[existingIndex]!,
type: 'node_complete',
content: nodeContent,
...buildReasoningMeta(event),
output: event.outputs,
status: nodeStatus,
};
}
updateAssistantMessage(msgId, {
reasoning_steps: [...currentSteps.value],
});
if (props.enableNodeEvents) {
emit('node-complete', {
node_id: event.node_id,
node_type: event.node_type,
status: event.status,
elapsed_time: event.elapsed_time,
tokens_used: event.tokens_used,
error_message: event.error_message,
outputs: event.outputs,
inputs: event.inputs,
loop_iteration: event.loop_iteration,
warnings: event.warnings,
});
}
break;
}
case 'parallel_complete': {
const existingIndex = currentSteps.value.findIndex(
(s) => s.node_id === event.node_id && s.type === 'parallel_start',
);
const branchResults = event.branch_results || {};
const nodeSucceeded = isSuccessEvent(event);
const content = nodeSucceeded
? `并行分支完成:${getItemCount(branchResults)} 个分支`
: `并行分支失败:${getStreamErrorMessage(event)}`;
const nodeStatus = nodeSucceeded ? 'completed' : 'failed';
if (existingIndex === -1) {
currentSteps.value.push({
type: 'parallel_complete',
content,
node_id: event.node_id,
node_type: 'parallel',
...buildReasoningMeta(event),
output: {
branch_results: branchResults,
total_tokens: event.total_tokens,
},
status: nodeStatus,
timestamp: new Date().toISOString(),
});
} else {
currentSteps.value[existingIndex] = {
...currentSteps.value[existingIndex]!,
type: 'parallel_complete',
content,
...buildReasoningMeta(event),
output: {
branch_results: branchResults,
total_tokens: event.total_tokens,
},
status: nodeStatus,
};
}
updateAssistantMessage(msgId, {
reasoning_steps: [...currentSteps.value],
});
break;
}
case 'parallel_start': {
const branches = event.branches || [];
currentSteps.value.push({
type: 'parallel_start',
content: `并行分支开始:${getItemCount(branches)} 个分支`,
node_id: event.node_id,
node_type: 'parallel',
...buildReasoningMeta(event),
params: { branches, branch_labels: event.branch_labels },
status: 'running',
timestamp: new Date().toISOString(),
});
updateAssistantMessage(msgId, {
reasoning_steps: [...currentSteps.value],
});
break;
}
case 'node_event': {
const nodeEvent = event.event;
if (nodeEvent?.type === 'message') {
const initialMsg = messages.value.find((m) => m.id === msgId);
if (initialMsg && !initialMsg.content?.trim()) {
updateAssistantMessage(msgId, {
status: 'completed',
content: nodeEvent.content || '',
...buildMessageMeta(event),
});
} else {
const newMsg: ChatMessage = {
id: generateId(),
role: 'assistant',
content: nodeEvent.content || '',
timestamp: new Date(),
status: 'completed',
...buildMessageMeta(event),
};
messages.value.push(newMsg);
}
} else if (nodeEvent) {
appendObservedStep(event, msgId);
}
break;
}
// ========== 节点事件 ==========
case 'node_start': {
const existingIndex = currentSteps.value.findIndex(
(s) => s.node_id === event.node_id,
);
if (existingIndex === -1) {
currentSteps.value.push({
type: 'node_start',
content: getNodeLabel(
event.node_id,
event.node_type,
event.node_label,
),
node_id: event.node_id,
node_type: event.node_type,
...buildReasoningMeta(event),
status: 'running',
timestamp: new Date().toISOString(),
});
}
updateAssistantMessage(msgId, {
reasoning_steps: [...currentSteps.value],
});
if (props.enableNodeEvents) {
emit('node-start', {
node_id: event.node_id,
node_type: event.node_type,
node_label: event.node_label,
inputs: event.inputs,
loop_iteration: event.loop_iteration,
});
}
break;
}
// ========== 通用事件 ==========
case 'start': {
currentRunId.value = event.run_id || '';
break;
}
case 'resume': {
currentRunId.value = event.run_id || currentRunId.value;
appendObservedStep(event, msgId, 'running');
break;
}
// ========== 对话流交互事件 ==========
case 'waiting_input': {
waitingForInput.value = true;
waitingConfig.value = event.waiting_config || event.config;
const configData = event.waiting_config || (event as any).config;
currentSteps.value.push({
type: 'waiting_input',
content:
configData?.title ||
configData?.question ||
event.node_label ||
'等待用户输入',
node_id: event.node_id,
node_type: event.node_type,
...buildReasoningMeta(event),
params: configData,
status: 'running',
timestamp: new Date().toISOString(),
});
if (configData?.type === 'design_preview') {
handleDesignPreview(event, msgId, configData);
} else {
handleDialogFlowInteraction(msgId, configData);
}
updateAssistantMessage(msgId, {
reasoning_steps: [...currentSteps.value],
});
running.value = false;
break;
}
default: {
if (
event.type &&
(event.content ||
event.message ||
event.error ||
event.error_message ||
event.event)
) {
appendObservedStep(
event,
msgId,
event.error || event.error_message ? 'failed' : 'completed',
);
}
break;
}
}
};
return {
handleStreamEvent,
updateAssistantMessage,
getWaitingPrompt,
generateId,
};
}