ストリーミングと Hooks
職務
① ストリーミング消費——agent の worker スレッドは同期的に stream_delta_callback(text) を呼ぶが、プラットフォーム送信は非同期。GatewayStreamConsumer が同期増分を非同期に橋渡し、同一のプラットフォームメッセージを段階的に編集し、ユーザーに「タイプライター」効果を見せる。② Hooks——メッセージライフサイクルのノード(受信、生成前、生成後、配信)にカスタムロジックを差し込む。セキュリティフィルタ、課金、Kanban 同期など。③ Relay——インスタンス/プロファイル間のメッセージ中継。
設計動機
agent の run_conversation は worker スレッドで走り、token を同期的に yield する——これは下層 SDK が決めることだ(LLM のストリーミングインターフェースは大抵 blocking iterator)。一方、プラットフォーム送信は asyncio コルーチンで、スレッドをまたいで直接呼ぶとイベントループを踏む。中間層がないと、回答全体を貯め終わってから送る(ユーザーは数十秒「入力中…」を見て内容が見えない)か、worker スレッドで毎回 asyncio.run_coroutine_threadsafe で future を作る(オーバーヘッドが大きく、エラー処理が難しい)の二択になる。GatewayStreamConsumer は queue.Queue でスレッド境界を作る——worker スレッドは queue.put するだけで、消費コルーチンが asyncio 側から get_nowait で取り出す。両側のクロックドメインが分離され、節流、マージ、フォールバック (fallback) はすべて消費側に足す。Hooks を独立させたのは、セキュリティ/課金のような横断関心事が主メッセージパイプを汚すべきでないからだ。
主要ファイル
class GatewayStreamConsumer:83-120— 非同期消費器本体、on_deltaがコールバックとして agent に渡されるclass StreamConsumerConfig:55-83— 単回消費の実行時設定(native draft ストリーミング戦略、最終編集遅延など)モジュールドキュメント:1-40— 同期→非同期橋渡しの原理と sentinel 完了シグナルの説明stream_dispatch— ストリーム (stream) イベントディスパッチ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 などを実行する。
on_delta は消費器の对外入口で、一つだけ——worker スレッドの token を queue.Queue に押し込む。
def on_delta(self, text: str) -> None:
"""Thread-safe callback — called from the agent's worker thread.
When *text* is ``None``, signals a tool boundary: the current message
is finalized and subsequent text will be sent as a new message so it
appears below any tool-progress messages the gateway sent in between.
"""
if text:
self._queue.put(text)
elif text is None:
self.on_segment_break()
def finish(self) -> None:
"""Signal that the stream is complete."""
self._queue.put(_DONE)None は tool boundary シグナル——以降の token は新メッセージで送り、tool 進度バブルの下に差し込む。_DONE は sentinel で、消費コルーチンが見たら最終編集をトリガーする。「値 = token、None = 境界、sentinel = 完了」で一つのキューがすべての制御フローを担う。
消費コルーチン run() は増分をバッチで取り、戦略に従って送る。
async def run(self) -> None:
"""Async task that drains the queue and edits the platform message."""
_len_fn = (
self.adapter.message_len_fn
if isinstance(self.adapter, _BasePlatformAdapter)
else len
)
_raw_limit = self._raw_message_limit()
_safe_limit = max(500, _raw_limit - _len_fn(self.cfg.cursor) - 100)
self._use_draft_streaming = self._resolve_draft_streaming()
try:
while True:
if not self._run_still_current():
return
got_done = False
while True:
try:
item = self._queue.get_nowait()
if item is _DONE:
got_done = True
breakいくつか細かい点:message_len_fn はアダプタ (adapter) から来る(Telegram は utf16 長)ので、溢出判定がプラットフォームの実際の制限と一致する。_safe_limit は cursor とリッチテキスト整形に 100 文字の余裕を残す。_run_still_current() を毎ループ检查し、session が /new や /stop されたら早めに抜け、古い token が画面に送られ続けるのを防ぐ。
GatewayEventDispatcher(gateway/stream_dispatch.py) はより構造的なラップだ——agent は typed events(MessageChunk、ToolCallChunk、Commentary…)を産出し、dispatcher は adapter 能力に応じてレンダリング方法を決める。プラットフォームが tool chrome をレンダリングできないときは None を返してイベントを飲み込む。これは「プラットフォームが純テキスト編集ではなくネイティブ UI(釘釘 AI Card、Slack Block Kit)を必要とする」ための拡張点だ。
class GatewayEventDispatcher:
"""Route typed stream events through an adapter onto a delivery sink."""
# Message/commentary/segment events flow into the consumer (native draft
# on Telegram DMs, edit-in-place elsewhere). Tool events are formatted
# by the adapter — which may return None to *eat* the event on platforms
# that can't render tool chrome — and the rendered line is enqueued onto
# the same tool progress queue the gateway already drains, so the two no
# longer race through independent code paths.境界と失敗
- ストリーム中断:session が
/newや/stopされると、消費コルーチンは次ループで_run_still_current()を检查して即 return する。既にqueue.putされた増分は失われるため、最終編集前に_flush_think_buffer()が残った途中の think タグを洗い流す。 - Flood control 貫通:
_MAX_FLOOD_STRIKES = 3、連続 3 回の 429 で非漸進式に永続降格し、本ストリームの残り token は一括で送る。 - プラットフォームが edit をサポートしない:Signal は
SUPPORTED_MESSAGE_EDITING = Falseで、消費器は非 edit フォールバックに回る——完全メッセージを一回送るだけで、ストリーミング中にタイプライターは見えない。これは能力ビットの安全網であってバグではない。 - sentinel 漏発:worker スレッドが
on_deltaの後にfinish()を呼ばずに落ちると、消費コルーチンはqueue.get()で詰まる。run_still_currentが安全網——session が片付けられても退出をトリガーする。
まとめ
ストリーミング消費器は「同期 agent コールバック ↔ 非同期プラットフォーム配信」のバッファと節流層で、二つのクロックドメインの不整合を解決する。Hooks は横断関心事(セキュリティ/課金/同期)の挿入点。両者により、ゲートウェイ (gateway) は agent 主幹を汚さずにリアルタイムタイプライター、セキュリティフィルタ、インスタンス間 relay をサポートできる。