DeerFlow Runtime:流式桥接、事件总线与状态持久化

Lead Agent 和中间件在执行时,会持续产生消息增量、工具调用和状态变更。前端需要打字机式输出和「正在搜索」一类中间态;关掉页面再回来,对话还要在;服务器重启后也不能整段失忆;调试和审计还要留住中间事件。这些若写进 agent 本体,会把业务逻辑淹没。DeerFlow 单独抽出 Runtime,把实时推送、状态持久化、事件存档做成管道层:Agent 只负责产生内容,管道负责播出去和记下来。

五个组件怎么分工

Runtime 核心组件

StreamBridgeruntime/stream_bridge/base.py)用生产者-消费者解耦 agent 与前端。publish(run_id, event, data) 入队,subscribe 异步迭代读出。默认 MemoryStreamBridge:每个 run 维护事件列表;约 15 秒无事件发心跳,避免反向代理掐断长连接;支持 Last-Event-ID 断线续读;缓冲有界(默认约 256),超出丢最旧,防止内存涨死。agent 生产速度和前端消费速度不再互相拖累。

RunManagerruntime/runs/manager.py)管 run 生命周期。RunRecordrun_idthread_idstatus(pending / running / success / error 等)、on_disconnect、后台 asyncio.Task。主索引按 run_id,辅索引按 thread_id,外加一把 asyncio.Lockcreate_or_reject 在同一把锁里检查同 thread 的 inflight 再创建,中间不 await,避免 TOCTOU 竞态导致同 thread 双 run;cancel 可保留胶片或回滚;cleanup / shutdown 负责收尾。shutdown 给 inflight 设 abort_event 并 cancel task,最多等约 5 秒,超时标 interrupted,避免硬杀损坏正在写入的 checkpoint。

RunJournalruntime/journal.py)继承 LangChain BaseCallbackHandler,把 on_llm_endon_tool_end 等回调标准化成 RunEvent。按 tag 把 token 分到 lead_agentsubagent:xxxmiddleware:xxx,便于看清费用落在主智能体、子代理还是摘要一类中间件。内部缓冲批量写库(阈值约 20 条),进度上报节流(约 5 秒),长时间 run 也能更新用量又不打爆数据库。

Checkpointerruntime/checkpointer/async_provider.py)经 make_checkpointer 按配置返回实现:memory 适合开发测试(重启丢),sqlite 适合单机,postgres 适合多实例并发。存的是 Agent 完整状态快照,用于断点续跑和回滚。

RunEventStore 存细粒度事件(消息、工具调用),按 thread_id + run_id 组织,带递增 seq,支持按 message 类目分页。和 Checkpointer、RunManager store 的分工可以对照成:胶片级状态、弹幕式事件流、场记级元数据(status、token 汇总)。写功能时先问自己「我要续跑、回放还是查汇总」,再选对存储,避免把事件日志塞进 checkpoint,或反过来指望 event store 恢复完整图状态。

用户侧大致对应这些能力:文字像打字一样冒出来,中间穿插工具调用提示,关掉网页再回来还能续聊。缺了 StreamBridge,前端只能轮询或干等整段结果;缺了 Checkpointer,进程一挂对话就断;缺了 EventStore,排障只能猜「模型到底调了几次工具」。管道抽出来之后,Lead Agent 的中间件链可以专心做决策相关的横切,不必同时当直播服务器。

一次 run 的生命周期

Gateway 收到发消息请求 → RunManager.create_or_reject → 后台 run_agentruntime/runs/worker.py)。Worker 侧关键步骤:若有 event_store 则建 RunJournal(可挂 progress_reporter 回写进度);set_status(running);执行前 aget_tuple 拍快照记下 pre_run_checkpoint_id 与深拷贝状态,供用户要求回滚;把 Journal 追加进 config["callbacks"];然后流式执行:

1
2
3
4
async for chunk in agent.astream(graph_input, config=runnable_config, stream_mode=lg_modes):
if record.abort_event.is_set():
break
await bridge.publish(run_id, _lg_mode_to_sse_event(mode), serialize(chunk, mode=mode))

finallyjournal.flush() 冲掉缓冲、update_run_completion 写 token 与末条消息等汇总、bridge.publish_end 通知前端结束。Gateway 的 sse_consumer 读请求头 Last-Event-ID,从断点继续 subscribe,刷新页面也不丢中间事件。

内存桥内部是 _RunStreamevents + asyncio.Condition + ended + start_offset。publish 追加事件(id 形如时间戳-序号)、必要时裁剪并 notify_all;subscribe 用 offset 推进,结束 yield END_SENTINEL,等待超时则 yield 心跳。断线重连时前端带上上次事件 id,桥从该点之后继续读,切换标签页或刷新不会从零丢失中间态。这是 SSE 单向推送加上可续读缓冲的组合,实现简单,却覆盖了大多数对话流场景。

再看并发安全:create_or_reject 必须原子完成「检查冲突 + 创建记录」。如果先检查再创建,中间插入 await,两个请求可能同时通过检查,同 thread 就会出现两个 running。把检查和创建放进同一把锁、且临界区内不做 IO,才能堵住这类竞态。优雅关闭同理:先发取消信号,留出短暂收尾窗口,再标记中断,比直接杀进程安全。

Agent 生产与前端消费互不阻塞;重启靠 Checkpointer,审计靠 EventStore,run 查询靠 Manager store。读代码建议从 worker.pyrun_agent 顺着 Journal → Bridge → Checkpointer 往下跟;Gateway 侧再补一眼 sse_consumer,就能看清「开拍 → 挂监听 → 流式推送 → 收尾」整条路径。

DeerFlow Runtime:流式桥接、事件总线与状态持久化

https://simonsu.net/2026/09/15/deerflow-03-runtime-stream-bridge/

Author

simonisacoder

Posted on

2026-09-15

Licensed under

Comments