Streaming y hooks
Responsabilidad
① Consumo de streaming — el hilo worker del agent llama síncronamente a stream_delta_callback(text), mientras que el envío a la plataforma es asíncrono; GatewayStreamConsumer puentea los incrementos síncronos al lado asíncrono, editando progresivamente el mismo mensaje de la plataforma para que el usuario vea un efecto de «máquina de escribir». ② Hooks — insertan lógica personalizada en los puntos del ciclo de vida del mensaje (recepción, pre-generación, post-generación, entrega): filtros de seguridad, facturación, sincronización con kanban. ③ Relay — reenvío de mensajes entre instancias/profiles.
Archivos clave
class GatewayStreamConsumer:83-120— consumidor asíncrono;on_deltase pasa como callback al agentclass StreamConsumerConfig:55-83— config en runtime de una sola ronda de consumo (política de streaming nativo draft, retardo de edición final, etc.)doc del módulo:1-40— explica el puente síncrono→asíncrono y la señal centinela de finalizaciónstream_dispatch— despacho de eventos de streamstream_events— definición de tipos de eventos de streamhooks— framework de registro de hooks de ciclo de vidadirectorio builtin_hooks/— implementaciones de hooks integradosdirectorio relay/— reenvío entre instanciasslash_commands— enrutado de comandos/authz_mixin— autorización (mezclado enGatewayRunner)profile_routing— enrutado multi-profile
Flujo de datos
GatewayRunnerarranca_start_stream_consumer:21218.- Construye un
GatewayStreamConsumery pasaconsumer.on_deltacomostream_delta_callbackaAIAgent(run_agent.py:400). - El agent llama síncronamente desde el hilo worker a
on_delta(text)→ el consumidor encola el incremento en una cola asíncrona. - La corutina del consumidor extrae incrementos de la cola y, según la política de
StreamConsumerConfig:55:auto/draft: prefiere streaming nativo draft (send_draft:2650), con fallback a edición normal- la edición final puede retrasarse (modelos de inferencia lenta) para evitar temblores de alta frecuencia
- Al acabar el stream, una señal centinela dispara la edición final, que se entrega al
delivery_ledger:155para su anotación. - Durante todo el proceso, los
hooksse disparan en cada punto (recepción/pre-generación/post-generación/entrega); los hooks integrados (gateway/builtin_hooks/) ejecutan seguridad, facturación, kanban, etc.
Resumen
El consumidor de streaming es la capa de buffer y throttling entre «callback síncrono del agent ↔ envío asíncrono a plataforma», resolviendo el desacople entre dos dominios de reloj. Los hooks son los puntos de inserción para preocupaciones transversales (seguridad/facturación/sincronización). Entre ambos permiten al gateway soportar escritura en tiempo real, filtrado de seguridad y relay entre instancias sin contaminar el tronco del agent.