ゲートウェイコア
職務
ゲートウェイ (gateway) は agent と外部世界の橋:各プラットフォーム(Telegram/Discord/Slack/Signal/微信…)からのメッセージを受け取り、統一 MessageEvent に変換して agent に渡し、agent のストリーミング出力を各プラットフォームの送信アクションに変換する。gateway/run.py がこの層のメインプロセス(GatewayRunner)で、アダプタ (adapter)、セッション (session)、ストリーミング消費、配信台帳 (delivery ledger)、hooks、cron 心拍など複数のサブシステムを協調させる。
設計動機
「プラットフォームと付き合う」ことと「agent ループを回す」ことを切り分けることで、agent は Telegram の長さ上限や Discord のスノーフレーク ID、Signal の signal-cli の RPC schema を気にしなくて済む。統一 MessageEvent + 抽象 BasePlatformAdapter により、一つの agent プロセスに十以上のプラットフォームを同時に掛けられ、新プラットフォームの追加はアダプタをいじるだけだ。さらに delivery_ledger を重ねるのは、ゲートウェイが長時間オンラインで、プロセスがいつ落ちてもよいからだ——agent が生成した最終返信が一緒に消えてはいけない。だからこの層は agent ではなくゲートウェイに属する。
主要ファイル
class GatewayRunner:3029— ゲートウェイメインクラス、認可/Kanban/Slash 能力をミックスイン_handle_message:9947— 入站メッセージディスパッチ (dispatch) 入口_handle_message_with_agent:11956— メッセージを agent に渡して一往復実行_start_one_profile_adapter / _start_secondary:9389-9947— 各 profile のプラットフォームアダプタを起動_start_stream_consumer:21218— ストリーミング消費器を起動_start_gateway_housekeeping:22246— 周期巡回(60s)_start_cron_ticker:22335— cron 心拍(60s)def main:22965— ゲートウェイプロセス入口例外と秘匿化:276-340—_gateway_loop_exception_handler/_redact_gateway_user_facing_secretsPlatformEntry:39-162— プラットフォーム登録項目データクラスclass PlatformRegistry:162-260— レジストリ (registry)(register / register_deferred / unregister)
データフロー
main()(gateway/run.py:22965) がGatewayRunnerを起動する。_start_one_profile_adapter:9389がプラットフォームアダプタを順に起動(プラットフォームアダプタ参照)。- アダプタがプラットフォーム生イベントを受け取り、統一
MessageEventに変換(gateway/platforms/base.py:1759参照)。 _handle_message(gateway/run.py:9947) がディスパッチ:認可チェック、セッション解析、slash コマンドルーティング、または_handle_message_with_agent(gateway/run.py:11956) に渡す。- 後者は agent の
run_conversation:588を呼び、ストリーミングコールバックをGatewayStreamConsumer:83に接続する。 - ストリーミング出力は消費器経由でプラットフォームメッセージを段階的に編集する。最終状態は
delivery_ledgerが記帳し、クラッシュ時にも最終返信を失わない。 - 巡回と cron 心拍はバックグラウンドで 60s 周期で走る(
gateway/run.py:22246/gateway/run.py:22335)。
GatewayRunner(gateway/run.py:3029) は多頭ミックスインクラスだ。
class GatewayRunner(GatewayAuthorizationMixin, GatewayKanbanWatchersMixin, GatewaySlashCommandsMixin):
"""Main gateway controller. Manages the lifecycle of all platform
adapters and routes messages to/from the agent."""
_busy_input_mode: str = "interrupt"
_draining: bool = False
def __init__(self, config: Optional[GatewayConfig] = None):
self.adapters: Dict[Platform, BasePlatformAdapter] = {}認可 / Kanban watcher / Slash コマンドはそれぞれ別の mixin が提供し、主クラスはライフサイクルとルーティングだけを管る——「薄い主クラス + 能力 mixin」により、各機能を独立にテストし、必要に応じて切り詰められる。
入站メッセージは _handle_message に流れ込み、パイプラインの冒頭で飛ばせないことをいくつかこなす——クロスセッションの ContextVar リセット、startup-restore 期の判定、scale-to-zero の時計の蓋、pre_gateway_dispatch プラグイン (plugin) フックの実行だ。
async def _handle_message(self, event: MessageEvent) -> Optional[str]:
source = event.source
# Cross-session leak guard: 新しい task は兄弟セッションの
# HERMES_SESSION_* ContextVar を受け継ぐ可能性がある。先に _UNSET に
# リセットしてから通常フローに入る。
try:
from gateway.session_context import reset_session_vars
reset_session_vars()
except Exception:
logger.debug("reset_session_vars failed at handler entry", exc_info=True)
if getattr(self, "_startup_restore_in_progress", False) and not getattr(event, "internal", False):
self._queue_startup_restore_event(event)
return Nonecreate_task() は spawning context をスナップショットする。並行メッセージが親タスク内で set_session_vars() を呼んでいた場合、新タスクは sibling の身分を受け継ぎ、subprocess-env 橋が誤った session を読む。
_handle_message_with_agent は session を探して作り、メッセージを agent に押し込みつつ、「死んだセッションの不正な復活」を防ぐ。
async def _handle_message_with_agent(self, event, source, _quick_key, run_generation):
"""Inner handler that runs under the _running_agents sentinel guard."""
session_entry = await self.async_session_store.get_or_create_session(source)
session_key = session_entry.session_key
pinned_session_id = str((getattr(event, "metadata", None) or {}).get("gateway_session_id") or "").strip()
if pinned_session_id and pinned_session_id != session_entry.session_id:
# Fail closed (#55578): spawning session は既に /new-reset されている
# 可能性があり、ユーザーが能動的に閉じたセッションを盲復活させてはいけない。
...非同期サブエージェントの回填時、spawning session が既に終了していたら、ゲートウェイは今回の注入を捨てても復活させない。
境界と失敗
- クロスセッション ContextVar の漏洩:
_handle_messageはcreate_task()の中で走り、親タスクのスナップショットを受け継ぐ。先にreset_session_vars()を呼ばないと、並行メッセージが sibling の session 身分を読み、subprocess 橋が誤った身分を環境に書き込む。 - 非同期回填が死んだセッションに注入する:サブエージェントが戻ったとき、spawning session は既に
/new-resetされているかもしれない。盲目的に切替えるとユーザーが能動的に閉じたセッションを復活させてしまう。fail closed で、結果は delegation records に残す。 - scale-to-zero の時計汚染:内部イベント(バックグラウンドプロセスの完了、startup-restore のリプレイ)はトラフィックとみなしてはいけない。さもないとアイドルなゲートウェイが保活されてしまう。コードは
is_internalで明示的に分流する。
まとめ
ゲートウェイコアは「イベントバス + ストリーミング橋渡し」:統一 MessageEvent 抽象でプラットフォーム差異を隠し、GatewayStreamConsumer が同期 agent コールバックを非同期プラットフォーム配信に橋渡し、delivery_ledger が最終状態を落とさない。agent とは疎結合——agent はループを回すことだけを考え、配信の詳細はすべてゲートウェイが担う。