Skip to content

チャット放送とイベント配送

源码版本v2026.6.11

責務

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 なしでテストできます。

主要ファイル

データフロー

createAgentEventHandler(createAgentEventHandler:299)は 3 つの放送プリミティブと 3 種の購読 registry を受け取り,イベント dispatcher (evt: AgentEventPayload) => void を返します。内部で throttle、lifecycle retry、spawnedBy cache などの状態を管理します。

broadcastbroadcastToConnIdscreateGatewayBroadcaster(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}
}

背圧は 2 段階:dropIfSlow: true(chat delta の高頻度で捨て可能ストリーム)は捨てつつ seq を更新し seq ずれを回避,dropIfSlow: false(重要 lifecycle、tool 結果)は直接 socket.close(1008, "slow consumer") で切断,切断しても重要イベントは捨てません。seq は非方向付き放送時だけ増分(isTargeted ? undefined : nextSeq)し,方向付き放送はフィルタ済みサブセットで,クライアント再整列に参加しません。

nodeSendToSessioncreateGatewayNodeSessionRuntime(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:globalglobal の 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 エラー再試行: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)。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