跳到主要内容

两套 WebSocket:daemon 通道、浏览器广播与多实例扇出

30 秒导读: Multica 是「托管型 agent 平台」——你在网页/桌面端派活,本地的 agent 进程(daemon)自动认领、执行、回报,你只管看进度。要做到这种「set it and forget it(设好就不用管)」的体验,页面不能靠轮询,必须服务端主动推。本章讲支撑这套推送的实时层:它由两套各司其职的 WebSocket 组成,中间用 Redis Stream 把多个服务端实例连成一张网。

本章在 Multica 全景中的位置(其余各章见 index):

本章只讲中间那条实时管道:线协议、两个 hub、多实例扇出、前端如何消费。


1. 这是什么(零基础也能懂)

一句话定义

实时层 = 一条从「服务端」通向「本地 daemon」的控制通道 + 一条从「服务端」通向「浏览器」的广播通道,两条都跑在 WebSocket 上,但用途完全不同。

为什么要两套,而不是一套

它们连接的对象、信任模型、消息方向都不一样,硬塞进一套只会互相拖累:

维度daemon 通道浏览器通道
连的是谁用户机器上的 daemon 进程网页 / 桌面 / 手机端
谁主动说话双向:服务端唤醒 daemon,daemon 发 RPC/心跳回来基本单向:服务端推,前端只发订阅/心跳
认证方式Authorization 头 + daemon token / PATCookie 或首帧 token(浏览器设不了自定义头)
核心用途「有活了,来认领」+ 认领 RPC + 心跳「这条 issue 变了,去刷新」
代码位置internal/daemonws/internal/realtime/

两个端点也是分开注册的(cmd/server/router.go:704 是浏览器 /wscmd/server/router.go:772 是 daemon 的 h.DaemonWebSocket)。

一句话直觉

  • daemon 通道像「工头对讲机」:工头(服务端)喊一声「3 号工位有活」,工人(daemon)听到后自己跑去领工单——喊话只是提醒,真正领活还要走正式流程(见 §4 的「best-effort 唤醒」)。
  • 浏览器通道像「广播喇叭」:办公室里所有人(浏览器标签)都听得到「A 项目进度更新了」,但喇叭只说「变了」,不念全文——听到的人自己去公告栏(数据库)取最新内容(见 §6 的「失效信号」)。

用起来什么样(一次真实往返)

你在网页点「让 agent 修这个 bug」,接下来这条链路全自动:

你点派活
→ 服务端把任务落库、判定该谁干(03 章)
→ 服务端通过 daemon 通道喊:"runtime X 有活了"(best-effort)
→ 你机器上的 daemon 听到,发起 tasks.claim RPC 认领
→ daemon 跑 agent,每有进展就 POST 一条 progress
→ 服务端把 progress 变成一条 task:progress 事件,走浏览器通道广播
→ 你和同事的页面收到,刷新看板上那张卡片的状态

你全程没刷新过页面。这就是本章要拆开讲的东西。


2. 顶层全景(它大概怎么转)

一张图看清两套通道

┌─────────────────────────────────────────┐
│ 服务端(可多实例) │
本地机器 │ │ 浏览器 / 桌面 / 手机
┌──────────┐ daemon │ ┌───────────────┐ 事件总线 ┌──────┐ │ 广播 ┌──────────────┐
│ daemon │◀────WS───▶│ │ daemonws.Hub │ │ bus │ │ WS │ realtime.Hub │◀──▶│ 浏览器标签 │
│ (agent) │ 控制通道 │ │ byRuntime/… │ └──┬───┘ │ ◀────▶│ rooms(scope) │ └──────────┘
└──────────┘ │ └──────┬────────┘ │ │ └──────┬───────┘
│ │ 唤醒/RPC/心跳 │广播 │ │ 房间扇出
│ │ ▼ │ │
│ │ ┌──────────────────┐ │
│ └─────────────▶│ Broadcaster │◀─────────┘
│ │ (DualWrite) │
└────────────────────────┴────────┬─────────┘
│ XADD / XREAD
┌──────▼──────┐
│ Redis Stream │ ← 多实例扇出:
│ relay │ 一个实例发,所有实例都收得到
└─────────────┘

怎么读这张图:左半边是 daemon 通道,右半边是浏览器通道,两者在服务端内部通过「事件总线 + Broadcaster」相连;最底下的 Redis Stream 让「多个服务端实例」表现得像一个。

部件一句话职责

部件干什么在哪
protocol.Message所有 WS 消息的统一信封(type + payloadpkg/protocol/messages.go:52
事件常量task:progressdaemon:rpc_request 等字符串枚举pkg/protocol/events.go
daemonws.Hubdaemon 连接注册表 + 唤醒推送 + RPC 分发 + 心跳internal/daemonws/hub.go:186
realtime.Hub浏览器连接的「房间」管理与扇出internal/realtime/hub.go:267
Broadcaster事件producer 只依赖的抽象接口internal/realtime/broadcaster.go:23
DualWriteBroadcaster本地即时扇出 + Redis 跨实例扇出,二者去重internal/realtime/redis_relay.go:502
ShardedStreamRelay固定分片的 Redis Stream 中继(当前默认)internal/realtime/sharded_stream_relay.go:79
RelayNotifier把 daemon 唤醒也塞进 Redis relay,让每个实例都能就近投递internal/daemonws/notifier.go:14
useRealtimeSync前端总入口:把事件翻译成缓存失效/打补丁packages/core/realtime/use-realtime-sync.ts:571

主线走一遍(不进代码)

一条事件的一生:producer 发布到事件总线 → registerListeners 决定发给谁(个人/工作区/daemon)→ Broadcaster 本地扇出并写 Redis → 别的实例从 Redis 读回来、就近扇出 → 浏览器收到 → 前端把它当失效信号刷新缓存。


3. 线协议:一个信封装下所有消息

本节讲「两端到底在网线上传什么」。答案是:都是同一个信封 protocol.Message,靠 type 字段区分。

3.1 统一信封

真源码就三行,是理解一切的地基:

// Message is the envelope for all WebSocket messages.
type Message struct {
Type string `json:"type"`
Payload json.RawMessage `json:"payload"`
}

pkg/protocol/messages.go:52MessagePayload 是延迟解码的原始 JSON——先看 Type,再决定拿什么结构体去解。这让协议天然向前兼容:不认识的 type 直接忽略(daemon 侧见 hub.go:703default 分支,浏览器侧见 hub.go:908)。

3.2 事件常量:一张「谁会飞过网线」的清单

所有 type 值集中定义在 pkg/protocol/events.go,按域分组。挑与本章最相关的几类:

事件常量字符串值方向含义
EventTaskDispatchtask:dispatch服务端→浏览器任务被 daemon 认领
EventTaskProgresstask:progress服务端→浏览器执行中的进度上报
EventTaskMessagetask:message服务端→浏览器单条 agent 消息(工具调用/文本)
EventTaskCompletedtask:completed服务端→浏览器任务完成
EventChatMessage / EventChatDonechat:message / chat:done服务端→浏览器聊天消息 / agent 回完
EventDaemonTaskAvailabledaemon:task_available服务端→daemon「有活了」唤醒提示
EventDaemonHeartbeat / _Ackdaemon:heartbeat / _ackdaemon↔服务端心跳与回执
EventDaemonRPCRequest / _Responsedaemon:rpc_request / _responsedaemon↔服务端通用 RPC(认领走这里)

任务事件的注释点破了一个设计取舍(events.go:28-41):前端按 task: 前缀订阅并整体刷新工作区任务快照,所以事件粒度是「用户想看到什么变化」,而不是「每一次内部状态翻转」。

3.3 RPC:在广播通道里做「请求-响应」

WebSocket 本身是「发了不管」的单向流,但认领任务需要「我领,你告诉我领没领到」的一问一答。Multica 用一对事件 + 一个关联 id 把 RPC 语义架在 WS 上(MUL-4257):

type RPCRequestPayload struct {
RequestID string `json:"request_id"` // 关联 id
Method string `json:"method"` // 如 "tasks.claim"
Body json.RawMessage `json:"body,omitempty"`
TimeoutMs int64 `json:"timeout_ms,omitempty"` // 服务端执行预算
}

pkg/protocol/messages.go:28RPCRequestPayload。响应 RPCResponsePayload:45)回显同一个 RequestID,并带一个仿 HTTP 的 Status——这样 daemon 处理 WS 结果和 HTTP 结果的代码可以完全一致,失败时(非 2xx)无缝回退到 HTTP 认领端点。

3.4 心跳与 pending update:一次往返干两件事

daemon 定期发 daemon:heartbeat,服务端回 daemon:heartbeat_ack。回执不是空的——它顺便夹带待办,让 daemon 少发几个 HTTP 请求:

type DaemonHeartbeatAckPayload struct {
RuntimeID string
Status string
ServerCapabilities []string
RuntimeGone bool // 运行时被删了,让 daemon 自己重注册
PendingUpdate *DaemonHeartbeatPendingUpdate // 该升级 CLI 了
PendingModelList *DaemonHeartbeatPendingModelList
PendingLocalSkills *DaemonHeartbeatPendingLocalSkills
// …还有 skill import 等
}

pkg/protocol/messages.go:262DaemonHeartbeatAckPayload。两个巧思:

  • RuntimeGone 代替 HTTP 404:运行时行被删除时不是断连报错,而是回一个 runtime_gone 状态,daemon 收到后清理本地状态并重新注册(否则死 UUID 会一直心跳到进程重启)。注释见 :254-261
  • pending 字段全用指针 + omitempty:老 daemon 不认识的新字段自动忽略,向前兼容不需要版本协商。

4. daemon 侧 hub:注册、唤醒、RPC 分发

本节讲 internal/daemonws/hub.go——服务端如何管理 daemon 连接、如何「喊话」、如何处理 daemon 发回来的 RPC。

4.1 三张索引表:按 runtime / workspace / user 找连接

一台机器上的 daemon 一条连接,可能持有多个 runtime(运行时)、隶属多个 workspace。Hub 用三张 map 维护反向索引,好按不同维度精准推送:

type Hub struct {
// …
clients map[*client]bool
byRuntime map[string]map[*client]bool // 按运行时找连接
byWorkspace map[string]map[*client]bool // 按工作区找连接
byUser map[string]map[*client]bool // 按用户找连接
// …
}

internal/daemonws/hub.go:186Hubregister:557)在连接建立时把 client 塞进所有三张表,unregister:601)反向清理。三种唤醒各走一张表:

唤醒方法找哪张表用途
NotifyTaskAvailablebyRuntime「这个运行时有活了」
NotifyRuntimeProfilesChangedbyWorkspace「工作区的运行时配置改了,去拉」
NotifyWorkspacesChangedbyUser「你的工作区成员关系变了,去对账」

4.2 唤醒是「提示」,不是「命令」

关键设计(Hub 顶部注释,hub.go:184):

消息是 best-effort 唤醒提示;daemon 仍然通过认领流程保证正确性。

也就是说 task_available 丢了也没关系——daemon 有自己的兜底轮询,唤醒只是让它更快响应,而非唯一触发。notifyFrame:420)用非阻塞发送体现了这一点:

select {
case c.send <- data:
delivered = true
default:
slow = append(slow, c) // 缓冲满 → 判为慢客户端,稍后踢掉
}

发送缓冲满就把该连接列入 slow 并在锁外 unregister + 关连接——宁可踢掉慢连接,也不让一个卡住的 daemon 阻塞整个推送循环

4.3 RPC 分发:goroutine + 限流 + 超时 + 回退

daemon 发来的 daemon:rpc_requesthandleRPCFramehub.go:714)处理,四道防线层层设防:

收到 rpc_request

├─ 没注册 handler ────────────▶ 回 503 → daemon 回退 HTTP

├─ 抢并发名额 rpcSem(容量 8)
│ └─ 抢不到 ─────────────▶ 回 429 → daemon 回退 HTTP

└─ 起独立 goroutine 执行
├─ 按 TimeoutMs 设执行预算(超了就取消,回滚,不留脏数据)
└─ 调 handler → 回 rpc_response(回显 request_id)

为什么每个 RPC 起独立 goroutine?注释说得明白(:709):一次认领是 DB 密集操作,不能卡住读取循环,否则下一个心跳都发不出去。为什么要 TimeoutMs 服务端预算?防止「daemon 已超时回退 HTTP、服务端却还在慢慢提交」造成重复认领——超时就取消并回滚(:742-750)。

RPC 目前只有一个方法 tasks.claim,它的实现有个极省事的巧思——直接复用 HTTP 认领处理器:

func (h *Handler) DaemonRPCHandler(ctx …, method string, body json.RawMessage) (int, json.RawMessage, error) {
switch method {
case "tasks.claim":
return h.rpcClaimTasks(ctx, identity, body)
// …
}
}

internal/handler/daemon_rpc.go:45rpcClaimTasks:54合成一个 HTTP 请求、用 rpcResponseCapture 捕获响应,再喂给现成的 h.ClaimTasksByRuntime。于是 WS 和 HTTP 两条路共用同一段认领逻辑,零重复。

4.4 心跳的一个反直觉细节:故意不加超时

handleHeartbeatFramehub.go:787)处理心跳时,注释特意警告不要给它包 WithTimeout:813-819):心跳会触达 Redis Lua 脚本 PopPending,脚本有副作用(ZREM + SET-running)中途取消无法安全回滚。所以它的自然边界是「读取循环的生命周期」(daemon 走了连接就关),而非人为的 per-call 超时。这类「什么时候该加超时」的判断,正是这套代码的老练之处。


5. 浏览器侧 hub 与多实例扇出

本节讲 internal/realtime/——面向前端的 hub、Broadcaster 抽象、以及让「多个服务端实例」协同工作的 Redis relay。

5.1 房间(scope)模型

浏览器 hub 把连接组织成「房间」,房间键是 {类型, id}。类型有五种(broadcaster.go:6):

scope含义谁能订阅
workspace工作区级广播连接时自动订阅本工作区
user用户级(跨工作区的个人事件)连接时自动订阅自己
task单个任务的高频流ScopeAuthorizer 授权
chat单个会话的高频流ScopeAuthorizer 授权
daemon_runtimedaemon 唤醒的 relay 传输通道只给 daemon hub 消费,不给浏览器

连接建立时 Runhub.go:311)自动把 client 加入 workspaceuser 两个房间。task/chat 这类敏感房间要显式发 subscribe 帧,并由 handleSubscribe:914)调 AuthorizeScope 查库确认「这个任务确实属于你的工作区」才放行。

重要现状(别被 scope 模型误导): 尽管房间机制、授权、Redis relay 都已就绪,per-resource 的 task/chat 路由当前并未启用cmd/server/listeners.go:170-184 的长注释说明:客户端还没落地「发 subscribe 帧 + 重连重放订阅」的那半边,此刻若真按 task/chat 房间路由会把消息全丢在地上。所以现在 task/chat 事件仍走 工作区广播,producer 已埋好 TaskID/ChatSessionID 提示,将来开关只是一行改动。

5.2 Broadcaster:producer 唯一依赖的抽象

所有事件 producer 不直接碰 *Hub,只依赖接口(broadcaster.go:23):

type Broadcaster interface {
BroadcastToScope(scopeType, scopeID string, message []byte)
BroadcastToWorkspace(workspaceID string, message []byte)
SendToUser(userID string, message []byte, excludeWorkspace ...string)
Broadcast(message []byte) // daemon:* 这类无工作区事件
}

这层抽象是「水平扩展计划」的地基:单机就用 *Hub,多机就换成 Redis relay,producer 代码一行不用改

registerListenerslisteners.go:24)就是把事件总线接到这个接口上:个人事件(inbox/邀请)走 SendToUser,其余走 BroadcastToWorkspace:186-192)。

5.3 本地即时 + Redis 跨实例:DualWriteBroadcaster

单实例时,一条事件本地扇出就行。多实例时,事件必须让别的实例上的连接也收到。DualWriteBroadcasterredis_relay.go:502)两件事一起干:

func (d *DualWriteBroadcaster) BroadcastToScope(scopeType, scopeID string, message []byte) {
id := ulid.Make().String()
frame := injectEventID(message, id)
d.local.BroadcastToScopeDedup(scopeType, scopeID, frame, id) // ① 本地立刻扇出
_ = d.relay.PublishWithID(scopeType, scopeID, "", message, id) // ② 写进 Redis
}

:521。关键是那个 id:本地扇出时给每个 client 标记「已见过 id」(markSeenhub.go:239),于是当同一条消息从 Redis 兜一圈回到本实例时,会被去重丢弃——本地连接不会收到两次。

5.4 两种 relay,两种权衡

Multica 有两套 Redis Stream relay 实现,靠环境变量切换(cmd/server/main.go:284):

relay一句话Redis 阻塞连接数适用
RedisRelay(legacy)每个活跃 scope 一条流、一个消费者,按需起停随活跃 scope 数增长早期实现
ShardedStreamRelay(默认)固定 N 条分片流,每 pod 每分片一个 XREAD 循环,本地按订阅过滤pod 数 × 分片数(有界)生产默认

ShardedStreamRelay 是为解决 legacy 的扩展性痛点而生的:活跃 scope 可能成千上万,每个都占一条阻塞的 Redis 连接会撑爆连接池。分片版把它压成固定 8 条(defaultShardedRelayShardssharded_stream_relay.go:17),scope 到分片用一致性哈希:

func (r *ShardedStreamRelay) shardFor(scopeType, scopeID string) int {
h := fnv.New32a()
h.Write([]byte(scopeType)); h.Write([]byte{0}); h.Write([]byte(scopeID))
return int(h.Sum32() % uint32(r.config.Shards))
}

:198。代价是每个 pod 会读到所有 scope 的消息,然后靠 deliverEnveloperedis_relay.go:109)根据本地订阅过滤——用「多读一点 + 本地过滤」换「连接数有界」。

还有个 pod 重启不丢消息的细节:分片读取从 now - ReplayGrace(默认 5 分钟)开始回放,而不是 $(只读新消息),配合下游幂等性(replayStartID:210)。

5.5 灰度切换用的 MirroredRelay

从 legacy 迁到 sharded 时怎么保证不出岔子?MirroredRelayrelay_lifecycle.go:26):同时启两套、都写、都读,同一 event_id 让回环投递幂等,还统计两边的分歧(RedisMirrorDivergenceTotal:99-109)。这是典型的「双写对账灰度」。

5.6 daemon 唤醒也搭 Redis relay 的便车

多实例下有个问题:daemon 连在实例 A,但「有活了」这个判定发生在实例 B。RelayNotifiernotifier.go:14)解决它——唤醒既发本地 daemon hub、也通过 relay 广播(用 ScopeDaemonRuntime 传输通道),于是每个实例都能尝试就近投递给自己持有的 daemon 连接:

func (n *RelayNotifier) NotifyTaskAvailable(runtimeID, taskID string) {
eventID := ulid.Make().String()
n.local.notifyTaskAvailable(runtimeID, taskID, eventID) // 本实例的 daemon
// …
n.relay.PublishWithID(realtime.ScopeDaemonRuntime, shardKey, "", frame, eventID) // 其它实例的 daemon
}

:23。relay 那头 DeliverDaemonRuntimehub.go:369)按帧类型把 scopeID 解释成 runtime / workspace / user,再投给对应的 daemon 连接。


6. 前端如何消费:事件是「失效信号」,不是「数据」

本节讲 packages/core/realtime/——收到事件后,前端绝不把 payload 镜像进本地 store,而是拿它当「该去刷新哪块缓存」的信号。这是全平台缓存一致性的纪律核心。

6.1 核心纪律:单一数据源 = TanStack Query 缓存

服务端数据(issue、任务、聊天…)只住在一个地方——TanStack Query 缓存。useRealtimeSync 顶部注释把这条纪律写死了(use-realtime-sync.ts:1078):

Single source of truth: the Query cache. No Zustand writes here —— 早先往 store 里镜像一份,导致「invalidate→refetch 窗口内缓存和 store 打架、UI 渲染出重复项」。

于是每条事件的处理只有两种动作:

  • invalidate(失效)qc.invalidateQueries(...) 把某块缓存标脏,触发按需 refetch——权威数据仍来自服务端;
  • patch(打补丁)qc.setQueryData(...) 直接把 payload 写进缓存,省一次往返——只用于「payload 自带全部所需字段」的低频场景。

6.2 两条订阅通道:粗粒度 onAny + 细粒度 on

useRealtimeSync 用一个 onAny 兜底 + 一堆精确 on 处理器(:789):

一条事件到达 WSClient.onmessage

├─▶ onAny:取前缀("task:progress" → "task")
│ └─ 命中 refreshMap[前缀] → 防抖 100ms → invalidate 一批缓存

└─▶ on("task:progress" …):若在 specificEvents 里则 onAny 跳过,
由专属 handler 做精确 setQueryData

onAnyrefreshMap:590),每个前缀对应一组失效动作,并按前缀防抖 100ms:751)——批量 issue 更新时合并成一次 refetch,避免请求风暴。specificEvents 集合(:764)列出「已有专属 handler、别再重复失效」的事件,比如所有 chat:*task:message

task:message 为什么被排除在前缀路径外?注释说(:777):它在长任务里每条流式消息都触发一次,若也走快照失效会「刷爆网络」。所以它单独用 setQueryData 把消息按 seq 合并进 taskMessages 缓存(:1089),让时间线原地更新。

6.3 单一 responder + self-initiated guard

CLAUDE 级纪律:一个事件只能有一个负责改状态的地方(单一 responder),并且「本客户端自己能触发这个事件」时要有自触发守卫,否则会和本地的乐观更新/导航抢跑。三个真实例子:

事件单一 responder 做什么自触发守卫
workspace:deleted清工作区存储 + 导航离开isWorkspaceDeletePending(id) 命中就直接 return——自己发起的删除由 mutation 收尾,避免和它的导航抢跑(:1002-1020
member:removed被踢者切换工作区只在 user_id === 自己 时动作(:1022
chat:session_deleted其它标签页同步删行 + 清空活动指针发起删除的那个标签已乐观删过,这里只服务「别处发起」的删除(:1359

注意 chat:session_deleted 里那句「清空客户端拥有的活动会话指针」——这正是 CLAUDE 规则允许的例外:事件绝不镜像服务端数据进 Zustand,但可以在单一 responder + 自触发守卫下清理客户端自己的指针(活动会话、选中项、当前工作区)。

6.4 一个安全红线:聚合态只能 invalidate,不能 setQueryData

聊天的 task:* 事件是工作区扇出——工作区每个成员都收得到,且 payload 不带创建者/可见性信息。所以跨会话的「待处理聚合」(驱动那个悬浮按钮的 has_pending 徽标)绝不能从这些事件乐观写入,否则「成员 B 起了个任务」会点亮「成员 A 的徽标」,绕过服务端权限过滤(:104-128:1104-1119):

// SECURITY: 这里是 invalidate,不是乐观 setQueryData —— 强制走
// /api/chat/pending-tasks 这个按创建者+可见 agent 过滤的权威端点。
export function refetchPendingChatAggregate(qc, wsId) {
qc.invalidateQueries({ queryKey: chatKeys.pendingTasks(wsId) });
}

:122。而 per-session 的 pendingTask 缓存可以直接写——它按会话 id 键控,只对「你有权打开的会话」渲染,不是跨用户聚合。「什么能 patch、什么必须 invalidate」在这里是一条安全边界,不只是性能选择。

6.5 重连与工作区切换:补回错过的事件

WebSocket 断线期间的事件是丢的。前端用两把「大扫除」找补:

  • 重连:1429):onReconnect 回调触发 invalidateWorkspaceScopedQueries:484)——把所有工作区级缓存标脏、按需 refetch,把断线期漏掉的都补回来。
  • 新 WSClient 实例:1447,即切换工作区):同样全量失效,但用 wsInstanceRef 跳过首次挂载,避免刚进页面就白刷一遍。

连接生命周期本身在 WSProviderprovider.tsx:49):它用 useSyncExternalStore 订阅当前工作区 slug,slug 一变就拆掉旧连接、开一条绑定新工作区的(:79-116),并只调一次 useRealtimeSync:121)。断线重连带指数退避 + 抖动(ws-client.ts:163)。


7. 端到端:一次进度上报如何点亮浏览器看板

把前面所有部件串起来,看「daemon 报一次 progress」怎么一路点亮你和同事的看板。这是本章的压轴流向图:

【本地机器】
daemon 跑 agent,产生一点进展
│ POST /api/daemon/tasks/progress

【服务端·实例 A】
handler.go:2824 ──▶ TaskService.ReportProgress(task.go:3701)

│ bus.Publish(EventTaskProgress, WorkspaceID=…)

事件总线 ──▶ registerListeners.SubscribeAll(listeners.go:152)
│ 非个人事件 + 有 WorkspaceID
│ → b.BroadcastToWorkspace(wsID, data) (listeners.go:188)

DualWriteBroadcaster.BroadcastToScope(redis_relay.go:521)
├───① 本地即时扇出:hub.BroadcastToScopeDedup(workspace, wsID)
│ → 实例 A 上订阅了该工作区的浏览器连接,立刻收到

└───② relay.PublishWithID → Redis XADD 到某分片流

【服务端·实例 B】每分片一个 XREAD BLOCK 循环
readShard(sharded_stream_relay.go:218)
→ deliverEnvelope(redis_relay.go:109)
→ hub.BroadcastToScopeDedup(workspace, wsID)
→ 实例 B 上的浏览器连接也收到(去重:event_id 已见过则丢)

【浏览器(连在 A 或 B 都行)】
WSClient.onmessage(ws-client.ts:102)
│ onAny → 前缀 "task"

refreshMap.task(use-realtime-sync.ts:697) 防抖 100ms
│ invalidateQueries(agentTaskSnapshot / workingAgents / …)

TanStack Query 按需 refetch → 看板那张卡片的状态刷新 ✔

三个值得记住的点:

  1. progress 走的是粗粒度失效,不是精确 patch——task:progress 不在 specificEvents 里,命中 refreshMap['task'],刷新的是「工作区任务快照」这类聚合缓存。
  2. 连在哪个实例都收得到,靠的正是 §5 的 Redis relay 扇出 + event_id 去重。
  3. 前端只被告知「变了」,真正的新进度是 refetch 从数据库取回来的——这就是 §6 的「失效信号」纪律。

8. 巧妙之处(可借鉴的技术)

  • RPC 复用 HTTP 处理器:WS 认领不重写逻辑,合成 HTTP 请求喂给现成 handler,一套认领逻辑两个传输共用(internal/handler/daemon_rpc.go:54)。
  • 心跳回执夹带待办:ack 顺手带 pending update / model list / skill import,省下多个独立 HTTP 轮询(pkg/protocol/messages.go:262)。
  • event_id 去重解决双写回环:本地扇出先 markSeen,Redis 兜回来的同 id 自动丢弃,双写不重复(redis_relay.go:521 + hub.go:239)。
  • 分片流换有界连接数pod × shard 条阻塞连接替代「每活跃 scope 一条」,多读一点本地过滤(sharded_stream_relay.go:198)。
  • 降级永远存在:daemon 通道任何一环失败(无 handler/限流满/缓冲满/超时)都回退 HTTP,唤醒丢了有兜底轮询——实时层是加速器,不是唯一路径hub.go:726/734)。
  • 前端 refetchType:"none" 的精修:时间线失效时不让已挂载的观察者立刻 refetch,避免 AI 流式输出时整棵评论树闪烁(use-realtime-sync.ts:898,MUL-1941)。

9. 边界与局限(诚实)

  • per-resource(task/chat)scope 路由尚未启用:机制全在,但客户端还没落地 subscribe 帧与重连重放,当前仍是工作区广播(cmd/server/listeners.go:170-184)。这意味着高频 task/chat 事件目前会发给整个工作区,而非只发给关心该资源的连接。
  • 唤醒不保证送达task_available 缓冲满就丢、慢连接直接踢,正确性完全依赖 daemon 的兜底认领流程(hub.go:420)。
  • 断线期事件靠「大扫除」补:重连不做增量重放,而是全量失效 + refetch(use-realtime-sync.ts:484)——简单可靠,但重连瞬间会有一波刷新流量。
  • MirroredRelay 是临时物:仅为 legacy→sharded 灰度存在(relay_lifecycle.go:26),不是长期架构。
  • ReportProgress 只发不落库:它只 bus.Publish 一条广播事件,看注释没有持久化 progress 本身(task.go:3701)——刷新看板靠的是 refetch 任务快照,不是回放 progress 流。

10. 代码地图(导航索引)

按符号名 grep 比按行号更抗漂移。

主题文件路径符号名
消息信封server/pkg/protocol/messages.goMessage
RPC 请求/响应载荷server/pkg/protocol/messages.goRPCRequestPayload / RPCResponsePayload
心跳回执 + pendingserver/pkg/protocol/messages.goDaemonHeartbeatAckPayload
事件常量server/pkg/protocol/events.goEventTaskProgress / EventDaemonRPCRequest
daemon hub / 索引表server/internal/daemonws/hub.goHub / register / notifyFrame
daemon 唤醒(best-effort)server/internal/daemonws/hub.goNotifyTaskAvailable
RPC 分发 + 限流/超时server/internal/daemonws/hub.gohandleRPCFrame
心跳处理(无超时)server/internal/daemonws/hub.gohandleHeartbeatFrame
RPC 认领复用 HTTPserver/internal/handler/daemon_rpc.goDaemonRPCHandler / rpcClaimTasks
daemon 唤醒搭 relayserver/internal/daemonws/notifier.goRelayNotifier
浏览器 hub / 房间server/internal/realtime/hub.goHub / subscribe / BroadcastToScopeDedup
订阅授权server/internal/realtime/hub.gohandleSubscribe
Broadcaster 抽象 + scope 常量server/internal/realtime/broadcaster.goBroadcaster / ScopeDaemonRuntime
双写广播 + event_id 去重server/internal/realtime/redis_relay.goDualWriteBroadcaster / injectEventID
按需 per-scope relayserver/internal/realtime/redis_relay.goRedisRelay / runConsumer
固定分片 relay(默认)server/internal/realtime/sharded_stream_relay.goShardedStreamRelay / shardFor / readShard
灰度双读双写server/internal/realtime/relay_lifecycle.goMirroredRelay
事件路由(个人/工作区/daemon)server/cmd/server/listeners.goregisterListeners
relay 模式装配server/cmd/server/main.gorelayMode switch,:284
进度上报 producerserver/internal/service/task.goReportProgress
前端总入口packages/core/realtime/use-realtime-sync.tsuseRealtimeSync
聚合态安全失效packages/core/realtime/use-realtime-sync.tsrefetchPendingChatAggregate
重连补数据packages/core/realtime/use-realtime-sync.tsinvalidateWorkspaceScopedQueries
连接生命周期packages/core/realtime/provider.tsxWSProvider
WS 客户端分发packages/core/api/ws-client.tsonmessage / onReconnect