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 ainvoke call advances a stateful thread, identified by its thread_id. The rest of this note follows from that.


Table of contents

  1. The checkpoint model
  2. Resume (partial replay)
  3. Time travel: history and forking
  4. The Postgres table layout
  5. Concurrency and reducers
  6. Human-in-the-loop with interrupt
  7. Multiple interrupts and id matching
  8. The re-run rule (why side effects double)
  9. Memory vs Postgres checkpointer
  10. Cheat sheet

1. The checkpoint model

When you compile a graph with a checkpointer, execution advances in discrete super-steps. Each super-step:

  1. runs the node(s) scheduled for this step,
  2. merges each node's returned dict into the state (via reducers, §5),
  3. writes a new checkpoint (a full snapshot plus a next pointer).

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 / dictnew run, starting from the entry node
Noneresume: continue the existing thread from its checkpoint
Command(resume=...)resume a thread that is paused at an interrupt (§6)

The thread_id (inside config) selects which thread. The input argument selects what to do with it. They are separate arguments, and the id always goes in config, never in input.

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:

usehowmutates?
resumeainvoke(None, cfg) from the latest checkpointadvances
inspect historyaget_state_history(cfg), iterating all snapshotsread-only
fork / what-ifaupdate_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.

tableholdskey idea
checkpointsone row per super-step: parent_checkpoint_id (tree structure) + a channel_versions pointer mapalmost no data, just which field points to which blob version
checkpoint_blobsthe actual field data, chunked by channel + versioncross-step de-dup: an unchanged field shares its version
checkpoint_writespending writes produced by each task within a super-steptask-level crash recovery
checkpoint_migrationsthe saver's own schema versionignore 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:

channelversionbytesnote
query...0002204stored once, referenced by every later checkpoint
options...00027never changes, so one row in total
partial...00042458written at step 2
result...00054630written at step 3
log...0002 → ...0005growsnew 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
fieldreducerparallel writes behaviour
mergeddict-mergeeach branch writes different keys → all kept
collectedoperator.addlists concatenated
singlenone (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 recoveryinterrupt / human-in-the-loop
why it stoppeda node raised (passive)a node called interrupt() (active)
where it stopsbefore the failed nodeat the interrupting node
how to continueainvoke(None, cfg)ainvoke(Command(resume=value), cfg)
underlyingthe same checkpointthe same checkpoint

What the caller must send to resume

Only two things:

  1. the thread_id, which locates the thread (in config);
  2. 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']
shapeidshow to resume
single interruptoneCommand(resume=value)
sequential in one nodesame id per position, one at a timeCommand(resume=value), repeat
parallel branchesdistinct per branchCommand(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 a and b in 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.

MemorySaverPostgresSaver
where checkpoints livein-process dictdatabase
in-process node failureresumableresumable
process / worker restartcheckpoints gonesurvives
resume across HTTP requests / workersno (may hit a different process)yes
time travel after the factno, lost on restartyes
durable human-in-the-loopnoyes
# 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 passmeans
a state dictnew run from the entry node
Noneresume 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

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

reducereffect
operator.addlist concat / append
{**l, **r}dict merge
none (default)overwrite; errors on parallel writes

Four rules

  1. Checkpoints are written between nodes, and they store each node's returned output, not its local variables mid-run.
  2. Resume re-runs next and everything after it; earlier checkpointed nodes are reused.
  3. Parallel writes need a reducer, or LangGraph raises.
  4. 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.