数据截至 (上游 commit 1f738cdeb7f5)
持久化、人在环路与时间旅行
30 秒导读: 一个长跑的多 agent 工作流,凭什么敢说自己"生产级"?靠三件事——崩了能接着跑(可恢复)、跑到一半能 停下等人拍板(人在环路)、能回到任意历史节点重来(时间旅行)。这三件事底层是同一个机制:在 Pregel 超步的每个边界,把工作流的完整状态拍成一张可序列化的快照(checkpoint)。本章讲这张快照里存了什么、怎么存、怎么恢复,以及"暂停等人"是怎么用同一套快照实现的。
本章紧接 03 章 Workflow 图引擎。超步(superstep)机制本身在那章讲透了,这里只用它的一个结论:超步边界 = 一个干净的、无并发在途的一致性时刻。快照就拍在这个时刻。
1. 这是什么(零基础也能懂)
一句话定义: 给工作流装上"存档/读档"能力——就像单机游戏的存档点,跑到一个安全点自动存盘,出事了从存档点读回来接着跑。
解决谁的什么问题? 三类生产痛点,对应三个卖点:
| 生产痛点 | 卖点 | 白话 |
|---|---|---|
| 工作流跑了 40 分钟,进程崩了 / 机器重启 | 可恢复(restartability) | 别从头再来,从最后一个存档点接着跑 |
| 流程走到一半需要人来批一下 / 补个信息 | 人在环路(human-in-the-loop) | 工作流停下、发出"我要问个问题"、挂起等人;人答完再继续 |
| 想调试"如果第 3 步换个输入会怎样" | 时间旅行(time travel) | 挑任意一张历史存档,从那儿重放 |
用起来什么样? 一个最小心智模型:构建工作流时给个存储后端,运行时框架自动在每个超步后存档;要恢复就把 checkpoint_id 交回去。
# 示意,非源码。演示三个卖点各自的入口调用
from agent_framework import WorkflowBuilder, FileCheckpointStorage
storage = FileCheckpointStorage("/var/checkpoints") # 落盘的存档柜
workflow = WorkflowBuilder(checkpoint_storage=storage).build() # 开启自动存档
# 卖点1 正常跑:每个超步边界自动存一张档
result = await workflow.run(message="分析这份合同")
# 卖点2 人在环路:跑到 request_info 会挂起,拿到待答请求
for req in result.get_request_info_events():
print(req.data) # "请人工确认:是否批准第 3 条?"
# 人答完,把答案按 request_id 交回去,工作流从挂起处继续
await workflow.run(responses={req.request_id: "批准"})
# 卖点3 时间旅行:挑一张历史档,从那儿重放
await workflow.run(checkpoint_id="某个历史 checkpoint 的 id")
一句话直觉: 把超步边界当"火车站台"——列车(工作流)只在站台完全停稳(无并发在途消息、状态已提交)时才允许拍照;这张照片信息完整,所以既能拿它重建列车(恢复),也能在照片这一刻往车厢里塞个乘客(人工输入)。
本节不出现底层细节。下面从全景开始逐层下钻。
2. 顶层全景(快照拍在哪、存了什么)
怎么读这张图: 从左到右是一个超步的生命周期;快照(★)总是拍在超步跑完、状态提交之后。
一个超步的生命周期(见 _runner.py:run_until_convergence)
┌──────────────────────────────────────────────────────────┐
│ ① 派发消息 ② 执行器并发跑 ③ 提交共享状态 ④ 拍快照★ │
│ drain_messages _run_iteration state.commit() create_ │
│ (超步边界) checkpoint │
└──────────────────────────── ──────────────────────────────┘
│
▼
┌───────────────────────────┐
│ WorkflowCheckpoint (一张档) │
│ ·在途消息 messages │
│ ·已提交状态 state │
│ ·待答人工请求 pending_... │
│ ·迭代数 iteration_count │
│ ·图指纹 graph_signature_hash │
└───────────────────────────┘
│ save()
▼
CheckpointStorage(存档柜:内存 / 文件)
关键时序:状态先提交,再拍快照。 _runner.py:163 先 self._state.commit()(把这一超步的写入固化),_runner.py:166 才 create_checkpoint_if_enabled()。所以快照里只有已提交状态,没有半途的脏写——这是"快照一致性"的根。
四个部件,一句话职责:
| 部件 | 干什么 | 在哪 |
|---|---|---|
WorkflowCheckpoint | 快照的数据结构:一张档里装什么 | _checkpoint.py:31 |
CheckpointStorage | 存档柜协议:save / load / list / get_latest | _checkpoint.py:129 |
State | 共享状态:pending 缓冲 + 超步边界 commit | _state.py:286 |
RequestInfoMixin / request_info | 人在环路:发请求挂起、按类型匹配 response_handler | _request_info_mixin.py:29、_workflow_context.py:403 |
三个卖点落到同一张快照:
┌──────────────┐
崩溃恢复 ─────►│ │ state + messages → 重建执行现 场
时间旅行 ─────►│ 一张快照 │ 挑任意历史 id → 从那超步重放
人在环路 ─────►│ │ pending_request_info_events → 挂起点
└──────────────┘
3. 检查点:快照里存了什么、怎么存
3.1 WorkflowCheckpoint——一张档的字段
先看数据结构,才知道"恢复"能恢复到什么程度。核心字段(_checkpoint.py:81-98):
| 字段 | 存的是 | 恢复时用来 |
|---|---|---|
messages | 超步之间在途未处理的消息(按源执行器分组) | 重建"下一超步该派发什么" |
state | 已提交的共享状态,含执行器自身状态(藏在保留键 _executor_state) | 重建全局变量与各执行器内部状态 |
pending_request_info_events | 尚未被回答的人工请求事件 | 恢复后仍知道"卡在等谁答话" |
iteration_count | 拍照时的超步序号 | 恢复后从这个序号接着数,不重头 |
graph_signature_hash | 工作流图拓扑的指纹(SHA-256) | 恢复前校验:图没变过才敢重放 |
previous_checkpoint_id | 上一张档的 id | 把历次快照串成链,形成可回溯的历史 |
一个刻意的设计:快照不绑定工作流实例。 文档字符串明说(_checkpoint.py:37-41):档只认"工作流定义"(靠 workflow_name + graph_signature_hash 识别),不记录是哪个运行实例产生的。好处: 同一个工作流定义的不同实例之间,快照可以互相共享、互相恢复。
快照成链 = 历史可回溯。 每张新档都用 previous_checkpoint_id 指向上一张(_runner.py:267 存完后 self._previous_checkpoint_id = checkpoint_id)。于是整段执行历史是一条单向链表——时间旅行就是"沿链挑一个节点跳回去"。
3.2 CheckpointStorage——存档柜协议
存档柜是一个 Protocol(_checkpoint.py:129,鸭子类型接口,任何实现了这些方法的类都算数)。六个方法:
| 方法 | 作用 | 定义 |
|---|---|---|
save(checkpoint) | 存一张 档,返回其 id | _checkpoint.py:132 |
load(checkpoint_id) | 按 id 取一张档 | _checkpoint.py:143 |
list_checkpoints(workflow_name) | 列出某工作流的所有档对象 | _checkpoint.py:157 |
delete(checkpoint_id) | 删一张档 | _checkpoint.py:168 |
get_latest(workflow_name) | 取最新一张档 | _checkpoint.py:179 |
list_checkpoint_ids(workflow_name) | 只列 id(轻量) | _checkpoint.py:190 |
框架自带两个实现:
InMemoryCheckpointStorage(_checkpoint.py:202)——给测试和开发用。 一个 dict 装档;save 时 copy.deepcopy(_checkpoint.py:211)防止外部后续修改污染已存的档。get_latest 靠比 timestamp 取最大(_checkpoint.py:240,max(..., key=lambda cp: datetime.fromisoformat(cp.timestamp)))。
FileCheckpointStorage(_checkpoint.py:249)——落盘持久化。 三个要点:
- 一张档 = 一个 JSON 文件,文件名是
{checkpoint_id}.json。JSON 结构人类可读,便于调试排查。 - 原子写(
_checkpoint.py:328_write_atomic): 先写.json.tmp,再os.replace(tmp, file)(_checkpoint.py:332)。os.replace在同一文件系统上是原子的——断电也不会留下半截损坏的档。 - 路径穿越防护(
_checkpoint.py:293_validate_file_path): 校验checkpoint_id解析出的路径确实落在存储目录内(is_relative_to),挡住有人用../../etc/xxx这类构造的 id 写到任意位置。
注意 get_latest 的代价差异: 文件版的 get_latest(_checkpoint.py:423)要先 list_checkpoints 把目录里所有档读出来反序列化再比时间戳,比内存版重得多。想只拿 id 用 list_checkpoint_ids(_checkpoint.py:439),它只 json.load 读顶层字段、不解码 pickle。
3.3 编码:JSON 骨架 + pickle/base64 填肉
要解决的小问题: 快照里的 state 可能装着任意 Python 对象(dataclass、自定义类、datetime),JSON 原生存不了。怎么既保持文件可读、又不丢对象保真度?
思路(_checkpoint_encoding.py 模块头 3-9 行):混合编码。 JSON 原生类型(str/int/float/bool/None)、以及 dict/list 这类容器——照原样递归写进 JSON,保持可读;其余一切(tuple、set、dataclass、自定义对象……)——pickle 序列化 + base64 编码成字符串,塞进 JSON 里一个带标记的小对象。
看 _encode(_checkpoint_encoding.py:307)的分支:
# 示意,非源码。重点看"能 JSON 就 JSON,不能就 pickle"的分流
def _encode(value):
if isinstance(value, (str, int, float, bool, type(None))):
return value # JSON 原生,直接过
if isinstance(value, dict):
return {str(k): _encode(v) for k, v in value.items()} # 递归
if isinstance(value, list):
return [_encode(x) for x in value] # 递归
return { # 其余:pickle + base64
"__pickled__": _pickle_to_base64(value),
"__type__": _type_to_key(type(value)), # 记下类型,解码时校验
}
解码有一道完整性检查。 _decode(_checkpoint_encoding.py:329)见到 __pickled__ + __type__ 标记就反序列化,然后 _verify_type(_checkpoint_encoding.py:362)比对"解出来的对象类型"是否等于"当初记下的类型"——不等就抛 WorkflowCheckpointException,提示档可能损坏或被篡改。注意这只是事后完整性检查,pickle.loads 那一刻代码早已执行(见下节安全模型)。
3.4 安全:RestrictedUnpickler 与"档是可信数据源"
pickle 是把双刃剑: 反序列化能执行任意代码。所以 _checkpoint_encoding.py 模块头(18-44 行)把安全模型讲得很硬:checkpoint 存储被当作【可信数据源】——绝不能把用户 HTTP 请求、消息体这类不可信输入喂给 decode_checkpoint_value;存档柜(文件系统 / Cosmos / Blob)必须做访问控制,当成数据库凭据一样看管。
纵深防御:受限反 pickle。 当传入 allowed_types 时,用 _RestrictedUnpickler(_checkpoint_encoding.py:157)。它重写 find_class(_checkpoint_encoding.py:204):只有类型键落在允许集内才放行,否则抛 UnpicklingError。允许集由四部分并起来:
| 来源 | 内容 | 依据 |
|---|---|---|
| 内置安全集 | 原语、datetime、uuid、Decimal、collections 等 | _BUILTIN_ALLOWED_TYPE_KEYS(_checkpoint_encoding.py:111) |
| 框架类型 | 所有 agent_framework. 开头的模块 | 前缀 _checkpoint_encoding.py:94 |
| OpenAI SDK 类型 | 所有 openai.types. 开头的模块 | 前缀 _checkpoint_encoding.py:97 |
| 调用方追加 | FileCheckpointStorage(allowed_checkpoint_types=[...]) 传入的 "模块:qualname" | _checkpoint.py:277 |
诚实的边界(模块头 21-25 行明说): 这个允许集是"减少攻击面的缓解",不是安全边界——某些必须放行的内置(如 getattr,用来重建枚举/具名元组)本身就有能力,拿不掉。真正的防线是"别让不可信数据进到这里"。
3.5 State:共享状态与超步边界提交
要解决的小问题: 一个超步里多个执行器并发跑,都要读写共享状态。怎么保证它们看到的是"这一超步开始时的一致快照",而不是彼此半途的脏写?
思路(_state.py:286 类文档):双缓冲 + 边界提交。 写不直接落到已提交状态,而是先进 pending 缓冲;读时先看 pending 再看 committed;直到超步边界由 Runner 调 commit() 一次性固化。
执行器 A ─set(k,v)─┐
执行器 B ─set(k,w)─┼──► _pending 缓冲(超步内)
│ │ 超步边界
│ ▼ Runner 调 state.commit() (_runner.py:163)
└──► _committed 已提交状态 ──► 进快照的就是这份
方法一览:
| 方法 | 行为 | 定义 |
|---|---|---|
set(k, v) | 写进 pending,不碰 committed | _state.py:310 |
get(k) | 先查 pending 再查 committed | _state.py:325 |
commit() | pending 全部固化进 committed,清空 pending | _state.py:370 |
discard() | 丢弃 pending,不提交 | _state.py:382 |
export_state() | 导出 committed 的副本(不含 pending) | _state.py:386 |
import_state(d) | 把字典合并进 committed | _state.py:393 |
两个细节:
- 删除靠哨兵。
delete(k)(_state.py:348)若键在 committed,就往 pending 塞一个_DeleteSentinel(_state.py:407)标记"提交时删掉";commit时见哨兵就pop。这样删除也遵守"边界才生效"的语义。 - 并发写:后写者胜。 同一超步内多个执行器写同一个键,都进同一个 pending 缓冲,
commit时最后一次写生效(_state.py:316-322文档)——与 .NET 版行为一致。
快照存的正是 export_state() 的结果。 拍档时 create_checkpoint 调 state.export_state()(_runner_context.py:436)——所以快照里永远是干净的已提交状态。
3.6 恢复:从档重放(时间旅行/可重启)
恢复入口是 Runner.restore_from_checkpoint(_runner.py:278)。 五步,顺序讲究:
restore_from_checkpoint(checkpoint_id)
│
① load 档 load_checkpoint / 外部 storage.load (_runner.py:302-305)
│
② 图指纹校验 ★ graph_signature_hash 不匹配就拒绝 (_runner.py:316)
│ "图变过了,请用原始工作流再恢复"
│
③ 重建共享状态 state.clear() → import_state(档.state) (_runner.py:327-328)
│ 先清后并,避免旧运行的残留键泄漏
│
④ 重建执行器状态 _restore_executor_states() (_runner.py:330)
│ 从 _executor_state 键逐个 on_checkpoint_restore
│
⑤ 应用到上下文 ctx.apply_checkpoint(档) (_runner.py:332)
│ 恢复在途消息 + 待答人工请求(并重发事件)
│
└─► _mark_resumed(档):iteration 跳回档.iteration_count (_runner.py:400-407)
第②步是时间旅行的安全 阀。 图拓扑指纹 graph_signature_hash 是把"起始执行器 + 各执行器签名 + 边组"规范化后做 SHA-256(_workflow.py:1112 _hash_graph_signature)。恢复前比对(_runner.py:313):图改过就拒绝重放,因为老档的消息/状态可能对不上新拓扑。这让"回到历史节点"是安全的,不会把状态灌进一个已经变形的图。
恢复后不重置迭代计数。 _mark_resumed(_runner.py:448)把 self._iteration 设回 checkpoint.iteration_count(_runner.py:454)——从档的超步序号接着数。reset_iteration_count 的文档(_runner.py:95-104)专门强调:从响应或检查点恢复时,迭代计数通常不重置。
"超步 0"也拍档。 入口档(捕获初始输入、任何执行器开跑之前的状态)由 Workflow 在把输入种入 start 执行器的内部自环后拍摄(_workflow.py:678-680 调 create_checkpoint_if_enabled);run_until_convergence 的 docstring 明说 runner 的职责收窄为「每个超步末拍档」,因为只有 Workflow 知道这次运行是全新输入、从档恢复还是从响应恢复(_runner.py:105-110)。这样连"刚开跑"这一刻也有存档点可回溯。
存档失败不拖垮工作流。 create_checkpoint_if_enabled(_runner.py:240)把整段 save 包在 try/except 里,失败只 logger.warning(_runner.py:268-276)、不抛——下一张成功的档会认上一张成功的档做父。存档是"尽力而为"的旁路,不阻断主流程。