Day 44 / 共 60 天 · 阶段 7 中断与人在环

时间旅行:翻看每一份存档、修改过去、从岔路重开

既然每个超步都存了 checkpoint(D33-40),那"回到过去"就是自然而然的能力:get_state_history() 把整条存档链读出来,update_state() 让你修改某个历史点的状态、甚至假装"是某个节点写的"。这就是时间旅行(time travel)——人在环里"人不仅能答,还能改、能倒带"的底座。

📍 阶段 7 · 中断与人在环(6 天)你在这里
D41 interrupt原理 D42 Command恢复 D43 静态断点 D44 时间旅行 D45 durability D46 replay幂等
L01

痛点:AI 走错一步,能不能倒回去重来?

🤔 痛点Agent 跑了 8 步,第 5 步模型选错了工具,导致后面全歪了。难道要从头 invoke 重跑(浪费前 4 步、而且模型可能又随机选错)?我想要的是:倒回到第 5 步之前,手动把状态改对,从那里继续——就像游戏读一个早期存档、改改装备再往下打。
💡 本质:存档链 + 三个只读/改写 API因为每个超步都落了一份 checkpoint、且每份都记着"父存档 id"(D35 见过),它们天然串成一条时间线。时间旅行就三件事:
get_state(cfg)——看当前状态(读最新档);
get_state_history(cfg)——看全部历史(把链倒序读出来);
update_state(cfg, values)——在某个历史点写一份新档,可作为分叉起点。
然后把某个历史 checkpoint 的 config(带它的 checkpoint_id)传给 invoke,图就从那个点重新往下跑
类比:checkpoint 链 = 游戏的存档列表;get_state_history = 打开存档列表;update_state = 读档后改角色属性再另存;带着那个档 invoke = 从那个存档继续游戏。
L02

StateSnapshot:一份存档读出来的样子

时间旅行的 API 都返回 StateSnapshot,定义在 types.py:643-659

# types.py:643
class StateSnapshot(NamedTuple):
    values: dict[str, Any] | Any        # 这一步的状态值
    next: tuple[str, ...]               # 下一步该跑哪些节点
    config: RunnableConfig              # 这份快照的坐标(含 checkpoint_id)
    metadata: CheckpointMetadata | None # 元数据(step/source/...)
    created_at: str | None              # 创建时间戳
    parent_config: RunnableConfig | None# 父快照坐标(往回走用)
    tasks: tuple[PregelTask, ...]       # 本步的任务(含中断/错误信息)
values那一刻的状态(各通道当前值)。你最常看的。
next下一步要跑的节点名。空元组 () = 图已结束;非空(如 ('n',))= 停在这、下一步是 n。判断"是否停住"就看它。
config这份快照的"门牌号"——里面的 checkpoint_id 就是它在时间线上的定位。把它传给 invoke 就能从这里重开。
parent_config父快照的门牌。一路顺着 parent 往回,就是整条历史(D35 讲的父子链在这兑现)。
tasks本步的任务列表,含 interrupts 字段——所以你能从快照里看出"这步是被哪个中断卡住的"。
💡 本质:StateSnapshot 是 checkpoint 的"用户友好视图"底层 checkpoint 是给引擎用的(含版本号、blob 引用等),StateSnapshot 是翻译给人看的:把通道值拼好(values)、算出下一步(next)、带上时间线坐标(config/parent_config)。它由 _prepare_state_snapshot() 从 CheckpointTuple 组装而来。
数据结构:StateSnapshot 的 7 个字段 values这一步的状态值(你最常看) next下一步跑哪些节点(空=结束,非空=停住) config本快照门牌(含 checkpoint_id,重开用) parent_config父快照门牌 → 顺它往回 = 整条历史 metadata/created_atstep/source/时间戳 tasks本步任务(含 interrupts → 看出卡在哪个中断)
图注:next 判"停没停",config/parent_config 是时间线坐标,tasks 带中断信息。
L03

get_state:读"最新那一份"

main.py:1392 起,先看它怎么拿 checkpointer、怎么处理子图:

# main.py:1392
def get_state(self, config, *, subgraphs=False) -> StateSnapshot:
    """Get the current state of the graph."""
    checkpointer = ensure_config(config)[CONF].get(
        CONFIG_KEY_CHECKPOINTER, self.checkpointer)
    if isinstance(checkpointer, BaseCheckpointSaver):
        checkpointer = self._apply_checkpointer_allowlist(checkpointer)
    if not checkpointer:
        raise ValueError("No checkpointer set")        # ① 没存档 = 没历史可看
    if (checkpoint_ns := config[CONF].get(CONFIG_KEY_CHECKPOINT_NS, "")) \
            and CONFIG_KEY_CHECKPOINTER not in config[CONF]:
        recast = recast_checkpoint_ns(checkpoint_ns)    # ② 子图:剥掉 task_id
        for _, pregel in self.get_subgraphs(namespace=recast, recurse=True):
            return pregel.get_state(...)                #    转发给对应子图
        else:
            raise ValueError(f"Subgraph {recast} not found")
    ...  # 否则读本图最新 checkpoint,组装成 StateSnapshot 返回
if not checkpointer: raise时间旅行的前提和中断一样:必须有 checkpointer。没存档,无历史可谈。
checkpoint_ns 非空 → 转发子图如果 config 指向某子图命名空间,get_state递归转发给那个子图的 Pregel 实例——所以你能查看子图内部的状态,不只是顶层。
recast_checkpoint_ns把命名空间里的 task_id 后缀剥掉("sub:abc123"→"sub"),才能按逻辑名找到子图。
返回 StateSnapshot最终读最新 checkpoint、拼成 L02 的快照返回。
🍼 一句话get_state = "现在什么状态、下一步跑啥"。它读的是最新档,是时间旅行里的"现在时刻"。
L04

get_state_history:把整条时间线读出来

main.py:1480-1531,核心是调 checkpointer.list()

# main.py:1480
def get_state_history(self, config, *, filter=None, before=None, limit=None):
    """Get the history of the state of the graph."""
    config = ensure_config(config)
    checkpointer = config[CONF].get(CONFIG_KEY_CHECKPOINTER, self.checkpointer)
    ...
    if not checkpointer:
        raise ValueError("No checkpointer set")
    # (子图转发逻辑,同 get_state) ...
    config = merge_configs(self.config, config, {...})
    # eagerly consume list() to avoid holding up the db cursor
    for checkpoint_tuple in list(                       # ① 一次性取出,别占着游标
        checkpointer.list(config, before=before, limit=limit, filter=filter)
    ):
        yield self._prepare_state_snapshot(             # ② 每份 → StateSnapshot
            checkpoint_tuple.config, checkpoint_tuple)
checkpointer.list(...)直接用 D34 学的 list 接口——它按时间倒序(新→旧)返回该 thread 的所有 checkpoint。时间旅行不过是 list 的一个上层封装。
before / limit / filter翻页与筛选:before 只看某个 checkpoint 之前的、limit 限条数、filter 按 metadata 过滤。适合历史很长时分页看。
list(checkpointer.list(...))注释点明:先一次性 list() 消费完再逐个 yield,避免生成器边迭代边占着数据库游标(长时间持有 DB cursor 是坑)。
_prepare_state_snapshot把每个 CheckpointTuple 翻译成 StateSnapshot(L02)。于是你 for s in app.get_state_history(cfg) 拿到一串快照。
💡 设计取舍①:为什么 get_state_history 只是 list() 的薄封装,而不自建历史存储?朴素想法:单独维护一份"历史列表"专供时间旅行。问题:那就有两份真相(执行落的 checkpoint + 历史列表),要同步、会不一致。源码做法:历史就是 checkpoint 链本身——执行时为了持久化/恢复本来就一份份存了,时间旅行只是"换个方向(倒序)读同一批数据"。零额外存储、绝不会和真实执行状态脱节。这就是为什么开了 checkpointer,时间旅行几乎"免费"送你。代价:历史的粒度 = checkpoint 的粒度(每超步一份),你没法看到超步内部的中间态——但那正是合理的观察边界。
L05

update_state:像"某个节点"那样写一份新档

update_statemain.py:2515,异步版 aupdate_state main.py:2528)不是"覆盖旧档",而是造一份新 checkpoint,模拟"某个节点执行并写入"。看它执行更新的核心 main.py:2413-2467

# main.py:2413
for as_node, values, provided_task_id in valid_updates:
    writers = self.nodes[as_node].flat_writers       # ① 拿"这个节点"的写入器
    if not writers:
        raise InvalidUpdateError(f"Node {as_node} has no writers")
    writes: deque[tuple[str, Any]] = deque()
    task = PregelTaskWrites((), as_node, writes, [INTERRUPT])
    task_id = provided_task_id or (... uuid5(UUID(checkpoint["id"]), INTERRUPT))
    run = RunnableSequence(*writers) if len(writers) > 1 else writers[0]
    await run.ainvoke(values, patch_config(config, ...  # ② 用你的 values 跑写入器
        configurable={ CONFIG_KEY_SEND: writes.extend, ... }))
# main.py:2468
apply_writes(checkpoint, channels, run_tasks,          # ③ 把写 apply 到通道
             checkpointer.get_next_version, self.trigger_to_nodes)
# main.py:2485
checkpoint = create_checkpoint(checkpoint, channels, step + 1, ...)  # ④ 造新档
next_config = await checkpointer.aput(checkpoint_config, checkpoint, ...) # ⑤ 存
writers = nodes[as_node].flat_writers关键:update_state 不是硬塞值,而是走 as_node 这个节点的写入器——于是你的 values 会像该节点的真实输出一样,经过 reducer(如 add_messages 的合并)处理。
run.ainvoke(values, ...)把你给的 values 喂给写入器执行,产出的写进 writes。等于"假装 as_node 刚跑完、输出了 values"。
apply_writes(...)D21 学的:把这批写合并进通道、涨版本号。此刻状态真的变了。
create_checkpoint + aputstep+1 的新 checkpoint 并存盘。所以 update_state 会在时间线上新增一份档,而不是改旧的。
返回 next_config返回新档的坐标——拿它 invoke,就从这个"人改过的点"继续。
💡 本质:改历史 = 追加一份"署名某节点"的新档时间线是只增不改(append-only)的。update_state 从不篡改旧 checkpoint,而是在你选的那个点之上再叠一份新档,署名 as_node。这既保住了"历史可回溯不被抹掉",又实现了"从某点分叉出新未来"。这和 Git 的理念一模一样:你不改历史提交,而是在某个 commit 上新建分支。
L06

as_node 的推断与"从哪叉出去"

你可以显式指定 as_node="某节点";不指定时,框架会main.py:2381-2397):

# main.py:2381
elif as_node is None:
    last_seen_by_node = sorted(
        (v, n)
        for n, seen in checkpoint["versions_seen"].items()
        if n in self.nodes
        for v in seen.values())
    # if two nodes updated the state at the same time, it's ambiguous
    if last_seen_by_node:
        if len(last_seen_by_node) == 1:
            as_node = last_seen_by_node[0][1]              # 只有一个候选,就是它
        elif last_seen_by_node[-1][0] != last_seen_by_node[-2][0]:
            as_node = last_seen_by_node[-1][1]             # 最后动的那个节点
if as_node is None:
    raise InvalidUpdateError("Ambiguous update, specify as_node")  # 猜不出→报错
versions_seen 找"最后动手的节点"借 D43 见过的 versions_seen 账本,找"看到的版本最高"的节点——通常就是最近执行的那个,默认让新更新"接在它后面"。
len==1 → 唯一只有一个节点动过,无歧义,直接用它。
[-1][0]!=[-2][0] → 有明确最后者最后两个节点版本不同 → 有明确的"最新节点",用它。
否则 raise Ambiguous并列第一时拒绝猜:两个节点同时更新、分不清谁最后,宁可报错让你显式指定 as_node,也不乱猜写错。
⚠️ 边界:并行节点后 update_state 不指定 as_node 会报错,不是 bug如果你的图有并行分支(两个节点在同一超步都写了状态),它们的 versions_seen 版本并列最高,框架故意InvalidUpdateError("Ambiguous update, specify as_node")。这不是缺陷,是安全设计:update_state 要"署名某个节点",署错名会让状态更新经过错误的 reducer、或从错误的节点位置分叉,后果隐蔽难查。与其猜错,不如逼你写清 as_node="a"。反模式:在并行图里偷懒不写 as_node。
控制流:时间旅行——回到过去、修改、分叉 c0 c1 c2 c3 c4 选中 c2 (回到过去) c2' update_state 造新档 c2' 从 c2' invoke → 新分支未来 c3/c4 是旧未来(虚线),未被删除,只是不再是当前线
图注:get_state_history 读出 c0..c4;选 c2、update_state 叠出 c2';从 c2' 继续 = 新分支。旧档不删(append-only)。
L07

子图的 replay:ReplayState + 今日小结

回到过去重跑时,一个棘手问题:图里如果有子图,子图也有自己的 checkpoint。重放到过去某点时,子图该读哪份存档?这归 ReplayState 管(_replay.py:14-73):

# _replay.py:14
class ReplayState:
    """Tracks which subgraphs have already loaded their pre-replay checkpoint."""
    __slots__ = ("checkpoint_id", "_visited_ns")
    def __init__(self, checkpoint_id: str) -> None:
        self.checkpoint_id = checkpoint_id      # 重放的"目标时间点"
        self._visited_ns: set[str] = set()      # 记哪些子图已加载过
    # _replay.py:52
    def get_checkpoint(self, checkpoint_ns, checkpointer, checkpoint_config):
        if self._is_first_visit(checkpoint_ns):             # ① 该子图第一次被访问
            for saved in checkpointer.list(
                checkpoint_config,
                before={"configurable": {"checkpoint_id": self.checkpoint_id}},
                limit=1):
                return saved                                 #    取"重放点之前"的档
            return None
        return checkpointer.get_tuple(checkpoint_config)     # ② 之后按正常最新档
_is_first_visit(ns)判断某子图是不是重放中第一次跑到_replay.py:34,会剥掉 task_id 后缀识别"同一个逻辑子图")。
第一次 → list(before=目标点)第一次访问时,取"重放目标时间点之前"的那份子图存档——保证子图也回到过去的正确状态,而不是用它现在的最新态。
之后 → get_tuple 正常同一子图在循环里被再次访问(如 for 循环里的子图),就用正常最新档——因为那些是重放过程中新产生的,该用新的。
共享一个实例类注释说明:一次父执行里所有派生 config 共享同一个 ReplayState(按引用传),才能全局记住"哪些子图已回到过去了"。
💡 设计取舍②:为什么子图重放要区分"第一次 vs 之后",不能一律取旧档?因为子图可能在循环里被多次调用。假设父图重放到过去,子图第一次跑当然要恢复"过去那份"存档(否则时间线不一致)。但如果子图在一个 for 循环里,第 2、3 次调用是重放过程中刚刚新产生的执行——它们的正确起点是"上一轮刚存的新档",而不是"远古那份旧档"。若一律取旧档,循环里每轮都从远古状态起跑,结果全错。_visited_ns 用一个集合记住"这个子图已经完成了它的一次性回溯",之后切回正常加载。用一个 set 的极小状态,精确处理了"回溯只发生一次、后续走正常"的微妙语义。

👶 小白:从过去某点重跑,前面已经花过的 token / 调过的 API 会重复消耗吗?

👨‍🏫 从你选的那个 checkpoint 往后的节点会真实重跑(重新调模型、花 token);但那个 checkpoint 之前的节点不会——它们的结果已经在你选中的存档里了,直接读回,不重算。所以时间旅行省的正是"重跑前半段"的钱。这也呼应明天 D46 的幂等:已完成的写会被复用而非重做。反过来注意:你选的点往后如果有副作用节点,重跑会再执行,这点和 D42 的幂等提醒一致。

🧠 今天你应该能回答

  • 时间旅行的三个 API 各干什么?(get_state 看当前 / get_state_history 翻历史 / update_state 改并分叉)
  • StateSnapshot 里判断"是否停住"看哪个字段?(next,空=结束、非空=停在此)
  • get_state_history 底层是什么?(checkpointer.list 的薄封装,倒序读同一批 checkpoint)
  • update_state 是覆盖旧档还是新增?(新增一份署名 as_node 的新档,append-only)
  • 为什么 update_state 要走节点的 writers?(让值经过 reducer 合并,像该节点真实输出)
  • 并行图不指定 as_node 为何报错?(版本并列、署名歧义,故意拒绝乱猜)
  • ReplayState 解决什么?(子图重放时"第一次回旧档、之后用新档")

✋ 10 分钟动手

# 1. 读三个 API 源码
sed -n '1480,1531p' libs/langgraph/langgraph/pregel/main.py     # get_state_history
sed -n '2413,2506p' libs/langgraph/langgraph/pregel/main.py     # update_state 核心
sed -n '14,73p'     libs/langgraph/langgraph/_internal/_replay.py
# 2. 亲手时间旅行
python - <<'PY'
from langgraph.graph import StateGraph, START, END
from langgraph.checkpoint.memory import InMemorySaver
from typing_extensions import TypedDict
class S(TypedDict):
    x: int
g=StateGraph(S)
g.add_node("a", lambda s:{"x":s["x"]+1})
g.add_node("b", lambda s:{"x":s["x"]+10})
g.add_edge(START,"a"); g.add_edge("a","b"); g.add_edge("b",END)
app=g.compile(checkpointer=InMemorySaver())
cfg={"configurable":{"thread_id":"t"}}
app.invoke({"x":0}, cfg)
hist=list(app.get_state_history(cfg))
print("历史档数:", len(hist))                       # 若干份(倒序)
for h in hist: print("  step", h.metadata["step"], "x=", h.values.get("x"), "next=", h.next)
# 挑一个 a 刚跑完的历史点,改 x,再从那分叉
target=[h for h in hist if h.next==("b",)][0]
app.update_state(target.config, {"x": 100}, as_node="a")
print("从改后的点继续:", app.invoke(None, {"configurable":{"thread_id":"t"}}))
PY
明天预告 · Day 45:时间旅行/中断都依赖"每步都落盘"。但落盘有成本。durability 三档 "sync"/"async"/"exit" 决定"什么时候写、写得多勤"——性能与安全的权衡。看 _loop.pydo_checkpoint_checkpointer_put_after_previous 怎么按档位决策。
← Day 43 静态断点 Day 45 · durability 模式 →