Skip to content

ストリーミングと Hooks

源码版本v2026.7.20

職務

ストリーミング消費——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 を作る(オーバーヘッドが大きく、エラー処理が難しい)の二択になる。GatewayStreamConsumerqueue.Queue でスレッド境界を作る——worker スレッドは queue.put するだけで、消費コルーチンが asyncio 側から get_nowait で取り出す。両側のクロックドメインが分離され、節流、マージ、フォールバック (fallback) はすべて消費側に足す。Hooks を独立させたのは、セキュリティ/課金のような横断関心事が主メッセージパイプを汚すべきでないからだ。

主要ファイル

データフロー

  1. GatewayRunner_start_stream_consumer:21218 を起動する。
  2. GatewayStreamConsumer を構築し、consumer.on_deltastream_delta_callback として AIAgent(run_agent.py:400) に渡す。
  3. agent は worker スレッドで同期的に on_delta(text) を呼ぶ → 消費器が増分を非同期キューに積む。
  4. 消費器コルーチンがキューから増分を取り出し、StreamConsumerConfig:55 戦略に従う:
    • auto/draft:ネイティブ draft ストリーミングを優先(send_draft:2650)、フォールバックは通常編集
    • 最終編集は遅延し得る(遅い推論モデルのケース)、高頻度ジッタを避ける
  5. ストリーム完了後、sentinel シグナルが最終編集をトリガーし、delivery_ledger:155 に記帳させる。
  6. 全程で各ノードが hooks をトリガー(受信/生成前/生成後/配信)、内蔵 hook(gateway/builtin_hooks/) がセキュリティ、課金、Kanban などを実行する。

on_delta は消費器の对外入口で、一つだけ——worker スレッドの token を queue.Queue に押し込む。

python
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() は増分をバッチで取り、戦略に従って送る。

python
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)を必要とする」ための拡張点だ。

python
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 をサポートできる。

非公式コミュニティ学習サイト。MIT ライセンスの NousResearch/hermes-agent ソースに基づく。