Skip to content

Núcleo del gateway

源码版本v2026.7.20

Responsabilidad

El gateway es el puente entre el agent y el mundo exterior: recibe mensajes de cada plataforma (Telegram/Discord/Slack/Signal/WeChat…), los convierte a un MessageEvent unificado que entrega al agent, y traduce la salida en streaming del agent en acciones de envío a cada plataforma. gateway/run.py es el proceso principal de esta capa (GatewayRunner), y coordina adaptadores (adapter), sesiones (session), consumo de streams, libro (ledger) de entrega, hooks y latido de cron.

Motivo de diseño

Si se separa «tratar con la plataforma» de «correr el bucle del agent», el agent no tiene que preocuparse por el límite de longitud de Telegram, los snowflake IDs de Discord o el esquema RPC de signal-cli. Un MessageEvent unificado + una BasePlatformAdapter abstracta permiten que un único proceso agent cuelgue a la vez una docena de plataformas; añadir una nueva toca solo el adaptador. Encima, delivery_ledger se pone porque el gateway está siempre activo y el proceso puede caer en cualquier momento — la respuesta final ya generada por el agent no puede perderse con el proceso; de eso se encarga esta capa y no el agent.

Archivos clave

Flujo de datos

  1. main() (gateway/run.py:22965) arranca GatewayRunner.
  2. _start_one_profile_adapter:9389 levanta uno a uno los adaptadores de plataforma (ver Adaptadores de plataforma).
  3. El adaptador recibe eventos nativos de la plataforma y los traduce a un MessageEvent unificado (ver gateway/platforms/base.py:1759).
  4. _handle_message (gateway/run.py:9947) despacha: comprobación de autorización, resolución de sesión, enrutado de slash commands, o entrega a _handle_message_with_agent (gateway/run.py:11956).
  5. Este último llama a run_conversation:588 del agent y conecta el callback de streaming al GatewayStreamConsumer:83.
  6. La salida en streaming edita progresivamente el mensaje de la plataforma a través del consumidor; el estado final lo anota delivery_ledger para que un crash no pierda la respuesta final.
  7. La inspección y el latido de cron corren en segundo plano cada 60s (gateway/run.py:22246 / gateway/run.py:22335).

GatewayRunner (gateway/run.py:3029) es una clase multi-mixin:

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] = {}

Auth / Kanban watcher / slash commands vienen de tres mixins separados; la clase principal solo se ocupa de lifecycle y enrutado — «clase principal delgada + mixins de capacidad» permite testear cada capacidad aislada y recortarla bajo demanda.

El mensaje entrante pasa por _handle_message; al principio del pipeline hace varias cosas que no se pueden saltar — resetear el ContextVar entre sesiones, detectar el periodo de startup-restore, marcar el reloj de scale-to-zero y lanzar el hook de plugin 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() snapshottea el contexto (context) desde el que se spawnea; si mensajes concurrentes hicieron set_session_vars() en la tarea padre, el nuevo task hereda la identidad sibling y el puente subprocess-env leería la sesión equivocada.

_handle_message_with_agent encuentra/crea la sesión y le entrega el mensaje al agent, evitando «revivir ilegalmente una sesión muerta»:

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,
        # 不能盲切回去复活用户主动关掉的会话。
        ...

Cuando un subagent asíncrono vuelve a rellenar y la sesión que lo spawneó ya terminó, el gateway prefiere perder esa inyección antes que revivirla.

Límites y fallos

  • Fuga de ContextVar entre sesiones: _handle_message corre dentro de un create_task() y hereda el snapshot del padre. Si no se llama reset_session_vars() antes, mensajes concurrentes pueden leer la identidad de sesión de un sibling; el puente subprocess-env la escribiría en el entorno equivocado.
  • Inyección asíncrona a una sesión muerta: cuando un subagent vuelve, la sesión que lo spawneó puede haber hecho /new-reset; seguir ciegamente reviviría una sesión que el usuario cerró a propósito. Fail closed deja el resultado en delegation records.
  • Contaminación del reloj scale-to-zero: los eventos internos (procesos en segundo plano terminando, replay de startup-restore) no deben contar como tráfico, si no un gateway ocioso se mantiene vivo. El código los separa explícitamente con is_internal.

Resumen

El núcleo del gateway es «bus de eventos + puente de streaming»: la abstracción unificada de MessageEvent oculta las diferencias entre plataformas, GatewayStreamConsumer puentea el callback síncrono del agent al envío asíncrono en cada plataforma, y delivery_ledger garantiza que el estado final no se pierda. Está desacoplado del agent: el agent solo corre el bucle; los detalles de entrega viven todos en el gateway.

Sitio de aprendizaje comunitario no oficial. Basado en el código fuente de NousResearch/hermes-agent (licencia MIT).