StateSnapshot: la instantánea de estado al final de cada superpaso
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 faltaread_channelspara leer los valores de canal (read_channels:1258). Una capa extraStateSnapshotdeja 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
CheckpointTuplesólo guarda los valores escritos;StateSnapshotllama aprepare_next_tasksal generarse (prepare_next_tasks:1178-1195) y calculanext=(node_name, ...), respondiendo directamente «qué va a pasar ahora». - Lista de tareas visible: cada
StateSnapshot.taskses la tupla dePregelTaskque se ejecutará en el paso actual (tasks:658), incluyendowrites/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_historydevuelve la lista de instantáneas en orden cronológico inverso, y cualquiera puede ser targeteada por suconfigpara hacer viaje en el tiempo. - Aplicación transparente de pending writes:
get_stateaplica por defecto lospending_writesque aún no se han persistido en un checkpoint antes de leer (apply pending writes:1239-1249); así, si llamasget_statejusto 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 campostatecontiene unStateSnapshotdel subgrafo o elRunnableConfigque lo localiza._prepare_state_snapshot:1145-1162— fallback con saved vacío: devuelve una instantánea vacía convalues={}/next=()/tasks=().step 计算 + prepare_next_tasks:1167-1195— obtiene el número de paso actual desdesaved.metadata["step"]y llama aprepare_next_taskspara calcular las tareas de este paso.task_states 递归:1199-1227— para cada tarea contask.name in subgraphs, compone elcheckpoint_nsy llama recursivamente alget_statedel subgrafo.apply pending writes:1229-1255— aplica temporalmente al canal las writes conNULL_TASK_IDy los pending writes normales, y luego calculatasks_with_writes.组装 StateSnapshot:1257-1266— ensambla elStateSnapshotfinal 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íacheckpoint_nsy decide siapply_pending_writes.get_state_history:1480— recorre todos los checkpoints históricos del mismothread_idy genera unStateSnapshotpara cada uno._aprepare_state_snapshot:1268— versión asíncrona, simétrica a la síncrona; sólo cambiachannels_from_checkpointporachannels_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.
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:
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]),
)Límites y fallos
- Sin checkpoint devuelve una instantánea vacía en vez de error (
空 saved:1152-1162), convalues={}/next=(); el llamador debe mirar sisnapshot.nextes la tupla vacía para distinguir «el grafo aún no corrió» de «el grafo terminó». - Llamar
get_statesin checkpointer lanzaValueError(no checkpointer:1401-1402) —No checkpointer set—. Un grafo puramente en memoria no tiene instantáneas históricas. checkpoint_nssin subgrafo que coincida lanza error (subgraph not found:1415-1416) —Subgraph {recast} not found—. Suele ocurrir cuando se ha concatenado mal el namespace conNS_SEP(|) oNS_END(:).apply_pending_writessólo se activa cuando no se pasacheckpoint_idexplí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
nextfiltra las tareas que ya tienen writes (next 过滤:1259): las tareas cont.writesno vacío no se cuentan como «next» —su input ya fue consumido en este paso, así quenextrefleja «las tareas que aún no han corrido». tasksincluyeinterrupts(interrupts 聚合:1265): se aplanan losInterruptde todas las tareas en una tupla, en el orden queCommand(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