实时通信

Agent 为什么需要实时通信

Agent 运行不是一次普通请求。它可能持续几十秒到几分钟,过程中会不断产生状态变化:正在构建上下文、正在连接 MCP、正在调用模型、正在执行工具、等待人工介入、完成、失败或被用户中止。

如果客户端只能靠轮询,就会出现三个问题:

  • 进度滞后:用户看不到 Agent 当前在做什么,只能等最终结果。
  • 交互变慢:Stop、HITL 回复、任务状态变化都要等下一次拉取。
  • 多端不同步:用户开着多个 tab,一个 tab 的任务完成了,另一个 tab 仍然显示旧状态。

所以我们的系统需要实时通信。它的目标不是“所有数据都靠 WebSocket 传”,而是让客户端在关键状态变化时被及时唤醒,再回到数据库和 HTTP API 读取权威状态。

两类实时流量

Agent 系统里的实时流量可以分成两类。

第一类是 任务内容流:模型 token、工具输出、状态 data part、UIMessageChunk。这类流量跟一次 run 强绑定,默认走 SSE,因为 SSE 路径支持基于 Redis Stream 的断连续传。WebSocket 也能承载同一套 chunk,但不是默认路径。

第二类是 任务生命周期通知:任务完成、失败、需要 HITL、被 abort、标题更新、列表刷新等。这类事件不一定属于当前打开的那个连接,而是属于用户、会话或某个房间。它们走 /ws 这个实时底座。

这两个层次要分清:

内容流:这一轮回答的每个 chunk 怎么送到当前读者
通知流:某个状态变化应该通知哪些连接

因此,实时通信文档讨论的是“消息怎么路由、连接怎么活、断线怎么补偿”;具体 UIMessageChunk 和 data part 协议见 流式传输架构。

/ws 的角色

/ws 是 Web 端共享的实时通道。每个浏览器 tab 建一条 WebSocket 连接,服务端用 wsHub 跟踪连接,并按消息语义选择分发方式。

我们保留三种分发模式:

模式用途例子
per-connection只回给当前连接点对点错误、当前 task stream 帧
per-room发给订阅同一房间的连接chat / conversation 广播
per-user发给该用户所有 tabtask:event、任务完成通知

最容易出错的是 task:event。任务生命周期事件不是“谁发起谁收到”,而是“这个用户所有打开的 tab 都应该知道”。所以它应该走 per-user,而不是 per-connection。

举例:用户开了两个 tab,一个在任务详情,一个在任务列表。任务在 worker 里完成后,如果只给发起连接回包,列表 tab 就永远不知道状态变了。per-user 能保证该用户所有活跃 tab 都收到刷新提示。

实时事件不是事实之源

WebSocket 不保证投递。连接断开时,正在路上的消息会丢;服务端也不会自动为每个客户端保存 WS 消息历史。

所以我们的 task:event 只做一件事:提示客户端去刷新权威数据。

{
  type: "task:event",
  taskId: "...",
  kind: "completed" | "hitl" | "aborted" | "errored"
}

客户端收到后 invalidate task 列表或任务详情,再通过 HTTP / DB-backed API 读取最新状态。事件本身不携带完整状态,也不作为最终事实。

这个约束很重要:

  • 任务完成事件丢了,重连后客户端要 refetch 补回来。
  • 任务失败事件丢了,reconciler 和列表刷新仍要能暴露最终状态。
  • 不能把余额、权限、任务完整内容这类权威数据只塞进 WS 事件。

如果某类业务未来要求“每条事件都不能丢”,那就不是普通 WS 通知了,需要 sequence number、ACK、重放存储,或者改走可靠队列 + HTTP 拉取。

连接一定会断

WebSocket 是长连接,长连接就一定会断。网络切换、VPN、浏览器休眠、服务端滚动升级、代理空闲超时,都可能让连接中断。更麻烦的是 TCP 半死:两端都以为连接还在,但实际消息已经到不了。

因此客户端和服务端都要把断线当作正常路径:

  • 客户端要有指数退避 + jitter 的重连策略,避免服务端重启时所有 tab 同时打回来。
  • 浏览器回到前台或网络恢复时,要能立刻探测重连,而不是傻等下一个退避窗口。
  • 服务端要用 ping / pong 和 sweeper 清理半死连接。
  • 客户端重连成功后,要主动 refetch 关键状态,因为断线期间的 WS 事件可能已经丢了。

实时系统的正确性不来自“连接永远在线”,而来自“连接断了也能恢复到正确状态”。

多实例必须跨进程 fan-out

单实例时,所有 WebSocket 连接都在本进程内存里,sendToUser(userId) 直接查本地连接表就能发。

生产多实例时情况不同:

用户连接在实例 A
任务完成发生在实例 B
实例 B 本地没有这个用户的 socket

如果只查本机内存,消息就丢了。多实例下,per-user 和 per-room 分发必须经过 Redis pub/sub 这类跨进程 fan-out。每个实例只负责把消息投递给自己持有的连接。

这也是为什么“单机开发正常”不能证明实时通信在生产正常。只要有滚动部署、负载均衡、多 worker、多 API 实例,就必须把跨实例投递作为默认路径设计。

Zombie 状态要靠 reconciler 收尾

实时通信还有一个相关问题:服务端可能在 Agent stream 中途死亡。

此时数据库里 task 可能还标记为 active,但实际已经没有进程在跑。用户看到的表现就是任务一直转圈,既不完成,也不失败。

解决方式不是指望 WebSocket 自己发现一切,而是让 reconciler 周期扫描 active task:

  1. Agent run 持有带 TTL 的 task lock,并定期续约。
  2. 进程死亡后,lock 心跳停止,TTL 自然过期。
  3. reconciler 发现 task 仍 active 但 lock 不存在。
  4. 它把 task 标记为 errored,并通过 task:event 通知用户刷新。

所以实时通信和任务锁、队列、reconciler 是一组配套机制。WS 负责“告诉客户端有变化”,reconciler 负责“把脏状态变成可见终态”。

协议要能演进

WebSocket 消息是客户端和服务端之间的公开协议。客户端版本可能滞后,不能假设所有用户都在最新代码上。

因此协议设计要守几条规则:

  • 顶层用字符串 type 做 discriminator,例如 task:event、task:stream:frame、chat:stream:frame。
  • 新增字段通常安全,删除字段和改变字段语义危险。
  • 新增 kind 要让旧客户端能忽略或 fallback。
  • 不要把所有消息塞进 { type: "notification", domain, payload } 这种弱类型大 envelope;短期省事,长期会丢掉 TypeScript 的 exhaustiveness。

对研发来说,新增 WS 事件就像新增一个 API。要考虑旧客户端、未知字段、未知 kind,以及事件丢失后的补偿路径。

写新实时事件前的自检

  • 这个事件是内容流,还是生命周期通知?
  • 它应该走 SSE、WS per-connection、per-room,还是 per-user?
  • 事件丢了会怎样?客户端重连后能否通过 refetch 补回来?
  • 事件是不是只做刷新提示,权威状态是否仍在数据库 / HTTP API?
  • 多 tab 是否都该收到?如果是,是否走 per-user?
  • 多实例下是否经过 Redis pub/sub,而不是只发本机连接?
  • 客户端收到未知字段或未知 kind 会不会崩?
  • 服务端进程中途死亡时,是否有 lock / reconciler 把状态收成终态?

详细 close codes、心跳参数、Redis pub/sub 接线、reconciler 时序和实现文件见 → 实时通信架构。

这页有帮助吗?