LangGraph Durable Execution
2026-08-04
This note explains, from the ground up, how LangGraph persists state, resumes work, travels through time, handles parallel writes, and supports human-in-the-loop. Every code sample is self-contained and uses dummy data. The outputs shown are representative: I wrote them to illustrate the behaviour, and they are not copied from a live run.
The short version: every
ainvokecall advances a stateful thread, identified by itsthread_id. The rest of this note follows from that.
Table of contents
- The checkpoint model
- Resume (partial replay)
- Time travel: history and forking
- The Postgres table layout
- Concurrency and reducers
- Human-in-the-loop with interrupt
- Multiple interrupts and id matching
- The re-run rule (why side effects double)
- Memory vs Postgres checkpointer
- Cheat sheet
1. The checkpoint model
When you compile a graph with a checkpointer, execution advances in discrete super-steps. Each super-step:
- runs the node(s) scheduled for this step,
- merges each node's returned dict into the state (via reducers, §5),
- writes a new checkpoint (a full snapshot plus a
nextpointer).
Checkpoints are written between nodes, never in the middle of one. A node either completes and gets checkpointed, or it did not run. There is no checkpoint for half a node.
super-step: -1 0 1 2 3
│ │ │ │ │
node run: input __start__ node_a node_b node_c → END
│ │ │ │ │
checkpoint: [c-1] [c0] [c1] [c2] [c3]
next: (__start__,) (node_a,) (node_b,) (node_c,) ()
state adds: - - query partial result
Read the next row as pointing one column to the right: each checkpoint's next is the
node that will run in the following step. So the checkpoint taken after node_a runs
(step 1) has next=(node_b,), because node_a is already done by then. next == ()
(step 3) means nothing is pending and the thread has finished.
Each checkpoint stores the merged full state at that point, plus that next pointer.
Sample: watch state accumulate
from typing import Annotated
from typing_extensions import TypedDict
import operator
from langgraph.graph import StateGraph, START, END
from langgraph.checkpoint.memory import MemorySaver
class State(TypedDict):
value: str
log: Annotated[list, operator.add] # append-only (see §5)
async def step_a(state): return {"value": "A", "log": ["ran A"]}
async def step_b(state): return {"value": "B", "log": ["ran B"]}
g = (StateGraph(State)
.add_node("a", step_a).add_node("b", step_b)
.add_edge(START, "a").add_edge("a", "b").add_edge("b", END)
.compile(checkpointer=MemorySaver()))
cfg = {"configurable": {"thread_id": "demo-1"}}
await g.ainvoke({"value": "", "log": []}, cfg)
async for snap in g.aget_state_history(cfg):
print(snap.metadata["step"], snap.next, snap.values.get("value"), snap.values["log"])
Representative output (newest first):
2 () B ['ran A', 'ran B']
1 ('b',) A ['ran A']
0 ('a',) []
-1 ('__start__',) []
log grows (append reducer) while value is overwritten at each step (default reducer).
§5 is about that difference.
2. Resume (partial replay)
Resuming means invoking the graph again with the same thread_id and continuing from the
last checkpoint instead of starting over. What happens depends on the first argument:
ainvoke first arg (input) | meaning |
|---|---|
| a state / dict | new run, starting from the entry node |
None | resume: continue the existing thread from its checkpoint |
Command(resume=...) | resume a thread that is paused at an interrupt (§6) |
The
thread_id(insideconfig) selects which thread. Theinputargument selects what to do with it. They are separate arguments, and the id always goes inconfig, never ininput.
The classic bug
# WRONG on the resume path: passing a fresh state re-runs from the top,
# even though a checkpoint exists.
result = await g.ainvoke(fresh_state, cfg)
# RIGHT: check for pending work, resume with None if interrupted.
snap = await g.aget_state(cfg)
graph_input = None if snap.next else fresh_state
result = await g.ainvoke(graph_input, config=cfg)
Scenario: a downstream step fails, then recovers
run 1: a ok → b ok → c raises
checkpoints saved: [after a] [after b] next = ('c',)
fix the cause, then:
run 2: ainvoke(None, cfg)
loads [after b], sees next=('c',), runs ONLY c
a and b are NOT re-run; their outputs come from the checkpoint
The work saved before the failure point is reused. Only next and everything after it
runs again.
3. Time travel: history and forking
The checkpointer keeps every checkpoint in the chain, and the same data supports three uses:
| use | how | mutates? |
|---|---|---|
| resume | ainvoke(None, cfg) from the latest checkpoint | advances |
| inspect history | aget_state_history(cfg), iterating all snapshots | read-only |
| fork / what-if | aupdate_state(old_cfg, {...}) then ainvoke(None, ...) | branches |
Forking
Pick an old checkpoint, write a change onto it, and run from there. The original branch stays as it was, and you grow a new one.
# 1. find an old checkpoint (say step 1)
target = None
async for snap in g.aget_state_history(cfg):
if snap.metadata["step"] == 1:
target = snap.config
break
# 2. write a modified value onto that historical point -> a new child checkpoint
new_cfg = await g.aupdate_state(target, {"value": "FORKED"})
# 3. run forward from the fork
await g.ainvoke(None, new_cfg)
The checkpoint chain becomes a tree, linked by parent_checkpoint_id:
c-1 ── c0 ── c1 ─┬─ c2 ── c3 (original branch)
│
└─ c2' ── c3' ── c4' (forked branch: value="FORKED")
^ source=update
Forking is cheap. The new checkpoint copies the pointer map and only writes a new blob for the fields you changed (see §4); unchanged fields are shared.
Typical uses are debugging (rewind, tweak the state, re-run from the bad node), human-in-the-loop edits, and what-if analysis.
4. The Postgres table layout
The Postgres saver creates four tables. The way they split the data explains both de-duplication and crash recovery.
| table | holds | key idea |
|---|---|---|
checkpoints | one row per super-step: parent_checkpoint_id (tree structure) + a channel_versions pointer map | almost no data, just which field points to which blob version |
checkpoint_blobs | the actual field data, chunked by channel + version | cross-step de-dup: an unchanged field shares its version |
checkpoint_writes | pending writes produced by each task within a super-step | task-level crash recovery |
checkpoint_migrations | the saver's own schema version | ignore it |
How a checkpoint stores data (dummy example)
A checkpoints row holds the pointer map, which is tiny:
{
"query": "0000...0002.xxxx", // never changed -> version 2 for the whole run
"options": "0000...0002.xxxx", // never changed -> shared
"partial": "0000...0004.xxxx", // written once at step 2
"result": "0000...0005.xxxx", // written at step 3
"log": "0000...0005.xxxx"
}
checkpoint_blobs holds the data, de-duplicated by version:
| channel | version | bytes | note |
|---|---|---|---|
| query | ...0002 | 204 | stored once, referenced by every later checkpoint |
| options | ...0002 | 7 | never changes, so one row in total |
| partial | ...0004 | 2458 | written at step 2 |
| result | ...0005 | 4630 | written at step 3 |
| log | ...0002 → ...0005 | grows | new version on each append |
Restoring a checkpoint means reading its pointer map, fetching each version's blob, and reassembling them. A run with N steps and M fields does not store N×M full copies, only the fields that changed, once per change.
checkpoint_writes and crash recovery
A super-step may run several tasks in parallel. Each task's outputs are written to
checkpoint_writes before they are merged into the next checkpoint:
node finishes → write its outputs to checkpoint_writes (persist "what I produced")
→ framework merges them into a new checkpoint (checkpoints + blobs)
→ pending writes are cleared
If the process crashes after the writes are recorded but before the checkpoint is formed, LangGraph reads the pending writes on restart instead of re-running that task. For parallel tasks A and B where A finished and B did not, only B re-runs, and A's result is read back. That makes recovery task-level instead of whole-step, so side-effecting calls in the finished tasks are not repeated.
5. Concurrency and reducers
A reducer is a merge function attached to a state field. When parallel branches each return a dict, the reducers decide how their writes combine.
class State(TypedDict):
merged: Annotated[dict, lambda l, r: {**(l or {}), **(r or {})}] # dict-merge
collected: Annotated[list, operator.add] # list-concat
single: str # NO reducer
| field | reducer | parallel writes behaviour |
|---|---|---|
merged | dict-merge | each branch writes different keys → all kept |
collected | operator.add | lists concatenated |
single | none (default = overwrite) | error if written by more than one parallel task |
Parallel fan-out with Send
from langgraph.types import Send
def route(state):
return [Send("worker", {"item": it}) for it in state["items"]]
async def worker(state):
return {"merged": {state["item"]: "done"}, "collected": [state["item"]]}
Representative output for items = ["x", "y", "z"]:
merged : {'x': 'done', 'y': 'done', 'z': 'done'} # all three kept
collected : ['x', 'y', 'z'] # deterministic order
Two properties matter here:
- The merge order is deterministic. It follows the fan-out order, not the order in which tasks complete, so the resulting state is the same however the tasks are scheduled. Checkpoints depend on that.
- Writing a field that has no reducer from more than one parallel task raises an error instead of silently letting the last write win:
InvalidUpdateError: At key 'single': Can receive only one value per step.
Use an Annotated key to handle multiple values.
That turns a class of silent data-loss races into an error you see on the first run. A field without a reducer is only safe when a single node writes it, after the merge.
6. Human-in-the-loop with interrupt
interrupt() makes a node stop on purpose, write a checkpoint, and hand control back to the
caller until an external value arrives.
from langgraph.types import interrupt, Command
async def gate(state):
decision = interrupt({"question": "approve this plan?", "plan": state["plan"]})
return {"approved": decision}
cfg = {"configurable": {"thread_id": "hitl-1"}}
# first invoke: runs until interrupt, then stops
r1 = await g.ainvoke({"plan": "deploy", "approved": ""}, cfg)
print(r1["__interrupt__"])
# [Interrupt(value={'question': 'approve this plan?', 'plan': 'deploy'}, id='a1b2c3...')]
snap = await g.aget_state(cfg)
print(snap.next) # ('gate',) -> paused at the node, awaiting input
# later (possibly a different request, hours later), resume with the decision
r2 = await g.ainvoke(Command(resume="APPROVED"), cfg)
print(r2["approved"]) # APPROVED
Interrupts use the same checkpoint machinery as crash recovery. Only the trigger is different:
| crash recovery | interrupt / human-in-the-loop | |
|---|---|---|
| why it stopped | a node raised (passive) | a node called interrupt() (active) |
| where it stops | before the failed node | at the interrupting node |
| how to continue | ainvoke(None, cfg) | ainvoke(Command(resume=value), cfg) |
| underlying | the same checkpoint | the same checkpoint |
What the caller must send to resume
Only two things:
- the
thread_id, which locates the thread (inconfig); - the decision value, as
Command(resume=value).
Business data such as the plan or earlier intermediate results is not sent again, because it is in the checkpoint. In production that looks like this:
POST /submit {plan_input...} -> {thread_id, question} # stash thread_id
POST /approve {thread_id, decision} -> {status: done, result} # resume by id
This needs a Postgres checkpointer. The two requests may reach different worker processes, so the paused thread has to be visible across processes (§9).
7. Multiple interrupts and id matching
Every interrupt has an id. The id is deterministic: it is derived from the interrupt's
position (which task or branch, and which interrupt within the node), not generated at
random. Two shapes behave differently:
(a) Sequential interrupts in one node: one at a time
async def two_gates(state):
first = interrupt({"ask": "A?"}) # stops here on run 1
second = interrupt({"ask": "B?"}) # stops here after first is resumed
return {"a": first, "b": second}
invoke -> stops at interrupt A
resume "A-VALUE" -> node re-runs; A returns "A-VALUE"; stops at B
resume "B-VALUE" -> node re-runs; A and B both return; node completes
One resume clears one interrupt. Two sequential interrupts need two resume calls, so three
ainvoke calls in total counting the initial invoke. Each resume supplies the value only for
the interrupt that is currently blocking. The node re-runs from the top, the interrupt that
already has an answer returns its stored value without stopping, and execution moves on to
the next unanswered interrupt, which stops again. N interrupts need N resumes.
LangGraph uses the interrupt id to know that an interrupt already has a resume value and
should not stop again. One consequence, covered in §8, is that the code before an interrupt
runs again.
The parallel case (b) below is different: it can clear every interrupt in one resume, because a dict keyed by id supplies all the values at once and no interrupt is left unanswered when the branches re-run.
(b) Parallel interrupts via Send: all at once, matched by id
def fan(state):
return [Send("review", {"item": it}) for it in state["items"]]
async def review(state):
d = interrupt({"review_item": state["item"]})
return {"approvals": [f"{state['item']}={d}"]}
The first invoke returns all the interrupts together, each with its own id:
__interrupt__ = [
Interrupt(id='id-x', value={'review_item': 'doc-X'}),
Interrupt(id='id-y', value={'review_item': 'doc-Y'}),
Interrupt(id='id-z', value={'review_item': 'doc-Z'}),
]
Resume with a dict keyed by id. One call routes each decision to its branch:
await g.ainvoke(Command(resume={
"id-x": "APPROVE",
"id-y": "REJECT",
"id-z": "APPROVE",
}), cfg)
# approvals -> ['doc-X=APPROVE', 'doc-Y=REJECT', 'doc-Z=APPROVE']
| shape | ids | how to resume |
|---|---|---|
| single interrupt | one | Command(resume=value) |
| sequential in one node | same id per position, one at a time | Command(resume=value), repeat |
| parallel branches | distinct per branch | Command(resume={id: value, ...}) |
In production, the /submit response has to return each interrupt's id so that
/approve can map the decisions back.
8. The re-run rule (why side effects double)
This is the behaviour most likely to cause a real bug. On resume, the interrupting node re-runs from the top, so code before the interrupt executes again.
async def node(state):
charge_customer() # SIDE EFFECT
d = interrupt({"ask": "ok?"}) # stops here
return {"decision": d}
first invoke : charge_customer() runs (1st time), then stops at interrupt
resume : node re-runs -> charge_customer() runs AGAIN (2nd time)
Why the checkpoint does not help here
A checkpoint stores what a node returns, not its local variables halfway through.
interrupt() raises a special exception that unwinds the node before it returns, so the
node has produced no output, there is nothing to checkpoint, and next still points at this
node. On resume the node counts as not completed, so it runs again. LangGraph cannot
serialise a half-executed function; only what you return into state survives.
Only the interrupting node re-runs. Nodes that already completed and were checkpointed are not touched, the same as
aandbin the §2 scenario.
Fix A: put the interrupt first, side effects after
async def node(state):
d = interrupt({"ask": "ok?"}) # stop first
charge_customer(d) # runs once, only after resume
return {"decision": d}
Code after the interrupt runs only once, on the resume pass.
Fix B: move the work into its own upstream node
async def prepare(state): # its own node -> its own checkpoint
charge_customer() # side effect
return {"prep": expensive()}
async def gate(state): # interrupt isolated here
d = interrupt({"ask": "ok?", "prep": state["prep"]})
return {"decision": d}
# edges: START -> prepare -> gate -> END
Representative counters:
prepare executed : 1 <- checkpointed after returning; not in the re-run set
gate executed : 2 <- the interrupting node re-runs (no side effect inside)
final prep : from the checkpoint, unchanged
run 1: prepare done (checkpointed) → gate paused (interrupt)
next = ('gate',)
resume: prepare NOT re-run (it's before `next`) → gate re-runs, completes
The design rule is that a node boundary is a checkpoint boundary and a re-run boundary.
Work that must happen once and survive pauses or crashes should get its own node and
return its result. Interrupts, and anything else that might re-run, go in a separate node
downstream.
9. Memory vs Postgres checkpointer
The API is the same, but what survives is different: MemorySaver recovers within the
process, and PostgresSaver recovers across restarts.
MemorySaver | PostgresSaver | |
|---|---|---|
| where checkpoints live | in-process dict | database |
| in-process node failure | resumable | resumable |
| process / worker restart | checkpoints gone | survives |
| resume across HTTP requests / workers | no (may hit a different process) | yes |
| time travel after the fact | no, lost on restart | yes |
| durable human-in-the-loop | no | yes |
# selection is config-driven, not a code change:
if DATABASE_URL:
async with AsyncPostgresSaver.from_conn_string(DATABASE_URL) as cp:
await cp.setup()
graph = build_graph(cp)
else:
graph = build_graph(MemorySaver())
Everything in §2 to §8 (resume, time travel, forking, durable interrupts) only works across
restarts and across workers when the checkpoints are in a shared store. With MemorySaver,
those capabilities last only as long as the original process.
10. Cheat sheet
Advancing a thread: the input argument of ainvoke(input, config)
| you pass | means |
|---|---|
| a state dict | new run from the entry node |
None | resume an interrupted/crashed thread |
Command(resume=v) | resume a thread paused at an interrupt |
Command(resume={id: v}) | resume parallel interrupts, matched by id |
Locating a thread
Always config = {"configurable": {"thread_id": ...}}. The id is the address, and the data
is in the checkpoint.
Key APIs
| call | does |
|---|---|
aget_state(cfg) | latest snapshot; .next = pending nodes, .values = state |
aget_state_history(cfg) | iterate all checkpoints (time travel) |
aupdate_state(cfg, {...}) | write onto a checkpoint → fork a branch |
interrupt(payload) | pause the node, await external input |
Reducers (Annotated[type, fn])
| reducer | effect |
|---|---|
operator.add | list concat / append |
{**l, **r} | dict merge |
| none (default) | overwrite; errors on parallel writes |
Four rules
- Checkpoints are written between nodes, and they store each node's returned output, not its local variables mid-run.
- Resume re-runs
nextand everything after it; earlier checkpointed nodes are reused. - Parallel writes need a reducer, or LangGraph raises.
- The interrupting node re-runs from the top, so keep side effects out of it (put the interrupt first, or move the side effects to an upstream node).
Node boundary = checkpoint boundary = re-run boundary.