跳到主要内容

第 12 章:事件总线、WebSocket 与实时一致性

12.1 实时系统的职责

Multica 的实时链路同时服务三类消费者:

  • 浏览器、Desktop、Mobile:Issue、Comment、Task、Inbox、Chat 等界面更新;
  • Daemon:任务 Wakeup、取消、Runtime 指令和 RPC;
  • Server 内部订阅者:通知、活动、Squad、Autopilot 等派生副作用。

它们不能共用一个无类型广播通道。浏览器需要 Workspace 隔离,Daemon 需要 Runtime 身份,内部订阅者还需要确定执行顺序。

12.2 同步进程内 Event Bus

server/internal/events/bus.goEvent 包含:

  • Type;
  • WorkspaceID;
  • ActorType / ActorID;
  • JSON 可序列化 Payload;
  • 可选 TaskID / ChatSessionID Scope Hint。

Bus 提供 Subscribe(type)SubscribeAll。Publish 时:

  1. 复制当前 Listener Slice;
  2. 按注册顺序同步调用类型 Listener;
  3. 再同步调用全局 Listener;
  4. 单个 Listener Panic 被 Recovery,不阻断后续 Listener。

12.3 为什么选择同步

同步分发提供一个简单而强的语义:

Handler 返回之前,当前进程中的事件订阅者已经按顺序看到了事件。

它便于维护:

  • Activity 先于全局 WebSocket Fanout;
  • 单元测试无需等待异步队列;
  • 同一请求的派生操作有确定顺序;
  • 不会为每个小事件启动无界 Goroutine。

代价是 Listener 不能做长时间阻塞工作。外部 I/O 应入队或异步化,而不是拖住原 HTTP 请求。

12.4 Panic 隔离不等于事务

Bus 会恢复 Listener Panic,但不会回滚之前已经完成的 Listener,更不会回滚数据库事务。

正确顺序通常是:

  1. 在事务中完成业务写入;
  2. Commit;
  3. Publish Domain Event;
  4. Listener 更新派生对象或 Fanout;
  5. 客户端收到信号后再读权威 API。

若在 Commit 前 Publish,客户端可能先收到一个数据库尚不可见的对象。

12.5 事件命名是跨层协议

常见事件族:

  • issue:
    comment:
  • task:*
  • chat:*
  • inbox:
    notification:
  • agent:
    runtime:
  • project:
    workspace:
  • autopilot:
    squad:

Event Type、Payload Shape、前端 Updater 与 Redis Relay 构成隐式公共协议。重命名不能只改 Server 常量。

12.6 浏览器 WebSocket Hub

浏览器连接先经过普通用户认证,再订阅可访问的 Workspace Room。Server 不接受客户端任意声明 Workspace 权限;身份与成员关系仍由 Server 校验。

Hub 将 Domain Event 编码为统一 Frame,携带 Event Type、Workspace 与 Payload。客户端对未知事件应忽略或降级 Invalidate,而不是让整个连接崩溃。

浏览器 WS 的目标是低延迟界面同步,不承载最终业务写入。

12.7 Daemon WebSocket 是另一条控制通道

入口在 server/internal/handler/daemon_ws.go,RPC Handler 在 daemon_rpc.go

Daemon WS 使用 Daemon 身份,并绑定它注册的 Runtime IDs。它承载:

  • Task Available Wakeup;
  • Heartbeat/可用性信号;
  • Cancel/Reconcile;
  • Model List、本地 Skill List/Import;
  • Runtime Update 等请求与响应。

它不能复用浏览器 WS Cookie,也不能让浏览器订阅 Runtime 控制 Room。

12.8 WS-first,不是 WS-only

Daemon 优先通过 WebSocket 收到 Wakeup 和 RPC,但任务正确性仍依赖 PostgreSQL Claim。

如果 WS 暂时断开:

  • Poller 继续通过 HTTP Batch Claim;
  • Heartbeat 可走兼容路径;
  • RPC 在安全条件下回退 HTTP;
  • Reconnect 后恢复低延迟通知。

消息丢失最多增加延迟,不应让任务永久丢失。这就是“WS 是提示,数据库是队列”。

12.9 Runtime 控制消息的幂等

Daemon 可能同时从 WS、HTTP Poll 和重连补偿看到同一动作。实现因此围绕持久 ID:

  • Task Claim 由数据库条件更新保证唯一;
  • Request ID 关联 Model、Skill、Update 响应;
  • Cancellation 读取 Task 当前状态;
  • Runtime Gone 触发可重复的 Re-register/Reconcile。

不能以“我只收到一次 WS Frame”为业务前提。

12.10 多实例 Server 与 Redis Relay

单实例时,内存 Hub 足够;多实例时:

  • 客户端可能连到 Server A;
  • 写请求可能落在 Server B;
  • B 的本地 Bus 无法直接通知 A 的 Socket。

Redis Relay 把事件写入按 Scope 分片的 Streams,再由各实例转发到本地 Hub。代码还保留对旧 Stream、双写或迁移形态的兼容,避免滚动升级期间一半节点收不到事件。

12.11 Scope 从 Workspace 向细粒度演进

Event 已有 TaskIDChatSessionID Hint,Fanout 层可据此选择更细 Room,而不重新解析 Payload。

但“Server 能发布细粒度 Scope”不等于“所有客户端已经订阅它”。当前前端主要仍以 Workspace 实时同步为主。扩展细粒度订阅时,需要同时核对:

  • Subscribe 请求;
  • 权限校验;
  • Redis Stream Key;
  • Hub Room;
  • Reconnect 后重订阅;
  • 旧客户端兼容。

12.12 前端 WS Client

共享前端实时入口可从 packages/core 中的 WebSocket 与 Realtime 模块追踪。客户端职责包括:

  • 依据 Api Base URL 构造 WS URL;
  • 携带现有认证上下文;
  • 校验 Frame 基本形状;
  • 指数退避并加入 Jitter;
  • 网络恢复时重连;
  • Workspace 切换时退订旧 Scope;
  • 将事件交给 Cache Updater。

移动端有独立实现,入口在 apps/mobile/data/realtime/ws-client.ts

12.13 React Query 是权威读模型

客户端不会把 WS Payload 当数据库完整镜像。Realtime Updater 通常选择:

  • Payload 足够且顺序明确:定点 Patch Query Cache;
  • 只知道某对象变化:Invalidate 对应 Query Key;
  • 影响多个聚合视图:Invalidate 列表、计数和详情;
  • 权限敏感聚合:只失效,不乐观注入未知对象。

这样可抵抗漏帧、乱序和版本偏差。

12.14 为什么 Query 可以长 Stale

部分 Query 使用很长甚至无限 staleTime,因为正常更新由事件驱动。它并不代表永远不刷新:

  • WebSocket Event 主动 Patch/Invalidate;
  • Reconnect 后 Refetch;
  • Workspace 切换产生新 Query Key;
  • Mutation Success/Settled 主动同步;
  • 用户重新聚焦或显式刷新可恢复。

长 Stale 减少每次页面切换的重复 GET,但把压力转移到事件覆盖测试。

12.15 乐观更新与实时回声

用户提交 Comment 时可能发生:

  1. Mutation 先把临时 Comment 放入 Cache;
  2. Server 创建真实 Comment;
  3. WebSocket 很快回送 comment:created
  4. Mutation Response 也返回实体。

若三路都 Append,就会短暂出现三条。Updater 需要按持久 ID 或 Client Correlation 去重;失败时撤销 Optimistic Snapshot,成功时用权威对象替换临时对象。

12.16 Workspace 切换的竞态

Workspace 是前端最重要的 Cache Namespace。切换时必须:

  • 更新 ApiClient 当前 Workspace;
  • 更新路由;
  • 切换 Query Keys;
  • 重建 WS Subscription;
  • 避免旧 Workspace 的迟到 Frame 写进新 Cache;
  • 对持久化 Zustand Store 使用 Workspace-aware Key。

仅在 URL 改一个 ID 而复用全局缓存,会造成跨工作区数据闪现,严重时形成隐私问题。

12.17 Inbox 与权限敏感事件

Inbox 是由 Issue、Comment、Assignment、Subscriber 关系派生的用户视图。同一个 Workspace Event 对不同成员的结果不同。

因此这类 Event 更适合作为“某个聚合可能变化”的信号,让客户端 Refetch 当前用户的 Inbox,而不是把 Workspace Payload 直接塞进每个人的列表。

12.18 断线一致性

WebSocket 不保证客户端离线期间逐条补齐所有 UI 事件。系统使用收敛模型:

  • 业务状态持久化在 PostgreSQL;
  • WS 提供在线低延迟;
  • Redis 解决节点间转发;
  • Reconnect/Refetch 重新读取当前事实;
  • Query Key 和更新时间帮助覆盖旧 Cache。

这是一种最终收敛的读模型,不是 Event Sourcing。

12.19 顺序与过期事件

同一对象可能快速发生 queued → dispatched → running → completed。跨网络 Frame 可能乱序。

客户端更新应优先依靠:

  • 数据库返回的状态;
  • updated_at 或单调版本;
  • 终态不回退规则;
  • Invalidate 后重新拉取。

若 Payload 没有可靠版本,安全策略是失效查询,而不是盲目覆盖。

12.20 背压与慢消费者

同步 Bus 不能让慢 Listener 阻塞太久;WebSocket 也不能让一个慢客户端拖住 Hub。

常见防线包括:

  • 有界发送队列;
  • 写超时与 Ping/Pong;
  • 队列满时断开并让客户端重连收敛;
  • Redis Consumer 独立处理;
  • 大对象不放 Event Payload,只发标识和摘要。

实时链路追求“尽快通知”,不是承诺把数据库所有内容可靠流送给每个客户端。

12.21 增加新事件的检查清单

  1. 业务写入是否先 Commit;
  2. Event Type 是否稳定且不碰撞;
  3. Payload 是否含最小必要 ID;
  4. Workspace、Task、Chat Scope 是否准确;
  5. 内部 Listener 顺序是否有依赖;
  6. Browser 与 Daemon 哪一类消费者需要它;
  7. 多节点是否经过 Redis;
  8. 前端应 Patch 还是 Invalidate;
  9. 是否会与 Mutation Response 产生重复;
  10. 断线重连能否通过 GET 收敛;
  11. 旧客户端遇到新事件是否安全忽略;
  12. 是否有跨 Workspace 泄漏测试。

12.22 排障路线

“数据库已更新但页面没变”按下面分层:

  1. Handler 是否在 Commit 后 Publish;
  2. Event WorkspaceID 是否正确;
  3. 类型 Listener 是否 Panic;
  4. 多实例 Relay 是否写到正确 Stream;
  5. 当前 Socket 是否订阅正确 Workspace;
  6. Frame 是否通过客户端 Schema;
  7. Updater 是否命中正确 Query Key;
  8. 是否被旧 Frame 覆盖;
  9. Reconnect Refetch 是否可恢复。

“Daemon 没被唤醒”则查 Daemon Hub、Runtime ID、Wakeup 与 Poll Fallback,不要沿浏览器 Hub 排查。

12.23 本章结论

Multica 的实时一致性不是靠一条永不丢包的 WebSocket,而是四层协作:

  • PostgreSQL 保存权威状态;
  • 同步 Event Bus 保证进程内顺序;
  • WebSocket/Redis 快速传递失效信号;
  • 客户端 Query Cache 通过 Patch、Invalidate 和 Refetch 收敛。

把 WS 当提示而非事实源,系统才能在断线、重连、多实例和版本漂移下继续正确工作。