Skip to content

Gateway-Kern

源码版本v2026.7.20

Verantwortung

Das Gateway (gateway) ist die Brücke zwischen Agent und Außenwelt: Es empfängt Nachrichten der Plattformen (Telegram/Discord/Slack/Signal/WeChat …), wandelt sie in ein einheitliches MessageEvent um und füttert den Agenten; den Streaming-Output des Agenten übersetzt es in Sende-Aktionen an die jeweilige Plattform-API. gateway/run.py ist der Hauptprozess dieser Schicht (GatewayRunner) und koordiniert Adapter (adapter), Sitzungen (sessions), Stream-Konsum (stream consumption), Delivery-Ledger (delivery ledger), Hooks, Cron-Heartbeat und weitere Subsysteme.

Designmotiv

«Plattformgeschäft» und «Agent-Schleife» zu trennen, befreit den Agenten davon, von Telegram-Längengrenzen, Discord-Schnee-IDs oder dem RPC-Schema von signal-cli wissen zu müssen. Ein einheitliches MessageEvent plus abstraktem BasePlatformAdapter erlaubt einem Agent-Prozess, über ein Dutzend Plattformen gleichzeitig zu bedienen; neue Plattformen berühren nur den Adapter. delivery_ledger wird darübergelegt, weil das Gateway lange online ist und jederzeit crashen kann — die bereits generierte Endantwort des Agenten darf nicht mitgerissen werden, deshalb liegt diese Schicht im Gateway, nicht im Agenten.

Schlüsseldateien

Datenfluss

  1. main() (gateway/run.py:22965) startet den GatewayRunner.
  2. _start_one_profile_adapter:9389 bringt nacheinander die Plattform-Adapter hoch (siehe Plattform-Adapter).
  3. Der Adapter empfängt das native Plattform-Event und wandelt es in ein einheitliches MessageEvent um (siehe gateway/platforms/base.py:1759) um.
  4. _handle_message (gateway/run.py:9947) dispatched: Autorisierungsprüfung, Sitzungs-Auflösung (session resolution), Slash-Command-Routing oder Übergabe an _handle_message_with_agent (gateway/run.py:11956).
  5. Letzteres ruft run_conversation:588 des Agenten auf und verdrahtet den Streaming-Callback an GatewayStreamConsumer:83.
  6. Der Streaming-Output wird vom Konsumenten inkrementell editiert auf der Plattform-Nachricht; der Endzustand wird vom delivery_ledger verbucht, damit ein Crash die finale Antwort nicht verliert.
  7. Inspektion und Cron-Heartbeat laufen im Hintergrund mit 60s-Periode (gateway/run.py:22246 / gateway/run.py:22335).

GatewayRunner (gateway/run.py:3029) ist eine Multi-Head-Mixin-Klasse:

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

Authorization/Kanban-Watcher/Slash-Command stammen aus drei separaten Mixins; die Hauptklasse verwaltet nur Lebenszyklus und Routing — «schlanke Hauptklasse + Fähigkeits-Mixin» macht jede Eigenschaft einzeln testbar und bei Bedarf abstreichbar.

Eingehende Nachrichten laufen durch _handle_message; am Anfang der Pipeline stehen ein paar unausweichliche Dinge — ContextVar reset, Prüfung auf Startup-Restore-Phase, Scale-to-Zero-Uhr aufziehen, pre_gateway_dispatch-Plugin (plugin)-Hook:

python
async def _handle_message(self, event: MessageEvent) -> Optional[str]:
    source = event.source
    # Cross-session leak guard: new task might inherit a sibling session's
    # HERMES_SESSION_* ContextVar; reset to _UNSET before going further.
    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() snapshottet den Spawning-Context; haben konkurrierende Nachrichten im Eltern-Task zuvor set_session_vars() aufgerufen, erbt der neue Task die Sibling-Identität, und die Subprocess-Env-Brücke liest die falsche Session.

_handle_message_with_agent sucht/legt die Session an und schiebt die Nachricht in den Agenten; es verhindert das «unautorisierte Wiederbeleben einer toten Sitzung»:

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 might have /new-reset,
        # don't blindly switch back and revive a user-closed session.
        ...

Kommt ein asynchroner Subagent zurück, während die Spawning-Session bereits beendet ist, verwirft das Gateway diese Injektion lieber, als die Sitzung wiederzubeleben.

Grenzen und Fehler

  • Cross-Session ContextVar-Leck: _handle_message läuft in einem create_task() und erbt den Snapshot des Eltern-Tasks. Ohne vorgeschaltetes reset_session_vars() lesen konkurrierende Nachrichten die Sibling-Session-Identität, und die Subprocess-Brücke schreibt die falsche Identität in die Umgebung.
  • Async-Rückfüllung in eine tote Session: Kommt der Subagent zurück, ist die Spawning-Session vielleicht schon per /new-reset beendet; blindes Umschalten würde eine vom User bewusst geschlossene Sitzung wiederbeleben. Fail-closed behält das Ergebnis in den Delegation-Records.
  • Scale-to-Zero-Uhr-Verschmutzung: Interne Events (Hintergrundprozess beendet, Startup-Restore-Replay) dürfen nicht als Traffic zählen, sonst bleibt ein idle Gateway ewig am Leben. Der Code zweigt sie über is_internal explizit ab.

Zusammenfassung

Der Gateway-Kern ist «Event-Bus + Stream-Brücke»: Die einheitliche MessageEvent-Abstraktion schirmt Plattform-Unterschiede ab, der GatewayStreamConsumer überbrückt den Graben zwischen dem synchronen Agent-Callback und der asynchronen Plattform-Zustellung, das delivery_ledger sichert den Endzustand. Der Agent ist entkoppelt – er dreht nur die Schleife, die Zustellungs-Details liegen voll im Gateway.

Inoffizielle Community-Lernseite. Basiert auf dem MIT-lizenzierten NousResearch/hermes-agent-Quellcode.