Batch B: settings popover 新增"扩展"tab,按知识库/Skills/MCP/Agent工具分面板 Batch C: knowledge-rag popover 新增"资料源管理/检索测试/检索参数"三tab Batch A: manifest.rs 按能力域拆为 10 个分组函数(context/doc/knowledge_rag/block/page/mindmap/office/onlyoffice/artifact) Batch D: 新增 agent-stream-event-router.js 共享模块(事件状态机 + SSE 帧解析) Batch E: settings popover 新增"仪表盘"tab — 知识库文件数/Skills数量/MCP连接状态/Agent run次数 836/838 Rust 测试通过(仅 2 个预存失败)
128 lines
4.2 KiB
JavaScript
128 lines
4.2 KiB
JavaScript
// AgentStreamEventRouter — 统一的 agent 流式事件状态机
|
|
// 供 Page AI sidebar runtime 和 document editor adapter runtime 共用
|
|
//
|
|
// 事件状态:
|
|
// init — 建立 request_id 绑定
|
|
// loading — message_delta / tool_call_delta → 增量更新
|
|
// stream_event — tool-started / tool-finished / agent_state
|
|
// finished — 正常终止
|
|
// interrupted — 中断(等待审批/恢复)
|
|
// error — 错误终止
|
|
// approval_required — 等待用户确认
|
|
|
|
export const AGENT_EVENT_STATES = {
|
|
INIT: 'init',
|
|
LOADING: 'loading',
|
|
STREAM_EVENT: 'stream_event',
|
|
FINISHED: 'finished',
|
|
INTERRUPTED: 'interrupted',
|
|
ERROR: 'error',
|
|
APPROVAL_REQUIRED: 'approval_required',
|
|
};
|
|
|
|
export const TERMINAL_STATES = new Set([
|
|
AGENT_EVENT_STATES.FINISHED,
|
|
AGENT_EVENT_STATES.INTERRUPTED,
|
|
AGENT_EVENT_STATES.ERROR,
|
|
]);
|
|
|
|
export function isTerminalState(status) {
|
|
return TERMINAL_STATES.has(status);
|
|
}
|
|
|
|
// 解析 SSE 帧为 eventName + payloadText
|
|
export function parseSSEFrames(buffer, lastBoundary) {
|
|
var frames = buffer.split('\n\n');
|
|
var remaining = frames.pop() || '';
|
|
var events = [];
|
|
frames.forEach(function(frame) {
|
|
var eventName = '';
|
|
var dataLines = [];
|
|
frame.split('\n').forEach(function(line) {
|
|
if (line.startsWith('event:')) eventName = line.slice(6).trim();
|
|
if (line.startsWith('data:')) dataLines.push(line.slice(5).trim());
|
|
});
|
|
var payloadText = dataLines.join('\n');
|
|
if (!eventName && payloadText) {
|
|
try {
|
|
var parsed = JSON.parse(payloadText);
|
|
eventName = parsed && parsed.event ? String(parsed.event) : '';
|
|
} catch (_) {}
|
|
}
|
|
if (eventName) events.push({ eventName: eventName, payloadText: payloadText });
|
|
});
|
|
return { events: events, remaining: remaining };
|
|
}
|
|
|
|
// 从 SSE 流读取事件
|
|
export async function readSSEStream(response, onFrame) {
|
|
if (!response.body || typeof response.body.getReader !== 'function') return;
|
|
var reader = response.body.getReader();
|
|
var decoder = new TextDecoder();
|
|
var buffer = '';
|
|
while (true) {
|
|
var chunk = await reader.read();
|
|
if (chunk.done) break;
|
|
buffer += decoder.decode(chunk.value, { stream: true });
|
|
var result = parseSSEFrames(buffer, '');
|
|
buffer = result.remaining;
|
|
result.events.forEach(function(evt) {
|
|
onFrame(evt.eventName, evt.payloadText);
|
|
});
|
|
}
|
|
}
|
|
|
|
// 创建标准事件路由器
|
|
// handlers: { onInit, onDelta, onToolCall, onToolResult, onAgentState, onTerminal, onApprovalRequired }
|
|
export function createAgentStreamRouter(handlers) {
|
|
var h = handlers || {};
|
|
|
|
return function routeEvent(eventName, payloadText) {
|
|
var payload = null;
|
|
try { payload = JSON.parse(payloadText || 'null'); } catch (_) {}
|
|
|
|
var status = (payload && payload.status) || eventName;
|
|
|
|
switch (status) {
|
|
case AGENT_EVENT_STATES.INIT:
|
|
if (h.onInit) h.onInit(payload);
|
|
break;
|
|
|
|
case AGENT_EVENT_STATES.LOADING:
|
|
if (h.onDelta) h.onDelta(payload);
|
|
if (h.onToolCall) {
|
|
var toolChunks = (payload && payload.tool_call_chunks) || (payload && payload.msg && payload.msg.tool_call_chunks);
|
|
if (toolChunks && toolChunks.length) h.onToolCall(payload);
|
|
}
|
|
break;
|
|
|
|
case AGENT_EVENT_STATES.STREAM_EVENT:
|
|
if (payload && payload.event === 'tool-finished') {
|
|
if (h.onToolResult) h.onToolResult(payload);
|
|
} else if (payload && payload.event === 'tool-started') {
|
|
// tool start — 可选处理
|
|
} else if (payload && payload.agent_state) {
|
|
if (h.onAgentState) h.onAgentState(payload.agent_state);
|
|
}
|
|
break;
|
|
|
|
case AGENT_EVENT_STATES.FINISHED:
|
|
case AGENT_EVENT_STATES.INTERRUPTED:
|
|
case AGENT_EVENT_STATES.ERROR:
|
|
if (h.onTerminal) h.onTerminal(status, payload);
|
|
break;
|
|
|
|
case AGENT_EVENT_STATES.APPROVAL_REQUIRED:
|
|
if (h.onApprovalRequired) h.onApprovalRequired(payload);
|
|
break;
|
|
|
|
default:
|
|
// 未知状态:尝试作为 loading 处理
|
|
if (payload && (payload.text || payload.delta || payload.content)) {
|
|
if (h.onDelta) h.onDelta(payload);
|
|
}
|
|
break;
|
|
}
|
|
};
|
|
}
|