Skip to content

StreamMux: dispatcher de eventos para streaming multi-modo

源码版本1.2.9

Responsabilidades

En el protocolo de streaming v3 de LangGraph, al invocar graph.stream_events(version="v3") el usuario ya no recibe una secuencia de tuplas raw (namespace, mode, payload), sino un objeto GraphRunStream que expone «proyecciones (projection)» como run.values / run.messages / run.updates / run.subgraphs / run.lifecycle(GraphRunStream:31-49). El componente que reparte estos eventos raw a las distintas proyecciones por modo, permitiendo al usuario consumirlas con un simple bucle for, es StreamMux — el centro de distribución de eventos. El código está en libs/langgraph/langgraph/stream/_mux.py (StreamMux:26).

StreamMux no implementa por sí mismo la lógica de ningún modo concreto; sólo hace tres cosas: registrar transformers, hacer push de eventos y cerrar. En cada push invoca StreamTransformer.process(event) en orden de registro(push:269-296); el transformer decide si incrusta el evento en su projection channel; a continuación el propio mux añade el evento al registro principal _events, que actúa como «copia de auditoría completa» del flujo raw. Las proyecciones abiertas por los transformers son instancias de StreamChannel, y los channel con nombre también son reenviados automáticamente por el mux al registro principal, de modo que los efectos del transformer se reflejan tanto en su propio channel como en el registro principal.

StreamMux es también la base del anidamiento de subgrafos. Al descubrir un subgrafo, SubgraphTransformer invoca mux._make_child(scope)(_make_child:193-225) para construir un mini-mux para el subgrafo; este mini-mux comparte el mismo conjunto de factories con el padre (por lo que las instancias de transformer son una por scope), y también comparte el binding del pump del mux padre — cuando el pump raíz avanza, todos los eventos de los subgrafos avanzan un paso.

Motivación de diseño

  • Pump dirigido por el llamador, sin hilos en background: el iterador de GraphRunStream es el pump(_pump_next:107-130). Cuando el usuario hace for event in run.values, cada iteración dispara _pump_next, que extrae un evento raw y lo pasa al mux; no se abre ningún hilo en background, y el consumo de memoria lo determina la tasa de consumo. bind_pump registra esta función pump en el mux(bind_pump:162-180), y luego el mini-mux del subgrafo la hereda automáticamente vía _make_child.
  • Transformers en serie por orden de registro: push usa for transformer in self._transformers: if not transformer.process(event): keep = False(push loop:288-292), de modo que los efectos secundarios del transformer anterior ocurren primero y el siguiente los puede ver; si un transformer devuelve False, sólo suprime el evento del registro principal — los transformers anteriores ya lo han procesado.
  • Distinción native vs extension: los transformers con _native = True tienen su proyección montada como atributo directo de GraphRunStream, como run.values(native attrs:78-79); los no nativos sólo entran en run.extensions. El conjunto native_keys(native_keys:109) es la whitelist de «atributos directos» que el protocolo v3 expone al SDK.
  • before_builtins cede paso a transformers de reescritura de contenido: transformers como los de filtrado PII o revisión de contenido deben ver el campo text original antes que los transformers internos como MessagesTransformer(before_builtins:94-109); el mux los reorganiza en dos lanes según before_builtins = True durante el registro(partition:131-148), pero el orden dentro de cada lane se conserva.
  • seq se asigna sólo en el mux raíz: _assign_seq es True por defecto; el child mux recibe False vía _make_child(_assign_seq:215-220). Así, los eventos del subgrafo reenviados directamente al registro raíz no ven alterado su envelope, manteniendo seq monótono y correlacionado con el orden de escritura raíz.

Archivos clave

  • StreamMux 类:26-46 — atributos y docstring del dispatcher central.
  • StreamMux.__init__:48-148 — recibe lista de transformers o factories, los registra en lanes según before_builtins.
  • _make_child:193-225 — construye un mini-mux al descubrir un subgrafo, hereda el binding del pump.
  • _register:227-267 — llama a transformer.init() para obtener el dict de projection, comprueba conflictos de key y enlaza los StreamChannel con nombre.
  • push:269-296 — ruta síncrona: llama a process en orden, y decide si entra en el registro principal.
  • apush:351-378 — ruta asíncrona: aprocess en serie con await, garantizando que los transformers siguientes ven los resultados async anteriores.
  • close:298-324 — cierre síncrono; llama a finalize de cada transformer, garantizando el cierre de todos los channel incluso si hay errores.
  • aclose:380-404 — cierre asíncrono; primero gather de todos los tasks programados, luego afinalize.
  • _collect_stream_modes:398-414 — deduce qué modos pedir al grafo a partir del required_stream_modes de los transformers registrados en el mux.
  • _pregel_stream_v3:3533-3558 — entrada síncrona v3, instancia StreamMux + GraphRunStream y pasa el stream_mode deducido por el mux al stream() subyacente.

Flujo de datos

El núcleo de StreamMux.push son unas pocas líneas — recorre los transformers en orden de registro, y si el evento sobrevive, lo mete en el registro principal:

python
def push(self, event: ProtocolEvent) -> None:
    keep = True
    for transformer in self._transformers:
        if not transformer.process(event):
            keep = False
    if keep:
        if self._assign_seq:
            self._seq += 1
            event["seq"] = self._seq
        self._events.push(event)

(push body:288-296) Parece simple, pero transformer.process puede hacer push a su propio StreamChannel, y un channel con nombre también es cableado automáticamente por el mux al registro principal(_register:227-267); por tanto un evento raw values puede acabar tanto en run.values por ValuesTransformer como en el registro principal tras ser transformado por un transformer personalizado. Los eventos fluyen en cadena entre transformers, conservando estrictamente el orden.

Límites y fallos

  • Un transformer async en modo sync lanza error al instante: _register invoca transformer_requires_async para comprobar(requires_async check:234-239); si el transformer reescribe aprocess pero no fija supports_sync=True, registrarlo en un mux síncrono lanza un error sugiriendo usar astream() en lugar de stream().
  • Conflicto de key de projection lanza ValueError: _register comprueba set(projection) & set(self.extensions)(conflict check:246-256); el error lista las keys en conflicto y el nombre del transformer propietario. Un transformer personalizado no puede reutilizar nombres de keys internos.
  • _make_child rechaza transformers pre-built: cuando el mux se construye con transformers=, no se pueden crear sub-mux(factory required:209-214), porque las instancias pre-built no pueden reinicializarse en un nuevo scope. SubgraphTransformer sólo funciona en mux construidos con factories=.
  • Errores de finalize no bloquean la limpieza: close usa first_error para conservar la primera excepción y continúa cerrando el resto de transformers y channels, lanzando al final(close error handling:312-324), para evitar fugas de recursos.
  • bind_pump debe llamarse antes de leer el channel: GraphRunStream.__init__ con wire_pump=True por defecto pasa por _wire_request_more(_wire_request_more:80-92); el sub-mux hereda el binding del padre vía _make_child; pero si el usuario construye manualmente una instancia de mux, olvidar bind pump hará que _request_more de StreamChannel.__iter__ sea siempre None(_request_more:191-192), y la iteración terminará de inmediato.
  • schedule(coro) sólo disponible en modo async: StreamTransformer.schedule requiere un event loop(schedule:233-263); el mux hace gather de todos los tasks programados en aclose y sólo tras terminar el async afinalize; invocar schedule en modo sync falla.

Resumen

StreamMux es el núcleo del protocolo de streaming v3: no implementa lógica de ningún modo, sólo programa transformers en orden, construye sub-mux por scope, y recoge tanto los eventos raw como las proyecciones de los transformers en el registro principal. Para ver cómo cada modo convierte los eventos en una forma amigable para el usuario, ver stream/transformers; para ver cómo los eventos raw se introducen desde PregelLoop, ver motor Pregel.

Véase la documentación oficial: Documentación de LangGraph · README