チャット放送とイベント配送
責務
server-chat.ts は agent メインループが吐き出したイベント (agent event) をクライアントに配送します。イベントは生産せず,分配型ディスパッチだけを行います:AgentEventPayload(delta、final、tool、lifecycle 等)を受け取り,イベントタイプ、セッション (session) 帰属、クライアント購読状況に応じて,すべての可視接続に放送 (broadcast) するか,一組の connId に方向付き配送 (broadcastToConnIds) するか,別ノードで該当 session を購読したプロセスへノードまたぎ配送 (nodeSendToSession) するかを決定します。
この層は 3 種の購読状態を管理します: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 (throttle)、あるイベントは遅いときに drop する必要があります。これらのポリシーをメインループに詰め込むと,ループ本体がカオスになります——メインループは agent 状態機械だけを気にし,放送ポリシーは独立した層として抽出すべきです。
3 つのプリミティブは createGatewayNodeSessionRuntime(createGatewayNodeSessionRuntime:971)が外層で注入し,createAgentEventHandler は消費するだけで実装を気にしません——chat モジュールは実際の WebSocket なしでテストできます。
主要ファイル
AgentEventHandlerOptions:255-293— 3 つの放送プリミティブ + 3 種の購読 registry 型宣言。createAgentEventHandler シグネチャ:299-317— 放送プリミティブを注入し,イベント dispatcher を返す。resolveNodeSessionDeliveryKeys:495— global session をagent:DEFAULT:global+globalの 2 つの 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)は 3 つの放送プリミティブと 3 種の購読 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}
}背圧は 2 段階:dropIfSlow: true(chat delta の高頻度で捨て可能ストリーム)は捨てつつ seq を更新し seq ずれを回避,dropIfSlow: false(重要 lifecycle、tool 結果)は直接 socket.close(1008, "slow consumer") で切断,切断しても重要イベントは捨てません。seq は非方向付き放送時だけ増分(isTargeted ? undefined : nextSeq)し,方向付き放送はフィルタ済みサブセットで,クライアント再整列に参加しません。
nodeSendToSession は createGatewayNodeSessionRuntime(nodeSendToSession:31)から来て,実装は 1 行 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」は,なぜノードまたぎ配送が直接 nodeRegistry.send ではなく subscription manager を経由するかを説明します——ノード再接続、明示的 unsubscribe 時に双方向インデックスが同期しなければならず,さもなくば切断済みノードに配送したり新規接続ノードを見落としたりします。
イベントエントリは返されるクロージャ(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)の分流ロジックは対称で,イベント名が "agent" に変わるだけです。
ノードまたぎ配送にもう一つ詳細が(resolveNodeSessionDeliveryKeys:495)あります:sessionKey が "global" のとき agent:DEFAULT:global と global の 2 つの 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)。3 層フォールバックでイベントは常に帰属 session を見つけられます。 - ノードまたぎ配送は seq を消費しない:
isTargeted ? undefined : nextSeq,方向付き放送も増分するとクライアントが seq 飛びを sees,イベントを失ったと誤認します。 - ノード購読同期:
nodeSendToSessionは直接nodeRegistry.sendではなく subscription manager を経由し,ノード再接続/切断/unsubscribe 時に双方向インデックスを同期。manager をバイパスするとノード再接続後に配送漏れが発生。
まとめ
放送層は agent イベント配送ポリシーをメインループから剥離し,3 つのプリミティブに抽象化し,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。