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.
Motivo de diseño
run_conversation del agent corre en un hilo worker; los tokens salen por yield síncrono — lo dicta el SDK subyacente (la interfaz de streaming de un LLM suele ser un blocking iterator). El envío a la plataforma es una corutina asyncio; llamarlo directamente desde otro hilo pisa el event loop. Sin una capa intermedia, o bien acumulas la respuesta entera antes de enviar (el usuario ve «escribiendo…» durante decenas de segundos sin contenido) o bien lanzas un asyncio.run_coroutine_threadsafe por cada token en el hilo worker (caro, manejo de errores enrevesado). GatewayStreamConsumer usa queue.Queue como frontera entre hilos — el hilo worker solo queue.put; la corutina consumidora hace get_nowait desde el lado asyncio. Los dos relojes quedan desacoplados; throttle, merge y fallback se añaden en el lado del consumidor. Los hooks se separan porque las preocupaciones transversales (seguridad, facturación) no deben contaminar el pipeline principal de mensajes.
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 (dispatch) 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.
on_delta es la entrada pública del consumidor y solo hace una cosa — meter en queue.Queue los tokens del hilo worker:
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 es la señal de tool boundary — los tokens siguientes van a un mensaje nuevo, debajo de las burbujas de progreso de tool. _DONE es el centinela que la corutina consumidora interpreta para disparar la edición final. «Valor = token, None = boundary, centinela = fin» permite que una sola cola cargue con todo el flujo de control.
La corutina consumidora run() extrae incrementos en lote y los despacha según la política:
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
breakUnos detalles: message_len_fn viene del adaptador (adapter) (Telegram usa longitud utf16) para que la detección de overflow coincida con el límite real de la plataforma; _safe_limit reserva 100 caracteres para cursor y formato de texto enriquecido; _run_still_current() se comprueba en cada vuelta — si la sesión (session) se fue con /new o /stop, se sale antes para no seguir empujando tokens stale a la pantalla.
GatewayEventDispatcher (gateway/stream_dispatch.py) es un envoltorio más estructurado — el agent emite typed events (MessageChunk, ToolCallChunk, Commentary…); el dispatcher decide cómo renderizarlos según las capacidades del adaptador. Si la plataforma no puede renderizar tool chrome, devuelve None y se come el evento. Es el punto de extensión para «plataformas con UI nativa (DingTalk AI Card, Slack Block Kit) en lugar de edición de texto plano»:
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.Límites y fallos
- Stream interrumpido: si la sesión se va con
/newo/stop, la corutina consumidora comprueba_run_still_current()en la próxima vuelta y sale. Los incrementos ya hechosqueue.putse pierden, así que antes de la edición final_flush_think_buffer()tira del buffer para vaciar medias etiquetasthinkque hayan quedado. - Flood control roto:
_MAX_FLOOD_STRIKES = 3; tres 429 seguidos degradan permanentemente a modo no progresivo y el resto de tokens de la ronda se envían en un solo lote. - Plataforma sin edición: Signal con
SUPPORTED_MESSAGE_EDITING = Falsecae al fallback no-edit — se envía un único mensaje completo y durante el stream no hay efecto máquina de escribir. Es el respaldo por capability bit, no un bug. - Centinela perdido: si el hilo worker muere tras
on_deltasin llamarfinish(), la corutina consumidora se queda colgada enqueue.get().run_still_currentes la red — cuando la sesión se limpia, también se fuerza la salida.
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.