Skip to content

Streaming et Hooks

源码版本v2026.7.20

Responsabilité

(1) Consommation en streaming — le thread worker de l'agent appelle synchrone stream_delta_callback(text), alors que l'envoi côté plateforme est asynchrone ; GatewayStreamConsumer fait le pont entre l'incrément synchrone et l'asynchrone, en éditant progressivement le même message de la plateforme pour donner à l'utilisateur un effet « machine à écrire ». (2) Hooks — insère une logique personnalisée aux nœuds du cycle de vie du message (réception, avant génération, après génération, livraison), comme filtrage de sécurité, facturation, synchro Kanban. (3) Relay — relais de messages entre instances / profils.

Mot de conception

Le run_conversation de l'agent tourne dans un thread worker ; les tokens sont yield de façon synchrone — c'est imposé par le SDK sous-jacent (l'interface streaming du LLM est typiquement un blocking iterator). Or l'envoi côté plateforme est une coroutine asyncio, l'appeler cross-thread depuis le worker piétinerait la event loop. Sans couche intermédiaire, soit on accumule toute la réponse avant d'envoyer (l'utilisateur regarde « en train d'écrire… » pendant des dizaines de secondes sans rien voir), soit on appelle asyncio.run_coroutine_threadsafe à chaque token depuis le worker (coût élevé, gestion d'erreur difficile). GatewayStreamConsumer utilise une queue.Queue comme frontière de thread — le worker ne fait que queue.put, la coroutine de consommation lit en get_nowait côté asyncio ; les deux horloges sont découplées, et throttle / coalescing / repli (fallback) se greffent côté consommation. Les hooks sont isolés parce que les préoccupations transverses (sécurité / facturation) ne doivent pas polluer le pipeline principal de message.

Fichiers clés

Flux de données

  1. GatewayRunner démarre _start_stream_consumer:21218.
  2. Construction d'un GatewayStreamConsumer ; consumer.on_delta est passé comme stream_delta_callback à AIAgent (run_agent.py:400).
  3. L'agent, dans son thread worker, appelle synchrone on_delta(text) → le consommateur pousse l'incrément dans une file asynchrone.
  4. La coroutine du consommateur retire les incréments de la file, selon la stratégie de StreamConsumerConfig:55 :
    • auto / draft : privilégie le streaming natif en draft (send_draft:2650), repli sur édition simple
    • l'édition finale peut être retardée (modèles à inférence lente), pour éviter les tremblements haute fréquence
  5. À la fin du flux (stream), un signal sentinel déclenche l'édition finale, confiée au delivery_ledger:155 pour la comptabilisation.
  6. Tout au long du flux, les hooks sont déclenchés à chaque nœud (réception / avant génération / après génération / livraison) ; les hooks intégrés (gateway/builtin_hooks/) exécutent sécurité, facturation, Kanban, etc.

on_delta est l'entrée exposée par le consommateur ; elle ne fait qu'une chose — pousser le token du thread worker dans une queue.Queue :

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 est le signal de tool boundary — les tokens suivants vont dans un nouveau message, inséré sous les bulles de progression outil (tool). _DONE est le sentinel qui déclenche l'édition finale côté coroutine de consommation. « Valeur = token, None = frontière, sentinel = fin » : une seule file porte tout le contrôle de flux.

La coroutine run() draine la file par lots et pousse selon la stratégie :

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

Quelques détails : message_len_fn vient de l'adaptateur (Telegram utilise la longueur utf16), pour que la détection de débordement colle à la vraie limite de la plateforme ; _safe_limit garde 100 caractères de marge pour le cursor et le formatage riche ; _run_still_current() est vérifié à chaque itération — si la session est /new ou /stop, on sort tout de suite pour éviter de pousser des tokens obsolètes à l'écran.

GatewayEventDispatcher (gateway/stream_dispatch.py) est un wrapper plus structuré — l'agent produit des événements typés (MessageChunk, ToolCallChunk, Commentary…), le dispatcher décide comment les rendre en fonction des capacités de l'adaptateur (adapter) ; si la plateforme ne sait pas rendre le chrome d'outils, il renvoie None et avale l'événement. C'est le point d'extension pour « plateformes qui veulent une UI native (DingTalk AI Card, Slack Block Kit) plutôt que de l'édition de texte brut » :

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.

Limites et échecs

  • Streaming interrompu : après /new ou /stop sur la session (session), la coroutine de consommation vérifie _run_still_current() à la prochaine itération et retourne. Les incréments déjà queue.put sont perdus, donc avant l'édition finale _flush_think_buffer() chasse les fragments résiduels de balises think.
  • Flood control transpercé : _MAX_FLOOD_STRIKES = 3 ; trois 429 consécutifs dégradent en permanence vers non-progressif, et les tokens restants du flux sont accumulés et envoyés en un seul lot.
  • Plateforme sans edit : Signal avec SUPPORTED_MESSAGE_EDITING = False, le consommateur suit le repli non-edit — un seul message complet est envoyé, l'utilisateur ne voit pas l'effet machine à écrire pendant le streaming. C'est le filet du bit de capacité, pas un bug.
  • Sentinel oublié : si le thread worker plante après on_delta sans appeler finish(), la coroutine de consommation reste bloquée sur queue.get(). run_still_current est le filet — la suppression de la session déclenche aussi la sortie.

Résumé

Le consommateur de streaming est la couche de buffer et de throttle entre « callback synchrone de l'agent ↔ livraison asynchrone à la plateforme », qui résout l'inadéquation des deux horloges. Les hooks sont le point d'insertion des préoccupations transverses (sécurité / facturation / synchro). Les deux permettent à la passerelle (gateway), sans polluer le tronc principal de l'agent, de supporter le typing en temps réel, le filtrage de sécurité et le relay inter-instances.

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