Skip to content

StateSnapshot : l'instantané d'état à la fin de chaque superpas

源码版本1.2.9

Responsabilités

StateSnapshot est l'instantané (snapshot) immuable que LangGraph expose à la frontière de chaque superpas (superstep) pour décrire « à quoi ressemble tout le graphe à la fin de cette étape ». Il empaquette cinq choses dans un NamedTuple : les valeurs courantes de chaque canal (channel) (values), les nœuds à exécuter ensuite (next), la config utilisée, les métadonnées du point de contrôle (checkpoint) correspondant (metadata + created_at), et la config du checkpoint parent (définition StateSnapshot:643-661). Lorsque vous appelez graph.get_state(config), graph.aget_state(config), ou itérez sur graph.get_state_history(config), c'est ce type que vous récupérez.

Dans l'architecture, il se situe « au-dessus du checkpoint, en dessous de la vue utilisateur ». Le checkpoint (CheckpointTuple) est le format interne sérialisé du moteur : il stocke directement des champs bas niveau comme channel_values / channel_versions / pending_writes, qui ne disent pas à quoi ressemble l'état métier ; StateSnapshot traduit cette couche en « dict lisible par le métier + étape suivante + liste de tâches » (_prepare_state_snapshot:1145-1162). L'état des sous-graphes passe aussi par lui — le champ task.state reçoit un StateSnapshot du sous-graphe lors d'une récursion subgraphs=True (PregelTask.state:200-201).

StateSnapshot lui-même est un NamedTuple, donc immuable ; pour « modifier l'état », il faut passer par graph.update_state(config, values) qui produit un pending write avec NULL_TASK_ID déposé dans le prochain checkpoint, puis génère un nouveau StateSnapshot.

Motivation de conception

Pourquoi ne pas exposer directement CheckpointTuple à l'utilisateur ?

  • Découpler la représentation interne : le format de checkpoint, pour sérialiser efficacement, stocke des dicts plats avec numéros de version comme channel_versions ; ce que l'utilisateur veut, c'est « mon dict d'état métier », ce qui nécessite une étape read_channels pour lire les valeurs des canaux (read_channels:1258). Une couche supplémentaire StateSnapshot permet au format de checkpoint d'évoluer librement sans casser l'API publique.
  • Porter la sémantique « étape suivante » : le moteur sait quels nœuds vont s'exécuter ensuite, mais CheckpointTuple lui-même ne stocke que les valeurs écrites ; StateSnapshot appelle prepare_next_tasks à la construction (prepare_next_tasks:1178-1195) pour calculer next=(node_name, ...), et répondre directement à « que va-t-il se passer ensuite ».
  • Visibilité de la liste de tâches : chaque StateSnapshot.tasks est le tuple des PregelTask à exécuter à cette étape (tasks:658) — avec writes / interrupts / état de sous-graphe — essentiel pour juger, en interaction humain-dans-la-boucle, « sur quel nœud est-on bloqué, quelle entrée attend-on ».
  • Immuabilité = comparabilité : les instantanés historiques forment une liste chaînée via parent_config (parent_config:656), get_state_history renvoie la liste des snapshots par ordre chronologique inverse, et n'importe lequel peut être ciblé à nouveau via son config pour du time travel.
  • Application transparente des pending writes : get_state applique par défaut les pending_writes non encore checkpointés dans les canaux avant de les lire (apply pending writes:1239-1249), donc si vous appelez get_state juste après l'écriture d'un nœud mais avant le checkpoint, vous voyez la dernière valeur, pas celle de l'étape précédente.

Fichiers clés

  • StateSnapshot NamedTuple:643-661 — 7 champs : values / next / config / metadata / created_at / parent_config / tasks / interrupts.
  • PregelTask:200-217 — objet tâche, dont le champ state contient en récursion soit le StateSnapshot du sous-graphe, soit un RunnableConfig le localisant.
  • _prepare_state_snapshot:1145-1162 — repli quand saved est vide : renvoie un snapshot vide avec values={} / next=() / tasks=().
  • calcul step + prepare_next_tasks:1167-1195 — déduit le numéro d'étape courant depuis saved.metadata["step"], appelle prepare_next_tasks pour calculer le dict des tâches à exécuter à cette étape.
  • récursion task_states:1199-1227 — pour chaque tâche dont task.name in subgraphs, assemble le checkpoint_ns puis appelle récursivement get_state sur le sous-graphe.
  • apply pending writes:1229-1255 — applique temporairement au canal les écritures NULL_TASK_ID et les pending writes régulières, puis calcule tasks_with_writes.
  • assemblage StateSnapshot:1257-1266 — assemble valeurs de canaux, tuple de nœuds next, config, metadata, ts, parent_config, tasks, interrupts dans le StateSnapshot final.
  • get_state:1392-1434 — entrée synchrone : résout le checkpointer, route vers le sous-graphe via checkpoint_ns, décide d'apply_pending_writes ou non.
  • get_state_history:1480 — itère tous les checkpoints historiques sous le même thread_id, génère un StateSnapshot pour chacun.
  • _aprepare_state_snapshot:1268 — version async, logique symétrique à la version sync, channels_from_checkpoint remplacé par achannels_from_checkpoint.

Flux de données

Voici le passage clé de _prepare_state_snapshot traduisant un checkpoint en StateSnapshot : déduction du numéro d'étape depuis saved.metadata, appel à prepare_next_tasks pour récupérer les tâches de l'étape, puis traitement de la récursion des sous-graphes.

python
step = saved.metadata.get("step", -1) + 1
stop = step + 2
channels, managed = channels_from_checkpoint(
    self.channels,
    saved.checkpoint,
    saver=self.checkpointer
    if isinstance(self.checkpointer, BaseCheckpointSaver) else None,
    config=saved.config,
)
next_tasks = prepare_next_tasks(
    saved.checkpoint,
    saved.pending_writes or [],
    self.nodes,
    channels,
    managed,
    saved.config,
    step,
    stop,
    for_execution=True,
    store=self.store,
    checkpointer=(self.checkpointer
                  if isinstance(self.checkpointer, BaseCheckpointSaver) else None),
    manager=None,
)

(appel prepare_next_tasks:1167-1195)

Puis, à partir des valeurs de canaux lues et en filtrant les noms de tâches ayant déjà des writes, assemblage du snapshot immuable :

python
return StateSnapshot(
    read_channels(channels, self.stream_channels_asis),
    tuple(t.name for t in next_tasks.values() if not t.writes),
    patch_checkpoint_map(saved.config, saved.metadata),
    saved.metadata,
    saved.checkpoint["ts"],
    patch_checkpoint_map(saved.parent_config, saved.metadata),
    tasks_with_writes,
    tuple([i for task in tasks_with_writes for i in task.interrupts]),
)

(assemblage:1257-1266)

Limites et échecs

  • Sans checkpoint, un snapshot vide est renvoyé plutôt qu'une erreur (saved vide:1152-1162) — values={} / next=(), l'appelant doit lui-même tester si snapshot.next est un tuple vide pour distinguer « graphe pas encore lancé » et « graphe terminé ».
  • Appeler get_state sans checkpointer lève ValueError (no checkpointer:1401-1402) — No checkpointer set ; un graphe purement en mémoire n'a pas de snapshot historique.
  • checkpoint_ns sans sous-graphe correspondant : erreur (subgraph not found:1415-1416) — Subgraph {recast} not found, typique quand l'assemblage du namespace se trompe sur NS_SEP (|) ou NS_END (:).
  • apply_pending_writes n'est activé que si checkpoint_id n'est pas donné explicitement (condition apply:1433) — autrement dit, un snapshot historique montre les valeurs telles qu'elles étaient à la persistance, sans pollution par un pending write ultérieur ; seul le snapshot « le plus récent » applique temporairement les pending.
  • Le champ next filtre les tâches ayant déjà des writes (filtrage next:1259) — une tâche dont t.writes n'est pas vide n'est pas comptée comme « next », car ses entrées ont déjà été consommées à cette étape ; next reflète donc « les tâches pas encore exécutées ».
  • tasks inclut interrupts (agrégation interrupts:1265) — aplatit les objets Interrupt de toutes les tâches en un tuple, dans l'ordre à respecter quand l'utilisateur appelle Command(resume=...).

Résumé

StateSnapshot est le type frontière par lequel LangGraph traduit son format interne de checkpoint en une vue lisible par le métier — il ne stocke pas de données, il assemble checkpoint + pending writes en un instantané immuable, augmenté d'une sémantique « étape suivante ». Toutes les API de reprise après interruption, de time travel et d'interaction humain-dans-la-boucle se construisent au-dessus de lui.

Il forme un binôme avec RunnableConfig : la config est la coordonnée « comment trouver ce snapshot », StateSnapshot est « ce qu'on voit une fois trouvé ». Continuer avec StateGraph pour voir comment est définie la forme d'état que l'on snapshotte, ou avec moteur Pregel pour voir comment _prepare_state_snapshot est appelé par get_state.

Voir la documentation officielle : LangGraph docs · README.