Streaming et Hooks
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
class GatewayStreamConsumer:83-120— consommateur asynchrone,on_deltaest branché comme callback par l'agentclass StreamConsumerConfig:55-83— config d'exécution pour une consommation (stratégie streaming du draft natif, délai d'édition finale, etc.)doc du module:1-40— explique le pont synchrone → asynchrone et le signal sentinel de finstream_dispatch— dispatch des événements de streamstream_events— définition des types d'événements de streamhooks— framework d'enregistrement des hooks de cycle de vierépertoire builtin_hooks/— implémentations de hooks intégréesrépertoire relay/— relais entre instancesslash_commands— routage des commandes/authz_mixin— autorisation (mixé dansGatewayRunner)profile_routing— routage multi-profils
Flux de données
GatewayRunnerdémarre_start_stream_consumer:21218.- Construction d'un
GatewayStreamConsumer;consumer.on_deltaest passé commestream_delta_callbackàAIAgent(run_agent.py:400). - L'agent, dans son thread worker, appelle synchrone
on_delta(text)→ le consommateur pousse l'incrément dans une file asynchrone. - 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
- À la fin du flux (stream), un signal sentinel déclenche l'édition finale, confiée au
delivery_ledger:155pour la comptabilisation. - Tout au long du flux, les
hookssont 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 :
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 :
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
breakQuelques 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 » :
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
/newou/stopsur la session (session), la coroutine de consommation vérifie_run_still_current()à la prochaine itération et retourne. Les incréments déjàqueue.putsont 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_deltasans appelerfinish(), la coroutine de consommation reste bloquée surqueue.get().run_still_currentest 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.