Skip to content

聊天广播与事件投递

源码版本v2026.6.11

职责

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 测试。

关键文件

数据流

createAgentEventHandler(createAgentEventHandler:299)接收三个广播原语和三类订阅 registry,返回事件 dispatcher (evt: AgentEventPayload) => void。内部维护 throttle、lifecycle retry、spawnedBy cache 等状态。

broadcastbroadcastToConnIds 来自 createGatewayBroadcaster(createGatewayBroadcaster:100),实现遍历 clients: Set<GatewayWsClient>,核心循环:

typescript
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:globalglobal 两个 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 错误重试:pendingTerminalLifecycleErrors Map 记录 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