Skip to content

Streaming und Hooks

源码版本v2026.7.20

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

Datenfluss

  1. GatewayRunner startet _start_stream_consumer:21218.
  2. Ein GatewayStreamConsumer wird erzeugt; consumer.on_delta wird als stream_delta_callback an AIAgent (run_agent.py:400) übergeben.
  3. Der Agent ruft im Worker-Thread synchron on_delta(text) auf → der Konsument stellt das Inkrement in eine asynchrone Queue.
  4. 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
  5. Nach Stream-Abschluss löst das Sentinel-Signal das finale Editieren aus und übergibt an das delivery_ledger:155 zur Verbuchung.
  6. An allen Knotenpunkten werden hooks getriggert (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:

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

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

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

python
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 /new oder /stop beendet, prüft die Konsum-Koroutine beim nächsten Durchlauf _run_still_current() und kehrt direkt zurück. Bereits per queue.put abgegebene Inkremente gehen verloren; deshalb spült _flush_think_buffer() vor dem finalen Editieren halbe think-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_delta ab, ohne finish() aufzurufen, hängt die Konsum-Koroutine im queue.get(). run_still_current ist 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.

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