聊天廣播與事件投遞
職責
server-chat.ts 把 agent 主循環吐出的事件 (agent event) 投遞回用戶端。它不生產事件,只做分派型分發:接 AgentEventPayload(delta、final、tool、lifecycle 等),按事件類型、工作階段 (session) 歸屬、用戶端訂閱情況,決定是廣播 (broadcast) 給所有可見連線、定向投遞 (broadcastToConnIds) 給一組 connId、還是跨節點投遞 (nodeSendToSession) 到另一台節點上訂閱了該 session 的進程。
這一層維護三類訂閱狀態:toolEventRecipients(runId 級工具事件訂閱者,用戶端發起 run 前註冊)、sessionEventSubscribers(session 級,operator UI 在 run 已在跑後附加)、sessionMessageSubscribers(session 級 chat 訊息訂閱,控制 UI 隱藏後定向投遞)。
設計動機
為什麼不在 agent 主循環裡直接 socket.send?同一個 session 可能被多個用戶端訂閱,用戶端可能跨多個閘道節點,agent 事件還要按 stream 類型(agent / chat / session.tool / sessions.changed)分流,有些事件要節流 (throttle)、有些要丟棄 (drop if slow)。把這些策略塞進主循環會讓循環體變成一鍋粥——主循環只該關心 agent 狀態機,廣播策略單獨抽一層。
三個原語由 createGatewayNodeSessionRuntime(createGatewayNodeSessionRuntime:971)在外層注入,createAgentEventHandler 只消費,不關心實作——chat 模組可以脫離真實 WebSocket 測試。
關鍵檔案
AgentEventHandlerOptions:255-293— 三個廣播原語 + 三類訂閱 registry 型別宣告。createAgentEventHandler 簽名:299-317— 注入廣播原語,返回事件 dispatcher。resolveNodeSessionDeliveryKeys:495— global session 展開成agent:DEFAULT:global+global兩個 delivery key。sendChatPayload:890— chat 事件投遞,控制 UI 可見走 broadcast,不可見走定向。sendAgentPayload:967— agent 事件投遞,同上分流。dispatcher 閉包:1136—(evt: AgentEventPayload) => {...}事件入口。createGatewayBroadcaster:100-209—broadcast/broadcastToConnIds真實實作,含 scope 過濾和背壓 (backpressure)。createGatewayNodeSessionRuntime:14-53—nodeSendToSession經由nodeSubscriptions路由。
資料流
createAgentEventHandler(createAgentEventHandler:299)接收三個廣播原語和三類訂閱 registry,返回事件 dispatcher (evt: AgentEventPayload) => void。內部維護 throttle、lifecycle retry、spawnedBy cache 等狀態。
broadcast 和 broadcastToConnIds 來自 createGatewayBroadcaster(createGatewayBroadcaster:100),實作遍歷 clients: Set<GatewayWsClient>,核心循環:
for (const c of params.clients) {
if (targetConnIds && !targetConnIds.has(c.connId)) continue;
if (!hasEventScope(c, event)) continue; // scope 過濾
const slow = c.socket.bufferedAmount > MAX_BUFFERED_BYTES;
if (slow && opts?.dropIfSlow) {
if (!isTargeted) clientSeq.set(c, (clientSeq.get(c) ?? 0) + 1); // 定向不消費 seq
continue; // 可丟就丟
}
if (slow) { c.socket.close(1008, "slow consumer"); continue; } // 不可丟就斷開
c.socket.send(frame); // {"type":"event","event":...,"payload":...,"seq":N}
}背壓分兩檔:dropIfSlow: true(chat delta 高頻可丟流)丟棄但更新 seq 避免 seq 錯位;dropIfSlow: false(關鍵 lifecycle、tool 結果)直接 socket.close(1008, "slow consumer") 斷連,寧可斷連也不丟關鍵事件。seq 只在非定向廣播時遞增(isTargeted ? undefined : nextSeq),定向廣播是已過濾子集,不參與用戶端重排。
nodeSendToSession 來自 createGatewayNodeSessionRuntime(nodeSendToSession:31),實作是單行 nodeSubscriptions.sendToSession(sessionKey, event, payload, nodeSendEvent)。原始碼註解「session fanout goes through the subscription manager so node reconnects and explicit unsubscribes keep both node->session indexes in sync」說明為什麼跨節點投遞要經過 subscription manager 而非直接 nodeRegistry.send——節點重連、顯式 unsubscribe 時雙向索引必須同步,否則會投遞到已斷開的 node,或漏掉新連上的 node。
事件入口是返回的閉包(dispatcher:1136),先解析 sessionKey、agentId、clientRunId(很多事件不帶這些欄位,要從 chatRunState.registry 反查),再按 stream 類型分流。sendChatPayload(sendChatPayload:890)按控制 UI 可見性分流:可見時 broadcast("chat", payload) 加 sendNodeSessionPayloadForAgent 跨節點投遞;隱藏時查 sessionMessageSubscribers.get(deliverySessionKey),有訂閱者就 broadcastToConnIds 定向投遞——通常是 operator 後台工具,主面板看不到這次 run。sendAgentPayload(sendAgentPayload:967)分流邏輯對稱,只是 event 名換成 "agent"。
跨節點投遞還有個細節(resolveNodeSessionDeliveryKeys:495):sessionKey 是 "global" 時展開成 agent:DEFAULT:global 和 global 兩個 delivery key。全域事件既要被訂閱「某個 agent 的 global」的用戶端收到,也要被訂閱「全域」的用戶端收到;但只有 agentId 就是預設 agent 時才同時推兩個,否則只推 agent-scoped——避免重複投遞。
邊界與失敗
- 慢消費者處理:
socket.bufferedAmount > MAX_BUFFERED_BYTES觸發背壓。dropIfSlow: true(串流 delta、chat.send_timing)跳過發送但更新 seq;dropIfSlow: false(關鍵 lifecycle、tool 結果)直接socket.close(1008, "slow consumer")斷連,不斷連會一直 buffer 拖垮進程(backpressure:157)。 - lifecycle 錯誤重試:
pendingTerminalLifecycleErrorsMap 記錄 end/error phase 事件在 grace period 內的重試狀態。lifecycleErrorRetryGraceMs控制視窗,新 lifecycle generation 到達時取消舊 timer。 - spawnedBy cache:
resolveSpawnedBy只對 subagent/ACP sessionKey 走儲存查詢,非 lineage 直接返回 null 不污染快取(spawnedByCache:383)。delta/flush/final 高頻路徑避免每次打 session 儲存。 - 重啟恢復 (restart recovery) 抑制:
resolveRestartRecoveryLifecycleState檢查 lifecycle 事件是否是 restart recovery 產物,是就suppress: true不廣播——避免重啟後被恢復的舊 run 事件重新彈給用戶端。 - sessionKey 解析回退:優先
chatLink.sessionKey,其次evt.sessionKey,最後resolveSessionKeyForRun(evt.runId)。三層回退保證事件總能找到歸屬 session。 - 跨節點投遞不消費 seq:
isTargeted ? undefined : nextSeq,定向廣播也遞增會讓用戶端看到 seq 跳號,誤以為丟事件。 - 節點訂閱同步:
nodeSendToSession走 subscription manager 而非直接nodeRegistry.send,讓 node 重連/斷開/unsubscribe 時雙向索引同步。繞過 manager 會在 node 重連後漏投。
小結
廣播層把 agent 事件投遞策略從主循環剝離,抽象成三個原語,由 createGatewayNodeSessionRuntime 在外層注入,createAgentEventHandler 只消費。真實廣播實作(createGatewayBroadcaster)處理 scope 過濾、背壓、seq 編號;跨節點投遞經由 subscription manager 路由,保證節點重連後索引同步。事件入口是 dispatcher 閉包,按 stream 類型分流到 chat / agent / session.tool / sessions.changed,每個 stream 又按控制 UI 可見性走全量或定向。
請求怎麼進來被分派,看 RPC 方法表與請求分派;廣播原語怎麼注入到整個閘道,看 閘道核心;session 訂閱狀態怎麼和 agent run 生命週期綁定,看 Session Manager。
對照官方資料:Gateway 文件 · README。