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

Command(resume=):把答案喂回去,图从原地复活

昨天看透了"停":节点抛 GraphInterrupt,现场落盘。今天看"续"——你用 Command(resume="yes") 再 stream 一次,图怎么就知道"这个 yes 是回答刚才那个 interrupt 的"?答案藏在一个叫 便签本(PregelScratchpad)的小对象里,和一个把恢复值"发牌"到便签本的函数 _scratchpad()

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

痛点:恢复不是"新一次调用",是"接着上次"

🤔 痛点图停在了"确认删库?"。现在我要说"yes"。但如果我像第一次那样 app.stream({"x":0}, cfg),那不就是重新跑一遍吗?框架怎么区分"我要重开"和"我要接着上次那个中断答 yes"?而且这个 "yes" 要精准送到那一个 interrupt 手里,不能塞错。
💡 本质:恢复 = 用 Command 携带答案 + 复用同一 thread你不传新的初始状态,而是传一个 信封app.stream(Command(resume="yes"), cfg)cfg 里的 thread_id 没变 → 框架去 checkpointer 把昨天冻结的现场读回来;Command.resume="yes" → 这就是"要喂给中断的答案"。框架把 "yes" 存成一条特殊的 RESUME 写,然后重跑那个节点。节点里的 interrupt() 这次一查便签本"有答案",直接返回 "yes",代码继续往下——就像从没停过。
类比:昨天的存档是"游戏退到桌面"。今天的 Command(resume=) 是"双击图标、点继续",而不是"新建游戏"。thread_id 就是那个存档槽位。
L02

Command 对象:一个信封装四种意图

Command 定义在 types.py:759-808,是个能装多种"命令"的数据类:

# types.py:759
@dataclass(...)
class Command(Generic[N], ToolOutputMixin):
    graph: str | None = None       # 发给哪张图(None=当前,PARENT=父图)
    update: Any | None = None      # 顺便更新一下状态
    resume: dict[str, Any] | Any | None = None   # ★ 今天的主角:恢复值
    goto: Send | Sequence[Send | N] | N = ()     # 顺便跳到某节点(D15 讲过)

    PARENT: ClassVar[Literal["__parent__"]] = "__parent__"   # types.py:808
resume本日核心。文档(types.py:768-772)说它两种形态:① 单个值——喂给"下一个待恢复的中断";② {中断id: 值} 映射——精确指定回答哪个中断(用昨天的 Interrupt.id)。
update恢复的同时还能改状态。比如人不仅点了"通过",还顺手改了金额。
goto恢复的同时指定下一步跳哪去(Command 本来就是 D15 的跨节点跳转对象,resume 只是它的一个字段)。
graph=PARENT子图里可以 Command(graph=Command.PARENT, ...) 把命令发给父图——人在环跨层时用得上。
💡 本质:resume 只是 Command 的一个字段,不是新概念LangGraph 没为"恢复"造新 API。Command 本来就是"给图下达指令"的通用信封(更新状态、跳转、发消息),resume 只是往这个信封里多加了一格"给中断的答案"。所以你能一次性 Command(resume="yes", update={"note":"已复核"}, goto="audit") ——恢复 + 改状态 + 指定去向,一个信封全办了。
L03

PregelScratchpad:任务级的"便签本"

昨天 interrupt() 里那个 scratchpad 究竟是什么?定义在 _scratchpad.py:8-19,短到一屏:

# _scratchpad.py:8
@dataclasses.dataclass(**_DC_KWARGS)
class PregelScratchpad:
    step: int
    stop: int
    # call
    call_counter: Callable[[], int]
    # interrupt
    interrupt_counter: Callable[[], int]     # ① 节点里 interrupt 的编号器
    get_null_resume: Callable[[bool], Any]   # ② 取"通用恢复值"的函数
    resume: list[Any]                        # ③ 这个任务的恢复值列表
    # subgraph
    subgraph_counter: Callable[[], int]
一个任务一本关键理解:便签本是任务级的(一次节点执行一本),不是全图共享。所以昨天文档说"恢复值 scoped 到具体任务、不跨任务共享"。
interrupt_counter昨天 L02 见过:给节点里的 interrupt 从 0 开始编号,让第 idx 个 interrupt 对上第 idx 个恢复值。
resume: list已认领的恢复值列表Command(resume=) 里针对这个任务的值会被放进来。interrupt 按 idx 从这里取。
get_null_resume取"通用恢复值"(没指定 id 的那个单值)的函数。谁先叫谁拿走(消费型)。L05 细看。
注意它全是 Callable(函数)而不是直接存数字/列表——比如 interrupt_counter 是个可调用的计数器。为什么?见 L07 的 LazyAtomicCounter:并发下 +=1 不安全,用函数封装原子操作。
L04

_scratchpad():从存档里"发牌"恢复值

便签本由谁填?由 _scratchpad()_algo.py:1280)在每次准备任务时构造。它从 pending_writes(还记得吗,昨天 RESUME 就存这里)里找恢复值:

# _algo.py:1280
def _scratchpad(parent_scratchpad, pending_writes, task_id,
                namespace_hash, resume_map, step, stop):
    if len(pending_writes) > 0:
        # find global resume value  —— 找"通用恢复值"(谁都能领的那份)
        for w in pending_writes:
            if w[0] == NULL_TASK_ID and w[1] == RESUME:   # ① 特殊 task_id + RESUME
                null_resume_write = w
                break
        else:
            null_resume_write = None
        # find task-specific resume value —— 找"点名给本任务的"恢复值
        for w in pending_writes:
            if w[0] == task_id and w[1] == RESUME:        # ② 就是本 task_id 的 RESUME
                task_resume_write = w[2]
                if not isinstance(task_resume_write, list):
                    task_resume_write = [task_resume_write]
                break
        else:
            task_resume_write = []
        # find namespace and task-specific resume value  —— 按中断 id 精确映射
        if resume_map and namespace_hash in resume_map:   # ③ Command(resume={id:值})
            mapped_resume_write = resume_map[namespace_hash]
            task_resume_write.append(mapped_resume_write)
    else:
        null_resume_write = None
        task_resume_write = []
NULL_TASK_ID + RESUME分支①:Command(resume="单个值") 会被存成 (NULL_TASK_ID, RESUME, 值)——不点名,谁的 interrupt 先要谁拿。这就是 get_null_resume 要取的那份。
w[0] == task_id分支②:点名给这个任务的恢复值(多值场景),直接进 task_resume_write 列表,成为便签本的 resume
resume_map[namespace_hash]分支③:Command(resume={中断id: 值}) 时,按命名空间哈希(=昨天 Interrupt.id 的来源)精确匹配到本任务,追加进恢复列表。
else: []没有 pending_writes(全新运行)→ 恢复列表空。于是 interrupt 走"抛异常"分支。首尾呼应昨天 L02。
💡 设计取舍①:为什么恢复值要分"通用(NULL_TASK_ID)"和"点名(task_id)"两种?通用值(NULL_TASK_ID)服务最常见场景:只有一个中断在等,你 Command(resume="yes") 给个单值,框架不用你操心它该给谁——"谁在等谁拿"。点名值(task_id / 按 id 映射)服务复杂场景:并行分支同时中断了多个(比如 Send 扇出的多个任务各自 interrupt),这时一个笼统的单值无法区分,必须用 {中断id: 值} 精确投递。框架用同一套 pending_writes 存两类,靠 task_id 是不是 NULL_TASK_ID 区分——简单场景零负担,复杂场景有精度。这是"常见路径极简 + 高级路径可用"的经典分层。
L05

get_null_resume:通用恢复值"取一次就没了"

_scratchpad() 内部定义了闭包 get_null_resume_algo.py:1320-1331),交给便签本:

# _algo.py:1320
def get_null_resume(consume: bool = False) -> Any:
    if null_resume_write is None:
        if parent_scratchpad is not None:
            return parent_scratchpad.get_null_resume(consume)   # ① 本层没有?问父层
        return None
    if consume:                                                  # ② 消费模式
        try:
            pending_writes.remove(null_resume_write)            #    从 pending 里删掉!
            return null_resume_write[2]                          #    返回值
        except ValueError:
            return None
    return null_resume_write[2]                                  # ③ 只看不消费
null_resume_write is None本任务没找到通用恢复值 → 向父便签本递归查(子图场景:恢复值可能是发给父图的)。父也没有则 None。
if consume: pending_writes.remove(...)核心:consume=True 时,取值的同时把它从 pending_writes 删除。昨天 interrupt() 里正是 get_null_resume(True)——领了就消费掉。
为什么要删保证"一个通用恢复值只被一个 interrupt 领走"。领完就不在池子里了,下一个 interrupt 再查就是 None → 该停还得停。避免一个 "yes" 被多个中断重复消费。
return ..[2](不删)consume=False 时只窥视不拿走——某些地方要"看看有没有"但不真正认领。
⚠️ 边界:一个节点里有 2 个 interrupt,只给 1 个 resume 值会怎样?假设节点里 a=interrupt("Q1"); b=interrupt("Q2"),第一次跑停在 Q1。你 Command(resume="A1") 恢复 → 节点重跑:Q1 这次 get_null_resume(True) 拿到 "A1" 并消费删除,返回;执行到 Q2 时再查,通用恢复值已被消费掉了 → None → 再次抛 GraphInterrupt 停在 Q2。所以多中断节点需要多轮恢复(一次 resume 只喂一个通用值),或者一次用 resume={id1:A1, id2:A2} 映射式喂全。反模式:以为一个 Command(resume=单值) 能同时答完一个节点里的所有 interrupt。
L06

闭环:interrupt 重跑时如何"直接返回答案"

现在把便签本填好后,回到昨天的 interrupt() 前半段(types.py:912-925)——这次它走的是另一条路

# types.py:912
scratchpad = conf[CONFIG_KEY_SCRATCHPAD]     # 拿到 _scratchpad() 刚填好的便签本
idx = scratchpad.interrupt_counter()          # 本 interrupt 编号
# ① 便签本里已有点名恢复值(task-specific)
if scratchpad.resume:
    if idx < len(scratchpad.resume):
        conf[CONFIG_KEY_SEND]([(RESUME, scratchpad.resume)])
        return scratchpad.resume[idx]         # 按编号取答案,直接返回!
# ② 或者有通用恢复值
v = scratchpad.get_null_resume(True)          # 消费型取值
if v is not None:
    assert len(scratchpad.resume) == idx, (scratchpad.resume, idx)
    scratchpad.resume.append(v)               # 记进便签本(下次同 idx 可复用)
    conf[CONFIG_KEY_SEND]([(RESUME, scratchpad.resume)])
    return v                                   # 返回,节点继续往下跑
if scratchpad.resume: return resume[idx]点名恢复值路径:便签本里第 idx 项就是答案,取出返回。多个 interrupt 各按自己的 idx 拿各自的值。
v = get_null_resume(True)通用恢复值路径:消费掉一个通用值。
scratchpad.resume.append(v)妙笔:把领到的通用值追加进便签本的 resume 列表。这样如果本任务因为别的原因再重跑一次,同一个 idx 就能从 resume[idx] 稳定取到,不用再去消费池子——幂等(D46 主题)。
CONFIG_KEY_SEND([(RESUME,...)])两条路都会把 RESUME 写回去——记录"这个恢复值已被用掉",随现场落盘,避免重复消费。
💡 本质:同一段 interrupt 代码,恢复时"改道"了昨天无恢复值 → 走到最后 raise。今天便签本里有值 → 在 if 里就 return 了,压根到不了 raise。你的节点函数一字没改,行为却从"停"变成"继续"——切换开关就是便签本里那几个恢复值,而它们由 _scratchpad()Command.resume 转存的 RESUME 写填充。整个"停—存—喂—续"闭环到此合拢。
L07

LazyAtomicCounter + 今日小结

最后补一个细节:为什么 interrupt_counter 是个函数而不是 self.idx += 1?看 _algo.py:1426-1439

# _algo.py:1426
class LazyAtomicCounter:
    __slots__ = ("_counter",)
    def __init__(self) -> None:
        self._counter = None
    def __call__(self) -> int:
        if self._counter is None:
            with LAZY_ATOMIC_COUNTER_LOCK:          # 加锁,双重检查
                if self._counter is None:
                    self._counter = itertools.count(0).__next__   # 惰性建计数器
        return self._counter()                      # itertools.count 的 __next__ 是原子的
itertools.count(0).__next__用 C 实现的 itertools.count,它的 __next__ 在 CPython 里是原子操作,天然线程安全——比手写 x+=1(读-加-写三步,会撞车)稳。
惰性(lazy)创建只有第一次真正调用计数时才建计数器。大多数节点根本没 interrupt,省掉无谓开销。
_scratchpad 里赋值回看 _algo.py:1334-1344interrupt_counter=LazyAtomicCounter()call_countersubgraph_counter 全用它——interrupt/call/subgraph 各一个独立原子计数器。
💡 设计取舍②:为什么计数器用 itertools.count 而非普通整数自增?因为一个节点可能并发执行子任务(Send 扇出、并行工具调用),多个执行流可能同时给同一个便签本的 interrupt 编号。普通 self.idx += 1 是"读旧值→加一→写回"三步,两个线程交错会拿到相同编号(竞态),导致恢复值配错 interrupt。itertools.count().__next__ 把"取号+自增"压成一个 C 层原子操作,杜绝撞号。代价是多一层函数调用和一个锁(还用了惰性初始化把锁的开销摊薄)。这是"用极小的并发安全成本,换恢复值配对的绝对正确"。
控制流:Command(resume) 恢复的返程 你 stream( Command(resume="yes")) 读回现场 同 thread_id 存档 存 RESUME 写 (NULL_TASK_ID,RESUME,值) _scratchpad() 填便签本 节点重跑:interrupt() 查便签本 有值 → return "yes"(不再 raise)
图注:Command.resume → RESUME 写 → _scratchpad 填便签本 → interrupt 查到值直接 return。
数据结构:恢复值的三种投递方式 通用值 Command(resume="yes") 存: NULL_TASK_ID 谁先要谁拿(消费) get_null_resume 点名值(task) 存: 具体 task_id 进 resume 列表 按 idx 取 并行多任务 映射值(id) resume={id: 值} 按 ns 哈希匹配 最精确 resume_map
图注:从"随便谁拿"到"精确到中断 id",投递精度递增,对应 _scratchpad 的三个分支。

👶 小白:恢复时整个节点重跑,那节点前半段有副作用(比如已经写了数据库)怎么办?会写两遍吗?

👨‍🏫 好问题,这正是幂等问题(D46 专讲)。原则:interrupt 之前的代码,恢复时会再执行一遍。所以带副作用的操作(写库、发消息、扣款)要么放在 interrupt 之后,要么自己做幂等保护。LangGraph 的 put_writes 层对"通道写"做了去重(D35 L05 见过),但你节点里手写的外部副作用它管不着。这是人在环编程必须记住的铁律:把 interrupt 想象成一道"存档线",线之前的代码可能重放。

🧠 今天你应该能回答

  • 怎么恢复一个中断的图?(stream/invoke 传 Command(resume=值),复用同一 thread_id)
  • Command.resume 的两种形态?(单值给"下一个中断";{id:值} 精确映射)
  • PregelScratchpad 是什么级别的?(任务级便签本,记 interrupt 编号和恢复值)
  • _scratchpad() 从哪找恢复值?(pending_writes 里的 RESUME 写,分 NULL_TASK_ID 通用 / task_id 点名 / resume_map 映射)
  • get_null_resume(True) 的"消费"是什么意思?(取值同时从 pending 删除,防重复领)
  • 一个节点两个 interrupt 只给一个 resume 会怎样?(第二个再次中断,需多轮恢复)

✋ 10 分钟动手

# 1. 读便签本定义与发牌逻辑
sed -n '8,19p'    libs/langgraph/langgraph/_internal/_scratchpad.py
sed -n '1280,1344p' libs/langgraph/langgraph/pregel/_algo.py
# 2. 完整跑一次"停→续"
python - <<'PY'
from langgraph.graph import StateGraph, START
from langgraph.checkpoint.memory import InMemorySaver
from langgraph.types import interrupt, Command
from typing_extensions import TypedDict
class S(TypedDict):
    x: int
def node(s):
    ans = interrupt("确认?")
    return {"x": 1 if ans=="yes" else 0}
g=StateGraph(S); g.add_node("n",node); g.add_edge(START,"n")
app=g.compile(checkpointer=InMemorySaver())
cfg={"configurable":{"thread_id":"t1"}}
print("第一次:", list(app.stream({"x":0}, cfg)))        # 停,出 __interrupt__
print("恢复后:", list(app.stream(Command(resume="yes"), cfg)))  # 续,节点跑完
print("最终:", app.get_state(cfg).values)               # {'x': 1}
PY
明天预告 · Day 43:除了在节点里手写 interrupt()(动态断点),LangGraph 还能在 compile/stream 时用 interrupt_before=["n"]静态断点——不改节点代码就能"在某节点前/后停"。看 should_interrupt() 怎么靠 versions_seen 判断该不该停。
← Day 41 interrupt 原理 Day 43 · 静态断点 →