Skip to content

Cœur de la passerelle

源码版本v2026.7.20

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

Flux de données

  1. main() (gateway/run.py:22965) démarre GatewayRunner.
  2. _start_one_profile_adapter:9389 lance un par un les adaptateurs de plateforme (voir Adaptateurs de plateforme).
  3. L'adaptateur reçoit l'événement brut de la plateforme et le transforme en MessageEvent unifié (voir gateway/platforms/base.py:1759).
  4. _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).
  5. Ce dernier appelle run_conversation:588 de l'agent, et branche le callback de streaming sur GatewayStreamConsumer:83.
  6. 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.
  7. 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 :

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

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 :

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

create_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 » :

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): 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_message tourne dans un create_task(), qui hérite du snapshot du parent. Si l'on ne reset_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.

Site d'apprentissage communautaire non officiel. Basé sur le code source de NousResearch/hermes-agent (licence MIT).