Cœur de la passerelle
Responsabilité
La passerelle (gateway) est le pont entre l'agent et le monde extérieur : elle reçoit les messages des plateformes (Telegram/Discord/Slack/Signal/WeChat…) et les transforme en un MessageEvent unifié qui alimente l'agent ; elle retransforme la sortie en streaming de l'agent en actions d'envoi propres à chaque plateforme. gateway/run.py est le processus principal de cette couche (GatewayRunner), qui coordonne plusieurs sous-systèmes : adaptateurs (adapters), sessions (sessions), consommation en streaming, registre de livraison (delivery ledger), hooks, heartbeat cron…
Mot de conception
Séparer « parler aux plateformes » de « faire tourner la boucle agent » permet à l'agent d'ignorer les limites de longueur Telegram, les snowflake IDs Discord, le schéma RPC signal-cli. Unifier MessageEvent + abstraire BasePlatformAdapter permet à un processus agent de servir une dizaine de plateformes simultanément, l'ajout d'une plateforme ne touchant que l'adaptateur. Par-dessus, delivery_ledger existe parce que la passerelle est longue-vie et peut crasher à tout moment — la réponse finale déjà générée par l'agent ne doit pas périr avec elle, donc cette responsabilité revient à la passerelle plutôt qu'à l'agent.
Fichiers clés
class GatewayRunner:3029— classe principale de la passerelle, mixe les capacités d'autorisation / Kanban / Slash_handle_message:9947— entrée de distribution des messages entrants_handle_message_with_agent:11956— confie un message à l'agent pour un tour_start_one_profile_adapter / _start_secondary:9389-9947— démarre les adaptateurs de plateforme pour chaque profil_start_stream_consumer:21218— démarre le consommateur de streaming_start_gateway_housekeeping:22246— ronde périodique (60 s)_start_cron_ticker:22335— heartbeat cron (60 s)def main:22965— entrée du processus passerelleexceptions et masquage:276-340—_gateway_loop_exception_handler/_redact_gateway_user_facing_secretsPlatformEntry:39-162— dataclass de l'entrée de plateformeclass PlatformRegistry:162-260— registre (register / register_deferred / unregister)
Flux de données
main()(gateway/run.py:22965) démarreGatewayRunner._start_one_profile_adapter:9389lance un par un les adaptateurs de plateforme (voir Adaptateurs de plateforme).- L'adaptateur reçoit l'événement brut de la plateforme et le transforme en
MessageEventunifié (voirgateway/platforms/base.py:1759). _handle_message(gateway/run.py:9947) distribue : vérification d'autorisation, résolution de session, routage des slash commands, ou délégation à_handle_message_with_agent(gateway/run.py:11956).- Ce dernier appelle
run_conversation:588de l'agent, et branche le callback de streaming surGatewayStreamConsumer:83. - La sortie en streaming édite progressivement le message de la plateforme via le consommateur ; l'état final est comptabilisé par
delivery_ledger, pour éviter qu'un crash ne perde la réponse finale. - La ronde de housekeeping et le heartbeat cron tournent en arrière-plan à 60 s (
gateway/run.py:22246/gateway/run.py:22335).
GatewayRunner (gateway/run.py:3029) est une classe 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] = {}Autorisation / Kanban watcher / Slash commands sont fournies par trois mixins séparés ; la classe principale ne gère que le cycle de vie et le routage — « classe principale fine + mixins de capacités » permet à chaque trait d'être testé isolément et taillé au besoin.
Les messages entrants passent par _handle_message, qui en tête de tube fait d'abord les choses qu'on ne peut pas sauter — reset des ContextVar cross-session, détection de la période startup-restore, pose de l'horloge scale-to-zero, exécution du hook plugin (plugin) pre_gateway_dispatch :
async def _handle_message(self, event: MessageEvent) -> Optional[str]:
source = event.source
# Cross-session leak guard: un nouveau task peut hériter de sessions sœurs
# HERMES_SESSION_* ContextVar, d'abord reset en _UNSET avant le flux normal.
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() snapshot le contexte (context) qui spawn ; si, dans une session (session) parente, on a fait set_session_vars(), la nouvelle task hérite de l'identité sibling et le pont subprocess-env lit alors la mauvaise session.
_handle_message_with_agent trouve/crée la session + pousse le message à l'agent, en se prémunissant contre la « résurrection illégitime d'une session morte » :
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): la spawning session a pu /new-reset,
# ne pas basculer aveuglément pour ressusciter une session que l'utilisateur a fermée.
...Au backfill d'un sous-agent asynchrone, si la spawning session a été terminée, la passerelle préfère jeter cette injection plutôt que de la ressusciter.
Limites et échecs
- Fuite ContextVar cross-session :
_handle_messagetourne dans uncreate_task(), qui hérite du snapshot du parent. Si l'on nereset_session_vars()pas d'abord, des messages concurrents lisent l'identité session d'un sibling, et le pont subprocess écrit la mauvaise identité dans l'environnement. - Injection backfill dans une session morte : au retour du sous-agent, la spawning session a pu être
/new-reset; un basculement aveugle ressusciterait une session que l'utilisateur a fermée. Fail closed laisse le résultat dans les delegation records. - Pollution de l'horloge scale-to-zero : les événements internes (fin d'un process background, replay startup-restore) ne doivent pas compter comme trafic, sinon la passerelle inactive est gardée vivante. Le code les aiguille explicitement via
is_internal.
Résumé
Le cœur de la passerelle est un « bus d'événements + pont de streaming » : l'abstraction unifiée MessageEvent masque les différences de plateforme, le GatewayStreamConsumer traduit le callback synchrone de l'agent en livraison asynchrone à la plateforme, le delivery_ledger garantit que l'état final n'est pas perdu. Il est découplé de l'agent — l'agent se contente de tourner sa boucle, les détails de livraison restent côté passerelle.