ストリーミングと Hooks
源码版本v2026.7.20
職務
① ストリーミング消費——agent の worker スレッドは同期的に stream_delta_callback(text) を呼ぶが、プラットフォーム送信は非同期。GatewayStreamConsumer が同期増分を非同期に橋渡し、同一のプラットフォームメッセージを段階的に編集し、ユーザーに「タイプライター」効果を見せる。② Hooks——メッセージライフサイクルのノード(受信、生成前、生成後、配信)にカスタムロジックを差し込む。セキュリティフィルタ、課金、Kanban 同期など。③ Relay——インスタンス/プロファイル間のメッセージ中継。
主要ファイル
class GatewayStreamConsumer:83-120— 非同期消費器本体、on_deltaがコールバックとして agent に渡されるclass StreamConsumerConfig:55-83— 単回消費の実行時設定(native draft ストリーミング戦略、最終編集遅延など)モジュールドキュメント:1-40— 同期→非同期橋渡しの原理と sentinel 完了シグナルの説明stream_dispatch— ストリームイベントディスパッチstream_events— ストリームイベント型定義hooks— ライフサイクル hook 登録フレームワークbuiltin_hooks/ ディレクトリ— 内蔵 hook 実装relay/ ディレクトリ— インスタンス間中継slash_commands—/コマンドルーティングauthz_mixin— 認可(GatewayRunnerにミックスイン)profile_routing— マルチ profile ルーティング
データフロー
GatewayRunnerが_start_stream_consumer:21218を起動する。GatewayStreamConsumerを構築し、consumer.on_deltaをstream_delta_callbackとしてAIAgent(run_agent.py:400) に渡す。- agent は worker スレッドで同期的に
on_delta(text)を呼ぶ → 消費器が増分を非同期キューに積む。 - 消費器コルーチンがキューから増分を取り出し、
StreamConsumerConfig:55戦略に従う:auto/draft:ネイティブ draft ストリーミングを優先(send_draft:2650)、フォールバックは通常編集- 最終編集は遅延し得る(遅い推論モデルのケース)、高頻度ジッタを避ける
- ストリーム完了後、sentinel シグナルが最終編集をトリガーし、
delivery_ledger:155に記帳させる。 - 全程で各ノードが
hooksをトリガー(受信/生成前/生成後/配信)、内蔵 hook(gateway/builtin_hooks/) がセキュリティ、課金、Kanban などを実行する。
まとめ
ストリーミング消費器は「同期 agent コールバック ↔ 非同期プラットフォーム配信」のバッファと節流層で、二つのクロックドメインの不整合を解決する。Hooks は横断関心事(セキュリティ/課金/同期)の挿入点。両者により、ゲートウェイは agent 主幹を汚さずにリアルタイムタイプライター、セキュリティフィルタ、インスタンス間 relay をサポートできる。