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"]
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 里的 database 或 run_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_name 或 model 里取真实模型名,写入 per-model accumulator。subagent 和 middleware 通过额外的 token records 并回父 run,最后一起落到 RunStore.update_run_completion()。因此 /threads/{thread_id}/token-usage 的 by_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 的持久化是多层契约,不是一个总开关。每一层都要问清楚“它存的是什么、谁消费、重启后应该恢复到什么程度”。