← 返回主线
09

第 09 站 · 已发布

持久化 · store · checkpointer

源码锚点 · @7e7f041
backend/app/gateway/deps.py @7e7f041 backend/app/gateway/services.py @7e7f041 backend/app/gateway/routers/threads.py @7e7f041 backend/app/gateway/routers/thread_runs.py @7e7f041 backend/packages/harness/deerflow/config/database_config.py @7e7f041 backend/packages/harness/deerflow/config/checkpointer_config.py @7e7f041 backend/packages/harness/deerflow/config/run_events_config.py @7e7f041 backend/packages/harness/deerflow/runtime/checkpointer/async_provider.py @7e7f041 backend/packages/harness/deerflow/runtime/store/async_provider.py @7e7f041 backend/packages/harness/deerflow/runtime/runs/manager.py @7e7f041 backend/packages/harness/deerflow/runtime/runs/worker.py @7e7f041 backend/packages/harness/deerflow/runtime/journal.py @7e7f041 backend/packages/harness/deerflow/runtime/events/store/base.py @7e7f041 backend/packages/harness/deerflow/runtime/events/store/db.py @7e7f041 backend/packages/harness/deerflow/persistence/engine.py @7e7f041 backend/packages/harness/deerflow/persistence/bootstrap.py @7e7f041 backend/packages/harness/deerflow/persistence/migrations/_helpers.py @7e7f041 backend/packages/harness/deerflow/persistence/migrations/versions/0001_baseline.py @7e7f041 backend/packages/harness/deerflow/persistence/migrations/versions/0002_runs_token_usage.py @7e7f041 backend/packages/harness/deerflow/persistence/run/model.py @7e7f041 backend/packages/harness/deerflow/persistence/run/sql.py @7e7f041 backend/packages/harness/deerflow/persistence/thread_meta/sql.py @7e7f041 backend/packages/harness/deerflow/runtime/serialization.py @7e7f041

读 DeerFlow 的持久化代码,不能只问“哪些东西被存起来”。在 agent 系统里,持久化更像一层事实边界:agent 跨轮次对话、调用工具、写文件、产生事件和副作用之后,系统必须知道哪些状态可以恢复,哪些记录可以查询,哪些事实可以审计。

所以这一站的目标是建立一条边界,不是背表名:

持久化给 agent 运行建立可恢复、可查询、可审计的事实层,范围远大于“存聊天记录”。Checkpointer 负责可恢复的图状态;Store 负责长期键值状态;RunStore、RunEventStore、ThreadMetaStore 负责 DeerFlow Gateway 面向查询、展示、恢复和审计的运行时记录。

flowchart TD
REQ["HTTP run request"] --> RM["RunManager / RunStore"]
REQ --> TM["ThreadMetaStore"]
RM --> WORKER["run_agent()"]
WORKER --> CP["LangGraph Checkpointer"]
WORKER --> STORE["LangGraph Store"]
WORKER --> JOURNAL["RunJournal"]
JOURNAL --> EV["RunEventStore"]
CP --> STATE["/state / history / wait"]
TM --> LIST["/threads/search"]
RM --> RUNS["/runs / token-usage"]
EV --> MSG["/messages / events"]
一次 run 会触碰多类存储。它们不是同一个层次,也不应该互相替代。

Gateway 启动时一次性装配运行时

持久化组件不是每个请求临时创建的。Gateway 启动时,langgraph_runtime() deps.py::langgraph_runtime )会用启动时的 AppConfig 快照创建一组长期对象:

make_stream_bridge(config)
init_engine_from_config(config.database)
make_checkpointer(config)
make_store(config)
RunRepository or MemoryRunStore
ThreadMetaRepository or MemoryThreadMetaStore
make_run_event_store(config.run_events)
RunManager(store=run_store)

这些对象会放到 app.state 上。请求进来后,路由不会重新建数据库连接池,也不会重新建 checkpointer;它们会从 app.state 取现成对象。

这解释了一个重要现象:DeerFlow 支持部分配置热加载,但数据库、checkpointer、store、run event store 这类基础设施是启动期绑定的。你修改 config.yaml 里的 databaserun_events,通常需要重启 Gateway 才会切换后端。

应用表 schema 也在启动期 bootstrap

持久化还包括数据库表结构怎么创建、怎么升级,不能只看“对象放进哪个 store”。应用表由 init_engine_from_config() 进入 bootstrap_schema() bootstrap.py::bootstrap_schema ),不会只做无条件 Base.metadata.create_all()

空库
  create_all()
  alembic stamp head

旧库:已有 DeerFlow 表,但没有 alembic_version
  补齐 baseline 表
  stamp 0001_baseline
  upgrade head

已版本化库:有 alembic_version
  alembic upgrade head

这套 hybrid bootstrap 的取舍是:空库仍然用 SQLAlchemy metadata 建表,避免手写一份容易漂移的 baseline;从 baseline 之后的变化交给 Alembic migration。0001_baseline 0001_baseline.py )主要是 Alembic 链路的根和旧库 stamp 目标,生产路径里新库通常是 create_all + stamp head,不会真的逐条执行 baseline DDL。

0002_runs_token_usage 0002_runs_token_usage.py )展示了 post-baseline 变更怎么做:给 runs 表补 token_usage_by_model。它通过 _helpers.safe_add_column() _helpers.py::safe_add_column )实现幂等添加:列已存在就不重复加,同时检查 nullable、default、类型是否和模型定义漂移。

并发上也分层处理:

Postgres
  用 advisory lock 串行化 bootstrap。

SQLite
  同进程用 asyncio.Lock。
  跨进程依赖 SQLite 文件锁 + 30 秒 busy_timeout,是 best-effort,不是强分布式锁。

migration 本身
  列级变更用幂等 helper,作为重试和并发边界的最后兜底。

这套启动期 bootstrap 解决的是一个很实际的问题:旧数据库升级时,不应该因为少一列或少一张 baseline 表而在第一个请求上才 500。schema bootstrap 是启动期责任,和下面讲的 checkpointer/store 分工是两条不同但相关的线。

五类存储分别存什么

先把几个名字摆清楚。

Checkpointer
  LangGraph 的 checkpoint 存储。

Store
  LangGraph 的长期 KV 存储,挂在 Runtime.store 上。

RunStore
  DeerFlow 的 run 元数据存储。

RunEventStore
  DeerFlow 的 run 事件流存储。

ThreadMetaStore
  DeerFlow 的线程列表和归属信息存储。

Checkpointer:图状态和恢复点

checkpointer 来自 LangGraph。异步 Gateway 使用 make_checkpointer() async_provider.py::make_checkpointer )创建它。

它保存的是 thread 的 checkpoint,里面最核心的是:

channel_values
  messages
  title
  thread_data
  sandbox
  其他 ThreadState 字段

metadata
  step
  source
  writes
  parents
  created_at / updated_at

tasks / pending_writes
  interrupt、错误、下一步任务等执行信息

它的主要消费者是:

/threads/{thread_id}/state
/threads/{thread_id}/history
/runs/wait 返回最终 channel_values
cancel(action=rollback)
run_agent finally 读取 title,同步到 thread_meta

所以,决定“下一轮对话能不能接上”的,是 checkpointer。消息历史、sandbox_id、thread_data 这些 ThreadState 字段能否恢复,也取决于 checkpointer 是否持久化。

Store:LangGraph 的长期 KV

Store 由 make_store() 创建( async_provider.py::make_store ),会传给:

Runtime(context=runtime_ctx, store=store)
agent.store = store

它不是 checkpoint,不负责保存每一步图状态。它是 LangGraph runtime 暴露给 agent 的长期键值存储能力。

DeerFlow 里还有一个兼容用法:当没有 SQL session factory 时,MemoryThreadMetaStore 会把 thread metadata 存到 Store 的 ("threads",) namespace 下。因此在 memory 模式里,线程列表可能借用 Store;但在 sqlite/postgres 模式里,线程列表主要走 threads_meta 表。

RunStore:一次 run 的外部生命周期

RunStore 的 SQL 实现是 RunRepository sql.py::RunRepository ),由 RunManager 使用。

它保存的是 run 这一层的外部记录:

run_id
thread_id
assistant_id
user_id
status
model_name
multitask_strategy
metadata / kwargs
error
created_at / updated_at
token usage
token_usage_by_model
message_count
first_human_message
last_ai_message

这些字段服务于:

列出某个 thread 的 runs
查询 run 当前状态
取消 run 时检查状态
聚合 token usage
重启后识别 orphaned inflight runs

它不保存完整对话,也不负责恢复 LangGraph state。

runs 表里的 token_usage_by_model model.py::RunRow )用来修正按模型统计的口径,不是重复字段:一次 run 里可能既有 lead agent 模型调用,也有 subagent 和 middleware 模型调用,不能只按 run 的 model_name 粗略归类。token_usage_by_model 保存每个真实模型的输入、输出和总 token。

RunEventStore:消息和执行事件

RunEventStore 是事件流接口( base.py::RunEventStore )。SQL 实现在 DbRunEventStore db.py::DbRunEventStore )。

一条 event 大致是:

thread_id
run_id
event_type
category
content
metadata
seq
created_at

category 很关键:

message
  前端可展示的消息。

trace / middleware / outputs / error
  调试、审计、生命周期记录。

seq 是同一个 thread 内递增的序号,用来分页和排序。SQL 后端通过唯一约束保证 (thread_id, seq) 不重复。

RunEventStore 的消费者是:

/threads/{thread_id}/messages
/threads/{thread_id}/runs/{run_id}/messages
/threads/{thread_id}/runs/{run_id}/events

它适合查询“发生过什么”,但不是图状态的权威来源。恢复图状态仍然看 checkpointer。

ThreadMetaStore:线程列表的外壳信息

ThreadMetaStore 的 SQL 实现是 ThreadMetaRepository sql.py::ThreadMetaRepository )。

它保存的是:

thread_id
assistant_id
user_id
display_name
status
metadata
created_at / updated_at

这不是对话正文,也不是 checkpoint。它是线程列表、权限检查和搜索所需的信息。

例如 /threads/search 不会扫描 checkpointer 的所有历史,它直接查 ThreadMetaStore。run 结束后,worker 会从最新 checkpoint 里读 title,再同步到 threads_meta.display_name,这样列表页能快速显示标题。

一次 run 如何写入这些存储

start_run() services.py::start_run )和 run_agent() worker.py::run_agent )串起了主流程。

1. start_run()
   RunManager.create_or_reject()
   -> 创建 RunRecord
   -> 写 RunStore pending 行

2. start_run()
   upsert thread_meta
   -> 不存在则 create
   -> 存在则 status=running

3. run_agent()
   RunManager.set_status(running)
   -> 更新 RunStore

4. run_agent()
   捕获 pre-run checkpoint
   -> 后续 rollback 使用

5. agent.astream(...)
   LangGraph 执行
   -> checkpointer 持续写 checkpoint
   -> StreamBridge 推 SSE

6. RunJournal callback
   -> 把 LLM / tool / middleware 回调转成 run events
   -> 聚合 token 和消息摘要

7. run_agent() 结束
   -> success / error / interrupted
   -> 更新 RunStore status

8. finally
   -> journal.flush()
   -> update_run_completion()
   -> checkpoint.title 同步到 thread_meta.display_name
   -> thread_meta.status 更新为 idle/error/interrupted

9. bridge.publish_end()
   -> SSE 结束
   -> bridge cleanup

这条链路说明三件事:

SSE
  当前连接看到的在线事件。

RunEventStore
  事后可查询的事件流。

Checkpointer
  下一轮继续执行所需的图状态。

它们可能都包含“消息”,但语义不同。

RunJournal:把回调转成可查询事件

RunJournal journal.py::RunJournal )是 LangChain callback handler。它不运行 agent,也不决定状态转移;它观察 LLM、tool、chain 回调,然后写入 RunEventStore。

它同时做了两个实用聚合:

事件流
  llm.human.input
  llm.ai.response
  llm.tool.result
  middleware:...
  run.start / run.end / run.error

run 摘要
  total_input_tokens
  total_output_tokens
  total_tokens
  llm_call_count
  lead_agent_tokens
  subagent_tokens
  middleware_tokens
  token_usage_by_model
  message_count
  first_human_message
  last_ai_message

为什么摘要要写回 RunStore?因为列表页和 token 统计不能每次都扫描完整事件流。RunEventStore 保存细节,RunStore 保存常用摘要,这是典型的读路径优化。

这里的统计口径要按真实模型来理解:RunJournal 不只累计总 token,也会从每次 LLM 响应的 response_metadata.model_namemodel 里取真实模型名,写入 per-model accumulator。subagent 和 middleware 通过额外的 token records 并回父 run,最后一起落到 RunStore.update_run_completion()。因此 /threads/{thread_id}/token-usageby_model 统计的是“这个 thread 下各模型实际产生了多少 token”,不是只看“这个 run 的主模型用了多少”。

RunRepository.aggregate_tokens_by_thread() sql.py::RunRepository.aggregate_tokens_by_thread )会优先从每行的 token_usage_by_model JSON 汇总;旧行没有这个字段内容时,才退回 model_name + total_tokens 的 legacy 口径。注意 by_model[model].runs 不是互斥计数:同一个 run 如果调用过多个模型,会同时计入多个模型的 runs。

拆开存

  • checkpoint 专注恢复图状态。
  • run row 适合列表和统计。
  • event stream 适合消息分页和审计。

必须同步

  • run 结束时要同步多个 surface。
  • 某个同步失败时,列表和 state 可能短暂不一致。

LangGraph state 和 DeerFlow runtime state

判断“重启后能不能恢复”,先要分清两类状态。

会进入 checkpointer 的,是 LangGraph state,即 graph 的 channel values:

messages
title
thread_data
sandbox
artifacts
promoted tools
其他 ThreadState 字段

这些状态用于:

下一轮对话
state/history API
rollback
工具运行时上下文

不会作为 ThreadState 进入 checkpoint,但 DeerFlow 仍然要管理的,是 Gateway runtime state:

RunRecord
  当前进程里的 asyncio.Task、abort_event、status。

RunRow
  可查询的 run 元数据和 token 摘要。

RunEventRow
  消息、执行事件和审计信息。

ThreadMetaRow
  线程列表、标题、归属用户和状态。

StreamBridge buffer
  当前 run 的 SSE 订阅和短期缓冲。

app.state singletons
  checkpointer、store、event_store、run_manager 等启动期对象。

这些状态用于取消、查询、展示、鉴权和在线传输。它们不应该混进 ThreadState。

重启后哪些东西应该恢复

不同后端下恢复能力不同。可以按 surface 看:

checkpointer 持久化
  能恢复 thread state、history、messages、title、thread_data、sandbox_id 等图状态。

run_store 持久化
  能恢复 run 列表、run status、token 摘要。

run_event_store 持久化
  db/jsonl 能恢复消息和事件;memory 会丢。

thread_meta 持久化
  SQL backend 能恢复线程列表和 owner。
  MemoryThreadMetaStore 是否能恢复,取决于它背后的 LangGraph Store 是否持久化。

StreamBridge
  进程内传输层,重启后不能恢复。

asyncio.Task / abort_event
  进程内对象,重启后不能恢复。

因此,一个持久化 run 的 row 还在,不代表那个 run 还在执行。执行中的 Python task 已经随进程消失了。

RunManager.reconcile_orphaned_inflight_runs() manager.py::RunManager.reconcile_orphaned_inflight_runs )处理的正是这个问题:

数据库里还有 pending/running 的 run
但当前进程没有对应 task
=> 标记为 error

这一步不会恢复执行;它把不确定状态变成明确失败,避免 UI 永远显示 running。

cancel 和 rollback 不是一回事

取消 run 时,RunManager.cancel() 只能取消当前进程里的 task:

record.abort_event.set()
record.task.cancel()
record.status = interrupted

这说明取消是进程本地能力。另一个进程即使能在 RunStore 里看到这条 run,也不能直接取消这个进程里的 asyncio.Task

rollback 处理的是图状态。run 开始前,run_agent() 会保存 pre-run checkpoint snapshot。如果用户选择 action=rollback,worker 会把 checkpoint 恢复到 run 开始前的状态。

rollback 这里不会回滚数据库事务;它会在 checkpointer 里写入一个新的恢复点:

pre-run checkpoint
  -> 复制内容
  -> 换新的 checkpoint id / ts
  -> 写回 checkpointer

所以 rollback 回滚的是 LangGraph state,不会撤销所有外部副作用。比如已经写到 sandbox 文件系统或外部服务里的副作用,需要其他机制处理。

API 输出还有一层 serialization

Checkpoint 里保存的是 LangChain / LangGraph 对象,不能直接返回给前端。runtime/serialization.py serialization.py::serialize_channel_values_for_api )负责把它们转成 JSON 结构。

这里还有一个实际保护:ViewImageMiddleware 可能把 base64 图片块放进 hidden message 供模型使用。REST 的 state/history 返回时,serialize_channel_values_for_api() 会把 hide_from_ui 消息里的 data: 图片块移除,避免把很大的内部上下文直接发给前端。

这说明 serialization 同时承担格式转换和 API 边界保护。

不要把 schema migration 和 runtime store 一致性混在一起

按当前代码看,应用表的创建和升级已经有明确路径:persistence.bootstrap、Alembic baseline、0002_runs_token_usage、幂等 column helper 都在 Gateway 启动期路径里。

但 schema bootstrap 只回答“表结构怎么创建和升级”。下面的 split-brain 风险回答的是另一件事:checkpointer、store、run store、run event store、thread metadata 是否都消费同一份持久化配置。它们是两类问题,不要混在一起。

当前的 split-brain 风险

代码里已经有新的统一配置:

database:
  backend: sqlite | postgres | memory

DatabaseConfig database_config.py::DatabaseConfig )的意图是:同一个 backend 同时约束 checkpointer 和 DeerFlow 应用数据。

但当前启动路径还没完全统一:

make_checkpointer(config)
  支持 legacy checkpointer
  也支持新的 database

init_engine_from_config(config.database)
  支持新的 database

make_store(config)
  只读取 legacy checkpointer
  没有 checkpointer 时退回 InMemoryStore

这会造成一个真实风险:

只配置 database.sqlite,没有配置 legacy checkpointer

checkpointer
  可能是 SQLite 持久化。

RunRepository / ThreadMetaRepository
  是 SQL 持久化。

LangGraph Store
  仍可能是 InMemoryStore。

这类情况可以称为 split-brain:同一次 Gateway 里,不同状态面按不同配置来源决定持久化后端。有些状态能恢复,有些状态不能恢复,但用户可能以为“我开了 database.sqlite,所以都持久化了”。

这个问题对只依赖 checkpoint 的对话恢复不一定立刻致命,但对 LangGraph Store 的长期 KV、memory、未来扩展和 embedded/sync 路径都会造成语义不一致。

最后一层判断:不要把“持久化”说成一个承诺

持久化不是一句“开了数据库就恢复一切”。更准确的说法应该按对象分开:

对话状态能否恢复?
  看 checkpointer。

线程列表能否恢复?
  看 ThreadMetaStore。

run 列表和 token 统计能否恢复?
  看 RunStore。

消息和审计事件能否恢复?
  看 RunEventStore。

正在执行的任务能否恢复?
  不能。asyncio.Task 是进程内对象。

SSE 连接能否恢复?
  不能。StreamBridge 是在线传输层。

这才是这章最重要的心智模型:DeerFlow 的持久化是多层契约,不是一个总开关。每一层都要问清楚“它存的是什么、谁消费、重启后应该恢复到什么程度”。