PregelRunner: el ejecutor de tareas dentro de un superpaso
Responsabilidades
PregelRunner hace una sola cosa: tomar el lote de PregelExecutableTask que preparó PregelLoop.tick, ejecutarlo de verdad, recolectar sus writes y devolverlos a PregelLoop. No le importa el aspecto del grafo, ni el checkpoint, ni cómo reducen los canales — solo se ocupa de «ejecución concurrente + manejo de fallos + persistencia de writes».
Su posición es muy estrecha: dentro del bucle principal de Pregel.astream(sync main loop:2964-2984), cuando while loop.tick() devuelve True, justo después viene for _ in runner.tick([t for t in loop.tasks.values() if not t.writes], ...): y, al terminar, loop.after_tick(). Es decir, la entrada del runner es «tareas calculadas por tick que aún no han escrito writes» y la salida es «los writes de esas tareas ya están en checkpoint_pending_writes»; no hay un segundo búfer intermedio.
El propio PregelRunner es muy fino: dos métodos públicos tick(tick:176) y commit(commit:574), más una versión asíncrona atick. El trabajo pesado está en los invocadores de bajo nivel _call / _acall(_call:700)/(_acall:789) — son los responsables de envolver la función del nodo en un future programable, gestionar la política de retry, el fan-out de subtask Send, las llamadas de streaming chunk y, finalmente, volcar los writes al objeto task.
Motivación de diseño
¿Por qué el runner es un generator y no un método normal?
- Yield en streaming: el runner tiene forma de generator
for _ in runner.tick(...)(yield return type:188); cada vez que termina una tarea hace yield y devuelve el control al llamador — que así puede emitir de inmediato el streaming chunk al cliente, sin esperar a que termine el superpaso entero. Esa es la base del streaming de LangGraph. - Camino rápido de tarea única: con
len(tasks) == 1 and timeout is None and get_waiter is Nonese hace una llamada síncrona directa, sin levantar thread pool(fast path:203-254). La mayoría de los escenarios agent de rama única caen aquí y se ahorran el overhead del thread pool. commitcomo callback de Future:FuturesDict.on_done(on_done:116) se dispara al completar cada future e internamente esweakref.WeakMethod(self.commit)—commites el gancho de persistencia de writes por tarea, no el cierre batch. Así el orden con que los writes entran enpending_writeses naturalmente el de finalización, y el checkpoint puede escribirse cuanto antes.- Encaminamiento al error handler:
_should_route_to_error_handler(_should_route_to_error_handler:171-174) comprueba si el nodo fallido tiene configurado un nodoerror_handler. Si acierta, añade el id de la excepción a_handled_exception_idsy programa una tarea handler en sustitución de la tarea original — así el fallo no hace panic inmediato, sino que se cierra por otro nodo. - Cancelación con
_should_stop_others: cuando cualquier tarea lanza una excepción que no esGraphBubbleUp, el runner cancela las demás tareas in-flight del mismo lote(_should_stop_others:616-634) — garantiza la semántica dentro del superpaso de «o todas, o al handler».GraphInterruptse excluye explícitamente, porque una interrupción no cuenta como fallo.
Archivos clave
class PregelRunner:135-138— definición de la clase; el docstring resume su responsabilidad en una frase: ejecutar tareas, commitar writes, ceder control y, si hace falta, interrumpir otras tareas.__init__:140-169— mantiene weakrefs desubmit/put_writes,node_error_handler_map, callbacksschedule_error_handler/aschedule_error_handlery_handled_exception_ids, el conjunto de dedup de excepciones que sobrevive entre ticks.FuturesDict:75-134— subclase de dict a medida; el callbackon_donedisparacommital completar cada future y el camposhould_stopretiene un partial de_should_stop_others.tick signature:176-188— la entrada esIterable[PregelExecutableTask]y la salidaIterator[None]; la signatura deja explícita la semántica de generator.fast path:203-254— con tarea única, sin timeout y sin waiter, se ejecuta síncrono directo; en fallo se programa el error handler si procede.schedule tasks:259-276— ruta multi-tarea: se lanza un future por tarea víaself.submit()(en realidadPregelLoop.submit).concurrent wait loop:282-323— bucle conconcurrent.futures.wait(FIRST_COMPLETED); cada vez que termina una tarea se hacecommit, se emite la salida y, si hace falta, se levanta una tarea handler.commit method:574-613— cierre por tarea: cancelled escribe ERROR,GraphInterruptescribe INTERRUPT, excepción común escribeERROR+ERROR_SOURCE_NODE, y una finalización normal escribetask.writes+ el marcadorNO_WRITES._should_stop_others:616-634— decide si cancelar otras tareas; excluye explícitamenteGraphBubbleUpy las excepciones ya handled._panic_or_proceed:650-697— función de cierre: cancela todos los futures in-flight, fusiona variosGraphInterrupten uno y, en timeout, lanzaTimeoutError._call:700-787— envoltorio de la llamada al nodo: retry, emisión de stream chunk, callbackschedule_taskpara el fan-out deSend, acumulación detask.writes._should_route_to_error_handler:171-174— devuelve True cuandotask.name in self.node_error_handler_map; el propio nodo handler no puede ser encaminado (evita recursión).
Flujo de datos
Al inicio de tick se construye FuturesDict, un dict mejorado que llama automáticamente a commit al completar cada future:
def tick(
self,
tasks: Iterable[PregelExecutableTask],
*,
reraise: bool = True,
timeout: float | None = None,
retry_policy: Sequence[RetryPolicy] | None = None,
get_waiter: Callable[[], concurrent.futures.Future[None]] | None = None,
schedule_task: Callable[
[PregelExecutableTask, int, Call | None],
PregelExecutableTask | None,
],
) -> Iterator[None]:
tasks = tuple(tasks)
futures = FuturesDict(
callback=weakref.WeakMethod(self.commit),
event=threading.Event(),
should_stop=partial(
_should_stop_others, handled_exception_ids=self._handled_exception_ids
),
future_type=concurrent.futures.Future,
)
# give control back to the caller
yieldEsto viene de tick 头部:176-199. El primer yield devuelve el control de inmediato — la primera iteración de for _ in runner.tick(...) del llamador solo «arranca»; la ejecución real de tareas ocurre en iteraciones posteriores. Es una forma simplificada de corrutina con generator, para no introducir la palabra clave async en la ruta síncrona.
commit es la función de cierre por tarea, que vuelca los writes del task a pending_writes:
def commit(
self,
task: PregelExecutableTask,
exception: BaseException | None,
) -> None:
if isinstance(exception, asyncio.CancelledError):
task.writes.append((ERROR, exception))
self.put_writes()(task.id, task.writes)
elif exception:
if isinstance(exception, GraphInterrupt):
if exception.args[0]:
writes = [(INTERRUPT, exception.args[0])]
if resumes := [w for w in task.writes if w[0] == RESUME]:
writes.extend(resumes)
self.put_writes()(task.id, writes)
elif isinstance(exception, GraphBubbleUp):
pass
else:
task.writes.append((ERROR, exception))
if self._should_route_to_error_handler(task) and not isinstance(
exception, GraphBubbleUp
):
task.writes.append((ERROR_SOURCE_NODE, task.name))
self._handled_exception_ids.add(id(exception))
self.put_writes()(task.id, task.writes)
else:
if self.node_finished and (
task.config is None or TAG_HIDDEN not in task.config.get("tags", [])
):
self.node_finished(task.name)
if not task.writes:
task.writes.append((NO_WRITES, None))
self.put_writes()(task.id, task.writes)Esto viene de commit:574-613. put_writes es una weakref a PregelLoop.put_writes(put_writes:415) que mete los writes en checkpoint_pending_writes y lanza el future del checkpointer para persistir. Hay que observar varios writes especiales:
CancelledError→ escribe(ERROR, exception), de modo que al cerrarafter_tickla tarea cuenta como «fallida pero registrada», sin panic.GraphInterrupt→ escribe(INTERRUPT, ...)+ cualquier writeRESUME; el valor de reanudación de la interrupción debe persistirse junto con ella y leerse junto al volver a resume.- Excepción común + error handler configurado → añade un write extra
(ERROR_SOURCE_NODE, task.name); esta marca la barreráPregelLoop._resume_error_handlers_if_applicable(ERROR_SOURCE_NODE scan:775) para programar el nodo handler correspondiente al reanudar. - Finalización normal sin writes → se añade
(NO_WRITES, None)para queprepare_next_tasksno la tome como «aún no corrida»(NO_WRITES marker:609-611).
El flujo de ejecución completo del runner se puede dibujar así:
Límites y fallos
_handled_exception_idspersiste entre ticks: en__init__este conjunto es de instancia, no por tick(_handled_exception_ids:169). Una vez que el id de un objeto excepción entra al conjunto,_should_stop_othersno lo vuelve a tratar como fallo, evitando que tras correr el error handler se reprotege la excepción original.GraphBubbleUpno cuenta como fallo:commithacepassexplícito sobreGraphBubbleUp, sin escribir nada(GraphBubbleUp pass:592-594). Estas excepciones (comoParentCommand) son señales al grafo padre, no errores reales.NO_WRITESevita que la tarea se repita: si una tarea termina con normalidad pero no escribió nada, también se añade(NO_WRITES, None)(NO_WRITES:609-611). Si no, trasafter_tick,prepare_next_tasksvería la tarea sin writes, la creería no corrida y la volvería a ejecutar al reanudar.schedule_error_handlerpuede devolver None: el propio nodo handler puede no estar configurado o no poder construirse, en cuyo caso la tarea sigue la vía de pánico común(handler optional:230-248). El writeERROR_SOURCE_NODEdecommitcoexiste con esta ruta;_resume_error_handlers_if_applicablelo intentará de nuevo al reanudar.- El timeout no lanza GraphInterrupt:
_panic_or_proceedhaceif inflight: raise timeout_exc_cls("Timed out")(timeout:691-697) — el timeout es una excepción real que burbujea fuera de Pregel, no se traga comoGraphInterrupt. - Recorte de traceback: antes del reraise se saltan los frames listados en
EXCLUDED_FRAME_FNAMES(tb trim:241-247), filtrando los stack frames internos de langgraph para que el traceback que ve el usuario apunte directamente a su código de nodo.
Resumen
PregelRunner es una cáscara fina: el generator tick programa la ejecución concurrente y el commit por tarea vuelca los writes a pending_writes, con encaminamiento al error handler y cancelación entre medias. Entenderlo es entender «cómo se ejecutan de verdad las tareas dentro de un superpaso». La capa de bucle PregelLoop se ve en /pregel/loop; las tres funciones de algoritmo prepare_next_tasks / apply_writes / should_interrupt en /pregel/algo; el ensamblaje integral en /pregel/pregel; cómo se procesan las salidas en streaming entre los yield del runner se ve en /stream/run-stream.
Véase la documentación oficial: documentación de LangGraph · README