流式 AI 界面的核心不是逐字打印,而是消费一组可排序、可去重、可恢复的领域事件。本文用一个可运行的 React 与 Node.js 示例,构建支持工具调用、请求取消、断线重连和异常恢复的消息状态机。

先把流式响应看成状态机

最简单的流式聊天实现,往往只有一行核心逻辑:text += chunk。它能演示打字效果,却无法准确表达真实 LLM 请求的生命周期。

一次请求可能先输出文本,再发起工具调用,拿到结果后继续生成;用户也可能主动取消,网络可能中断,服务端还可能返回可展示的业务错误。若所有内容都被压成字符串,前端很难回答这些问题:当前还能否取消?工具是否仍在执行?重连后哪些片段已经处理过?error 是模型失败,还是连接暂时断开?

更合适的做法是把每个回答建模为状态机:

阶段含义是否终态
idle尚未开始
streaming正在接收文本或工具事件
reconnecting连接中断,准备续传
cancelling已请求取消,等待服务端确认
completed正常完成
cancelled服务端确认取消
failed业务错误或重试耗尽

状态改变只能由明确事件驱动,例如 text_deltatool_starttool_resulterrorcancelleddone。这样,文本只是状态的一部分,而不是整个协议。

设计可恢复的 SSE 事件协议

浏览器原生 EventSource 只建立 GET 请求,因此示例把流程拆成三个接口:

  1. POST /sessions 创建请求并返回会话 ID;
  2. GET /sessions/:id/stream?after=事件ID 订阅或续传事件;
  3. 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 终态收敛;
  • 生产环境应将会话日志放入共享存储,并设置过期、鉴权和资源上限。

当协议能够明确表达生命周期后,打字动画只是渲染细节;真正重要的是,无论正常完成、工具执行、用户取消还是网络抖动,界面都能落入一个可解释、可恢复的状态。