StreamMux: dispatcher de eventos para streaming multi-modo
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
GraphRunStreames el pump(_pump_next:107-130). Cuando el usuario hacefor 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_pumpregistra 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:
pushusafor 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 devuelveFalse, sólo suprime el evento del registro principal — los transformers anteriores ya lo han procesado. - Distinción native vs extension: los transformers con
_native = Truetienen su proyección montada como atributo directo deGraphRunStream, comorun.values(native attrs:78-79); los no nativos sólo entran enrun.extensions. El conjuntonative_keys(native_keys:109) es la whitelist de «atributos directos» que el protocolo v3 expone al SDK. before_builtinscede 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 comoMessagesTransformer(before_builtins:94-109); el mux los reorganiza en dos lanes segúnbefore_builtins = Truedurante el registro(partition:131-148), pero el orden dentro de cada lane se conserva.seqse asigna sólo en el mux raíz:_assign_seqes True por defecto; el child mux recibeFalsevía_make_child(_assign_seq:215-220). Así, los eventos del subgrafo reenviados directamente al registro raíz no ven alterado su envelope, manteniendoseqmonó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únbefore_builtins._make_child:193-225— construye un mini-mux al descubrir un subgrafo, hereda el binding del pump._register:227-267— llama atransformer.init()para obtener el dict de projection, comprueba conflictos de key y enlaza los StreamChannel con nombre.push:269-296— ruta síncrona: llama aprocessen orden, y decide si entra en el registro principal.apush:351-378— ruta asíncrona:aprocessen serie con await, garantizando que los transformers siguientes ven los resultados async anteriores.close:298-324— cierre síncrono; llama afinalizede cada transformer, garantizando el cierre de todos los channel incluso si hay errores.aclose:380-404— cierre asíncrono; primerogatherde todos los tasks programados, luegoafinalize._collect_stream_modes:398-414— deduce qué modos pedir al grafo a partir delrequired_stream_modesde 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 alstream()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:
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:
_registerinvocatransformer_requires_asyncpara comprobar(requires_async check:234-239); si el transformer reescribeaprocesspero no fijasupports_sync=True, registrarlo en un mux síncrono lanza un error sugiriendo usarastream()en lugar destream(). - Conflicto de key de projection lanza
ValueError:_registercompruebaset(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_childrechaza transformers pre-built: cuando el mux se construye contransformers=, no se pueden crear sub-mux(factory required:209-214), porque las instancias pre-built no pueden reinicializarse en un nuevo scope.SubgraphTransformersólo funciona en mux construidos confactories=.- Errores de finalize no bloquean la limpieza:
closeusafirst_errorpara 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_pumpdebe llamarse antes de leer el channel:GraphRunStream.__init__conwire_pump=Truepor 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_moredeStreamChannel.__iter__sea siempreNone(_request_more:191-192), y la iteración terminará de inmediato.schedule(coro)sólo disponible en modo async:StreamTransformer.schedulerequiere un event loop(schedule:233-263); el mux hacegatherde todos los tasks programados enaclosey sólo tras terminar el asyncafinalize; 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