Skip to content

Streaming y hooks

源码版本v2026.7.20

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

Flujo de datos

  1. GatewayRunner arranca _start_stream_consumer:21218.
  2. Construye un GatewayStreamConsumer y pasa consumer.on_delta como stream_delta_callback a AIAgent (run_agent.py:400).
  3. El agent llama síncronamente desde el hilo worker a on_delta(text) → el consumidor encola el incremento en una cola asíncrona.
  4. 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
  5. Al acabar el stream, una señal centinela dispara la edición final, que se entrega al delivery_ledger:155 para su anotación.
  6. Durante todo el proceso, los hooks se 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:

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 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:

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

Unos 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»:

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.

Límites y fallos

  • Stream interrumpido: si la sesión se va con /new o /stop, la corutina consumidora comprueba _run_still_current() en la próxima vuelta y sale. Los incrementos ya hechos queue.put se pierden, así que antes de la edición final _flush_think_buffer() tira del buffer para vaciar medias etiquetas think que 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 = False cae 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_delta sin llamar finish(), la corutina consumidora se queda colgada en queue.get(). run_still_current es 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.

Sitio de aprendizaje comunitario no oficial. Basado en el código fuente de NousResearch/hermes-agent (licencia MIT).