把任务生命周期和 HTTP 连接生命周期解耦,让 SSE 断线后能从断点续读。
最开始暴露出来的问题很普通:前端页面偶尔显示“当前服务异常,请稍后重试”。
但诊断信息里不是一个明确的业务错误,而是 AGT_STREAM_ERROR、network error、Failed to fetch 这类传输层异常。更麻烦的是,它经常出现在长任务里:Agent 已经搜索完、生成完一部分内容,接下来还要跑工作流、等图像模型、写分镜。连接断了,后台不一定真的失败。
这类问题如果只从前端看,很容易得出一个朴素结论:给 SSE 加重试。
后来真正改的时候,我发现更核心的问题不是“有没有重试”,而是后端架构里把两件事绑得太紧:
- Agent run:后台任务的生命周期
- SSE connection:前端订阅事件的传输连接
用户网络抖一下,SSE 可以断;页面刷新,SSE 可以重连;浏览器切后台,连接也可能被代理或网关影响。可是 Agent run 不能因为这条 HTTP 连接断了就被杀。
新的原则很简单:run 继续跑,事件先进 Redis Stream,SSE endpoint 只是 reader。
后端先拆成三层
以前最容易写出来的是一个大的 /chat 流式接口:请求进来,创建 Agent,边跑边往 response 写。
这会让“连接是否还在”和“任务是否还在”混在一起。连接断了以后,后端到底是取消任务、继续任务、还是等前端重连,就会变得很含糊。
现在后端拆成三层。
第一层是 Run API。
POST /v1/agent/runs 只负责创建 run,返回 run_id。它不负责把整个回答流完。
这里有几个约束:
client_request_id做幂等,避免前端重复点击或重试创建多个 run。- 同一个 session 同时只允许一个 active run。
- 如果已经有 active run,就返回
409 ACTIVE_RUN_EXISTS和active_run_id,前端直接订阅旧 run。 - 如果用户在
POST /runs返回前点取消,用cancel-by-request按session_id + client_request_id预取消,避免后台 run 刚创建就继续跑。
第二层是 Background Agent Producer。
它和 SSE 连接无关。run 创建成功以后,后台任务继续执行 Agent loop、tool call、workflow,然后把事件写进 Redis Stream。
第三层是 SSE Events Endpoint。
GET /v1/agent/runs/{run_id}/events 不创建任务,也不执行 Agent。它只读 Redis Stream,把事件转成标准 SSE 推给前端。
这个拆法一旦成立,很多边界都会自然变清楚:连接只是订阅,run 才是任务。
Redis Stream 是运行期日志
Redis Stream 不是最终聊天数据库。
它更像一段短期、可恢复的运行期事件日志。每个 run 一条 stream,事件形态是:
id: 1778494100696-0
event: chat.message.delta
data: {"text":"hello"}后端写事件时用 XADD。Redis 返回的 stream id 就是全局恢复游标。
redis_id = await redis.xadd(
f"agent:run:{run_id}:events",
{
"event": event_name,
"data": json.dumps(data, ensure_ascii=False),
},
maxlen=2000,
approximate=True,
)这样做有两个好处。
第一,事件天然有序。前端不用自己发明 seq,也不用猜哪条 delta 在前。
第二,断线后可以补。前端重连时带 Last-Event-ID,后端从这个 id 后面继续 XREAD。
Redis Stream 只保短期窗口,比如 24 小时,或者每个 run 最近 2000 条事件。窗口之外的完整聊天历史仍然进 MongoDB。Redis 负责“运行中恢复”,MongoDB 负责“最终历史”。
如果前端拿着一个太旧的 Last-Event-ID 回来,Redis 里已经没有对应窗口,就不要假装能恢复。后端应该返回 410 RESUME_WINDOW_EXPIRED,前端停止重试,提示刷新会话历史。
/events 不是 Agent 协程
这个点很容易误解。
前端发起的 /v1/agent/runs/{run_id}/events 请求,确实会在后端创建一个处理这条 HTTP 长连接的协程。但这个协程不是 Agent 本体,也不负责生成内容。
它只是一个 SSE handler:
- 校验用户能否访问这个 run。
- 读取
Last-Event-IDheader,缺失时用"0-0"。 - 调 Redis
XREAD,从 cursor 后面读事件。 - 把 Redis event 格式化成标准 SSE。
- 看到
run.completed、run.failed、run.cancelled后关闭连接。
简化后像这样:
RUN_TERMINAL_EVENTS = {"run.completed", "run.failed", "run.cancelled"}
async def event_stream(run_id: str, last_event_id: str | None):
cursor = last_event_id or "0-0"
while True:
events = await redis_store.read_events(
run_id,
last_event_id=cursor,
block_ms=15000,
)
if not events:
yield format_sse_comment("heartbeat")
continue
for item in events:
cursor = item.redis_id
yield format_sse_event(
event_id=item.redis_id,
event_name=item.event,
data=item.data,
)
if item.event in RUN_TERMINAL_EVENTS:
return心跳也不写进 Redis Stream。它只是 SSE comment:
: heartbeat这类 comment 是保活连接用的,不应该进入业务事件处理。
标准 SSE 形态
后端输出标准 SSE:
id: 1778494100696-0
event: chat.message.delta
data: {"session_id":"...","message_id":"...","text":"hello"}
三个字段各自有职责:
id是 Redis Stream id,也是恢复游标。event是业务事件名,比如chat.message.delta。data是 JSON payload。
多行 JSON、中文、换行都不要靠散落在各处的字符串拼接处理。后端应该有统一 formatter。
def format_sse_event(event_id: str, event_name: str, data: dict) -> str:
payload = json.dumps(data, ensure_ascii=False, separators=(",", ":"))
lines = [f"id: {event_id}", f"event: {event_name}"]
for line in payload.splitlines() or [""]:
lines.append(f"data: {line}")
return "\n".join(lines) + "\n\n"这个 helper 看起来小,但它避免了很多很难查的协议问题。
agent.* 和 run.* 分开
这次改造里,最重要的状态语义是把业务状态和生命周期状态拆开。
agent.* 是业务/UI 事件:
agent.completed
agent.failed
agent.stopped它告诉前端:Agent 内容层面完成了、失败了、停止了。前端可以更新消息状态,但 SSE handler 不一定立刻关闭。
真正让 SSE 关闭的是 run.*:
run.completed
run.failed
run.cancelled事件顺序固定下来以后,前后端都少很多猜测。
成功:
chat.message.completed
agent.completed
run.completed失败:
agent.failed
run.failed取消:
agent.stopped
run.cancelled为什么不让 agent.completed 直接关连接?
因为 agent.completed 只是业务内容完成。后端可能还要写最终消息、清理 active run、持久化状态。用 run.completed 作为唯一生命周期终态,前端就知道什么时候停止重连,后端也知道什么时候清理 run。
HITL 还有一个细节:run 可以完成这一段输出,但 session 状态进入 waiting_hitl。这时 SSE 可以因为 run.completed 关闭,页面上仍然保留等待用户选择的状态。下一次用户选择后,再创建下一轮 run。
active-run 只用于恢复
active-run 不能当轮询接口。
如果每次 React 状态变化都查一次 active-run,Network 面板里会出现一串重复请求。更糟的是,发消息前先查 active-run,再创建 run,会把并发窗口拉长。
现在更合适的策略是:
- 页面刷新、进入 session、历史加载完成后,查一次 active-run。
- 如果有
active_run_id,直接订阅/runs/{run_id}/events。 - 发送消息前不主动查 active-run。
- 发送消息直接
POST /runs。 - 如果后端返回
409 ACTIVE_RUN_EXISTS,前端移除本次 optimistic 消息,订阅后端返回的active_run_id。
也就是说,active-run 是“页面恢复入口”,不是“状态轮询接口”。
前端把传输层交给库
前端原来自己维护 response.body.getReader()、SSE parser、\n\n 分割、重连循环。这个复杂度很容易和业务逻辑搅在一起。
现在 /events 这一条连接交给 @microsoft/fetch-event-source。
前端只保留几件事:
onopen校验 HTTP 状态和content-type: text/event-stream。onmessage读取msg.id、msg.event、msg.data。- 维护最近 500 个 SSE id 的去重窗口,防御重复渲染。
- 网络错误、
429、5xx指数退避重连。 400、401、403、404、410不重试。agent.failed、run.failed是业务失败,不当网络错误重试。
第一次 /events 请求看不到 Last-Event-ID 是正常的,因为前端还没有收到任何事件。
只有断线重连时,请求头里才会出现:
Last-Event-ID: 1778494100696-0后端拿这个 id 继续读 Redis Stream。
前端还有一个很实际的坑:不要在同一条 chat.message.delta 上再做异步“打字机切片”。后端已经把 delta 切好了,前端如果再 await sleep() 逐字追加,多个 delta 的处理可能并发交错,最后出现 URL 被拼断、文本错位。这里应该让后端 delta 成为最小显示单位,前端直接 append。
为什么不是 Pub/Sub
Pub/Sub 更快,但它没有恢复窗口。
订阅者断线期间发布的消息就丢了。为了补消息,你还得再做一份事件缓存。最后系统会变成 Pub/Sub + Stream 或 Pub/Sub + DB,复杂度反而更高。
这个场景要解决的不是最低延迟广播,而是:
- 长任务不断线也能继续跑。
- 断线后能从上次事件继续。
- 同一个 run 的事件有序。
- 前端不重复渲染。
- 恢复窗口过期时明确失败。
Redis Stream 正好卡在这个位置。
单副本和多副本的边界
当前如果是单副本,进程内还能知道某个 background task 是否存在。
但这不是长期架构。
一旦 Railway 多实例、进程重启、或者 SSE 重连打到另一台实例,本进程没有 task 并不代表 run 已经死了。更稳的做法是 runner lease:
- run 启动时写
agent:run:{run_id}:lease - value 里放
instance_id - TTL 比如 30 秒
- background task 每 10 秒刷新一次
- 只有 run 仍是
running,并且 lease 过期,才追加agent.failed+run.failed
这不是为了让 SSE 更复杂,而是为了避免“连接重连到另一台机器”被误判成“任务中断”。
单副本阶段可以先不做,但要知道问题边界在哪里。
这套架构真正改变了什么
它把“用户看见的流”从“任务本身”降级成“任务事件的订阅视图”。
Agent run 继续在后台跑。事件先写 Redis Stream。SSE endpoint 从 Redis 读,再输出标准 id/event/data。前端断了就用 Last-Event-ID 重连。最终聊天历史再落 MongoDB。
这套东西不是为了让代码更酷。
它只是承认一个事实:长任务里,网络连接一定会断。真正要保证的不是连接永远不断,而是连接断了以后,任务还在,事件不丢,用户能继续看见结果。