PregelLoop: el driver del bucle de superpasos
Responsabilidades
PregelLoop es el «metrónomo» del runtime de LangGraph. Pregel mismo solo se ocupa de ensamblar recursos y exponer las entradas invoke / astream; lo que de verdad corre el grafo vuelta a vuelta es una instancia de PregelLoop: mantiene el punto de control (checkpoint) actual, el conjunto de canales (channel), pending_writes, el contador step, la máquina de estados status y un diccionario tasks — todo el contexto de un único superpaso (superstep). El patrón exterior while loop.tick(): ... loop.after_tick() se sostiene sobre él.
Sus dos métodos centrales son tick(tick:599) y after_tick(after_tick:683). El primero cubre «antes de correr»: comprueba si se excede el límite de pasos, llama a prepare_next_tasks para cargar en self.tasks los nodos que ejecutará este paso, verifica interrupciones, atiende drain_requested y reintegra los pending_writes dejados por una reanudación previa sobre los nodos exitosos. El segundo cubre «después de correr»: recoge todos los task.writes, llama a apply_writes para fusionarlos en los canales, limpia los pending, persiste el checkpoint y luego decide si aplica interrupt_after. Ambos, uno antes y otro después, dejan al ciclo BSP «leer entrada → programar → ejecutar → volcar writes → persistir» encajado entre medias.
Dicho de otro modo, PregelLoop no ejecuta nodos por sí mismo — ese es el trabajo de PregelRunner; solo decide «a quién miro en este paso, cómo cierro tras correr y si sigo con el paso siguiente».
Motivación de diseño
¿Por qué partir la lógica del bucle en tick y after_tick en lugar de resolverlo con un único método step()?
- Ceder la ejecución al llamador: cuando
tickdevuelveTrue, el llamador (el cuerpo dePregel.astream, versync main loop:2964-2984) se encarga de impulsarrunner.tickpara que termine de correr ese lote de tareas y luego vuelve a llamar aloop.after_tick. AsíPregelLoopno necesita retener una referencia aPregelRunner; la relación entre ambos es de colaboración acoplada de forma laxa, no de anidamiento. - Máquina de estados explícita: el campo
statustoma siete valores"input" / "pending" / "done" / "draining" / "interrupt_before" / "interrupt_after" / "out_of_steps"(status Literal:256-264). Cada estado corresponde a un motivo de salida; el código externo se apoya en él para decidir si resume, si raise o si break, sin tener que adivinar la semántica del valor de retorno. - Semántica clara de reanudación:
is_replayingsolo es verdadero en el primer tick — indica que este paso está «reproduciendo» tareas que ya habían terminado antes de la interrupción, y se apaga cuandoafter_ticktermina(is_replaying = False:716). Junto con_reapply_writes_to_succeeded_nodes, que repega los writes exitosos rescatados del checkpoint sobre las tareas en memoria, se consigue que los nodos exitosos no se repitan y los fallidos pasen al error handler. - Drain es la vía de intervención externa:
RunControl.drain_requestedlo fija el consumidor del stream (por ejemplo, el cliente quiere parar a mitad); en cuantoticklo ve pasa al estadodrainingy devuelveFalse, de modo que elwhileexterno sale de forma natural, en lugar de matar el nodo a la fuerza a mitad de ejecución. - Escritura delta en modo exit y orden temporal del checkpoint:
_delta_write_futs(_delta_write_futs:207) recolecta los futures deput_writesde todos los canales delta;_checkpointer_put_after_previousdebe drenarlos antes de escribir el siguiente checkpoint — garantiza el orden causal de «los writes producen el checkpoint, no que el checkpoint tape los writes».
Archivos clave
class PregelLoop:158— definición de la clase; los campos son casi todos attributes desnudos, sin envolver con property, para que las subclasesSyncPregelLoop/AsyncPregelLooplos puedan sobreescribir directamente.status Literal:256-264— siete estados;"out_of_steps"se dispara al alcanzar el límite de pasos,"draining"por petición externa de parada,"done"es cierre normal.tick method:599-681— toda la lógica de un tick: comprobación de pasos,prepare_next_tasks, rama done para tareas vacías, drain,_reapply_writes_to_succeeded_nodes,should_interrupt.out_of_steps:607-609— cuandoself.step > self.stopentra en estadoout_of_stepsy devuelveFalse; elwhileexterno sale.done branch:653-655— cuandoself.tasksestá vacío entra en estadodone; es la «muerte natural» del grafo.draining branch:657-659— cuandocontrol.drain_requestedes verdadero entra en estadodraining.reapply + resume handlers:662-664— al reanudar, primero se reaplican los pending writes sobre los nodos exitosos y luego se llama a_resume_error_handlers_if_applicablepara preparar tareas handler en los nodos fallidos.interrupt_before:667-671— llama ashould_interruptpara comprobar si el nodo cae en la listainterrupt_before; si acierta, raiseGraphInterrupt.after_tick method:683-726— recoger writes →apply_writes→ emit values → limpiar pending →_put_checkpoint→ comprobarinterrupt_after._reapply_writes_to_succeeded_nodes:736-749— reaplica los pending writes sobre las tareas en memoria, pero saltando las cuatro señales de controlERROR / ERROR_SOURCE_NODE / INTERRUPT / RESUME._resume_error_handlers_if_applicable:751-816— recorre la marcaERROR_SOURCE_NODEy construye tareas handler para los nodos fallidos, de modo que el runner salte la tarea original y corra directamente el handler.put_writes:415-508— entrada de los writes que produce una tarea; dedup del canal de control, acumula null task y lanza el futureput_writesdel checkpointer.SyncPregelLoop:1469/AsyncPregelLoop:1722— las dos subclases, que sobreescriben__enter__/__exit__,accept_push,schedule_error_handlery otras piezas que diferencian síncrono de asíncrono.
Flujo de datos
La entrada de tick hace primero la comprobación del límite de pasos y luego llama a prepare_next_tasks para calcular «qué tareas corren en el paso siguiente». Este paso es la semilla del superpaso entero — todas las acciones posteriores trabajan sobre ese lote de tareas.
def tick(self) -> bool:
"""Execute a single iteration of the Pregel loop.
Returns:
True if more iterations are needed.
"""
# check if iteration limit is reached
if self.step > self.stop:
self.status = "out_of_steps"
return False
# prepare next tasks
self.tasks = prepare_next_tasks(
self.checkpoint,
self.checkpoint_pending_writes,
self.nodes,
self.channels,
self.managed,
self.config,
self.step,
self.stop,
for_execution=True,
manager=self.manager,
store=self.store,
checkpointer=self.checkpointer,
trigger_to_nodes=self.trigger_to_nodes,
updated_channels=self.updated_channels,
retry_policy=self.retry_policy,
cache_policy=self.cache_policy,
)Esto viene de tick 头部:599-629. self.step > self.stop es la condición de salida por «límite de pasos» — stop se fija en PregelLoop.__init__ a recursion_limit; si la profundidad de recursión se excede, se para en out_of_steps. prepare_next_tasks toma el checkpoint actual, pending_writes, la tabla de nodos y la de canales y calcula qué tareas ejecutar en este paso, volcándolas en self.tasks. for_execution=True indica que estas tareas se van a correr de verdad (no es un dry-run para vista previa en streaming).
A continuación vienen tres condiciones de salida en paralelo — tareas vacías / drain / reanudar reaplicando writes — y luego should_interrupt:
# if no more tasks, we're done
if not self.tasks:
self.status = "done"
return False
if self.control is not None and self.control.drain_requested:
self.status = "draining"
return False
# if there are pending writes from a previous loop, apply them
if not self.is_replaying and self.checkpoint_pending_writes:
self._reapply_writes_to_succeeded_nodes(self.tasks)
self._resume_error_handlers_if_applicable()
# before execution, check if we should interrupt
if self.interrupt_before and should_interrupt(
self.checkpoint, self.interrupt_before, self.tasks.values()
):
self.status = "interrupt_before"
raise GraphInterrupt()Esto viene de tick 中段:652-671. La comprobación de interrupt_before se hace solo después de reaplicar los pending_writes en la reanudación — porque al reanudar la versión de los canales ya refleja los writes anteriores y should_interrupt puede decidir correctamente «si hubo nuevas actualizaciones desde el último interrupt».
after_tick se llama cuando el runner ha terminado de correr todas las tareas y se encarga del cierre:
def after_tick(self) -> None:
# finish superstep
writes = [w for t in self.tasks.values() for w in t.writes]
self._delta_channels_with_overwrite.update(
ch
for ch, v in writes
if isinstance(self.specs.get(ch), DeltaChannel) and _get_overwrite(v)[0]
)
# all tasks have finished
self.updated_channels = apply_writes(
self.checkpoint,
self.channels,
self.tasks.values(),
self.checkpointer_get_next_version,
self.trigger_to_nodes,
)Esto viene de after_tick 头部:683-698. Primero se aplanan los writes de todas las tareas y se anotan qué canales delta han sufrido overwrite (afecta al valor inicial del sparse replay); luego se llama a apply_writes para fusionar esos writes en los canales; el valor de retorno, updated_channels, lo usará prepare_next_tasks en el paso siguiente para acelerar «encontrar el siguiente nodo a disparar». Al final de after_tick también se vacía checkpoint_pending_writes, se pone is_replaying a False, se persiste el punto de control con _put_checkpoint({"source": "loop"}) y luego se ejecuta la comprobación de interrupt_after(after_tick 尾部:714-724).
El ciclo de vida completo de un superpaso se puede dibujar así:
Límites y fallos
out_of_stepsno es un error:step > stopdevuelveFalsedirectamente ystatusse fija enout_of_steps; elwhileexterno sale de forma natural. Más arriba,Pregeldecide según el status si lanzarRecursionErroro cerrar en silencio — es semántica de «límite blando», no una vía de excepción(out_of_steps:607-609).drainingno puede emitir lifecycle events:_emit_graph_lifecycle_eventrechaza explícitamente que se lo llame en estadodraining(draining guard:384-385), porque drain es una parada iniciada por el usuario y no cuenta como interrupción / reanudación.- Al reanudar, los pending writes con señales de control deben saltarse:
_reapply_writes_to_succeeded_nodestiene que saltarERROR / ERROR_SOURCE_NODE / INTERRUPT / RESUME(skip control signals:746-747); si no, las tareas que fallaron se tratarían como exitosas y los handler ya corridos se reactivarían. - Reprogramación del error handler al reanudar:
_resume_error_handlers_if_applicablesolo aplica a los nodos conerror_handler_nodeconfigurado(handler_node check:791-793); los nodos fallidos sin handler se dejan tal cual y siguen la vía de pánico por_should_stop_othersdel runner. - Semántica de acumulación del null task en
put_writes: los writes conNULL_TASK_IDno se sobreescriben, se acumulan(null task accumulate:422-431); sirve para input writes — writes «que no pertenecen a ningún nodo pero deben entrar en el checkpoint». - Los futures de los canales delta deben persistir antes que el siguiente checkpoint:
_delta_write_futsrecolecta los futuresput_writesde todos los canales delta y se drena antes de persistir el siguiente checkpoint(_delta_write_futs:201-207) — invertir el orden haría que el sparse replay no viera los writes que produjeron el checkpoint.
Resumen
PregelLoop es un metrónomo de superpasos dirigido por máquina de estados que envuelve «preparar tareas → correr → cerrar» en torno a PregelRunner sin tocar la ejecución. Una vez entendidos los dos tramos tick / after_tick, los siete estados de status y el mecanismo de reanudación con is_replaying + _reapply_writes_to_succeeded_nodes, se tiene el eje principal del modelo de ejecución de LangGraph. Cómo hace el runner para correr de verdad las tareas se ve en /pregel/runner; los detalles internos de la capa de algoritmo prepare_next_tasks / apply_writes / should_interrupt en /pregel/algo; el ensamblaje integral en /pregel/pregel. La semántica de escritura en canales se ve en /channel/base-channel.
Véase la documentación oficial: documentación de LangGraph · README