实时通信

单一 /ws 端点 + 三种消息分发模式(per-connection / per-room / per-user)、Redis pub/sub 跨实例递送、 四层可靠性(server ping、client watchdog、sweeper、带 jitter 的重连),以及进程在 stream 中途死亡时清理 zombie 状态的 task reconciler。

一个 /ws 端点扛起所有实时流量

/ws 是 chat 对话、task streaming、任务生命周期通知共同的运行时基础设施——每个浏览器 tab 一条 WebSocket 连接,所有实时流量都走它。真正的难点不在“发消息”,而在三件事:一条消息该路由给谁(广播到 room?推给某个用户的所有 tab?还是回给单条连接?)、连接怎么在网络抖动里活下来、以及进程半路死掉后留下的 zombie 任务状态怎么自动清掉。本页把这三件讲透。

要先分清一件事:这里讲的是实时传输的运行时底座;流过它的内容协议(data parts、AI SDK chunks)是另一回事,见 流式架构

单端点 + 三种分发模式

每个浏览器 tab 一条 WebSocket 连接,所有实时流量走它。服务端 wsHub 单例跟踪每个连接,按以下三种模式之一路由消息:

sequenceDiagram participant C1 as Client tab A · userX participant C2 as Client tab B · userX participant C3 as Client tab C · userY participant Hub as wsHub participant Redis Note over C1,Redis: ① 三种分发模式共享同一连接 C1->>Hub: send_message Hub->>Hub: send(connectionId, error) Note right of Hub: per-connection · 点对点回复 Hub->>Redis: publish room channel Redis-->>Hub: deliver Hub-->>C1: chat:stream:frame Hub-->>C3: chat:stream:frame Note right of Hub: per-room · 显式 join,多订阅者 Hub->>Redis: publish user channel Redis-->>Hub: deliver Hub-->>C1: task:event Hub-->>C2: task:event Note right of Hub: per-user · 自动,该用户所有 tab
模式API用途
Per-connectionwsHub.send(connectionId, msg)HITL 回复、错误响应、点对点发给调用者的 task stream 帧
Per-room(显式 join)wsHub.broadcast(roomId, msg) + join/leaveChat 对话:多个参与者订阅同一 conversationId
Per-user(隐式)wsHub.sendToUser(userId, msg)Task 生命周期事件:用户开着的任意 tab 都收到,无需订阅

per-user 模式是通知系统的基石。任务在实例 A 完成、用户在实例 B 上开着两个 tab 时,两个 tab 都能收到事件——因为 userId 索引在每个实例上维护,且通过 Redis pub/sub 串起来。

连接生命周期

sequenceDiagram participant C as Client participant W as /ws route participant Auth as verifyJwt (jose) participant Hub as wsHub participant T as Heartbeat-Sweeper Note over C,T: ① 建立连接 C->>W: WebSocket upgrade · ?token=JWT W->>Auth: verifyJwt(Bearer token) Auth-->>W: AuthUser or null Note right of W: 未授权 → close 4001(终态)<br/>已授权 → register W->>Hub: register(connectionId, ws, userId) Hub->>Hub: 维护 userToConnections 索引 Hub-->>C: open · status connected rect rgba(52, 211, 153, 0.18) Note over C,T: ② 稳态心跳(每 25 s) T->>C: ping(JSON 应用层) C->>W: pong W->>Hub: touch(connectionId) · 刷新 lastActivityAt end Note over C,T: ③ 异常检测 — 三条独立路径 Hub->>C: 优雅关闭(1000 / 1006 / 1011 / 1008 / 4002) C->>C: watchdog · lastMessageAt 超过 35 s · 主动 close T->>Hub: sweeper · lastActivityAt 超过 90 s Hub->>C: close 4002(stale) Note over C,T: ④ 重连策略 C->>C: 检查 close code Note right of C: 终态(1000 / 1008 / 4001)→ 停止重连<br/>瞬态(1006 / 1011 / 4002 / 其它)→ jitter 退避 C->>W: 重连尝试(1, 2, 4, 8, 16, 32 s 封顶 · full jitter) Note over C: visibilitychange / online 事件绕过退避立即重连

Close codes

Code含义客户端反应
1000Normal closure停止重连(用户主动断开)
1006Abnormal closure退避重连(最常见)
1008Policy violation停止重连(协议错误)
1011Internal server error退避重连
4001Unauthorized(自定义)停止重连;由上游 auth 流程处理重定向
4002Stale connection(自定义)退避重连(被 sweeper 关掉)

为什么 server-driven 心跳

服务端每 25 s 发应用层 {type: "ping"},客户端回 {type: "pong"},两端各自跟踪 lastActivityAt。这套与 Socket.IO / SignalR 默认一致。server 驱动优于 client 驱动 的理由:

  1. 集中调参,无需重新发布客户端
  2. server ping 抵达本身证明 server 还活着——client 驱动只能证明 client 活着
  3. 集中速率限制 / 滥用控制

未来若 WS adapter 暴露 RFC 6455 原生 ping,可在传输层加一层补充:某些中间代理会剥离应用帧但保留协议级控制帧。当前 JSON ping 是唯一心跳。

浏览器事件触发的立即重连

两个 window 事件绕过退避计时器,触发立即重连探测:

  • visibilitychange(visible)——笔记本唤醒 / tab 切回前台
  • online——网络恢复 / VPN 重连

这两个事件比 watchdog 超时更早探测到“网络回来了”,尤其在长时间休眠后 OS 还没察觉断开的场景。

多实例 · Redis pub/sub

每个 server 进程只拥有落到它身上的连接。要把 task:event 推给该用户所有 tab(可能分散在多个实例上),所有 send 都走 Redis pub/sub:

Redis pub/sub:跨实例投递 两个实例都订阅用户 channel · 发布者自身的 subscriber 也会收到 fanout Instance A userToConnections userX → conn-A1, conn-A2 已订阅 ws:user:userX Redis pub/sub channel zapvol:ws:user:userX 投递给所有 subscriber Instance B userToConnections userX → conn-B1 已订阅 ws:user:userX ① publish ② fanout ② fanout(发布者自身 subscriber 也收到) ③ handler · localSendToUser → conn-A1, conn-A2 ③ handler · localSendToUser → conn-B1

订阅是惰性的:每个实例仅在至少有该用户一个连接时才订阅 zapvol:ws:user:{userId},最后一个连接断开时取消订阅。订阅集随实例上的活跃用户线性变化,而不是总用户数。

Redis 不可用时(本地开发)sendToUser 退化为本进程内本地分发。单进程开发模式正常工作;多实例生产必须配 Redis

room 频道(zapvol:ws:room:{roomId})走同一套机制,由 wsHub.broadcast 调用。

可靠性的四层防御

四种机制协作,确保系统对故障忠实:

可靠性:4 层防御 每层捕捉不同故障 · server 主导(indigo)与 client 主导(teal)交替 1 · Server-driven 心跳(每 25 s) ping/pong 维护 lastActivityAt · 从 server 侧捕捉 client 端 TCP 半死 2 · Client watchdog(35 s 超时) lastMessageAt 超过 35 s · 主动 close · 从 client 侧捕捉 server 端 TCP 半死 3 · Server sweeper(90 s 阈值) 每 30 s 扫描 · 心跳路径被代理剥离或第 1 层静默失败时的兜底清理 4 · 重连策略(backoff + 浏览器事件) 瞬态 → jitter 退避 · 终态(1000 / 1008 / 4001)→ 停止 · visibilitychange / online 绕过计时器

每层独立、互补——它们不互相替代,而是各自覆盖不同的故障模式。从故障到层的完整映射:

故障形态首次发现兜底
TCP 半死(server 端)server sweeper 90 s——
TCP 半死(client 端)client watchdog 35 sonclose
网络切换 / NAT rebindonline eventwatchdog 35 s
笔记本唤醒visibilitychangewatchdog 35 s
Server 重启onclose 1006退避重连
鉴权过期close 4001停止重连
代理剥离心跳路径sweeper 90 s——

一个真正生产级的 WS 实现需要每一行。

Task reconciler · 清理 zombie 状态

与 WS 独立但同源(都涉及多实例问题):server 进程在 stream 中途死亡时,任务在数据库里 isActive = true,但实际无 compute 在跑。不干预的话,sidebar 永远显示该任务在转圈。

sequenceDiagram participant A as Instance A participant B as Instance B participant Redis participant DB as Postgres A->>DB: task X · isActive=true A->>Redis: SET task-lock:X · TTL 90 s loop task X 跑着的每 30 s A->>Redis: SET XX task-lock:X · TTL 90 s end Note over A: 进程被杀(SIGKILL · OOM · 滚动升级) Note over Redis: 心跳停 · TTL 在 90 s 内自然过期 loop 每个实例每 2 min B->>DB: listActive() DB-->>B: 包含 X 的行 B->>Redis: EXISTS task-lock:X Redis-->>B: 0(锁已不在) B->>DB: update X · isActive=false B->>Redis: publish zapvol:ws:user:userX(errored 事件) end

reconciler 在每个实例运行。操作幂等,多实例并发跑安全—— reconciler 自身不需要分布式锁。zombie 到清理的最大延迟:lock TTL(90 s) + sweep 周期(2 min) ≈ 3.5 min

Redis 不可用时 isTaskLocked 返回 false,reconciler 退化到时间戳判定:updatedAt < now - 5 min 视为 zombie。单实例开发模式正确工作;此 fallback 仅在没有 Redis 的场景被用到。

reconciler 也是“server 重启”的用户体验恢复路径:zombie 变成 kind: "errored" 事件 → 红色 toast(可选桌面通知)告诉用户任务没扛住部署。

协议形态

每条消息顶层 type 是 discriminator,采用 {domain}:{subkind} 命名空间。例如:

chat:stream:frame    chat:stream:end    chat:stream:error
task:stream:frame    task:stream:end    task:stream:error
task:event           agent:state        message:new          message:updated
typing               presence           ping                 error

task:event 是 task 生命周期通知的标准载体:

interface WsTaskEvent {
  type: "task:event";
  taskId: string;
  kind: "created" | "completed" | "hitl" | "aborted" | "errored";
  finishReason?: string;
  errorMessage?: string;
}

客户端收到 → invalidate task 列表 → 按 kind 触发 toast / sidebar 徽章 / 桌面通知。无需 prev/current diff,事件本身权威。

面向未来扩展

新 domain 需要通知时(schedule 触发、credit 警告等),按同一命名规范新增类型:

interface WsScheduleEvent {
  type: "schedule:event";
  scheduleId: string;
  kind: "fired" | "failed" | "missed";
  runId?: string;
}

interface WsCreditEvent {
  type: "credit:event";
  kind: "warning" | "exhausted";
  remaining: number;
}

每个 domain 保持强类型。WsServerMessage union 线性增长,客户端 discriminator switch 在 TypeScript narrowing 下保持 exhaustive。

故意做成通用的 {type: "notification", domain, event, payload} envelope ——类型擦除以短期简单换长期成本,跨 domain 代码失去编译期 discriminator 后维护成本陡升。

这个子系统不做什么

  • 不做保证投递 / 重放——客户端断线期间发出的事件丢失。客户端的补偿是:每次 disconnected → connected 都 invalidate task 列表,权威地重取状态。对 task 生命周期足够;若是更细粒度的事件流(每键位的光标位置),需要 sequence-number 重放协议。
  • 不做客户端发起任务创建——任务创建走 HTTP POST /api/tasks。WS 只承担通知和已有 HTTP 入口的双向流。
  • 不做 WS 侧速率限制——应用层处理。server 进程依赖 Redis pub/sub fanout 上限和 OS 文件描述符限制保护;激进的滥用处理超出范围。
  • 不做协议版本控制——消息是带 type 辨识的 JSON。前向兼容靠客户端忽略未知 type。硬破坏性变更(改字段名、删 kind)目前未版本化,需要协调客户端发布。

关键参数

理由
端点/ws所有实时流量共用一个端点
鉴权JWT 经 ?token= query 参数浏览器在 WS 升级时无法设置 header,token 走 query string → 合成 Authorization: BearerverifyJwt(jose)
Server ping 周期25 s对齐 Socket.IO / SignalR 默认值
Client watchdog 超时35 s> ping 周期 + 一次往返;触发主动重连
连接 sweeper 周期30 s服务端僵尸连接扫描频率
连接僵尸阈值90 s三次心跳错过;sweeper 周期的两倍
Task lock TTL90 s与连接 sweeper 阈值一致,lock-as-truth 统一
Lock heartbeat 周期30 sTTL 的 1/3;允许漏一次 Redis 调用
Reconciler 启动延迟60 s让同时启动的其它实例先完成 boot
Reconciler 周期2 minzombie 到清理的最大延迟:进程死亡后约 3.5 分钟
客户端重连基础延迟1 s指数退避,32 s 封顶,full jitter
最大重连尝试数10总等待约 17 分钟后放弃

实现位置

关注点文件
连接生命周期apps/server/src/routes/ws.ts
Hub 状态 + 路由apps/server/src/lib/ws-hub.ts
心跳 / sweeper同上——pingTimer, sweeperTimer in ensureSweeper()
Lock + heartbeatapps/server/src/lib/task-lock.ts
Reconcilerapps/server/src/services/task-reconciliation.ts
Boot 接线apps/server/src/index.ts
客户端重连packages/app/src/hooks/use-websocket.ts
通知监听器packages/app/src/components/app-notification-listener.tsx
浏览器原生通知packages/app/src/hooks/use-browser-notification.ts
偏好存储packages/app/src/stores/preferences-store.ts
Wire 类型packages/common/src/types/ws.ts

延伸阅读

  • 流式架构——流过 WS 的内容、data part 协议、resumable SSE
  • 任务编排——一个任务从 POST 到完成的全过程;锁获取 + 心跳是其生命周期的一部分
  • 生产部署——Redis 配置、环境变量、滚动部署机制
这页有帮助吗?