Skip to content

网关核心

源码版本v2026.7.20

职责

网关是 agent 与外部世界的桥梁:接收各平台(Telegram/Discord/Slack/Signal/微信…)消息,转成统一的 MessageEvent 喂给 agent;把 agent 的流式输出转成各平台的发送动作。gateway/run.py 是这层的主进程(GatewayRunner),协调适配器、会话、流式消费、投递账本、hooks、cron 心跳等多条子系统。

设计动机

把「跟平台打交道」和「跑 agent 循环」切开,agent 才不必关心 Telegram 长度上限、Discord 雪花 ID、Signal signal-cli 的 RPC schema。统一 MessageEvent + 抽象 BasePlatformAdapter 让一个 agent 进程同时挂十几个平台,新增平台只动适配器。再叠 delivery_ledger,是因为网关长时在线,进程随时可能崩——agent 已生成的最终回复不能跟着丢,所以这层归网关而不是 agent。

关键文件

数据流

  1. main()(gateway/run.py:22965)启动 GatewayRunner
  2. _start_one_profile_adapter:9389 逐个拉起平台适配器(见平台适配器)。
  3. 适配器把平台原始事件转成统一 MessageEvent(gateway/platforms/base.py:1759)。
  4. _handle_message(gateway/run.py:9947)分发:授权、会话、slash,或交给 _handle_message_with_agent(gateway/run.py:11956)。
  5. 后者调 agent run_conversation:588,流式回调接到 GatewayStreamConsumer:83
  6. 流式输出经消费器渐进式编辑平台消息;最终态由 delivery_ledger 记账。
  7. 巡检与 cron 心跳在后台 60s 周期跑(gateway/run.py:22246 / gateway/run.py:22335)。

GatewayRunner(gateway/run.py:3029)是个多头混入类:

python
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 插件钩子:

python
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 None

create_task() 会快照 spawning context,并发消息若在父任务里 set_session_vars() 过,新任务会继承 sibling 身份,subprocess-env 桥就读到错误 session。

_handle_message_with_agent 找/建 session + 把消息塞进 agent,并防「非法复活已死会话」:

python
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 replay)不应算流量,否则空闲网关被保活。代码用 is_internal 显式分流。

小结

网关核心是「事件总线 + 流式桥接」:统一 MessageEvent 抽象屏蔽平台差异,GatewayStreamConsumer 把同步 agent 回调桥到异步平台投递,delivery_ledger 保证最终态不丢。它和 agent 解耦——agent 只管跑循环,投递细节全在网关。

非官方社区学习站,内容以 MIT 许可的 NousResearch/hermes-agent 源码为依据。