Appearance
sse-agent-orchestration
1. 实验目标
用纯 Node(http 模块,零依赖)实现最小 SSE + Agent 编排(对应 docs/08-ai-engineering/03、05、08):
- SSE 流式:
GET /api/tasks/:id/stream推送step/done/cancelled事件 - 任务状态机:
pending → running → completed / failed / cancelled - 中断恢复:
POST /api/tasks/:id/cancel置状态,SSE 推cancelled;客户端重连时先收snapshot恢复当前状态 - 三层架构:本服务即"后端编排层",状态存内存(真实场景落 DB/Redis)
2. 工程要点
tasksMap 承载任务状态(真实场景是存储层)runAgent用setInterval模拟 Agent 逐步执行(每 1s 一步),向所有 SSE 连接广播- SSE 帧格式:
data: {json}\n\n(空行分隔);客户端重连先收snapshot对齐状态 broadcast遍历任务的clientsSet,连接断开时req.on('close')移除
3. 关键实现速览
3.1 SSE 端点 + 快照恢复
为什么看这段:SSE 的关键是 text/event-stream + 帧格式 + 断线恢复(先发 snapshot)。
js
res.writeHead(200, {
'Content-Type': 'text/event-stream',
'Cache-Control': 'no-cache',
Connection: 'keep-alive',
});
task.clients.add(res);
res.write(`data: ${JSON.stringify({ type: 'snapshot', state: task.state })}\n\n`); // 重连恢复
req.on('close', () => task.clients.delete(res));3.2 任务状态机 + 中断
为什么看这段:状态变更先落存储(这里改 task.state),再广播——"先写后推"保证一致性。
js
function runAgent(taskId, plan) {
let step = 0;
const timer = setInterval(() => {
const t = tasks.get(taskId);
if (t.cancelled) { // 中断检查:状态机决策
t.state = 'cancelled';
broadcast(taskId, { type: 'cancelled' });
return clearInterval(timer);
}
if (step >= plan.length) { // 完成
t.state = 'completed';
return broadcast(taskId, { type: 'done' }), clearInterval(timer);
}
broadcast(taskId, { type: 'step', step, desc: plan[step] });
step++;
}, 1000);
}4. 运行方式
bash
cd labs/node/ai-engineering/sse-agent-orchestration
node server.js # 终端 1:启动服务(端口 3000)
node client.js # 终端 2:完整跑完 5 步
node client.js --cancel-after 2500 # 终端 3:2.5 秒后中断,观察 cancelled5. 预期现象
- 服务启动打印三个端点说明
- client 完整模式:依次收到
snapshot→step 0..4→done - 中断模式:2.5 秒后收到
cancelled,任务状态变cancelled,不再推进
6. 常见误区
| 误区 | 实际情况 |
|---|---|
| SSE 是轮询 | 是服务端主动推送的单向长连接;客户端断开不影响其他连接 |
| 任务状态放前端 | 页面刷新即丢;必须后端存储 + 状态机推进(见 docs 08) |
| 中断是前端行为 | 前端只发事件,状态转移由后端决策(本 demo 的 cancel 接口) |
7. 对应知识库文档
- 理论主文档:SSE 流式传输、网络中断恢复与大模型打字机打点
- Agent 编排:Agent Tool Calling (Function Calling) 架构与安全防护
- 三层架构:Agent 前端 / 后端 / 存储基本架构