Skip to content

StateSnapshot: la instantánea de estado al final de cada superpaso

源码版本1.2.9

Responsabilidades

StateSnapshot es la instantánea inmutable que LangGraph expone en cada frontera de superpaso (superstep) para responder «cómo está el grafo justo después de correr este paso». Empaqueta cinco cosas en un NamedTuple: los valores actuales de cada canal (channel) (values), qué nodos se ejecutarán a continuación (next), el config usado, los metadatos del checkpoint correspondiente (metadata + created_at) y el config del checkpoint padre (StateSnapshot 定义:643-661). Lo que devuelven graph.get_state(config), graph.aget_state(config) o lo que recorre graph.get_state_history(config) es siempre de este tipo.

En la arquitectura se sitúa «por encima del checkpoint, por debajo de la vista de usuario». El checkpoint (CheckpointTuple) es el formato interno de serialización del motor: guarda campos de bajo nivel como channel_values / channel_versions / pending_writes, que no permiten leer cómo es el estado de negocio; StateSnapshot traduce esa capa en «un dict legible por el negocio + el siguiente paso + la lista de tareas» (_prepare_state_snapshot:1145-1162). Los estados de subgrafos también lo usan: el campo task.state, en la recursión con subgraphs=True, aloja un StateSnapshot del subgrafo (PregelTask.state:200-201).

StateSnapshot es un NamedTuple inmutable; para «cambiar el estado» hay que pasar por graph.update_state(config, values), que el motor convierte en un pending write con NULL_TASK_ID y persiste en el siguiente checkpoint, generando luego un nuevo StateSnapshot.

Motivación de diseño

¿Por qué no exponer CheckpointTuple directamente al usuario?

  • Desacoplar la representación interna: el formato del checkpoint, para serializar de forma eficiente, guarda dicts planos con números de versión como channel_versions; al usuario le importa «mi dict de estado de negocio», y entre medias hace falta read_channels para leer los valores de canal (read_channels:1258). Una capa extra StateSnapshot deja evolucionar el formato del checkpoint sin romper la API externa.
  • Semántica de «siguiente paso»: el motor sabe qué nodos se ejecutarán a continuación, pero CheckpointTuple sólo guarda los valores escritos; StateSnapshot llama a prepare_next_tasks al generarse (prepare_next_tasks:1178-1195) y calcula next=(node_name, ...), respondiendo directamente «qué va a pasar ahora».
  • Lista de tareas visible: cada StateSnapshot.tasks es la tupla de PregelTask que se ejecutará en el paso actual (tasks:658), incluyendo writes / interrupts / estado del subgrafo —clave en flujos humano-en-el-bucle para decidir «en qué nodo estamos atascados y qué input esperamos».
  • Inmutable = comparable: las instantáneas históricas se encadenan con parent_config (parent_config:656); get_state_history devuelve la lista de instantáneas en orden cronológico inverso, y cualquiera puede ser targeteada por su config para hacer viaje en el tiempo.
  • Aplicación transparente de pending writes: get_state aplica por defecto los pending_writes que aún no se han persistido en un checkpoint antes de leer (apply pending writes:1239-1249); así, si llamas get_state justo después de que un nodo escribió pero antes de que se persista el checkpoint, ves el valor más reciente, no el valor viejo del paso anterior.

Archivos clave

  • StateSnapshot NamedTuple:643-661 — 7 campos: values / next / config / metadata / created_at / parent_config / tasks / interrupts.
  • PregelTask:200-217 — objeto de tarea; en recursión el campo state contiene un StateSnapshot del subgrafo o el RunnableConfig que lo localiza.
  • _prepare_state_snapshot:1145-1162 — fallback con saved vacío: devuelve una instantánea vacía con values={} / next=() / tasks=().
  • step 计算 + prepare_next_tasks:1167-1195 — obtiene el número de paso actual desde saved.metadata["step"] y llama a prepare_next_tasks para calcular las tareas de este paso.
  • task_states 递归:1199-1227 — para cada tarea con task.name in subgraphs, compone el checkpoint_ns y llama recursivamente al get_state del subgrafo.
  • apply pending writes:1229-1255 — aplica temporalmente al canal las writes con NULL_TASK_ID y los pending writes normales, y luego calcula tasks_with_writes.
  • 组装 StateSnapshot:1257-1266 — ensambla el StateSnapshot final con los valores de canal, la tupla de nodos next, config, metadata, ts, parent_config, tasks e interrupts.
  • get_state:1392-1434 — entrada síncrona: resuelve el checkpointer, enruta al subgrafo vía checkpoint_ns y decide si apply_pending_writes.
  • get_state_history:1480 — recorre todos los checkpoints históricos del mismo thread_id y genera un StateSnapshot para cada uno.
  • _aprepare_state_snapshot:1268 — versión asíncrona, simétrica a la síncrona; sólo cambia channels_from_checkpoint por achannels_from_checkpoint.

Flujo de datos

Este es el tramo clave en el que _prepare_state_snapshot traduce el checkpoint a un StateSnapshot: a partir de saved.metadata obtiene el número de paso, llama a prepare_next_tasks para traer las tareas del paso y luego gestiona la recursión de subgrafos.

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

(prepare_next_tasks 调用:1167-1195)

Después, con los valores de canal leídos y filtradas las tareas que ya tienen writes pendientes, se ensambla la instantánea inmutable:

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

(组装:1257-1266)

Límites y fallos

  • Sin checkpoint devuelve una instantánea vacía en vez de error (空 saved:1152-1162), con values={} / next=(); el llamador debe mirar si snapshot.next es la tupla vacía para distinguir «el grafo aún no corrió» de «el grafo terminó».
  • Llamar get_state sin checkpointer lanza ValueError (no checkpointer:1401-1402) —No checkpointer set—. Un grafo puramente en memoria no tiene instantáneas históricas.
  • checkpoint_ns sin subgrafo que coincida lanza error (subgraph not found:1415-1416) —Subgraph {recast} not found—. Suele ocurrir cuando se ha concatenado mal el namespace con NS_SEP (|) o NS_END (:).
  • apply_pending_writes sólo se activa cuando no se pasa checkpoint_id explícito (apply 条件:1433); es decir, las instantáneas históricas muestran los valores persistidos en su momento, sin contaminación por pending writes posteriores; sólo la instantánea «más reciente» aplica los pending temporalmente.
  • El campo next filtra las tareas que ya tienen writes (next 过滤:1259): las tareas con t.writes no vacío no se cuentan como «next» —su input ya fue consumido en este paso, así que next refleja «las tareas que aún no han corrido».
  • tasks incluye interrupts (interrupts 聚合:1265): se aplanan los Interrupt de todas las tareas en una tupla, en el orden que Command(resume=...) debe respetar al rellenar.

Resumen

StateSnapshot es el tipo de frontera con el que LangGraph traduce el formato interno de checkpoint en una vista legible por el negocio —no almacena datos, sólo ensambla el checkpoint + los pending writes con la semántica de «siguiente paso» en una instantánea inmutable. Todas las APIs de reanudación tras interrupción, viaje en el tiempo y humano-en-el-bucle se construyen sobre él.

Forma pareja con RunnableConfig: el config es la coordenada «cómo encontrar esta instantánea» y StateSnapshot es «qué se ve al encontrarla». Continúa en StateGraph para ver cómo se define la forma del estado que se instantánea, o en motor Pregel para ver cómo _prepare_state_snapshot acaba siendo invocado desde get_state.

Véase la documentación oficial: documentación de LangGraph · README