Gateway-Kern
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
class GatewayRunner:3029— Gateway-Hauptklasse, mischt Authorization/Kanban/Slash ein_handle_message:9947— Dispatch-Eingang (dispatch entry) eingehender Nachrichten_handle_message_with_agent:11956— übergibt die Nachricht an den Agenten für einen Turn (turn)_start_one_profile_adapter / _start_secondary:9389-9947— startet die Plattform-Adapter der Profile_start_stream_consumer:21218— startet den Stream-Konsumenten (stream consumer)_start_gateway_housekeeping:22246— zyklische Inspektion (60s)_start_cron_ticker:22335— Cron-Heartbeat (60s)def main:22965— Gateway-ProzesseinstiegExceptions und Redaktion:276-340—_gateway_loop_exception_handler/_redact_gateway_user_facing_secretsPlatformEntry:39-162— Datenklasse eines Plattform-Registry-Eintragsclass PlatformRegistry:162-260— Registry (register / register_deferred / unregister)
Datenfluss
main()(gateway/run.py:22965) startet denGatewayRunner._start_one_profile_adapter:9389bringt nacheinander die Plattform-Adapter hoch (siehe Plattform-Adapter).- Der Adapter empfängt das native Plattform-Event und wandelt es in ein einheitliches
MessageEventum (siehegateway/platforms/base.py:1759) um. _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).- Letzteres ruft
run_conversation:588des Agenten auf und verdrahtet den Streaming-Callback anGatewayStreamConsumer:83. - Der Streaming-Output wird vom Konsumenten inkrementell editiert auf der Plattform-Nachricht; der Endzustand wird vom
delivery_ledgerverbucht, damit ein Crash die finale Antwort nicht verliert. - 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:
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:
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 Nonecreate_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»:
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_messageläuft in einemcreate_task()und erbt den Snapshot des Eltern-Tasks. Ohne vorgeschaltetesreset_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-resetbeendet; 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_internalexplizit 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.