niuzj
Command Palette

Search for a command to run...

Blog

AI Agent SSE 可恢复事件流设计:任务与连接生命周期解耦

把任务生命周期和 HTTP 连接生命周期解耦,让 SSE 断线后能从断点续读。

最开始暴露出来的问题很普通:前端页面偶尔显示“当前服务异常,请稍后重试”。

但诊断信息里不是一个明确的业务错误,而是 AGT_STREAM_ERRORnetwork errorFailed to fetch 这类传输层异常。更麻烦的是,它经常出现在长任务里:Agent 已经搜索完、生成完一部分内容,接下来还要跑工作流、等图像模型、写分镜。连接断了,后台不一定真的失败。

这类问题如果只从前端看,很容易得出一个朴素结论:给 SSE 加重试。

后来真正改的时候,我发现更核心的问题不是“有没有重试”,而是后端架构里把两件事绑得太紧:

  • Agent run:后台任务的生命周期
  • SSE connection:前端订阅事件的传输连接

用户网络抖一下,SSE 可以断;页面刷新,SSE 可以重连;浏览器切后台,连接也可能被代理或网关影响。可是 Agent run 不能因为这条 HTTP 连接断了就被杀。

新的原则很简单:run 继续跑,事件先进 Redis Stream,SSE endpoint 只是 reader。

Agent SSE 可恢复流架构

后端先拆成三层

以前最容易写出来的是一个大的 /chat 流式接口:请求进来,创建 Agent,边跑边往 response 写。

这会让“连接是否还在”和“任务是否还在”混在一起。连接断了以后,后端到底是取消任务、继续任务、还是等前端重连,就会变得很含糊。

现在后端拆成三层。

第一层是 Run API

POST /v1/agent/runs 只负责创建 run,返回 run_id。它不负责把整个回答流完。

这里有几个约束:

  • client_request_id 做幂等,避免前端重复点击或重试创建多个 run。
  • 同一个 session 同时只允许一个 active run。
  • 如果已经有 active run,就返回 409 ACTIVE_RUN_EXISTSactive_run_id,前端直接订阅旧 run。
  • 如果用户在 POST /runs 返回前点取消,用 cancel-by-requestsession_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:

  1. 校验用户能否访问这个 run。
  2. 读取 Last-Event-ID header,缺失时用 "0-0"
  3. 调 Redis XREAD,从 cursor 后面读事件。
  4. 把 Redis event 格式化成标准 SSE。
  5. 看到 run.completedrun.failedrun.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.idmsg.eventmsg.data
  • 维护最近 500 个 SSE id 的去重窗口,防御重复渲染。
  • 网络错误、4295xx 指数退避重连。
  • 400401403404410 不重试。
  • agent.failedrun.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。

这套东西不是为了让代码更酷。

它只是承认一个事实:长任务里,网络连接一定会断。真正要保证的不是连接永远不断,而是连接断了以后,任务还在,事件不丢,用户能继续看见结果。

Command Palette

Search for a command to run...