聊天广播与事件投递
职责
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。