Streaming und Hooks
Verantwortung
(1) Stream-Konsum (stream consumption) – der Worker-Thread des Agenten ruft synchron stream_delta_callback(text) auf, während das Plattform-Senden asynchron ist; GatewayStreamConsumer übersetzt die synchronen Inkremente in die asynchrone Welt und editiert inkrementell dieselbe Plattform-Nachricht, sodass der Nutzer eine «Schreibmaschine»-Wirkung sieht. (2) Hooks – an Lebenszyklus-Punkten (Empfang, Pre-Generierung, Post-Generierung, Zustellung) wird eigene Logik eingefügt, etwa Sicherheitsfilter, Abrechnung, Kanban-Sync. (3) Relay – nachrichtenweiterleitung über Instanzen/Profiles hinweg.
Designmotiv
run_conversation des Agenten läuft in einem Worker-Thread; Token werden synchron yield-weise herausgereicht — das gibt der untere SDK vor (LLM-Streaming-Schnittstellen sind meist blockierende Iteratoren). Aber das Senden an die Plattform ist eine asyncio-Koroutine; ein direkter threadübergreifender Aufruf tritt in die Event-Loop. Ohne Mittelschicht bleibt nur: entweder die gesamte Antwort aufsammeln und dann senden (der Nutzer sieht «tippt …» für zig Sekunden ohne Inhalt) oder im Worker-Thread pro Token asyncio.run_coroutine_threadsafe mit neuem Future (teuer, Fehlerbehandlung schwer). GatewayStreamConsumer zieht eine queue.Queue als Thread-Grenze — der Worker-Thread macht nur queue.put; die Konsum-Koroutine holt sie auf asyncio-Seite per get_nowait; die Taktbereiche sind entkoppelt, Drosselung, Merging und Backoff sitzen auf der Konsum-Seite. Hooks sind separat geführt, weil Querschnitts-Themen wie Sicherheit/Abrechnung die Hauptnachrichten-Pipe nicht verschmutzen dürfen.
Schlüsseldateien
class GatewayStreamConsumer:83-120— der asynchrone Konsument;on_deltaist der an den Agenten übergebene Callbackclass StreamConsumerConfig:55-83— Laufzeitkonfiguration einer Konsumtion (Native-Draft-Streaming-Strategie, End-Edit-Verzögerung usw.)Modul-Dokumentation:1-40— erklärt die Sync→Async-Brücke und das Sentinel-Abschlusssignalstream_dispatch— Dispatch (dispatch) von Stream-Events (stream events)stream_events— Definition der Stream-Event-Typenhooks— Lebenszyklus-Hook-Registrierungbuiltin_hooks/-Verzeichnis— eingebaute Hook-Implementierungenrelay/-Verzeichnis— instanzübergreifendes Relayslash_commands—/-Command-Routingauthz_mixin— Autorisierung (inGatewayRunnereingemischt)profile_routing— Multi-Profile-Routing
Datenfluss
GatewayRunnerstartet_start_stream_consumer:21218.- Ein
GatewayStreamConsumerwird erzeugt;consumer.on_deltawird alsstream_delta_callbackanAIAgent(run_agent.py:400) übergeben. - Der Agent ruft im Worker-Thread synchron
on_delta(text)auf → der Konsument stellt das Inkrement in eine asynchrone Queue. - Die Konsumenten-Koroutine holt Inkremente aus der Queue und folgt der
StreamConsumerConfig:55-Strategie:auto/draft: bevorzugt natives Draft-Streaming (send_draft:2650), Fallback (fallback) auf normales Editieren- das finale Editieren kann verzögert sein (bei langsamen Inferenzmodellen), um hochfrequentes Flackern zu vermeiden
- Nach Stream-Abschluss löst das Sentinel-Signal das finale Editieren aus und übergibt an das
delivery_ledger:155zur Verbuchung. - An allen Knotenpunkten werden
hooksgetriggert (Empfang / Pre-Gen / Post-Gen / Zustellung); eingebaute Hooks (gateway/builtin_hooks/) übernehmen Sicherheit, Abrechnung, Kanban usw.
on_delta ist der Außeneingang des Konsumenten; es macht nur eins — Token aus dem Worker-Thread in eine queue.Queue stellen:
def on_delta(self, text: str) -> None:
"""Thread-safe callback — called from the agent's worker thread.
When *text* is ``None``, signals a tool boundary: the current message
is finalized and subsequent text will be sent as a new message so it
appears below any tool-progress messages the gateway sent in between.
"""
if text:
self._queue.put(text)
elif text is None:
self.on_segment_break()
def finish(self) -> None:
"""Signal that the stream is complete."""
self._queue.put(_DONE)None ist das Tool-Boundary-Signal — nachfolgende Token gehen als neue Nachricht und erscheinen unter den Tool-Fortschritts-Blasen. _DONE ist ein Sentinel; die Konsum-Koroutine löst bei Sicht das finale Editieren aus. «Wert = Token, None = Boundary, Sentinel = fertig» lässt eine einzige Queue den gesamten Kontrollfluss tragen.
Die Konsum-Koroutine run holt Inkremente batchweise und liefert sie nach Strategie aus:
async def run(self) -> None:
"""Async task that drains the queue and edits the platform message."""
_len_fn = (
self.adapter.message_len_fn
if isinstance(self.adapter, _BasePlatformAdapter)
else len
)
_raw_limit = self._raw_message_limit()
_safe_limit = max(500, _raw_limit - _len_fn(self.cfg.cursor) - 100)
self._use_draft_streaming = self._resolve_draft_streaming()
try:
while True:
if not self._run_still_current():
return
got_done = False
while True:
try:
item = self._queue.get_nowait()
if item is _DONE:
got_done = True
breakEin paar Details: message_len_fn kommt vom Adapter (adapter) (Telegram nutzt UTF-16-Länge), damit die Überlaufprüfung mit dem realen Plattform-Limit übereinstimmt; _safe_limit lässt 100 Zeichen Reserve für Cursor und Rich-Text-Formatierung; _run_still_current() prüft jeden Schleifendurchlauf — wurde die Session (session) per /new oder /stop beendet, wird direkt abgebrochen, damit keine veralteten Token weiter auf den Bildschirm geschoben werden.
GatewayEventDispatcher (gateway/stream_dispatch.py) ist eine strukturiertere Hülle — der Agent erzeugt typisierte Events (MessageChunk, ToolCallChunk, Commentary …); der Dispatcher entscheidet nach Adapter-Fähigkeit, wie gerendert wird, und wenn eine Plattform Tool-Chrome nicht rendern kann, gibt er None zurück und schluckt das Event. Das ist der Erweiterungspunkt für «Plattformen, die native UI wollen (DingTalk AI Card, Slack Block Kit) statt reinem Text-Editieren»:
class GatewayEventDispatcher:
"""Route typed stream events through an adapter onto a delivery sink."""
# Message/commentary/segment events flow into the consumer (native draft
# on Telegram DMs, edit-in-place elsewhere). Tool events are formatted
# by the adapter — which may return None to *eat* the event on platforms
# that can't render tool chrome — and the rendered line is enqueued onto
# the same tool progress queue the gateway already drains, so the two no
# longer race through independent code paths.Grenzen und Fehler
- Streaming unterbrochen: Wird die Session per
/newoder/stopbeendet, prüft die Konsum-Koroutine beim nächsten Durchlauf_run_still_current()und kehrt direkt zurück. Bereits perqueue.putabgegebene Inkremente gehen verloren; deshalb spült_flush_think_buffer()vor dem finalen Editieren halbethink-Tags weg. - Flood-Control durchbrochen:
_MAX_FLOOD_STRIKES = 3; drei 429er in Folge degradieren dauerhaft auf nicht-inkrementell; die restlichen Token dieses Streams werden als Batch auf einmal gesendet. - Plattform ohne Edit: Signal hat
SUPPORTED_MESSAGE_EDITING = False; der Konsument geht den Nicht-Edit-Fallback — einmalige komplette Nachricht, während des Streamings gibt es keine Schreibmaschine. Das ist ein Fähigkeitsbit-Auffangwehr, kein Bug. - Sentinel vergessen: Bricht der Worker-Thread nach
on_deltaab, ohnefinish()aufzurufen, hängt die Konsum-Koroutine imqueue.get().run_still_currentist die Auffangwehr — wird die Session weggeräumt, löst das ebenfalls den Ausstieg aus.
Zusammenfassung
Der Stream-Konsument ist die Puffer- und Drosselungsschicht zwischen «synchronem Agent-Callback ↔ asynchroner Plattform-Zustellung» und löst das Sync-/Async-Mismatch. Hooks sind Einhängpunkte für Querschnitts-Themen (Sicherheit/Abrechnung/Sync). Beide lassen das Gateway Echtzeit-Tippen, Sicherheitsfilter und instanzübergreifendes Relay unterstützen, ohne den Agent-Stamm zu verschmutzen.