
前言
在普通 AI 对话中,前端通常发起一次请求,后端调用大模型,并通过 SSE 将生成结果逐步推送给浏览器。
用户发送消息
↓
后端创建 Agent Run
↓
Agent 开始执行
↓
SSE 持续推送事件
↓
前端实时更新对话界面
只要用户一直停留在当前页面,这套流程并不复杂。
真正棘手的问题出现在以下场景:
- Agent 正在输出时,用户切换到了另一个会话;
- 用户进入其他页面,稍后又返回;
- 浏览器刷新;
- 网络短暂断开;
- 前端重新部署导致页面重载;
- 同一个会话在多个标签页打开;
- Agent 后台执行完了,但原来的 SSE 连接已经不存在。
例如,Agent 已经执行到一半:
用户:分析这个项目的技术架构
Agent:
1. 正在读取项目文件
2. 正在分析服务依赖
3. 正在调用代码分析工具
4. 正在生成架构总结……
这时用户刷新页面。
如果系统只依赖当前 SSE 连接,页面刷新后可能出现:
- 已生成内容消失;
- 消息一直显示“生成中”;
- 工具调用卡在运行状态;
- Agent 已经完成,但前端不知道;
- 重新连接后内容重复;
- 同一个回答出现两份;
- 无法恢复中间执行步骤。
因此,Agent 对话中的 SSE 设计,不能只解决“如何流式输出”,还必须解决:
当连接随时可能中断时,如何重建完整对话状态,并从中断位置继续接收事件。
一、先理解 SSE 在 Agent 系统中的定位
SSE,全称是 Server-Sent Events,是一种基于 HTTP 的服务器单向推送机制。
浏览器建立连接后,服务器持续返回:
Content-Type: text/event-stream
服务器可以不断发送事件:
event: assistant.delta
id: 101
data: {"text":"你好"}
event: assistant.delta
id: 102
data: {"text":",我是"}
event: assistant.delta
id: 103
data: {"text":"智能助手"}
标准 SSE 事件可以包含 event、data、id 和 retry 等字段;浏览器的 EventSource 会保持长连接,并在连接异常时尝试重新连接。
但是必须认识到:
SSE 只是事件传输通道,不应该成为 Agent 状态的唯一来源。
如果所有状态只存在于当前连接中:
Agent 输出
↓
SSE 连接
↓
浏览器内存
那么页面刷新后,浏览器内存和 SSE 连接都会消失。
更可靠的架构应该是:
持久化状态
+
可回放事件
+
实时 SSE
也就是:
数据库快照
负责恢复完整页面
事件日志
负责补发断线期间的变化
SSE
负责推送最新实时事件
二、核心设计:Snapshot + Event Log + Live Stream
一个可恢复的 Agent 对话系统,建议采用三层结构。
第一层:Snapshot
当前会话和消息的完整快照
第二层:Event Log
Agent 执行过程中产生的有序事件
第三层:Live Stream
通过 SSE 推送最新事件
1. Snapshot:状态快照
Snapshot 是当前会话已经确定的业务状态,例如:
{
"conversationId": "conv_1001",
"messages": [
{
"id": "msg_user_1",
"role": "user",
"content": "分析一下这个项目",
"status": "completed"
},
{
"id": "msg_assistant_1",
"role": "assistant",
"content": "项目采用了微服务架构……",
"status": "streaming",
"runId": "run_9001"
}
],
"activeRun": {
"id": "run_9001",
"status": "running",
"lastEventId": "1058"
}
}
它用于页面首次进入或刷新时恢复整个界面。
2. Event Log:事件日志
Event Log 记录 Agent Run 执行过程中产生的变化:
1051 run.started
1052 assistant.message.created
1053 assistant.delta
1054 assistant.delta
1055 tool.started
1056 tool.completed
1057 assistant.delta
1058 run.progress
它用于补发用户离开期间错过的事件。
3. Live Stream:实时流
当历史事件补发完成后,SSE 连接继续等待新事件:
历史事件补发
↓
追上当前最新位置
↓
阻塞等待新事件
↓
实时推送
因此,一个完整的恢复过程应该是:
读取 Snapshot
↓
获得当前 lastEventId
↓
补发 lastEventId 之后的事件
↓
进入实时监听
三、Agent 对话中的几个核心数据对象
建议至少明确区分以下四类对象。
1. Conversation
Conversation 表示一段完整会话。
conversation
├── id
├── user_id
├── title
├── status
├── created_at
└── updated_at
2. Message
Message 表示用户或 Agent 的最终消息。
message
├── id
├── conversation_id
├── role
├── content
├── status
├── run_id
├── created_at
└── updated_at
消息状态可以设计为:
pending
streaming
completed
failed
cancelled
3. Run
Run 表示 Agent 的一次执行过程。
run
├── id
├── conversation_id
├── input_message_id
├── output_message_id
├── status
├── last_event_id
├── started_at
├── finished_at
└── error
常见状态:
queued
running
waiting_tool
waiting_user
completed
failed
cancelled
4. Event
Event 表示 Run 中发生的一次状态变化。
event
├── id
├── conversation_id
├── run_id
├── sequence
├── type
├── payload
└── created_at
这四个概念不要混在一起:
Conversation:一段对话
Message:用户或 Agent 的消息
Run:Agent 的一次执行
Event:执行过程中发生的变化
一段 Conversation 可以有很多 Message。
一次用户提问通常会创建一个 Run。
一个 Run 又会产生很多 Event。
四、不要只发送文本 Token,要发送结构化事件
很多系统最初只发送一种 SSE 消息:
{
"text": "你好"
}
这种设计只能处理简单的大模型文本输出。
但 Agent 对话中通常还包含:
- 思考状态;
- 计划;
- 工具调用;
- 子 Agent;
- 文件解析;
- 进度变化;
- 审批;
- 错误;
- 重试;
- 最终结果。
因此应该设计结构化事件。
推荐事件类型包括:
run.created
run.started
run.progress
run.completed
run.failed
run.cancelled
message.created
assistant.delta
assistant.message.completed
tool.started
tool.progress
tool.completed
tool.failed
subagent.started
subagent.completed
subagent.failed
approval.required
approval.resolved
heartbeat
统一事件结构可以设计为:
{
"eventId": "1058",
"sequence": 1058,
"type": "tool.completed",
"conversationId": "conv_1001",
"runId": "run_9001",
"blockId": "block_tool_3",
"timestamp": "2026-08-04T16:10:00+08:00",
"payload": {
"toolName": "search_documents",
"status": "completed",
"resultSummary": "找到 12 个相关文档"
}
}
对应 SSE:
event: tool.completed
id: 1058
data: {"eventId":"1058","sequence":1058,"type":"tool.completed","runId":"run_9001","blockId":"block_tool_3","payload":{"toolName":"search_documents","status":"completed"}}
这里的 id 非常重要。
它不是为了展示,而是为了:
- 事件排序;
- 前端去重;
- 断线补发;
- 判断缺失事件;
- 恢复消费位置。
五、三种页面中断场景应该分别如何处理
场景一:切换到另一个会话
例如用户从:
/conversations/1001
切换到:
/conversations/1002
此时前端应该主动关闭原会话的 SSE:
eventSource.close();
关闭 SSE 只表示:
当前页面不再实时接收这个 Run 的事件。
它不应该取消后台 Agent Run。
正确关系是:
关闭页面订阅
≠
取消 Agent 执行
Agent 仍然在服务器后台运行,继续:
- 调用模型;
- 执行工具;
- 写入事件;
- 更新消息;
- 保存最终状态。
当用户再次进入会话 1001 时,再重新恢复。
场景二:切换回来重新进入会话
用户重新进入会话时,不要直接创建一个空页面然后只等待新 SSE。
正确顺序应该是:
1. 请求会话 Snapshot
2. 渲染已有消息
3. 检查是否存在 activeRun
4. 获取 lastEventId
5. 建立 SSE 连接
6. 补发遗漏事件
7. 继续实时接收
例如:
GET /api/conversations/conv_1001
返回:
{
"conversationId": "conv_1001",
"messages": [],
"activeRun": {
"id": "run_9001",
"status": "running",
"lastEventId": "1058"
}
}
然后前端建立:
GET /api/runs/run_9001/events?after=1058
如果 Agent 在用户离开期间又产生了:
1059 assistant.delta
1060 tool.started
1061 tool.completed
1062 assistant.delta
后端先补发这些事件,再继续等待 1063 之后的新事件。
场景三:浏览器刷新
刷新与普通路由切换的区别是:
React/Vue 内存状态消失
SSE 对象消失
当前 lastEventId 可能消失
未完成消息的本地文本也可能消失
因此页面刷新后的恢复不能依赖前端内存。
应该重新执行:
刷新页面
↓
查询 Conversation Snapshot
↓
恢复 Messages 和 Blocks
↓
识别 activeRun
↓
重新建立 SSE
↓
从服务端游标继续
浏览器的 EventSource 在同一个对象发生临时断线时,可以自动重连,并在重新建立连接时使用最近收到的事件 ID;相关标准定义了 Last-Event-ID 请求头。
但是完整页面刷新后,原来的 EventSource 对象已经销毁。
新页面不能只依赖旧对象的自动恢复能力,因此仍然需要显式保存和恢复游标,例如:
服务端 Snapshot 中的 lastEventId
或者:
GET /runs/{runId}/events?after={lastEventId}
六、推荐的后端接口设计
可以将“发送消息”和“订阅事件”拆成两个接口。
1. 提交消息
POST /api/conversations/{conversationId}/messages
请求:
{
"content": "请分析这个项目的技术架构"
}
响应:
{
"messageId": "msg_user_1001",
"runId": "run_9001",
"status": "queued"
}
这个接口负责:
不要让这个接口持续保持几十分钟的 HTTP 请求。
2. 获取会话快照
GET /api/conversations/{conversationId}
负责返回:
- 历史消息;
- 已完成工具调用;
- 当前运行状态;
- Assistant 当前已生成文本;
- 最近事件游标。
3. 订阅 Run 事件
GET /api/runs/{runId}/events?after={eventId}
响应类型:
text/event-stream
这个接口负责:
七、后端推荐执行流程
整体流程可以设计成:
用户发送消息
↓
API 创建 Message 和 Run
↓
后台 Agent Runtime 开始执行
↓
每发生一个动作就写入 Event Log
↓
SSE Gateway 读取 Event Log
↓
推送给浏览器
关键点是:
Agent Runtime 不应该直接依赖某一条浏览器连接。
错误架构:
Agent Runtime
↓
直接向当前 HTTP Response 写数据
这种架构一旦连接断开,后续事件就无处可去。
推荐架构:
Agent Runtime
↓
写入统一事件总线
↓
Event Log
├── 数据库投影
├── SSE Gateway
├── 日志系统
└── 监控系统
这样,即使当前没有用户在线,Agent 仍然可以继续运行并保存事件。
八、使用 Redis Streams 保存可回放事件
中小型 Agent 系统可以使用 Redis Streams 作为短期事件日志。
Redis Stream 是一种有序、可追加的日志结构,支持时间有序 ID、历史读取、阻塞读取和保留策略。Redis 官方也将 XREAD 描述为适合 UI 实时流、调试器和类似 tail -f 的只读消费场景。
每个 Run 可以对应一个 Stream:
agent:run:run_9001:events
写入事件:
XADD agent:run:run_9001:events *
type assistant.delta
data {"text":"项目采用"}
Redis 返回事件 ID:
1722768000000-0
SSE Gateway 可以从指定 ID 后开始读取:
XREAD BLOCK 15000
STREAMS agent:run:run_9001:events 1722768000000-0
其行为是:
存在历史事件
↓
立即返回
没有新事件
↓
阻塞等待
有新事件写入
↓
立即返回并推送
对于浏览器订阅,不建议使用一个共享 Consumer Group 将事件分给不同浏览器。
因为每个打开该会话的客户端通常都应该看到完整事件流,而不是多个浏览器共同瓜分事件。
因此:
任务 Worker 消费:
可以使用 XREADGROUP
浏览器 SSE 实时订阅:
通常使用 XREAD
九、Snapshot 和 Redis Stream 应如何配合
Redis Stream 不应该成为长期保存所有对话内容的唯一数据库。
推荐职责划分:
MySQL/PostgreSQL
保存 Conversation
保存最终 Message
保存 Run 状态
保存工具调用结果
保存当前 Assistant 内容
Redis Stream
保存短期实时事件
支持断线补发
支持 SSE 实时监听
Agent 输出 Token 时,不需要每个 Token 都写一次数据库。
可以采用批量策略:
模型不断产生 delta
↓
每个 delta 写入 Redis Stream
↓
前端实时显示
↓
每 300~1000 毫秒批量更新数据库
↓
完成时强制写入最终完整内容
例如:
assistant.delta:实时写事件流
assistant.message.completed:保存最终 Message
这样同时兼顾:
- 实时性;
- 数据库压力;
- 页面恢复;
- 最终一致性。
十、SSE 服务端示例
下面使用接近 FastAPI 的伪代码展示核心逻辑。
FastAPI 当前提供 SSE 响应支持,可以通过生成器持续 yield 事件,并设置 event、id、data、retry 等字段。
import asyncio
import json
from collections.abc import AsyncGenerator
from fastapi import APIRouter, Depends, Request
from fastapi.sse import EventSourceResponse, ServerSentEvent
router = APIRouter()
@router.get("/api/runs/{run_id}/events")
async def subscribe_run_events(
run_id: str,
request: Request,
after: str | None = None,
current_user=Depends(get_current_user),
) –> EventSourceResponse:
await check_run_permission(
run_id=run_id,
user_id=current_user.id,
)
async def event_generator() –> AsyncGenerator[ServerSentEvent, None]:
cursor = after or "0-0"
while True:
if await request.is_disconnected():
break
events = await event_store.read_after(
run_id=run_id,
cursor=cursor,
block_ms=15_000,
count=100,
)
if not events:
yield ServerSentEvent(
comment="heartbeat",
)
continue
for event in events:
cursor = event.id
yield ServerSentEvent(
event=event.type,
id=event.id,
data=json.dumps(
event.to_dict(),
ensure_ascii=False,
),
retry=3000,
)
if event.type in {
"run.completed",
"run.failed",
"run.cancelled",
}:
return
return EventSourceResponse(
event_generator(),
headers={
"Cache-Control": "no-cache",
"X-Accel-Buffering": "no",
},
)
这里有几个重要设计。
1. after
前端告诉服务端:
我最后处理到了哪个事件
2. 先补历史事件
如果 Redis Stream 中已经有新事件,先返回历史事件。
3. 再阻塞等待
历史事件追平后,等待新事件。
4. 心跳
长时间没有事件时发送:
: heartbeat
SSE 中以冒号开头的内容是注释,不会作为普通业务事件处理,但可以用于维持连接。
5. 终态事件后结束
收到以下事件后可以主动结束当前流:
run.completed
run.failed
run.cancelled
十一、前端进入会话时的标准流程
前端不要一进入页面就直接连接 SSE。
更合理的流程是:
async function enterConversation(conversationId: string) {
const snapshot = await conversationApi.getConversation(
conversationId,
);
conversationStore.replaceSnapshot(snapshot);
const activeRun = snapshot.activeRun;
if (!activeRun) {
return;
}
subscribeRun({
runId: activeRun.id,
after: activeRun.lastEventId,
});
}
这里首先用 Snapshot 恢复页面,再订阅后续事件。
否则可能出现:
SSE 事件先到
↓
对应 Message 和 Block 还没有创建
↓
前端不知道把事件应用到哪里
十二、前端事件必须使用 Reducer 统一处理
不要在不同监听器中随意修改 UI:
source.addEventListener("tool.started", () => {
// 随意修改某个组件
});
source.addEventListener("assistant.delta", () => {
// 再修改另一份状态
});
更推荐把全部事件交给统一 Reducer:
function applyAgentEvent(
state: ConversationState,
event: AgentEvent,
): ConversationState {
if (state.processedEventIds.has(event.eventId)) {
return state;
}
const nextState = reduceAgentEvent(state, event);
nextState.processedEventIds.add(event.eventId);
nextState.lastEventId = event.eventId;
return nextState;
}
事件处理逻辑示例:
function reduceAgentEvent(
state: ConversationState,
event: AgentEvent,
): ConversationState {
switch (event.type) {
case "assistant.delta":
return appendAssistantText(
state,
event.blockId,
event.payload.text,
);
case "tool.started":
return createToolBlock(
state,
event.blockId,
event.payload,
);
case "tool.completed":
return completeToolBlock(
state,
event.blockId,
event.payload,
);
case "run.completed":
return completeRun(
state,
event.runId,
);
case "run.failed":
return failRun(
state,
event.runId,
event.payload.error,
);
default:
return state;
}
}
这样可以保证:
- 实时事件和补发事件使用同一逻辑;
- 页面恢复逻辑一致;
- 事件可以测试;
- 重复事件不会重复渲染;
- 状态流转更加清晰。
十三、EventSource 前端示例
let currentEventSource: EventSource | null = null;
interface SubscribeRunOptions {
runId: string;
after?: string;
}
function subscribeRun({
runId,
after,
}: SubscribeRunOptions): void {
currentEventSource?.close();
const params = new URLSearchParams();
if (after) {
params.set("after", after);
}
const url = `/api/runs/${runId}/events?${params.toString()}`;
const source = new EventSource(url, {
withCredentials: true,
});
currentEventSource = source;
const eventTypes = [
"run.started",
"run.progress",
"assistant.delta",
"assistant.message.completed",
"tool.started",
"tool.completed",
"tool.failed",
"run.completed",
"run.failed",
"run.cancelled",
];
for (const eventType of eventTypes) {
source.addEventListener(eventType, (rawEvent) => {
const messageEvent = rawEvent as MessageEvent;
const event = JSON.parse(
messageEvent.data,
) as AgentEvent;
conversationStore.applyEvent(event);
if (messageEvent.lastEventId) {
conversationStore.setLastEventId(
runId,
messageEvent.lastEventId,
);
}
if (
event.type === "run.completed" ||
event.type === "run.failed" ||
event.type === "run.cancelled"
) {
source.close();
}
});
}
source.onerror = () => {
conversationStore.markReconnecting(runId);
};
}
离开页面时:
function leaveConversation(): void {
currentEventSource?.close();
currentEventSource = null;
}
这里关闭的只是当前页面连接,不是后台 Run。
十四、为什么必须去重
SSE 和分布式事件系统一般需要按照“至少一次”的思路设计。
以下场景都可能导致事件重复:
事件已发送给浏览器
↓
连接突然中断
↓
浏览器没有来得及保存游标
↓
重新连接后再次补发
或者:
前端已经处理事件 1058
↓
服务端恢复时又发送一次 1058
因此前端必须根据 eventId 去重:
if (processedEventIds.has(event.eventId)) {
return;
}
仅仅根据文本内容去重是不可靠的。
例如大模型完全可能连续生成:
哈哈
哈哈
两段文本内容相同,但它们可能是两个合法 delta。
正确去重依据应该是:
runId + eventId
或者:
runId + sequence
十五、顺序错乱如何处理
每个 Run 的事件应该有严格递增的序号:
1051
1052
1053
1054
前端收到 1054 时,如果当前只处理到了 1052,就能发现:
缺少 1053
此时不应盲目继续应用,而应触发恢复:
检测到事件断档
↓
暂停当前流
↓
请求 after=1052
↓
补发 1053、1054
↓
继续处理
事件结构中可以同时保留:
{
"eventId": "1722768000000-0",
"sequence": 1054
}
其中:
- eventId 用于事件存储和读取;
- sequence 用于业务顺序检测。
十六、页面恢复时以谁为准
页面恢复时通常会同时存在:
- 数据库 Snapshot;
- Redis Event Log;
- 浏览器本地状态。
推荐优先级:
数据库 Snapshot
作为业务基础状态
Event Log
补充 Snapshot 之后的新变化
浏览器本地缓存
只用于优化体验
不要让 LocalStorage 成为最终状态来源。
因为它可能:
- 被清理;
- 过期;
- 多标签页冲突;
- 与服务端不一致;
- 保存了错误的半成品状态。
本地缓存可以保存:
当前会话 ID
最近事件游标
未发送输入框内容
滚动位置
折叠状态
但消息和 Run 的最终状态仍然应以服务端为准。
十七、Snapshot 和 Event Log 之间的竞争问题
一个典型问题是:
前端开始请求 Snapshot
↓
Snapshot 查询完成前,Agent 又产生两个事件
↓
前端随后连接 SSE
如果处理不正确,中间两个事件可能丢失。
解决方式是让 Snapshot 返回一个一致的游标:
{
"messages": [],
"activeRun": {
"id": "run_9001",
"lastEventId": "1058"
}
}
这个响应表达:
当前 Snapshot 已经包含了截至事件 1058 的状态。
前端随后请求:
/events?after=1058
这样事件 1059 以后一定会被补发。
因此 Snapshot 中的 lastEventId 不能随便读取,它需要与 Snapshot 状态保持一致。
可以在同一个事务或投影更新过程中:
应用事件
↓
更新 Message Snapshot
↓
更新 Run.last_event_id
十八、终态事件非常重要
Agent 执行结束时,一定要发出明确的终态事件:
run.completed
run.failed
run.cancelled
不要让前端根据“很久没收到 Token”猜测任务是否完成。
成功事件:
{
"eventId": "1100",
"type": "run.completed",
"runId": "run_9001",
"payload": {
"outputMessageId": "msg_assistant_1001",
"finishReason": "stop"
}
}
失败事件:
{
"eventId": "1100",
"type": "run.failed",
"runId": "run_9001",
"payload": {
"errorCode": "TOOL_TIMEOUT",
"errorMessage": "代码分析工具执行超时",
"retryable": true
}
}
前端收到终态事件后:
十九、EventSource 还是 Fetch Streaming
使用 EventSource
浏览器原生 EventSource 的优势是:
- API 简单;
- 原生支持 SSE 格式;
- 自动重连;
- 支持事件类型;
- 支持 Last-Event-ID;
- 适合 GET 长连接。
但它也有一些限制:
- 只能使用 GET;
- 不方便自定义任意请求头;
- 不能直接携带复杂请求体;
- 手动控制重连策略的能力较弱。
使用 Fetch Streaming
如果需要:
- POST 请求;
- Bearer Token Header;
- 请求体;
- 更灵活的取消;
- 自定义重连逻辑;
可以使用:
const response = await fetch(url, {
method: "GET",
headers: {
Authorization: `Bearer ${token}`,
},
signal: abortController.signal,
});
const reader = response.body?.getReader();
Fetch 的响应体可以通过 ReadableStream 增量读取,而不必等待整个响应结束。
不过使用 Fetch Streaming 时,需要自己处理:
- SSE 文本解析;
- 断线重连;
- 游标保存;
- AbortController;
- 错误重试;
- 心跳超时。
对于普通 Cookie 鉴权的 Agent 对话,优先选择 EventSource 会更简单。
对于需要自定义 Header 的系统,可以使用 Fetch Streaming。
二十、Nginx 和网关配置
SSE 在本地正常,部署后“每隔几十秒一次性出来”,通常不是前端问题,而是代理层缓冲。
Nginx 默认可能缓冲上游响应。关闭缓冲后,数据会在接收到时直接转发给客户端;也可以通过响应头 X-Accel-Buffering: no 控制。
配置示例:
location /api/runs/ {
proxy_pass http://agent_backend;
proxy_http_version 1.1;
proxy_set_header Connection "";
proxy_buffering off;
proxy_cache off;
proxy_read_timeout 3600s;
proxy_send_timeout 3600s;
gzip off;
}
后端响应头建议包含:
Content-Type: text/event-stream
Cache-Control: no-cache
Connection: keep-alive
X-Accel-Buffering: no
同时建议每隔 15~30 秒发送心跳:
: heartbeat
避免:
- Nginx 认为连接空闲;
- 网关关闭连接;
- 负载均衡器超时;
- 浏览器长时间无法感知断线。
二十一、多实例部署时不能只使用进程内队列
单实例开发阶段,可能这样实现:
run_queues: dict[str, asyncio.Queue] = {}
Agent 将事件写入本地 asyncio.Queue,SSE 接口从中读取。
这在单进程中可以工作。
但扩容为多个实例后:
Agent Runtime 在实例 A
SSE 请求被负载均衡到实例 B
实例 B 的内存中没有实例 A 的事件。
因此多实例环境不能只依赖:
- 进程内 Queue;
- 本地 EventEmitter;
- 单机内存;
- 本地文件。
应该引入共享事件基础设施:
Redis Streams
Kafka
NATS JetStream
数据库事件表
对于大多数中小型 Agent 平台:
业务状态:PostgreSQL/MySQL
短期事件:Redis Streams
浏览器推送:SSE Gateway
通常已经能够满足需求。
二十二、推荐的完整架构
┌────────────────────┐
│ Web 前端 │
│ Snapshot + SSE恢复 │
└─────────┬──────────┘
│
├── GET Conversation Snapshot
│
└── GET Run Event Stream
│
┌─────────────────────▼─────────────────────┐
│ Agent API │
│ │
│ Conversation API SSE Gateway │
│ │ │ │
└─────────┼──────────────────┼───────────────┘
│ │
▼ ▼
┌────────────────┐ ┌──────────────────────┐
│ MySQL/Postgres │ │ Redis Streams │
│ │ │ │
│ Conversation │ │ Run Event Log │
│ Message │ │ Replay + Live Tail │
│ Run Snapshot │ │ │
└───────▲────────┘ └──────────▲───────────┘
│ │
└──────────┬────────────┘
│
┌────────▼────────┐
│ Agent Runtime │
│ │
│ LLM │
│ Tool │
│ Subagent │
│ Memory │
└─────────────────┘
其中:
Agent Runtime
负责执行
Redis Streams
负责实时事件和短期回放
数据库
负责完整业务状态
SSE Gateway
负责向浏览器推送
前端 Reducer
负责重建界面
二十三、常见错误设计
错误一:SSE 断开就取消 Agent
用户切换页面并不代表用户要求取消任务。
应区分:
unsubscribe:停止订阅
cancel:取消后台执行
错误二:只保存最终回答
如果只在 Agent 完成时保存回答,刷新时会丢失当前已经生成的内容。
应定期保存流式中的 Message Snapshot。
错误三:只保存文本,不保存工具状态
刷新后文本恢复了,但工具调用全部消失。
工具调用、子 Agent、计划和审批也应该是可持久化 Block。
错误四:重新进入后只连接最新事件
如果不补发历史事件,用户离开期间发生的工具调用和输出会丢失。
错误五:不设置事件 ID
没有事件 ID,就无法可靠去重、排序和恢复。
错误六:一个 Token 对应一次数据库写入
这会造成大量小事务和数据库压力。
应对 delta 做批量持久化。
错误七:仅依赖 Redis Pub/Sub
Pub/Sub 更适合在线广播。订阅者离线期间无法自然回放已经错过的消息。
需要恢复能力时,应使用可回放事件日志,例如 Redis Streams。
错误八:页面刷新后直接创建新的 Run
这可能让同一个用户问题执行两遍。
页面恢复应该连接已有 Run,而不是重新提交用户消息。
二十四、建议的恢复算法
最终可以把前端恢复过程总结为以下算法。
进入 Conversation 页面
↓
关闭上一个页面的 SSE
↓
获取 Conversation Snapshot
↓
用 Snapshot 替换本地状态
↓
是否存在未完成 Run?
├── 否:结束
└── 是
↓
获取 Snapshot.lastEventId
↓
连接 /events?after=lastEventId
↓
服务端补发遗漏事件
↓
前端按 eventId 去重
↓
按 sequence 检查顺序
↓
Reducer 更新 Message、Tool 和 Run
↓
继续接收实时事件
↓
收到 completed/failed/cancelled
↓
关闭 SSE
服务端恢复算法:
收到 SSE 请求
↓
校验 Run 访问权限
↓
读取 after / Last-Event-ID
↓
从 Event Log 查询之后的事件
↓
依次补发
↓
阻塞等待新事件
↓
持续推送
↓
发送心跳
↓
遇到终态事件后关闭
二十五、测试清单
这类功能不能只测试“正常生成”。
至少应该覆盖以下场景:
页面操作
- Agent 输出中切换会话;
- 切换回来;
- Agent 输出中刷新页面;
- Agent 完成后重新进入;
- 同一会话打开两个标签页;
- 快速来回切换会话。
网络异常
- SSE 短暂断线;
- 断网后恢复;
- Nginx 重启;
- 后端实例重启;
- Redis 短暂不可用。
事件一致性
- 同一事件重复发送;
- 事件中间缺失;
- 事件乱序;
- Event Log 已被清理;
- Snapshot 比事件流更新;
- Snapshot 比事件流落后。
Agent 状态
- 文本生成中刷新;
- 工具执行中刷新;
- 子 Agent 执行中刷新;
- 等待用户审批时刷新;
- Agent 失败时刷新;
- Agent 取消时刷新。
最终验收标准应该是:
无论用户何时离开、切换或刷新,再次进入时都能看到正确的完整状态,并且不会重复生成、不会丢事件、不会永远卡在运行中。
总结
Agent 对话中的 SSE 恢复问题,本质上不是一个简单的前端重连问题。
它是一个完整的分布式状态恢复问题。
可靠方案需要同时具备:
Conversation Snapshot
+
Run 状态机
+
有序 Event Log
+
事件游标
+
断线补发
+
前端幂等 Reducer
+
实时 SSE
最重要的设计原则是:
SSE 负责传输事件,但数据库和事件日志才负责保存事实。
切换页面时,可以关闭 SSE,但不要停止 Agent。
重新进入时,先加载 Snapshot,再从事件游标之后补发事件。
刷新页面时,不依赖浏览器内存,而是通过服务端状态重新构建 UI。
收到重复事件时,根据 eventId 去重。
发现事件缺口时,根据 sequence 重新补发。
Agent 完成时,必须发送明确的终态事件,并保存最终消息。
最终形成的标准链路是:
提交消息
↓
创建 Run
↓
Agent 后台执行
↓
事件写入 Event Log
↓
状态投影到数据库
↓
SSE 实时推送
↓
页面断开
↓
Snapshot 恢复
↓
事件补发
↓
继续实时接收
当这套机制建立以后,Agent 对话才真正具备生产级的可恢复性,而不只是一个“页面不刷新时能够流式输出”的演示系统。




