流式 AI 界面的核心不是逐字打印,而是消费一组可排序、可去重、可恢复的领域事件。本文用一个可运行的 React 与 Node.js 示例,构建支持工具调用、请求取消、断线重连和异常恢复的消息状态机。
先把流式响应看成状态机
最简单的流式聊天实现,往往只有一行核心逻辑:text += chunk。它能演示打字效果,却无法准确表达真实 LLM 请求的生命周期。
一次请求可能先输出文本,再发起工具调用,拿到结果后继续生成;用户也可能主动取消,网络可能中断,服务端还可能返回可展示的业务错误。若所有内容都被压成字符串,前端很难回答这些问题:当前还能否取消?工具是否仍在执行?重连后哪些片段已经处理过?error 是模型失败,还是连接暂时断开?
更合适的做法是把每个回答建模为状态机:
| 阶段 | 含义 | 是否终态 |
|---|---|---|
idle | 尚未开始 | 否 |
streaming | 正在接收文本或工具事件 | 否 |
reconnecting | 连接中断,准备续传 | 否 |
cancelling | 已请求取消,等待服务端确认 | 否 |
completed | 正常完成 | 是 |
cancelled | 服务端确认取消 | 是 |
failed | 业务错误或重试耗尽 | 是 |
状态改变只能由明确事件驱动,例如 text_delta、tool_start、tool_result、error、cancelled 和 done。这样,文本只是状态的一部分,而不是整个协议。
设计可恢复的 SSE 事件协议
浏览器原生 EventSource 只建立 GET 请求,因此示例把流程拆成三个接口:
POST /sessions创建请求并返回会话 ID;GET /sessions/:id/stream?after=事件ID订阅或续传事件;DELETE /sessions/:id请求取消。
每个 SSE 消息都有单调递增的 id,数据部分采用统一信封:
{
"id": "4",
"type": "tool_result",
"toolCallId": "weather-1",
"output": "上海 18°C,多云"
}
服务端暂存事件日志。连接断开后,客户端携带最后成功处理的 ID 重连,服务端只重放更新的事件。前端仍要保留 seen 集合进行去重,因为重放、代理缓存边界和本地处理时机都可能造成重复投递。
这里还要区分两类错误:type: error 是服务端已经确认的业务终态;EventSource.onerror 只是传输层异常,应先进入 reconnecting,而不是立即把回答标记为失败。
Node.js 服务端:记录、重放与取消
下面使用真实存在的 Express、CORS 和 SSE API。模型输出由定时器模拟,便于独立运行;接入具体 LLM 时,只需把 produce 替换为模型流,并继续调用 emit。
先创建服务端:
mkdir sse-server && cd sse-server
npm init -y
npm install express cors
保存为 server.mjs:
import express from "express";
import cors from "cors";
import { randomUUID } from "node:crypto";
const app = express();
app.use(cors());
app.use(express.json());
const sessions = new Map();
const terminalTypes = new Set(["done", "cancelled", "error"]);
function emit(session, type, payload = {}) {
if (session.terminal) return;
const event = {
id: String(++session.sequence),
type,
...payload
};
session.events.push(event);
const frame =
`id: ${event.id}\n` +
`event: ai\n` +
`data: ${JSON.stringify(event)}\n\n`;
for (const response of session.clients) response.write(frame);
if (terminalTypes.has(type)) {
session.terminal = true;
clearInterval(session.timer);
for (const response of session.clients) response.end();
session.clients.clear();
setTimeout(() => sessions.delete(session.id), 60_000).unref();
}
}
function produce(session) {
if (session.terminal) return;
emit(session, "start");
const steps = session.prompt.includes("失败")
? [
() => emit(session, "text_delta", { delta: "正在处理请求……" }),
() => emit(session, "error", { message: "这是示例业务错误" })
]
: [
() => emit(session, "text_delta", { delta: "我先查询天气。" }),
() => emit(session, "tool_start", {
toolCallId: "weather-1",
name: "get_weather"
}),
() => emit(session, "tool_result", {
toolCallId: "weather-1",
output: "上海 18°C,多云"
}),
() => emit(session, "text_delta", {
delta: " 上海当前 18°C,多云,适合外出。"
}),
() => emit(session, "done")
];
let index = 0;
session.timer = setInterval(() => {
if (session.terminal) return;
steps[index++]();
if (index === steps.length) clearInterval(session.timer);
}, 700);
}
app.post("/sessions", (req, res) => {
const prompt = String(req.body?.prompt ?? "").trim();
if (!prompt) return res.status(400).json({ error: "prompt is required" });
const session = {
id: randomUUID(),
prompt,
sequence: 0,
events: [],
clients: new Set(),
terminal: false,
timer: null
};
sessions.set(session.id, session);
setTimeout(() => produce(session), 50);
res.status(201).json({ id: session.id });
});
app.get("/sessions/:id/stream", (req, res) => {
const session = sessions.get(req.params.id);
if (!session) return res.sendStatus(404);
res.set({
"Content-Type": "text/event-stream",
"Cache-Control": "no-cache, no-transform",
Connection: "keep-alive"
});
res.flushHeaders();
const after = Number(req.query.after || req.get("Last-Event-ID") || 0);
for (const event of session.events) {
if (Number(event.id) <= after) continue;
res.write(
`id: ${event.id}\nevent: ai\ndata: ${JSON.stringify(event)}\n\n`
);
}
if (session.terminal) return res.end();
session.clients.add(res);
req.on("close", () => session.clients.delete(res));
});
app.delete("/sessions/:id", (req, res) => {
const session = sessions.get(req.params.id);
if (!session) return res.sendStatus(404);
if (!session.terminal) emit(session, "cancelled");
res.sendStatus(202);
});
app.listen(3000, () => console.log("SSE server: http://localhost:3000"));
运行 node server.mjs。终态事件保留一分钟,使短暂断线仍能恢复;生产环境通常应改用带过期时间的 Redis、数据库或消息日志,不能依赖单进程内存。
React 客户端:归约事件而非拼接响应
用 Vite 创建项目并替换 src/App.jsx:
npm create vite@latest web -- --template react
cd web
npm install
npm run dev
import { useEffect, useReducer, useRef, useState } from "react";
const API = "http://localhost:3000";
const initial = {
phase: "idle",
text: "",
tools: [],
error: null,
seen: {}
};
function reducer(state, action) {
if (action.kind === "reset") return { ...initial, phase: "streaming" };
if (action.kind === "reconnecting") return { ...state, phase: "reconnecting" };
if (action.kind === "cancelling") return { ...state, phase: "cancelling" };
if (action.kind === "transport_failed") {
return { ...state, phase: "failed", error: action.message };
}
const event = action.event;
if (!event || state.seen[event.id]) return state;
const next = { ...state, seen: { ...state.seen, [event.id]: true } };
switch (event.type) {
case "start":
return { ...next, phase: "streaming" };
case "text_delta":
return { ...next, phase: "streaming", text: next.text + event.delta };
case "tool_start":
return {
...next,
tools: [...next.tools, {
id: event.toolCallId,
name: event.name,
status: "running"
}]
};
case "tool_result":
return {
...next,
tools: next.tools.map(tool =>
tool.id === event.toolCallId
? { ...tool, status: "completed", output: event.output }
: tool
)
};
case "done":
return { ...next, phase: "completed" };
case "cancelled":
return { ...next, phase: "cancelled" };
case "error":
return { ...next, phase: "failed", error: event.message };
default:
return next;
}
}
export default function App() {
const [prompt, setPrompt] = useState("上海天气如何?");
const [state, dispatch] = useReducer(reducer, initial);
const sourceRef = useRef(null);
const sessionRef = useRef(null);
const lastIdRef = useRef(0);
const generationRef = useRef(0);
const retryTimerRef = useRef(null);
useEffect(() => () => {
generationRef.current++;
sourceRef.current?.close();
clearTimeout(retryTimerRef.current);
}, []);
async function start() {
sourceRef.current?.close();
clearTimeout(retryTimerRef.current);
const generation = ++generationRef.current;
lastIdRef.current = 0;
dispatch({ kind: "reset" });
try {
const response = await fetch(`${API}/sessions`, {
method: "POST",
headers: { "Content-Type": "application/json" },
body: JSON.stringify({ prompt })
});
if (!response.ok) throw new Error(`创建请求失败:${response.status}`);
const { id } = await response.json();
if (generation !== generationRef.current) return;
sessionRef.current = id;
let retries = 0;
const connect = () => {
if (generation !== generationRef.current) return;
const source = new EventSource(
`${API}/sessions/${id}/stream?after=${lastIdRef.current}`
);
sourceRef.current = source;
source.onopen = () => { retries = 0; };
source.addEventListener("ai", message => {
const event = JSON.parse(message.data);
if (Number(event.id) <= lastIdRef.current) return;
lastIdRef.current = Number(event.id);
dispatch({ kind: "event", event });
if (["done", "cancelled", "error"].includes(event.type)) {
source.close();
}
});
source.onerror = () => {
source.close();
if (generation !== generationRef.current) return;
dispatch({ kind: "reconnecting" });
retries += 1;
if (retries > 5) {
dispatch({
kind: "transport_failed",
message: "连接多次失败,请重新发起请求"
});
return;
}
retryTimerRef.current = setTimeout(connect, Math.min(1000 * retries, 5000));
};
};
connect();
} catch (error) {
dispatch({ kind: "transport_failed", message: error.message });
}
}
async function cancel() {
if (!sessionRef.current) return;
dispatch({ kind: "cancelling" });
await fetch(`${API}/sessions/${sessionRef.current}`, { method: "DELETE" });
}
const active = ["streaming", "reconnecting", "cancelling"].includes(state.phase);
return (
<main>
<h2>流式消息状态机</h2>
<input value={prompt} onChange={e => setPrompt(e.target.value)} />
<button onClick={start}>开始</button>
<button onClick={cancel} disabled={!active}>取消</button>
<p>状态:{state.phase}</p>
<p>{state.text}</p>
{state.tools.map(tool => (
<pre key={tool.id}>{JSON.stringify(tool, null, 2)}</pre>
))}
{state.error && <p role="alert">{state.error}</p>}
</main>
);
}
可在浏览器开发者工具中临时切换离线状态。恢复网络后,客户端使用 after 续传,服务端重放缺失事件,归约器再用事件 ID 去重。连续五次失败才进入 failed,用户随后点击“开始”即可创建全新的状态机,实现异常后的显式恢复。
总结
流式 AI UI 应以事件协议和状态机为边界,而不是围绕一个不断增长的字符串组织代码。关键要点包括:
- 为文本、工具调用、错误、取消和完成定义独立事件;
- 使用单调事件 ID、服务端日志和客户端去重实现可恢复投递;
- 区分业务错误与暂时的传输错误,避免断网即失败;
- 取消必须由服务端确认,并以
cancelled终态收敛; - 生产环境应将会话日志放入共享存储,并设置过期、鉴权和资源上限。
当协议能够明确表达生命周期后,打字动画只是渲染细节;真正重要的是,无论正常完成、工具执行、用户取消还是网络抖动,界面都能落入一个可解释、可恢复的状态。
评论