Skip to content

PregelLoop : pilote de la boucle de superpas

源码版本1.2.9

Responsabilités

PregelLoop est le « métronome » du runtime LangGraph. Pregel lui-même ne fait qu'assembler les ressources et exposer les entrées invoke / astream ; ce qui fait tourner le graphe, cercle après cercle, est une instance de PregelLoop : elle détient le point de contrôle (checkpoint) courant, l'ensemble des canaux (channels), les pending_writes, le compteur step, la machine à états status, ainsi qu'un dictionnaire tasks — tout le contexte d'un seul superpas (superstep). Le schéma enveloppant while loop.tick(): ... loop.after_tick() s'appuie sur elle.

Ses deux méthodes centrales sont tick (tick:599) et after_tick (after_tick:683). La première gère le « avant exécution » : vérifier si l'on a dépassé la limite de pas, appeler prepare_next_tasks pour charger dans self.tasks les nœuds à exécuter pour ce pas, examiner les interruptions, traiter drain_requested, et recoller les pending_writes d'une reprise précédente sur les nœuds déjà réussis. La seconde gère le « après exécution » : collecter tous les task.writes, appeler apply_writes pour les fusionner dans les canaux, nettoyer les pending, poser le checkpoint, puis décider des interrupt_after. L'une et l'autre encadrent le cycle BSP « lire l'entrée → planifier → exécuter → réécrire → persister ».

Autrement dit, PregelLoop n'exécute pas elle-même les nœuds — c'est le travail de PregelRunner ; elle décide seulement « qui regarder à ce pas, comment finaliser après exécution, et s'il faut continuer au pas suivant ».

Motivation de conception

Pourquoi séparer la logique de boucle en tick et after_tick plutôt que d'embrayer tout dans une méthode step() ?

  • Rendre l'exécution à l'appelant : une fois que tick renvoie True, l'appelant (le corps de Pregel.astream, voir sync main loop:2964-2984) pilote lui-même runner.tick pour faire tourner ce batch de tâches, puis rappelle loop.after_tick. PregelLoop n'a ainsi pas besoin de détenir une référence à PregelRunner ; les deux objets collaborent en couplage lâche, sans imbrication.
  • Machine à états explicite : le champ status prend sept valeurs — "input" / "pending" / "done" / "draining" / "interrupt_before" / "interrupt_after" / "out_of_steps" (status Literal:256-264) ; chaque état correspond à une raison de sortie, le code externe s'en sert pour décider s'il faut resume, raise ou break, sans devoir deviner la sémantique de la valeur de retour.
  • Sémantique claire de reprise après rupture : is_replaying n'est vrai qu'au premier tick — il indique que ce pas « rejoue » les tâches déjà achevées avant la dernière interruption ; dès que after_tick termine, il est désactivé (is_replaying = False:716). Combiné à _reapply_writes_to_succeeded_nodes qui recolle les writes réussis du checkpoint sur les tâches mémoire, cela permet aux nœuds réussis de ne pas être relancés et aux nœuds en échec d'entrer dans le handler d'erreur.
  • Drain est la porte d'intervention externe : RunControl.drain_requested est positionné par le consommateur du stream (par ex. un client qui veut s'arrêter en cours) ; dès que tick le voit, il passe à l'état draining et renvoie False, le while externe sort naturellement, sans avoir à tuer un nœud en plein milieu.
  • Écriture en mode delta et timing du checkpoint : _delta_write_futs (_delta_write_futs:207) collecte les futures put_writes de tous les canaux delta ; _checkpointer_put_after_previous doit d'abord les drainer avant d'écrire le checkpoint suivant — pour garantir l'ordre causal « l'écriture produit le checkpoint, pas le checkpoint qui écrase l'écriture ».

Fichiers clés

  • class PregelLoop:158 — définition de classe ; les champs sont quasiment tous des attributs bruts, sans property, pour que les sous-classes SyncPregelLoop / AsyncPregelLoop puissent les surcharger directement.
  • status Literal:256-264 — énumération à sept états ; "out_of_steps" est déclenché par la limite de pas, "draining" par une requête d'arrêt externe, "done" est la fin normale.
  • tick method:599-681 — toute la logique d'un tick complet : vérification du pas, prepare_next_tasks, branche done sur tâche vide, drain, _reapply_writes_to_succeeded_nodes, should_interrupt.
  • out_of_steps:607-609 — quand self.step > self.stop, passe à l'état out_of_steps et renvoie False, le while externe sort.
  • done branch:653-655 — quand self.tasks est vide, passe à l'état done, c'est la « mort naturelle » du graphe.
  • draining branch:657-659 — quand control.drain_requested est vrai, passe à l'état draining.
  • reapply + resume handlers:662-664 — à la reprise, on colle d'abord les pending writes sur les nœuds réussis, puis on appelle _resume_error_handlers_if_applicable pour préparer des tâches de handler pour les nœuds en échec.
  • interrupt_before:667-671 — appelle should_interrupt pour vérifier si un nœud tombe dans la liste interrupt_before ; si oui, lève GraphInterrupt.
  • after_tick method:683-726 — collecte les writes → apply_writes → émet values → nettoie pending → _put_checkpoint → vérification interrupt_after.
  • _reapply_writes_to_succeeded_nodes:736-749 — recolle les pending writes sur les tâches mémoire, mais saute les quatre signaux de contrôle ERROR / ERROR_SOURCE_NODE / INTERRUPT / RESUME.
  • _resume_error_handlers_if_applicable:751-816 — scanne le marqueur ERROR_SOURCE_NODE, construit une tâche de handler d'erreur pour le nœud en échec, afin que le runner saute la tâche originelle et exécute directement le handler.
  • put_writes:415-508 — entrée des writes produites par une tâche ; déduplique les canaux de contrôle, accumule les null task, lance les futures put_writes du checkpointer.
  • SyncPregelLoop:1469 / AsyncPregelLoop:1722 — deux sous-classes, surchargent __enter__/__exit__, accept_push, schedule_error_handler et les autres parties différant entre sync et async.

Flux de données

L'entrée de tick effectue d'abord la vérification de la limite de pas, puis appelle prepare_next_tasks pour calculer « quelles tâches exécuter au pas suivant ». C'est la graine de tout le superpas — toutes les opérations ultérieures travaillent ce batch de tâches.

python
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,
    )

Extrait de tick 头部:599-629. self.step > self.stop est la condition de sortie « limite de pas » — stop est fixé à recursion_limit dans PregelLoop.__init__ ; si la profondeur de récursion est dépassée, on s'arrête sur out_of_steps. prepare_next_tasks prend le checkpoint courant, les pending_writes, la table des nœuds et la table des canaux pour calculer quelles tâches exécuter à ce pas et les ranger dans self.tasks. for_execution=True indique que ces tâches doivent réellement s'exécuter (et non un dry-run pour aperçu streaming).

Viennent ensuite trois conditions de sortie en parallèle — tâche vide / drain / reprise collant les writes — puis should_interrupt :

python
# 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()

Extrait de tick 中段:652-671. La vérification interrupt_before n'a lieu qu'après re-collage des pending_writes en cas de reprise — car à la reprise les versions de canaux reflètent déjà les écritures précédentes, ce qui permet à should_interrupt de décider correctement « y a-t-il eu de nouvelles mises à jour depuis la dernière interruption ».

after_tick est appelé quand le runner a terminé toutes les tâches ; il gère le nettoyage :

python
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,
    )

Extrait de after_tick 头部:683-698. On aplatit d'abord les writes de toutes les tâches, en notant au passage quels canaux delta ont subi un overwrite (ce qui affecte la valeur de départ du replay sparse) ; puis on appelle apply_writes pour fusionner ces writes dans les canaux, et la valeur renvoyée updated_channels servira de hint au prepare_next_tasks suivant pour accélérer « trouver le prochain nœud à déclencher ». La fin de after_tick vide aussi checkpoint_pending_writes, passe is_replaying à False, pose le checkpoint via _put_checkpoint({"source": "loop"}), puis exécute le test interrupt_after (after_tick 尾部:714-724).

Le cycle de vie complet d'un superpas se dessine ainsi :

Limites et échecs

  • out_of_steps n'est pas une erreur : step > stop renvoie directement False, le status est mis à out_of_steps, le while exterme sort naturellement. L'étage supérieur de Pregel décide selon le status s'il faut lever RecursionError ou s'arrêter silencieusement — c'est une sémantique de « plafond souple », pas un chemin d'exception (out_of_steps:607-609).
  • draining ne doit pas émettre d'événements lifecycle : _emit_graph_lifecycle_event refuse explicitement d'être appelé à l'état draining (draining guard:384-385), car drain est un arrêt volontaire de l'utilisateur, et n'a pas la sémantique d'interrupt / resume.
  • À la reprise, les pending writes contenant des signaux de contrôle doivent être sautés : _reapply_writes_to_succeeded_nodes doit ignorer ERROR / ERROR_SOURCE_NODE / INTERRUPT / RESUME (skip control signals:746-747), sinon les tâches qui ont échoué seraient considérées comme réussies, et les handlers déjà joués seraient redéclenchés.
  • Re-planification du handler d'erreur au resume : _resume_error_handlers_if_applicable ne s'applique qu'aux nœuds disposant d'un error_handler_node (handler_node check:791-793) ; les nœuds en échec sans handler restent en l'état, et le _should_stop_others du runner les emmène sur le chemin panic.
  • Sémantique d'accumulation du null task dans put_writes : les writes de NULL_TASK_ID ne s'écrasent pas mais s'accumulent (null task accumulate:422-431), pour les écritures de type input qui « n'appartiennent à aucun nœud mais doivent entrer dans le checkpoint ».
  • Les futures du canal delta doivent être flushés avant le prochain checkpoint : _delta_write_futs collecte les futures put_writes de tous les canaux delta ; avant d'écrire le checkpoint suivant, on les draine (_delta_write_futs:201-207) — inverser l'ordre ferait que le replay sparse ne verrait pas les writes qui ont produit le checkpoint.

Résumé

PregelLoop est un métronome de superpas piloté par une machine à états, qui encadre le triptyque « préparer les tâches → exécuter → nettoyer » autour de PregelRunner sans toucher lui-même à l'exécution. Comprendre les deux segments tick / after_tick, les sept états de status, et le mécanisme de reprise is_replaying + _reapply_writes_to_succeeded_nodes, c'est saisir l'axe principal du modèle d'exécution de LangGraph. Pour la façon dont le runner fait réellement tourner les tâches, voir /pregel/runner ; pour les détails internes des fonctions algorithmiques prepare_next_tasks / apply_writes / should_interrupt, voir /pregel/algo ; pour l'assemblage global, voir /pregel/pregel. La sémantique d'écriture des canaux est dans /channel/base-channel.

Voir la documentation officielle : LangGraph docs · README.