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.
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, 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.
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, 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.