Skip to content
LLM 应用状态管理:消息流、任务状态与数据同步
概述
LLM API 本身是无状态的。以 chat completion 为例,客户端向接口发送一个包含完整消息历史的请求,服务端返回一次响应;服务端不会为客户端保存会话记录。这意味着应用层必须自己维护状态,才能实现连续对话、流式渲染、任务追踪和多端同步。
在 LLM 应用中,状态模型通常由三个核心对象构成:
- 会话(Session):一次对话的聚合单位,包含会话 ID、消息列表和相关元数据。
- 消息(Message):会话中一次交互的最小单元,对应
user输入、assistant输出或工具调用结果。 - 任务(Task):一次模型调用过程的执行单元,记录调用从开始到结束所处的状态。
三个对象通过事件相互关联。用户发送消息时,客户端创建一条 user 消息,同时创建一个任务;服务端开始执行任务后,产生 assistant 消息的增量内容;每个增量内容更新消息状态,消息完成后任务进入终态。消息流、任务状态和数据同步不是三个独立模块,而是同一状态模型的不同侧面。
与普通 Web 应用的“请求-响应”模式相比,LLM 应用的状态管理有两个显著差异。第一,响应是流式的:模型按 token 序列输出,客户端需要持续追加内容,消息因此存在“未完成”的中间状态。第二,调用是异步的:一次模型调用可能耗时数秒到数十秒,期间可能发生网络中断、API 错误、用户取消,或者等待外部工具执行结果,应用需要显式跟踪任务的执行进度。
基本概念
消息(Message)
消息是会话中可渲染的最小状态单元。从状态管理的角度,一条消息只需要暴露渲染和更新所需的字段。消息角色表示消息的来源:user 是用户输入,assistant 是模型输出,tool 是工具执行后返回的结果。
typescript
type MessageRole = 'user' | 'assistant' | 'tool';
interface Message {
id: string;
role: MessageRole;
content: string;
status: 'streaming' | 'completed' | 'error';
createdAt: number;
updatedAt?: number;
}消息 ID 由客户端或服务端显式生成,不建议使用数组索引作为 ID。流式增量到达时,客户端需要通过 ID 定位目标消息,并原地更新其 content。
任务(Task)
一次 LLM 调用,尤其是一次涉及工具调用的多步任务,无法用一个布尔标志表示“完成或未完成”。任务需要经历创建、排队、执行、暂停、结束等多个阶段,这些阶段组成一个状态机。
typescript
type TaskStatus =
| 'pending'
| 'running'
| 'succeeded'
| 'failed'
| 'cancelled'
| 'needs-action';
interface Task {
id: string;
status: TaskStatus;
attempt: number;
maxAttempts: number;
createdAt: number;
updatedAt: number;
}状态含义:
pending:任务已创建,等待被调度执行。running:任务正在执行,可能正在流式输出,也可能正在等待内部子步骤。succeeded:任务正常完成,结果已产出。failed:任务执行失败,可以进入重试,也可以直接结束。cancelled:任务被用户或系统取消。needs-action:任务执行到某个点,需要外部输入才能继续。例如等待工具执行结果、等待用户确认。
会话(Session)
会话是消息和任务的聚合边界。会话本身不参与流式增量计算,但它是消息列表的读取单位,也是数据同步中快照的粒度。一个会话包含一条按时间排序的消息列表,以及当前正在执行的任务集合。
消息流建模
消息对象与角色
消息是流式输出的最小可渲染单元。status 表示消息内容是否完整。SSE 流式输出中,模型内容不是一个完整字符串一次到达的,而是以增量 token 逐步累积。当消息正在接收增量时,status 为 streaming;增量接收完成后改为 completed;如果流中断或处理出错,改为 error。UI 可以根据这三个状态渲染不同的视觉效果:
tsx
function MessageItem({ message }: { message: Message }) {
return (
<div className={`message message-${message.role}`}>
{message.status === 'streaming' && <span className="cursor">▍</span>}
{message.content}
{message.status === 'error' && <span>(错误)</span>}
</div>
);
}工具调用场景中,一条 assistant 消息可能携带对某个工具的调用参数,而工具结果则以 tool 角色追加到会话中。不同模型提供方可能增加额外角色,例如用于携带系统指令的角色;这里是消息流建模的基础集合,具体实现以所用 API 的文档为准。
增量事件与流式协议
LLM 流式响应最常用的传输协议是 SSE(Server-Sent Events)。OpenAI 兼容的流式接口在请求参数中设置 stream: true 后,返回 text/event-stream 格式。每个事件是一个以 data: 开头的数据行,事件之间用空行分隔。流结束时返回一个 data: [DONE] 标记 [2]。以下是一个典型的响应片段:
text
data: {"id":"chatcmpl-123","object":"chat.completion.chunk","choices":[{"delta":{"role":"assistant"},"index":0}]}
data: {"id":"chatcmpl-123","object":"chat.completion.chunk","choices":[{"delta":{"content":"你"},"index":0}]}
data: {"id":"chatcmpl-123","object":"chat.completion.chunk","choices":[{"delta":{"content":"好"},"index":0}]}
data: [DONE]每个数据块是一个 JSON 对象。choices[0].delta 表示这一帧的增量,delta.content 是新增的文本片段。流式响应的核心语义是:delta 是“增量”,不是“完整消息”。客户端必须把每个 delta 追加到已有内容上,而不是用 delta 替换整条消息。
流式响应结束前,最后一帧可能带有 choices[0].finish_reason,用于标识结束原因,例如 stop、length 或 tool_calls。如果请求参数中设置了 stream_options: { "include_usage": true },流末尾还会返回一个包含 usage 字段的帧,提供 prompt_tokens、completion_tokens 和 total_tokens 三个计数 [3]。
读取与解析 SSE
在 Node.js 中读取 SSE 流通常使用 fetch 返回的 ReadableStream:
typescript
const response = await fetch('https://api.example.com/v1/chat/completions', {
method: 'POST',
headers: { 'Content-Type': 'application/json' },
body: JSON.stringify({
model: 'your-model-id',
messages,
stream: true,
}),
});
if (!response.body) {
throw new Error('response has no body');
}
const reader = response.body.getReader();
const decoder = new TextDecoder();
let buffer = '';
while (true) {
const { done, value } = await reader.read();
if (done) break;
buffer += decoder.decode(value, { stream: true });
const lines = buffer.split('\n');
buffer = lines.pop() ?? '';
for (const line of lines) {
if (!line.startsWith('data:')) continue;
const data = line.slice(5).trim();
if (data === '[DONE]') return;
const chunk = JSON.parse(data);
const content = chunk.choices?.[0]?.delta?.content;
if (content) {
process.stdout.write(content);
}
}
}这里有两个需要注意的问题。第一,decoder.decode(value, { stream: true }) 告诉 TextDecoder 输入可能是不完整的字节序列。如果不设置这个选项,当一个 UTF-8 多字节字符恰好被拆到两个网络分块中时,第二个分块会被解码为乱码。第二,reader.read() 每次返回的字节长度不等,一帧 SSE 数据可能被拆到多次 read() 中,因此需要维护一个 buffer,按 \n 切分,并把最后一个不完整的行留在 buffer 中。
工具调用与多角色增量
当模型响应包含工具调用时,增量字段不只是 content。OpenAI 兼容接口中,delta 可能携带 tool_calls 数组,数组元素通过 index 区分不同工具调用,调用的参数文本也以字符串片段形式逐帧返回,需要按 index 累积。此外,在工具调用的多轮链中,流式响应的角色边界可能比内容边界更复杂。由于 delta.role 在部分实现中并不总是 assistant,按角色切割消息流的客户端需要处理这种情况 [1]。
WebSocket 场景
SSE 是单向通道,适合服务端向客户端推送事件。如果客户端在流式过程中需要发送控制信号(例如取消、切换模型),可以使用 WebSocket。WebSocket 没有标准的事件格式,需要自行定义。常见做法是让每个消息带一个 type 字段,例如:
json
{ "type": "delta", "id": "msg_123", "content": "你" }
{ "type": "done", "id": "msg_123", "finish_reason": "stop" }
{ "type": "error", "id": "msg_123", "message": "timeout" }WebSocket 的接入代价高于 SSE,但它允许客户端中途发送取消指令,而不需要额外建立一个控制连接。
异步任务状态机
状态转换
允许的转换关系如下:
text
pending -> running | cancelled
running -> succeeded | failed | cancelled | needs-action
needs-action -> running | cancelled
failed -> pending | cancelled
succeeded -> (终态)
cancelled -> (终态)needs-action 是“暂停但未结束”的状态。任务从 running 进入 needs-action,外部输入到达后回到 running;如果外部输入永远不来,任务可以被取消。
状态转换可以用一个纯函数实现,职责是校验合法性并生成新状态:
typescript
const ALLOWED_TRANSITIONS: Record<TaskStatus, TaskStatus[]> = {
pending: ['running', 'cancelled'],
running: ['succeeded', 'failed', 'cancelled', 'needs-action'],
'needs-action': ['running', 'cancelled'],
failed: ['pending', 'cancelled'],
succeeded: [],
cancelled: [],
};
function transition(task: Task, next: TaskStatus): Task {
if (!ALLOWED_TRANSITIONS[task.status].includes(next)) {
throw new Error(
`invalid task transition: ${task.status} -> ${next}`,
);
}
return { ...task, status: next, updatedAt: Date.now() };
}这个纯函数不包含副作用。实现业务逻辑时,状态转换会伴随发起请求、写入结果、安排重试等动作。把它单独抽取出来,是为了让状态流转的合法性可以在单元测试中直接断言。
任务编排
状态机定义了任务“可以怎么走”,任务编排决定任务“实际怎么走”。
重试
进入 failed 的任务可以选择重试。重试时必须更新 attempt 计数。下面的执行器把一次尝试封装在 execute 回调中:
typescript
const sleep = (ms: number) =>
new Promise<void>((resolve) => setTimeout(resolve, ms));
async function runWithRetry(
task: Task,
execute: (task: Task) => Promise<void>,
): Promise<Task> {
let current: Task = { ...task, status: 'pending' };
while (current.attempt < current.maxAttempts) {
current = transition(current, 'running');
try {
await execute(current);
return transition(current, 'succeeded');
} catch (err) {
current = transition(current, 'failed');
current = {
...current,
attempt: current.attempt + 1,
updatedAt: Date.now(),
};
if (current.attempt >= current.maxAttempts) {
return current;
}
current = transition(current, 'pending');
await sleep(Math.min(1000 * 2 ** current.attempt, 10000));
}
}
return current;
}attempt 的递增位置是关键。如果重试时忘记更新 attempt,while 条件永远不会满足退出条件,执行器会无限重试。指数退避的延迟时间基于更新后的 attempt 计算:第一次重试延迟 1 秒,第二次 2 秒,第三次 4 秒,上限 10 秒。最后一次尝试失败后,attempt 等于 maxAttempts,函数返回处于 failed 状态的任务,不会继续循环。
超时与取消
超时控制可以用 AbortController 实现。把 controller.signal 传给 fetch,超时后调用 abort():
typescript
async function runWithTimeout(
task: Task,
timeoutMs: number,
execute: (signal: AbortSignal) => Promise<void>,
): Promise<Task> {
const controller = new AbortController();
const timer = setTimeout(() => controller.abort(), timeoutMs);
let current = transition(task, 'running');
try {
await execute(controller.signal);
return transition(current, 'succeeded');
} catch (err) {
if (err instanceof Error && err.name === 'AbortError') {
return transition(current, 'cancelled');
}
return transition(current, 'failed');
} finally {
clearTimeout(timer);
}
}AbortError 是 abort() 触发的错误类型名。取消和超时都会触发同一个错误,业务上把主动取消标记为 cancelled,把超时标记为 failed 还是 cancelled,取决于产品语义。上面的示例把超时视为取消。
取消流式任务的关键是让底层请求响应取消信号,而不是只改任务状态。如果只把任务改为 cancelled 而不中止 fetch,服务端仍然会继续生成内容,浪费资源。
幂等
幂等用于解决“请求发出后没有收到响应,于是再次请求”的重复问题。客户端为每个请求生成幂等键,服务端保存已处理的幂等键,遇到相同键时直接返回已有结果:
typescript
const taskStore = new Map<string, Task>();
function createTask(request: {
idempotencyKey: string;
messages: Message[];
}): Task {
const existing = taskStore.get(request.idempotencyKey);
if (existing && existing.status !== 'failed') {
return existing;
}
const task: Task = {
id: crypto.randomUUID(),
status: 'pending',
attempt: 0,
maxAttempts: 3,
createdAt: Date.now(),
updatedAt: Date.now(),
};
taskStore.set(request.idempotencyKey, task);
return task;
}这个幂等只覆盖任务层。如果重试导致 LLM API 被调用了两次,应用层无法撤销第一次调用产生的费用。应用层幂等的作用是让客户端看到一致的任务状态,而不是从底层消除重复计费。
数据同步策略
快照、增量与乐观更新
会话状态同时存在于客户端和服务端。客户端需要快速渲染,服务端需要持久化。数据同步的三个基本策略是快照、增量和乐观更新。
快照是某一时刻的完整状态。客户端刷新页面后,不可能只靠增量事件重建整个会话,需要先拉取一个快照作为基线。快照的结构通常是消息列表加上正在执行的任务:
json
{
"session": { "id": "sess_1" },
"messages": [
{
"id": "msg_1",
"role": "user",
"content": "你好",
"status": "completed"
}
],
"tasks": {
"task_1": { "id": "task_1", "status": "running", "attempt": 0 }
}
}增量是在快照之上应用的变化。流式输出是典型的增量场景:每个 delta 只改变一条消息的 content。用 reducer 处理增量事件时,delta 事件应该携带增量文本,并在 reducer 内部执行追加:
typescript
interface ChatSnapshot {
messages: Message[];
tasks: Record<string, Task>;
}
type ChatEvent =
| { type: 'message.created'; message: Message }
| { type: 'message.updated'; id: string; patch: Partial<Message> }
| { type: 'message.delta'; id: string; delta: string }
| { type: 'task.updated'; task: Task };
function chatReducer(state: ChatSnapshot, event: ChatEvent): ChatSnapshot {
switch (event.type) {
case 'message.created':
return { ...state, messages: [...state.messages, event.message] };
case 'message.updated':
return {
...state,
messages: state.messages.map((m) =>
m.id === event.id ? { ...m, ...event.patch } : m,
),
};
case 'message.delta':
return {
...state,
messages: state.messages.map((m) =>
m.id === event.id && m.status === 'streaming'
? { ...m, content: m.content + event.delta }
: m,
),
};
case 'task.updated':
return {
...state,
tasks: { ...state.tasks, [event.task.id]: event.task },
};
}
}这里的 message.delta 不能用 message.updated 的 patch 来替代。如果服务端把“当前完整 content”放进事件,客户端在乱序或事件丢失时很难判断是否应该覆盖;而 message.delta 是追加语义,客户端只需要保证事件顺序。
乐观更新是指客户端在服务端确认之前先渲染预期状态。用户发送消息时,如果等待服务端回包再显示,界面会让用户感觉“卡住”。通常的流程是:客户端创建一条 user 消息,立即写入本地状态并渲染,同时向服务端发送请求:
typescript
function sendMessage(
store: { dispatch: (e: ChatEvent) => void },
text: string,
): void {
const optimisticMessage: Message = {
id: crypto.randomUUID(),
role: 'user',
content: text,
status: 'completed',
createdAt: Date.now(),
};
store.dispatch({ type: 'message.created', message: optimisticMessage });
// 发送到服务端,服务端会以相同 id 去重或确认
fetch('/api/chat', {
method: 'POST',
headers: { 'Content-Type': 'application/json' },
body: JSON.stringify({ message: optimisticMessage }),
});
}乐观更新会让客户端先显示用户消息,然后等待 assistant 消息流式到达。如果服务端返回的消息 ID 与客户端不同,需要在收到服务端确认后更新本地 ID,否则后续增量无法定位到消息。
冲突处理
当多个客户端或服务端同时修改状态时,需要冲突处理机制。聊天状态中,消息本身以追加写为主,冲突概率较低,但任务状态和消息顺序仍然可能出现竞争。
冲突处理的主要机制包括:
- 客户端 ID:消息 ID 由客户端生成,服务端按 ID 去重。客户端重试或网络重放时,不会插入重复消息。
- 版本号:每条消息或任务保存版本号,更新时递增。提交方携带期望的旧版本号,服务端发现版本不匹配则拒绝更新。
- 事件日志:所有变更以追加事件的形式记录,客户端从某个游标开始顺序消费。事件日志天然提供全序关系,且支持重建状态。
最终一致性
最终一致性的含义是:如果停止新的写入,经过一段时间后,所有副本会收敛到相同状态。事件日志是实现最终一致的常见基础。客户端首次加载时拉取快照,之后订阅增量事件;重新连接时带上游标,服务端返回游标之后的事件:
typescript
interface SyncState {
cursor: number;
snapshot: ChatSnapshot;
}
async function applyEvents(
state: SyncState,
baseUrl: string,
): Promise<SyncState> {
const res = await fetch(`${baseUrl}/events?cursor=${state.cursor}`);
const events = await res.json();
let next = state;
for (const event of events) {
next = {
cursor: event.seq,
snapshot: chatReducer(next.snapshot, event.payload),
};
}
return next;
}这里 cursor 是服务端分配的事件序号,不是客户端时间戳。不同设备的时钟可能不同,时间戳不能作为事件顺序的依据。事件序号要求单调递增,可以由数据库自增 ID 或分布式序列服务生成。
状态存储选型
状态存储需要保存的对象包括会话元数据、消息内容、任务状态和同步事件日志。不同对象的访问模式不同,因此在实际系统中通常会混合使用多种存储。
内存存储
内存存储适合开发阶段的原型或单进程应用。内存读写最快,但进程重启后状态全部丢失,也不支持多实例共享。
Redis
Redis 适合保存运行中的任务状态和最近会话快照。任务状态可以用 Hash 结构保存,以任务 ID 为字段,以 JSON 字符串为值:
typescript
import { createClient } from 'redis';
const redis = createClient();
// 需要先 await redis.connect()
async function saveTask(task: Task): Promise<void> {
await redis.hSet('tasks', task.id, JSON.stringify(task));
}
async function loadTask(taskId: string): Promise<Task | null> {
const raw = await redis.hGet('tasks', taskId);
return raw ? (JSON.parse(raw) as Task) : null;
}Redis 的 TTL 机制适合清理过期任务,但需要注意:正在运行的长任务如果 TTL 太短,状态会在运行中途被清除。运行中的任务应该使用滑动过期,即每次状态更新时重新设置 TTL,而不是设置固定 TTL。
关系型数据库
关系型数据库(PostgreSQL、MySQL 等)适合保存完整的历史消息和任务记录。消息表通常按会话 ID 查询,任务表按任务 ID 查询。相比 Redis,关系型数据库提供事务和 SQL 查询,但高频增量写入需要合理设计表和索引。
事件日志
事件日志是一种追加型存储,可以建立在数据库表或专门的消息队列上。每条日志记录一个状态变更事件,带单调递增序号。事件日志的价值在于可以回溯和重建状态:如果某个客户端落后,可以从上游游标开始重放事件;如果发现数据错误,可以追溯导致错误的变更记录。事件日志的缺点是存储量增长快,因此通常配合定期快照使用:快照提供基线,日志提供快照之后的增量。
消息内容与任务状态建议分开存储。消息按会话顺序读取,访问模式是“按会话查询列表”;任务按 ID 查询,访问模式是“按 ID 获取最新状态”。把它们放在同一张表或同一个 Key 空间里,会导致查询互相干扰。
示例:React 与 Node.js 协同实现
前端状态管理与服务端状态管理可以共享同一个事件模型。下面给出一个最小可运行的结构:React 负责渲染,Node.js 负责转发上游 SSE。
创建框架无关的 Store
typescript
type Listener = () => void;
class ChatStore {
private state: ChatSnapshot;
private listeners = new Set<Listener>();
constructor(initialState: ChatSnapshot) {
this.state = initialState;
}
getSnapshot = (): ChatSnapshot => this.state;
subscribe = (listener: Listener): (() => void) => {
this.listeners.add(listener);
return () => this.listeners.delete(listener);
};
dispatch(event: ChatEvent): void {
this.state = chatReducer(this.state, event);
this.listeners.forEach((listener) => listener());
}
}getSnapshot 返回 this.state,只有 dispatch 产生新对象时引用才变化。这满足 useSyncExternalStore 对快照稳定性的要求。
在 React 中订阅 Store
React 组件用 useSyncExternalStore 订阅 store:
tsx
import { useSyncExternalStore } from 'react';
function ChatView({ store }: { store: ChatStore }) {
const state = useSyncExternalStore(store.subscribe, store.getSnapshot);
return (
<div>
{state.messages.map((message) => (
<div key={message.id}>
[{message.role}] {message.content}
{message.status === 'streaming' && <span>▍</span>}
</div>
))}
</div>
);
}useSyncExternalStore 要求 getSnapshot 返回的引用在状态没有变化时保持稳定。上面的 ChatStore 总是返回 this.state,因此符合要求。
发送消息并处理流式响应
typescript
async function sendMessage(store: ChatStore, text: string): Promise<void> {
const previousMessages = store.getSnapshot().messages;
const userMessage: Message = {
id: crypto.randomUUID(),
role: 'user',
content: text,
status: 'completed',
createdAt: Date.now(),
};
store.dispatch({ type: 'message.created', message: userMessage });
const assistantMessage: Message = {
id: crypto.randomUUID(),
role: 'assistant',
content: '',
status: 'streaming',
createdAt: Date.now(),
};
store.dispatch({ type: 'message.created', message: assistantMessage });
const task: Task = {
id: crypto.randomUUID(),
status: 'running',
attempt: 0,
maxAttempts: 3,
createdAt: Date.now(),
updatedAt: Date.now(),
};
store.dispatch({ type: 'task.updated', task });
try {
const events = fetchChatCompletion([...previousMessages, userMessage]);
for await (const event of events) {
if (event.type === 'text') {
store.dispatch({
type: 'message.delta',
id: assistantMessage.id,
delta: event.text,
});
} else if (event.type === 'finish') {
store.dispatch({
type: 'message.updated',
id: assistantMessage.id,
patch: { status: 'completed' },
});
store.dispatch({
type: 'task.updated',
task: { ...task, status: 'succeeded', updatedAt: Date.now() },
});
}
}
} catch {
store.dispatch({
type: 'message.updated',
id: assistantMessage.id,
patch: { status: 'error' },
});
store.dispatch({
type: 'task.updated',
task: { ...task, status: 'failed', updatedAt: Date.now() },
});
}
}fetchChatCompletion 是上一节中的 SSE 解析函数,事件类型是上层自定义的 { type: 'text' } 和 { type: 'finish' }。
服务端 SSE 转发
浏览器端不能直接连接需要 API 密钥的 LLM 服务,通常由 Node.js 服务端转发。服务端解析上游 SSE,并把 text 和 finish 事件以 SSE 格式转发给客户端:
typescript
import http from 'node:http';
http.createServer(async (req, res) => {
res.writeHead(200, {
'Content-Type': 'text/event-stream',
'Cache-Control': 'no-cache',
Connection: 'keep-alive',
});
const messages = JSON.parse(/* 从请求体中读取 */);
const events = fetchChatCompletion(messages);
for await (const event of events) {
res.write(`data: ${JSON.stringify(event)}\n\n`);
}
res.end();
}).listen(3000);客户端使用 fetch 读取服务端转发流:
typescript
const response = await fetch(`/api/chat/${sessionId}/events`, {
method: 'POST',
headers: { 'Content-Type': 'application/json' },
body: JSON.stringify({ messages }),
});
const reader = response.body!.getReader();
const decoder = new TextDecoder();
let buffer = '';
while (true) {
const { done, value } = await reader.read();
if (done) break;
buffer += decoder.decode(value, { stream: true });
const lines = buffer.split('\n');
buffer = lines.pop() ?? '';
for (const line of lines) {
if (!line.startsWith('data:')) continue;
const data = line.slice(5).trim();
if (data === '[DONE]') return;
const event = JSON.parse(data);
store.dispatch(event);
}
}实际路由还需要处理请求体解析和错误响应,这里只给出了 SSE 转发的核心循环。事件应尽快写入,不要在服务端批量缓冲再一次性发送,否则客户端感知不到流式效果。
可观测性与调试
状态流转和流式事件都是异步的,出现问题后很难通过一次调用栈定位。可观测性应该覆盖三个层面:状态转换日志、事件追踪和状态流转监控。
状态转换日志
状态转换日志记录每次任务状态迁移。日志至少要包含任务 ID、旧状态、新状态、attempt 计数和时间戳:
typescript
function logTransition(task: Task, next: TaskStatus): void {
console.info({
event: 'task.transition',
taskId: task.id,
from: task.status,
to: next,
attempt: task.attempt,
timestamp: Date.now(),
});
}attempt 是区分“第一次尝试”和“重试”的关键字段。缺少 attempt 的日志无法判断一次失败是用户操作导致还是不断重试导致。
消息增量日志可以在调试时记录每次 delta 的内容长度和累积长度。内容本身较长,不需要完整输出,打印长度和消息 ID 就足够定位问题。
事件追踪
追踪用于串联一次完整请求中的多个环节:客户端事件、服务端任务、上游 LLM 调用。可以在请求入口生成一个 traceId,通过 AsyncLocalStorage 保存,在日志中输出。这样,当任务未按预期完成时,可以根据 traceId 搜索出这次请求在哪个环节中断。
状态流转监控
状态流转监控包括两类指标:非法转换发生的次数,以及状态停留时长。非法转换说明程序走到了设计之外的分支,通常是一个缺陷。状态停留时长可以用于判断上游 LLM 的响应速度:从 running 到 succeeded 的时长就是模型生成耗时。
测试辅助
调试异步逻辑时,SSE mock 可以让测试不依赖真实 API。本地起一个返回固定 SSE 片段的服务器,前端代码不需要修改就可以验证增量累积逻辑:
typescript
http.createServer((req, res) => {
const chunks = ['你', '好', ',', '世', '界'];
res.writeHead(200, { 'Content-Type': 'text/event-stream' });
chunks.forEach((text, i) => {
setTimeout(() => {
res.write(`data: ${JSON.stringify({ delta: { content: text } })}\n\n`);
if (i === chunks.length - 1) {
res.end();
}
}, i * 100);
});
}).listen(3000);重试逻辑中的指数退避会让测试时间过长。使用假计时器(如 Sinon 的 fake timers 或 Vitest 的 vi.useFakeTimers())可以在不等待真实时间的情况下验证重试过程。
状态机的合法转换可以用单元测试直接断言:
typescript
assert.throws(() => transition(pendingTask, 'succeeded'));
assert.doesNotThrow(() => transition(pendingTask, 'running'));非法转换应该在开发阶段被发现,而不是等用户界面表现出异常后才开始排查。
小结
LLM 应用的状态管理由三个模型构成。消息流模型处理模型输出的增量到达,定义消息对象、角色和流式协议。任务状态机模型跟踪一次模型调用的执行过程,定义生命周期和状态转换。数据同步模型在客户端与服务端之间传递快照和增量事件,处理乐观更新和冲突。
三者通过事件相互驱动:任务推进产生消息增量,消息增量完成驱动任务状态迁移,同步事件把这两类变化传播到所有客户端。设计状态存储时,运行中任务、消息历史和事件日志的访问模式各不相同,可以分别选择内存、Redis、关系型数据库或事件日志。实现层面,React 的 useSyncExternalStore 与框架无关的 reducer/store 可以组合成一致的消息渲染管线;Node.js 负责把上游 SSE 解析为业务事件,再通过 SSE 或 WebSocket 转发给前端。可观测性是状态模型的一部分,状态转换日志和事件追踪需要从一开始就嵌入到状态流转代码中。
