第 2 章
一次任务的完整执行
从提交、模型与工具循环,一直走到文件交付
这一章只做一件事:沿着运行时间向前走。浏览器先通过 thread.submit 提交任务,Gateway 创建 RunRecord,worker 调用 run_agent;随后 Agent 在 agent.astream 中反复使用模型和工具。中途需要哪个机制,我们就在那个位置停下来解释。
先用四个阶段定位,再在正文中沿真实函数逐层进入;旁支只挂在对应调用点上。
- 01Submit
浏览器提交 thread 输入
- 02Manage
Gateway 建立运行身份
- 03Execute
worker 驱动 Agent 图
- 04Deliver
文件登记并结束运行
submit → run → loop → deliver先看完整流程:主链只有七个调用点
先只读流程图中间一列:thread.submit 提交任务,start_run 创建并安排 Run,run_agent 准备运行环境,agent_factory 装配 Agent,agent.astream 驱动模型与工具,present_files 登记交付物,worker 最后写入 RunStatus.success。这七个调用点组成主链。
流程图侧面的 Middleware、Memory、Sub-agent、Sandbox、Checkpoint 和事件后端都不是下一步。它们只在主链经过某个调用点时提供能力。例如 Sandbox 挂在文件工具调用上,Checkpoint 挂在 run_agent 的图状态上。把侧枝暂时遮住,主链仍然能够从上走到下。
下面每个标题前都有一枚位置标记。“主链”表示控制流向下一层移动;“主链内部”表示仍在当前函数里;“挂在”表示临时解释支撑机制,读完会回到标出的函数。
中间是必须依次经过的主链;右侧短句是只在某一步介入的机制。先顺着中间读到底,再回头看分支。
- 01Browser
thread.submitinput + config + context
- 附件已上传,只传路径与 metadata
- 02Gateway
start_runRunRecord + background task
- SSE 只订阅这条 Run
- StreamBridge 传递实时事件
- 03Worker
run_agentRunnableConfig + RunContext
- Checkpoint 提供图状态
- RunEventStore 接收持久事件
- 04Assembly
agent_factorymodel + tools + middleware + state
- 模式决定模型与工具
- Middleware 包住模型和工具调用
- 05Agent loop
agent.astreamAIMessage ↔ ToolMessage
- Memory 在模型调用前召回
- Sub-agent 通过 task 工具委派
- Sandbox 在文件工具首次使用时取得
- 06Delivery
present_filesThreadState.artifacts
- 只有 outputs 中明确选中的文件会交付
- 07Settlement
RunStatus.success终态 + receipt + 可恢复结果
- Checkpoint 留状态
- RunEventStore 留历史
- 页面可重新查询 Artifact
thread.submitthread.submit 把任务交给 Gateway
主链从前端 hooks.ts 中的 thread.submit 开始。教学命令把 thread、mode、附件和目标压在一行,真实网页则把它们拆成结构化数据:第一项是图输入,第二项是 threadId、config 与 context。
这里最重要的区别是“用户目标”和“运行选项”。用户目标告诉 Agent 要做什么;mode、thinking_enabled、subagent_enabled 等选项告诉 Harness 怎样装配这次执行。两者会在 worker 中重新相遇,但不会混成一段无法校验的提示词。
await thread.submit(
{
messages: buildThreadSubmitMessages({
text,
additionalKwargs: options?.additionalKwargs,
additionalInputMessages: options?.additionalInputMessages,
filesForSubmit,
}),
},
{
threadId: threadId,
// No streamSubgraphs: subtask progress arrives via root-namespace
// custom events, while subgraph frames would leak a delegated
// subagent's values/messages into the thread view (#4399).
streamResumable: true,源码中第一个对象包含 messages,第二个对象包含 threadId、streamResumable、config 和 context。上传动作此前已经完成,这里传的是 filesForSubmit 生成的文件信息,不是把 PDF 二进制再次塞进模型消息。
LangGraph SDK 随后请求 Gateway 的 /runs/stream 路由。到这里仅仅完成提交:模型还没有被调用,sales-review.md 也不存在。
stream_run → start_runstream_run 进入统一的 start_run
thread.submit 的流式请求落到 stream_run。这个路由只取出 StreamBridge 与 RunManager,然后调用 start_run;普通创建接口和等待接口也复用同一个 start_run,所以 DeerFlow 没有为 SSE 另写一套 Agent 执行逻辑。
start_run 返回 record 后,路由把 sse_consumer 包进 StreamingResponse。浏览器由此开始观察这条 Run,但响应流本身不负责执行模型。
bridge = get_stream_bridge(request)
run_mgr = get_run_manager(request)
record = await start_run(body, thread_id, request)
return StreamingResponse(
sse_consumer(bridge, record, request, run_mgr),
media_type="text/event-stream",
headers={
"Cache-Control": "no-cache",
"Connection": "keep-alive",
"X-Accel-Buffering": "no",
# LangGraph Platform includes run metadata in this header.
# The SDK uses a greedy regex to extract the run id from this path,
# so it must point at the canonical run resource without extra suffixes.
"Content-Location": f"/api/threads/{thread_id}/runs/{record.run_id}",
},
)源码中的关键顺序是 record = await start_run(...) 在前,StreamingResponse 在后。也就是说,系统先建立运行身份,再开放事件流。下一节进入 start_run 内部,看 RunRecord 和后台 worker 怎样出现。
create_or_reject → create_taskstart_run 先保存 RunRecord,再安排 worker
start_run 先规范化输入与配置,再调用 run_mgr.create_or_reject。这个调用把 thread_id、input、metadata、模型和多任务策略写进可查询的 RunRecord;若同一 thread 的并发策略不允许新任务,它也在这里拒绝。
只有 durable admission 成功以后,代码才构造 worker,并用 asyncio.create_task 挂到 record.task。顺序不能倒过来:否则后台已经开始跑,系统却还没有一条能取消、查询或接管的运行记录。
record = await run_mgr.create_or_reject(
thread_id,
body.assistant_id,
on_disconnect=disconnect,
metadata=body.metadata or {},
# Persist a secret-redacted copy of the config: the run record is
# written to runs.kwargs_json and echoed by the run API, so a
# request-scoped secret (#3861) must not ride along. The live
# config built above keeps the secrets for the actual run.
kwargs={"input": body.input, "config": redact_config_secrets(body.config)},
multitask_strategy=body.multitask_strategy,
model_name=model_name,
user_id=owner_user_id,
idempotency_key=idempotency_key,
)
if record.idempotency_reused:
return record
worker = run_after_metadata(record)
try:
# No await is allowed between durable admission and task
# attachment. Metadata setup runs inside the attached
# worker so a pending cancellation can bypass stalled
# thread-store IO and still reach run_agent's startup
# barrier / stream finalization.
record.task = asyncio.create_task(worker)源码把两个动作放在同一临界区:create_or_reject 产生管理事实,create_task 产生进程内执行任务。start_run 随即返回 record,不等待报告完成。
worker 实际执行的是 run_after_metadata(record)。它完成 thread metadata 准备后,下一跳就是 await run_agent(...)。
sse_consumerSSE 只是观察 start_run 创建的同一条 Run
主链现在停在 start_run 已经返回 RunRecord 的位置。sse_consumer 从 bridge 读取这条 record 的事件并发送给浏览器,不负责调用模型,也不拥有 Agent 的生命周期。
因此,“前端没收到流”与“后台没执行”不是同一个故障。前者看 SSE 和 StreamBridge,后者看 RunRecord 与 worker。连接断开后是否取消由 on_disconnect 策略决定,不由 TCP 连接替系统猜测。
这条旁支到此结束。控制流回到 record.task 中的 worker,继续进入 run_agent。
run_agentworker 调用 run_agent 进入 Harness
start_run 中的后台协程最终执行 await run_agent(...)。这行是 Gateway 与 Harness 的明确边界:左边仍在管理 HTTP 请求和 RunRecord,右边开始准备 LangGraph 运行环境。
调用参数也把职责说清了。record 提供本次运行身份,graph_input 是规范化后的模型输入,config 是运行配置,ctx 是 RunContext,提供 checkpointer、event_store、thread_store 等基础设施,bridge 负责实时事件。
await run_agent(
bridge,
run_mgr,
record,
ctx=run_ctx,
agent_factory=agent_factory,
graph_input=graph_input,
config=config,
stream_modes=stream_modes,
stream_subgraphs=body.stream_subgraphs,
interrupt_before=body.interrupt_before,
interrupt_after=body.interrupt_after,
)run_agent 不接收浏览器组件,也不返回一篇最终文字。它托管 Agent 图的整个生命周期:取得运行所有权、装配 Agent、消费图输出、处理取消、计算交付结果,最后写入终态。
下面先留在 run_agent 内部,看输入配置怎样变成真正可执行的 Agent 图。
RunnableConfigrun_agent 先准备本次运行配置
前端提交的模式、模型选择与功能开关,需要进入 Agent 工厂;LangGraph 自己也有 configurable、context、metadata 等运行字段。DeerFlow 不让各个模块随意从请求对象取值,而是在装配入口先把兼容字段合并。
这一步发生在真正创建模型和 Middleware 之前。结果决定是否启用深度思考、Sub-agent、特定模型,以及当前用户身份怎样传给后续运行组件。
def _get_runtime_config(config: RunnableConfig) -> dict:
"""Merge legacy configurable options with LangGraph runtime context."""
cfg = dict(config.get("configurable", {}) or {})
context = config.get("context", {}) or {}
if isinstance(context, dict):
cfg.update(context)
return cfg_get_runtime_config 先复制 configurable,再用 context 中的键覆盖。这个很短的函数说明运行配置存在明确合并点;排查某个开关失效时,应先确认它是否进入 cfg,而不是直接修改模型提示词。
合并后的配置只决定“这次怎么装配”,不会自动成为 ThreadState 的业务字段。下一步,lead agent 用它挑选模型、工具和 Middleware。
agent_factoryrun_agent 调用 agent_factory 构建图
RunnableConfig 准备好以后,worker 把它放进 agent_factory_kwargs,并调用前面由 Gateway 解析出的 agent_factory。默认情况下,这个工厂就是 assemble_lead_agent。
这里不要把 agent_factory 理解成又一次模型调用。它做的是组装:根据本次配置选择模型、工具、Middleware、系统提示和 ThreadState schema,返回可以运行的 LangGraph 图。
agent_factory_kwargs: dict[str, Any] = {"config": initial_runnable_config}
if ctx.app_config is not None and _agent_factory_supports_app_config(agent_factory):
agent_factory_kwargs["app_config"] = ctx.app_config
from deerflow.extensions import bind_agent_build_extensions
with bind_agent_build_extensions(extensions):
agent = _agent_graph(agent_factory(**agent_factory_kwargs))源码中的 agent = _agent_graph(agent_factory(...)) 同时说明两件事:工厂允许返回带描述信息的 LeadAgentAssembly,worker 只取其中真正可执行的 graph;具体装配细节仍封装在 lead_agent/agent.py。
下一节继续进入 agent_factory,看模型和工具怎样接到同一张图。
create_agentagent_factory 把模型、工具和 Middleware 接到同一张图
所谓“创建 Agent”,实际是把四类东西接到同一张图:聊天模型、可调用工具、系统提示和 Middleware。ThreadState 则规定图运行期间共享的数据形状。
对报告任务来说,工具集合可能包含网页搜索、文件读写、shell 命令、present_files 和 task;配置决定其中哪些真正可用。Middleware 再为这些调用补上摘要、记忆、进度与错误处理。
graph = create_agent(
model=create_chat_model(name=model_name, thinking_enabled=thinking_enabled, reasoning_effort=reasoning_effort, app_config=resolved_app_config, attach_tracing=False, model_overrides=agent_model_overrides),
tools=final_tools,
middleware=normalize_middleware_state_schemas(middlewares, mode),
system_prompt=system_prompt,
state_schema=get_thread_state_schema(mode),
)create_agent 接收 model、final_tools、normalize_middleware_state_schemas 的结果、system_prompt 和 get_thread_state_schema(mode)。源码把它们并排传入,比任何“智能体平台”口号都更能说明边界:模型只是装配件之一。
Agent 创建完成后,控制权回到前面的 agent.astream。图读取当前状态,把消息交给模型;模型若发出 tool call,图执行对应工具并把 ToolMessage 放回状态,然后开始下一轮。
agent.astreamagent.astream 驱动模型与工具循环
Agent 图装配完成后,run_agent 调用 agent.astream(input_payload, config, stream_mode)。从这一行开始,模型与工具才真正工作。它不是等待一个最终字符串,而是持续产出消息、状态更新和自定义事件。
每个 chunk 返回 worker 后,worker 检查 abort_event,再通过 bridge.publish(run_id, ...) 发布。取消和事件发送包在图循环外面,因此不会改变模型下一步选择哪个工具。
模型选择工具,工具结果写回状态;没有新的 tool call 时,图才结束本轮 Run。
single_mode = lg_modes[0]
async for chunk in agent.astream(input_payload, config=stream_config, stream_mode=single_mode):
if record.abort_event.is_set():
logger.info("Run %s abort requested — stopping", run_id)
break
llm_error_fallback_message = llm_error_fallback_message or _extract_llm_error_fallback_message(chunk, pre_existing_message_ids)
sse_event = _lg_mode_to_sse_event(single_mode)
await bridge.publish(run_id, sse_event, serialize(chunk, mode=single_mode))
if single_mode == "custom":
await subagent_events.add(chunk)图内部则反复执行同一个闭环:模型读取 ThreadState,产生普通回答或 tool call;工具执行后把 ToolMessage 或 Command.update 写回状态;模型再读新状态。没有新的 tool call 时,图才准备结束。
报告任务的第一批工具通常包括 read_file、搜索和写文件。下面的 Sandbox、Middleware、Memory 与 Sub-agent 都挂在这段 agent.astream 循环上;读完任何一条旁支,都要回到这里继续下一轮。
read_fileread_file 第一次访问文件时取得 Sandbox
Sandbox Middleware 不一定在 Run 开始时立刻创建环境。第一次 read_file、写文件或执行命令需要文件系统时,ensure_sandbox_initialized 才会按 thread_id 和 user_id 获取环境。纯文本任务因此不必提前承担 Sandbox 成本。
按 thread 获取还有另一个效果:同一 thread 的后续 Run 可以继续访问此前工作文件。第二次让系统修改报告时,不必把所有中间材料重新塞进消息。
附件位于可读输入区,中间文件留在工作区,准备交付的结果进入 outputs。
def _acquire_sandbox(self, thread_id: str, *, user_id: str) -> str:
provider = get_sandbox_provider()
sandbox_id = provider.acquire(thread_id, user_id=user_id)
logger.info(f"Acquiring sandbox {sandbox_id}")
return sandbox_id_acquire_sandbox 从 provider.acquire 取得 sandbox_id,并记录当前环境。具体 provider 可以映射到本地目录、容器或远程 Sandbox;上层工具只依赖统一的文件与命令接口。
环境准备好以后,read_file 读取 sales-notes.pdf,把访谈内容作为工具结果交回 ThreadState。模型发现“客户更关注价格”还只是内部判断,于是继续调用搜索工具核对行业数据;稍后写报告时也复用这份工作区。主线现在回到 agent.astream 的循环。
wrap_*_callMiddleware 包在 agent.astream 的调用边界上
主线现在停在一次工具调用。工具真正执行前后,多个 Middleware 会像函数包装器一样套在外面。读取后再写、进度上报和错误处理看似互不相干,但包装顺序会改变谁能看到原始异常、谁能看到规范化结果。
源码在构造 tail 时特意留下顺序说明:ToolProgressMiddleware 要位于 ToolErrorHandlingMiddleware 外面,才能观察后者已经写入 metadata 的结果。
外层记录完整过程,内层规范失败;返回时结果按相反方向穿回去。
tool_progress_config = app_config.tool_progress
if tool_progress_config.enabled:
from deerflow.agents.middlewares.tool_progress_middleware import ToolProgressMiddleware
tail.append(ToolProgressMiddleware.from_config(tool_progress_config))
tail.append(ToolErrorHandlingMiddleware(app_config=app_config))ReadBeforeWriteMiddleware 先进入列表,ToolProgressMiddleware 随后追加,ToolErrorHandlingMiddleware 最后成为内层。调用向内经过这些 wrapper,返回时反向穿出;所以“列表先后”不只是视觉排序。
对于读者,记住一个判断方法就够了:Middleware 决定某一步怎样被执行和记录,模型仍决定是否调用这一步。插叙结束,工具结果回到 ThreadState,agent.astream 继续下一轮。
SummarizationMiddleware消息过长时,摘要发生在下一次模型调用前
报告任务经过多轮搜索与文件读取后,messages 会越来越长。模型上下文有长度和成本限制,不能把所有历史原封不动地带进每一次调用。Summarization Middleware 会在达到条件时压缩较早内容,同时保留继续执行所需的事实。
摘要不是 Memory。它服务当前 thread 的上下文窗口,目标是让这条任务能继续;Memory 则用于跨 thread 保留用户或 Agent 的长期经验。混淆两者,会把“这次对话太长”和“下次还要记得”当成同一个问题。
DeerFlowSummarizationMiddleware 在工厂中接收 hooks、模型名称、app_config 和 extensions。当前注册的 memory_flush_hook 会在旧消息从上下文移除前,把即将压缩的消息送进长期 Memory 队列;它不是一个通用的附件或任务保存器。
摘要完成后,新的 summary 与最近消息继续交给模型,主循环没有换成另一条 Run。插叙结束,我们仍停在同一个 agent.astream 中,只是下一次模型输入更短。
get_memory_contextMemory 在模型调用前取回长期经验
假设用户过去明确说过“销售分析优先比较环比,不要把未确认线索写成结论”。这些偏好不一定在当前 thread 消息里,却会影响报告质量。Dynamic Context Middleware 可以在模型调用前按 agent_name 与 user_id 取回相关 Memory。
检索结果不是把整个记忆库倾倒进提示。配置决定是否注入,Memory manager 负责召回,Middleware 再把结果整理成有限长度的上下文块,与当前日期等动态信息一起进入模型。
摘要压缩当前 thread;Memory 按用户和 Agent 检索跨 thread 经验,并在 Agent 结束时安排更新。
injection_enabled = self._app_config.memory.injection_enabled if self._app_config else True
memory_context = (
_get_memory_context(
self._agent_name,
app_config=self._app_config,
user_id=resolve_runtime_user_id(runtime),
)
if injection_enabled
else ""
)
current_date = _format_current_date()
date_reminder = _format_current_date_reminder(current_date)
memory_block = memory_context.strip() if memory_context else None
return date_reminder, memory_blockinjection_enabled 控制召回路径;_get_memory_context 返回 memory_context,随后整理为 memory_block。若关闭注入,模型照常运行,只是不会看到长期经验。
这说明 Memory 是可选增强,不是 ThreadState 的替代品。当前报告的消息、附件和 artifacts 仍在图状态里;长期偏好只在需要时成为额外上下文。
manager_class: openvikingOpenViking 是 Memory 后端,不是另一套 Agent
DeerFlow 的 Memory manager 可以有不同实现,OpenViking 是当前支持的一种后端。选择它不会改变 agent.astream 的主循环;变化发生在召回与写入长期经验的接口后面。
配置同时区分总开关、注入开关、manager_class 和 backend_config。部署者还要决定服务地址、API key、启动失败策略以及检索阈值。
memory:
enabled: true
injection_enabled: true
shutdown_flush_timeout_seconds: 30
manager_class: openviking
mode: middleware
backend_config:
base_url: http://127.0.0.1:1933
owner_user_id: default
api_key_env: OPENVIKING_API_KEY
startup_policy: fail_fast
failure_policy:
read: fail_open
write: log_and_drop
retrieval:
top_k: 8
score_threshold: 0.25
max_injection_chars: 12000
content_mode: overview
injection_query: >-
user profile preferences important entities events ongoing goals
constraints and prior decisionsmanager_class: openviking 选择实现,retrieval.top_k 和 score_threshold 控制候选数量与相关性,max_injection_chars 限制进入模型的字符量。failure_policy.read 设为 fail_open,意味着记忆服务读取失败时任务仍可继续。
因此不要把“启用了 OpenViking”理解成模型自动拥有完美记忆。召回质量取决于写入内容、用户隔离、检索配置和注入预算;Memory 出故障时是否阻断主任务,也是明确的工程选择。
tasktask 工具把独立支线交给 Sub-agent
如果主 Agent 同时调查价格、竞品和区域差异,上下文会迅速变乱。它可以调用 task 工具,把“核对华东区竞品降价”作为边界清楚的支线交给 Sub-agent,自己保留整份报告的统筹责任。
这不是偷偷再开一条用户请求。task 本身是一个普通工具调用,有 description、prompt、subagent_type 和 tool_call_id;结果最终仍回到主 Agent 的工具消息。
@tool("task", parse_docstring=True)
async def task_tool(
runtime: Runtime,
description: str,
prompt: str,
subagent_type: str,
tool_call_id: Annotated[str, InjectedToolCallId],
) -> str | Command:task_tool 用 subagent_type 选择执行配置,用 prompt 传递明确任务。返回类型可以是字符串或 Command,说明委派结果既可能直接返回文本,也可能带状态更新。
Sub-agent 的价值是上下文隔离与职责分解,而不是数量越多越聪明。任务太小会增加延迟,边界不清则会重复搜索;只有能独立验收的支线才值得委派。
subagent stateSub-agent 拿到的是新状态,不是整条主对话
委派时最容易误解的一点,是以为 Sub-agent 自动继承主 Agent 的全部消息。实际执行器会组装自己的 system prompt,再把具体 task 作为 HumanMessage 放进新的初始状态。
需要共享的文件、技能或上下文要由运行时明确装配;不相关的主对话不会无条件复制。这样既控制 token,也减少支线被无关讨论干扰。
messages 先放可选的 SystemMessage,再追加表示实际任务的 HumanMessage,最后写入 state。这个新 state 属于支线执行;完成后返回经过选择的结果,而不是把两份消息历史粗暴合并。
当竞品调研结果返回,主 Agent 把它与已有证据一起纳入报告,随后调用写文件工具把成稿放进 Sandbox 的 outputs。Sub-agent 插叙结束;文件此时已经存在,下一步是明确把它交给用户。
present_filespresent_files 把选定文件登记为 Artifact
Agent 调用 present_files,并传入准备交付的输出路径。工具先规范化路径、检查文件是否位于允许目录,再构造状态更新。只有这一步成功,前端才有稳定、受控的结果列表。
此时文件名第一次成为交付事实:/mnt/user-data/outputs/sales-review.md。此前即使正文提过这个目标名称,也不能说明文件已经生成并可供下载。
# The merge_artifacts reducer will handle merging and deduplication
return Command(
update={
"artifacts": normalized_paths,
"messages": [ToolMessage("Successfully presented files", tool_call_id=tool_call_id)],
},
)Command.update 把 normalized_paths 写入 artifacts,并附上一条 ToolMessage。ThreadState 的 reducer 负责合并和去重;页面把这份 artifacts 列表解释为 Artifact,而不是扫描 Sandbox。
sales-review.md 因而同时拥有两种事实:它是文件系统里的字节,也是状态里被明确登记的交付物。前一个事实让工具能继续修改,后一个事实让用户能看见、下载和追踪。
after_agentAgent 结束时,Memory 才安排写回
present_files 返回以后,模型确认任务已经满足,Agent 生命周期准备结束。现在才轮到 Memory 的捕获阶段:系统从本次对话中筛选值得长期保存的内容,而不是每产生一个 chunk 就同步修改记忆库。
Memory Middleware 的 after_agent 先解析 thread_id、messages、user_id 与 trace_id;如果当前状态不满足写入条件,就直接返回,不影响主任务结果。注意,此处是 Agent 已经停止选择新工具之后,但 worker 还没有完成 Run 的最终结算。
@override
def after_agent(self, state: MemoryMiddlewareState, runtime: Runtime) -> dict | None:
"""Queue conversation for memory update after agent completes."""
add_args = self._resolve_add_args(state, runtime)
if add_args is None:
return None
thread_id, messages, user_id, trace_id = add_args
# Hand raw messages to the manager; the backend filters to user + final-AI
# turns, validates, detects correction/reinforcement, and enqueues.
get_memory_manager().add(
thread_id,
messages,
agent_name=self._agent_name,
user_id=user_id,
trace_id=trace_id,
)
return Noneafter_agent 调用 get_memory_manager().add 安排更新,随后返回 None。源码中的 Queue conversation 表明长期记忆写入与主响应解耦;它不会再回到模型工具循环,也不应该让次要的 Memory 写入阻塞报告交付。
Memory 写回安排完成后,Agent 图产生最后的状态,worker 才离开 agent.astream。主线接下来进入 Run 结算。
set_status_if_not_cancelledworker 验证交付后写入 RunStatus.success
present_files 返回以后,结果先回到 Agent 图。模型确认任务已经满足,agent.astream 才结束;worker 随后计算本次 Run 真正产生的输出路径,并检查交付 receipt 是否完整。
所以“工具调用成功”和“Run 成功”是两个时刻。只有没有 delivery_error,worker 才把候选终态设为 RunStatus.success;若输出校验失败,同一处代码会写入 RunStatus.error。
delivery_error = _delivery_error(delivery_content)
cancel_action = await run_manager.set_status_if_not_cancelled(
run_id,
RunStatus.error if delivery_error else RunStatus.success,
error=delivery_error,
stop_reason=stop_reason,
**terminal_status_kwargs,
)
if cancel_action is not None:
await _finish_cancellation(cancel_action)set_status_if_not_cancelled 还会在落终态前处理并发取消,避免成功结果覆盖已经收到的取消请求。源码中的条件表达式正是正常交付与错误交付的最后分叉。
到这里主链结束:RunRecord 有终态,artifacts 中有交付路径。Checkpoint、RunEventStore 和 StreamBridge 接下来解释的是这些事实怎样保存或送达,不是新的执行步骤。
checkpointerCheckpoint 保存 agent.astream 推进到哪里
主线已经到达 Run 收尾。现在回头解释故障恢复需要的三类事实。第一类是 Checkpoint:LangGraph 在节点边界保存状态快照,使同一 thread 的图能够从已知状态继续,而不必只依赖进程内变量。
Checkpoint 关心 messages、artifacts、todos 等图状态;它不等于事件日志,也不负责把实时 chunk 推给网页。部署可以选择不同持久化实现,因此运行时通过 provider 创建 checkpointer。
make_checkpointer 是 async context manager,返回 AsyncIterator[Checkpointer]。上下文管理器负责建立和关闭后端资源,Agent 工厂只接收统一接口,不必知道它来自内存还是数据库。
有了 Checkpoint,系统知道 Agent 图保存到哪里;但用户还需要查询“这条 Run 发生过什么”,这属于下一类事实。
RunEventStoreRunEventStore 保存可查询的运行历史
RunEventStore 按 thread_id 与 run_id 保存模型、工具、进度、错误和终态事件。页面刷新后,它可以重新列出历史;运维排查时,也能把同一 thread 中的多条 Run 分开。
持久事件有自己的顺序 seq。查询可以按事件类型、task_id 和 after_seq 过滤,适合断点续读或只查看某类运行事实。
@abc.abstractmethod
async def list_events(
self,
thread_id: str,
run_id: str,
*,
event_types: list[str] | None = None,
task_id: str | None = None,
limit: int = 500,
after_seq: int | None = None,
) -> list[dict]:list_events 同时接收 thread_id、run_id 与 after_seq,并返回事件记录列表。这个接口说明持久历史的游标是序列号,不应与实时流后端的消息 ID 随意比较。
Checkpoint 回答“图状态是什么”,RunEventStore 回答“执行过程发生了什么”。两者都能帮助恢复,却保存不同语义;最后还剩实时交付这第三类事实。
bridge.publishStreamBridge 只负责把正在发生的事件送出去
agent.astream 产生 chunk 后,worker 通过 StreamBridge 发布;SSE consumer 再从 Bridge 消费。Bridge 追求低延迟,不承担全部长期历史。单进程开发可用内存队列,多实例部署则需要共享后端。
选择内存还是 Redis 不改变 worker 的 publish 接口,却会改变跨进程可见性。worker 与 Gateway 若不在同一进程,MemoryStreamBridge 无法让另一个进程读到事件。
管理状态、图快照、持久事件和实时流各有用途;恢复时不要把它们当成一份数据。
messages、todo、sandbox、artifacts
status、owner、lease、终止摘要
按 seq 查询的 message 与 trace
SSE chunk、重放窗口与 StreamGap
if config is None or config.type == "memory":
from deerflow.runtime.stream_bridge.memory import MemoryStreamBridge
maxsize = config.queue_maxsize if config is not None else 256
bridge = MemoryStreamBridge(queue_maxsize=maxsize)
logger.info("Stream bridge initialised: memory (queue_maxsize=%d)", maxsize)
try:
yield bridge
finally:
await bridge.close()
return
if config.type == "redis":
from deerflow.runtime.stream_bridge.redis import RedisStreamBridge
redis_url = _resolve_redis_url(config)
bridge = RedisStreamBridge(
redis_url=redis_url,
queue_maxsize=config.queue_maxsize,
max_connections=config.max_connections,
stream_ttl_seconds=config.stream_ttl_seconds,
)MemoryStreamBridge 在上下文结束时 close;RedisStreamBridge 则接收 redis_url、queue_maxsize、连接数与 TTL。两种实现共享抽象,但可靠性边界不同。
现在三类恢复事实齐了:Checkpoint 保存图状态,RunEventStore 保存可查询历史,StreamBridge 传递正在发生的事件;再加上 RunRecord 的管理终态,页面刷新后才能既恢复结果,又不把旧任务误判为仍在执行。
query Run + artifacts回到用户:刷新页面以后,什么还在
让我们把插叙全部收起。主 Agent 已经读取 sales-notes.pdf,必要时召回 Memory、调用 Sub-agent 补充调查,在 Sandbox 中写好 sales-review.md,并通过 present_files 将它登记为 Artifact。worker 刷新日志,RunRecord 进入 success,实时流发送结束信号。
如果用户此时刷新页面,旧 SSE 连接当然消失,但 thread 与 Run 仍可重新查询;Checkpoint 保留图状态,RunEventStore 提供历史事件,artifacts 让页面重新渲染下载入口。只有未持久化的进程内队列不应被当成唯一事实来源。
若结果不完整,可以沿同一条链排查:没有 RunRecord,看提交与 start_run;Run 一直 running,看 worker 和 agent.astream;模型缺少旧偏好,看 Memory 召回;Sub-agent 没结果,看 task 工具事件;磁盘有文件但页面没有,看 present_files 与 artifacts;刷新后事件丢失,看 RunEventStore 和 StreamBridge 配置。
到这里,教学命令中的每个部分都有了真实落点:thread 选择持续工作空间,mode 影响运行配置,file 指向线程附件,目标进入图输入;最终文件经过 Sandbox 与 Artifact 边界返回界面。这就是 DeerFlow 执行一项任务的完整骨架。
9 个真实调用点沿真实文件再走一遍核心调用栈
缩进表示调用进入下一层;点击任意一行可回到对应源码。
- 01frontend/src/core/threads/hooks.ts
thread.submit(...)把消息、附件信息和运行选项交给 Gateway - 02backend/app/gateway/routers/thread_runs.py
start_run(...)流式路由和普通路由共用同一个启动入口 - 03backend/app/gateway/services.py
create_or_reject(...)先得到可查询的 RunRecord - 04backend/app/gateway/services.py
asyncio.create_task(worker)再把实际工作放进后台任务 - 05backend/app/gateway/services.py
run_agent(...)把 RunRecord、配置和运行后端交给 Harness - 06backend/packages/harness/deerflow/runtime/runs/worker.py
agent_factory(...)装配出包含模型、工具和 Middleware 的图 - 07backend/packages/harness/deerflow/runtime/runs/worker.py
agent.astream(...)驱动模型与工具循环,并持续产出事件 - 08backend/packages/harness/deerflow/tools/builtins/present_file_tool.py
present_files(...)把 outputs 中选定的文件登记进 artifacts - 09backend/packages/harness/deerflow/runtime/runs/worker.py
set_status_if_not_cancelled(..., RunStatus.success)确认交付无误后写下 Run 的终态
这张调用栈只保留会把控制权交给下一层的代码。缩进表示调用深入,回退表示内部工作完成后回到 worker;文件名则告诉你该去仓库哪里继续读。
Middleware、Memory、Sub-agent、Sandbox、Checkpoint 与事件系统没有从调用栈中消失。它们分别挂在 agent_factory、agent.astream、工具调用和 worker 收尾处,但不改变主链从 thread.submit 走到 RunStatus.success 的方向。