Núcleo del gateway
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
class GatewayRunner:3029— clase principal del gateway, mezcla capacidades de auth/Kanban/Slash_handle_message:9947— entrada de despacho (dispatch) de mensajes entrantes_handle_message_with_agent:11956— entrega el mensaje al agent para correr un turno (turn)_start_one_profile_adapter / _start_secondary:9389-9947— arranca los adaptadores de plataforma de cada profile_start_stream_consumer:21218— arranca el consumidor de streams_start_gateway_housekeeping:22246— inspección periódica (60s)_start_cron_ticker:22335— latido de cron (60s)def main:22965— entrada del proceso gatewayexcepciones y redacción:276-340—_gateway_loop_exception_handler/_redact_gateway_user_facing_secretsPlatformEntry:39-162— dataclass de registro (registry) de plataformaclass PlatformRegistry:162-260— registro (register / register_deferred / unregister)
Flujo de datos
main()(gateway/run.py:22965) arrancaGatewayRunner._start_one_profile_adapter:9389levanta uno a uno los adaptadores de plataforma (ver Adaptadores de plataforma).- El adaptador recibe eventos nativos de la plataforma y los traduce a un
MessageEventunificado (vergateway/platforms/base.py:1759). _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).- Este último llama a
run_conversation:588del agent y conecta el callback de streaming alGatewayStreamConsumer:83. - La salida en streaming edita progresivamente el mensaje de la plataforma a través del consumidor; el estado final lo anota
delivery_ledgerpara que un crash no pierda la respuesta final. - 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:
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:
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() 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»:
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_messagecorre dentro de uncreate_task()y hereda el snapshot del padre. Si no se llamareset_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.