网关核心
职责
网关是 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。
关键文件
class GatewayRunner:3029— 网关主类(混入授权/Kanban/Slash)_handle_message:9947— 入站消息分发入口_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— 注册表
数据流
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 插件钩子:
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 replay)不应算流量,否则空闲网关被保活。代码用
is_internal显式分流。
小结
网关核心是「事件总线 + 流式桥接」:统一 MessageEvent 抽象屏蔽平台差异,GatewayStreamConsumer 把同步 agent 回调桥到异步平台投递,delivery_ledger 保证最终态不丢。它和 agent 解耦——agent 只管跑循环,投递细节全在网关。