跳到正文
DEERFLOW BOOK
02 / 02EN
02

第 2 章

一次任务的完整执行

从提交、模型与工具循环,一直走到文件交付

这一章只做一件事:沿着运行时间向前走。浏览器先通过 thread.submit 提交任务,Gateway 创建 RunRecord,worker 调用 run_agent;随后 Agent 在 agent.astream 中反复使用模型和工具。中途需要哪个机制,我们就在那个位置停下来解释。

CONTROL FLOW七个核心调用落在四个执行阶段

先用四个阶段定位,再在正文中沿真实函数逐层进入;旁支只挂在对应调用点上。

  1. 01Submit

    浏览器提交 thread 输入

  2. 02Manage

    Gateway 建立运行身份

  3. 03Execute

    worker 驱动 Agent 图

  4. 04Deliver

    文件登记并结束运行

全局图submit → run → loop → deliver

先看完整流程:主链只有七个调用点

先只读流程图中间一列:thread.submit 提交任务,start_run 创建并安排 Run,run_agent 准备运行环境,agent_factory 装配 Agent,agent.astream 驱动模型与工具,present_files 登记交付物,worker 最后写入 RunStatus.success这七个调用点组成主链。

流程图侧面的 MiddlewareMemorySub-agentSandboxCheckpoint 和事件后端都不是下一步。它们只在主链经过某个调用点时提供能力。例如 Sandbox 挂在文件工具调用上,Checkpoint 挂在 run_agent 的图状态上。把侧枝暂时遮住,主链仍然能够从上走到下。

下面每个标题前都有一枚位置标记。“主链”表示控制流向下一层移动;“主链内部”表示仍在当前函数里;“挂在”表示临时解释支撑机制,读完会回到标出的函数。

DEERFLOW MAP · EXECUTION FLOW一次 DeerFlow 任务的完整执行流程

中间是必须依次经过的主链;右侧短句是只在某一步介入的机制。先顺着中间读到底,再回头看分支。

  1. 01
    Browserthread.submit

    input + config + context

    • 附件已上传,只传路径与 metadata
  2. 02
    Gatewaystart_run

    RunRecord + background task

    • SSE 只订阅这条 Run
    • StreamBridge 传递实时事件
  3. 03
    Workerrun_agent

    RunnableConfig + RunContext

    • Checkpoint 提供图状态
    • RunEventStore 接收持久事件
  4. 04
    Assemblyagent_factory

    model + tools + middleware + state

    • 模式决定模型与工具
    • Middleware 包住模型和工具调用
  5. 05
    Agent loopagent.astream

    AIMessage ↔ ToolMessage

    • Memory 在模型调用前召回
    • Sub-agent 通过 task 工具委派
    • Sandbox 在文件工具首次使用时取得
  6. 06
    Deliverypresent_files

    ThreadState.artifacts

    • 只有 outputs 中明确选中的文件会交付
  7. 07
    SettlementRunStatus.success

    终态 + receipt + 可恢复结果

    • Checkpoint 留状态
    • RunEventStore 留历史
    • 页面可重新查询 Artifact
主链 01 / 07thread.submit

thread.submit 把任务交给 Gateway

主链从前端 hooks.ts 中的 thread.submit 开始。教学命令把 thread、mode、附件和目标压在一行,真实网页则把它们拆成结构化数据:第一项是图输入,第二项是 threadId、config 与 context。

这里最重要的区别是“用户目标”和“运行选项”。用户目标告诉 Agent 要做什么;mode、thinking_enabledsubagent_enabled 等选项告诉 Harness 怎样装配这次执行。两者会在 worker 中重新相遇,但不会混成一段无法校验的提示词。

frontend/core/threads/hooks.ts · thread.submit
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 也不存在。

主链 02 / 07stream_run → start_run

stream_run 进入统一的 start_run

thread.submit 的流式请求落到 stream_run。这个路由只取出 StreamBridge 与 RunManager,然后调用 start_run;普通创建接口和等待接口也复用同一个 start_run,所以 DeerFlow 没有为 SSE 另写一套 Agent 执行逻辑。

start_run 返回 record 后,路由把 sse_consumer 包进 StreamingResponse。浏览器由此开始观察这条 Run,但响应流本身不负责执行模型。

thread_runs.py · stream_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 怎样出现。

主链 03 / 07create_or_reject → create_task

start_run 先保存 RunRecord,再安排 worker

start_run 先规范化输入与配置,再调用 run_mgr.create_or_reject。这个调用把 thread_id、input、metadata、模型和多任务策略写进可查询的 RunRecord;若同一 thread 的并发策略不允许新任务,它也在这里拒绝。

只有 durable admission 成功以后,代码才构造 worker,并用 asyncio.create_task 挂到 record.task。顺序不能倒过来:否则后台已经开始跑,系统却还没有一条能取消、查询或接管的运行记录。

gateway/services.py · Run admission and worker 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(...)。

挂在 start_run 之后sse_consumer

SSE 只是观察 start_run 创建的同一条 Run

主链现在停在 start_run 已经返回 RunRecord 的位置。sse_consumer 从 bridge 读取这条 record 的事件并发送给浏览器,不负责调用模型,也不拥有 Agent 的生命周期。

因此,“前端没收到流”与“后台没执行”不是同一个故障。前者看 SSE 和 StreamBridge,后者看 RunRecord 与 worker。连接断开后是否取消由 on_disconnect 策略决定,不由 TCP 连接替系统猜测。

这条旁支到此结束。控制流回到 record.task 中的 worker,继续进入 run_agent

主链 04 / 07run_agent

worker 调用 run_agent 进入 Harness

start_run 中的后台协程最终执行 await run_agent(...)。这行是 Gateway 与 Harness 的明确边界:左边仍在管理 HTTP 请求和 RunRecord,右边开始准备 LangGraph 运行环境。

调用参数也把职责说清了。record 提供本次运行身份,graph_input 是规范化后的模型输入,config 是运行配置,ctx 是 RunContext,提供 checkpointer、event_storethread_store 等基础设施,bridge 负责实时事件。

gateway/services.py · enter run_agent
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 图。

主链内部 · run_agentRunnableConfig

run_agent 先准备本次运行配置

前端提交的模式、模型选择与功能开关,需要进入 Agent 工厂;LangGraph 自己也有 configurable、context、metadata 等运行字段。DeerFlow 不让各个模块随意从请求对象取值,而是在装配入口先把兼容字段合并。

这一步发生在真正创建模型和 Middleware 之前。结果决定是否启用深度思考、Sub-agent、特定模型,以及当前用户身份怎样传给后续运行组件。

lead_agent/agent.py · _get_runtime_config
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

主链内部 · run_agentagent_factory

run_agent 调用 agent_factory 构建图

RunnableConfig 准备好以后,worker 把它放进 agent_factory_kwargs,并调用前面由 Gateway 解析出的 agent_factory。默认情况下,这个工厂就是 assemble_lead_agent

这里不要把 agent_factory 理解成又一次模型调用。它做的是组装:根据本次配置选择模型、工具、Middleware、系统提示和 ThreadState schema,返回可以运行的 LangGraph 图。

runtime/runs/worker.py · build agent graph
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,看模型和工具怎样接到同一张图。

主链内部 · agent_factorycreate_agent

agent_factory 把模型、工具和 Middleware 接到同一张图

所谓“创建 Agent”,实际是把四类东西接到同一张图:聊天模型、可调用工具、系统提示和 MiddlewareThreadState 则规定图运行期间共享的数据形状。

对报告任务来说,工具集合可能包含网页搜索、文件读写、shell 命令、present_files 和 task;配置决定其中哪些真正可用。Middleware 再为这些调用补上摘要、记忆、进度与错误处理。

lead_agent/agent.py · create_agent
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_toolsnormalize_middleware_state_schemas 的结果、system_promptget_thread_state_schema(mode)。源码把它们并排传入,比任何“智能体平台”口号都更能说明边界:模型只是装配件之一。

Agent 创建完成后,控制权回到前面的 agent.astream。图读取当前状态,把消息交给模型;模型若发出 tool call,图执行对应工具并把 ToolMessage 放回状态,然后开始下一轮。

主链 05 / 07agent.astream

agent.astream 驱动模型与工具循环

Agent 图装配完成后,run_agent 调用 agent.astream(input_payload, config, stream_mode)。从这一行开始,模型与工具才真正工作。它不是等待一个最终字符串,而是持续产出消息、状态更新和自定义事件。

每个 chunk 返回 worker 后,worker 检查 abort_event,再通过 bridge.publish(run_id, ...) 发布。取消和事件发送包在图循环外面,因此不会改变模型下一步选择哪个工具。

DEERFLOW MAP · AGENT LOOPagent.astream 内部不是一次回答,而是反复循环

模型选择工具,工具结果写回状态;没有新的 tool call 时,图才结束本轮 Run。

graph_input.messages → agent.astreamModeltool_callsToolMessagestate updateAgent loop一圈 = 一次新的观察无 tool call → final
runtime/runs/worker.py · graph stream
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、搜索和写文件。下面的 SandboxMiddlewareMemorySub-agent 都挂在这段 agent.astream 循环上;读完任何一条旁支,都要回到这里继续下一轮。

挂在 agent.astream 的工具调用read_file

read_file 第一次访问文件时取得 Sandbox

Sandbox Middleware 不一定在 Run 开始时立刻创建环境。第一次 read_file、写文件或执行命令需要文件系统时,ensure_sandbox_initialized 才会按 thread_iduser_id 获取环境。纯文本任务因此不必提前承担 Sandbox 成本。

按 thread 获取还有另一个效果:同一 thread 的后续 Run 可以继续访问此前工作文件。第二次让系统修改报告时,不必把所有中间材料重新塞进消息。

DEERFLOW MAP · FILESYSTEM同一 thread 的文件工作区

附件位于可读输入区,中间文件留在工作区,准备交付的结果进入 outputs。

THREAD SANDBOX · /mnt/user-data
uploads/sales-notes.pdf · 只作为输入
workspace/调查笔记、临时脚本
outputs/sales-review.md · 可交付
ThreadState.artifactsabsolutePath + metadata
Frontend预览或下载入口
sandbox/middleware.py · acquire
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_sandboxprovider.acquire 取得 sandbox_id,并记录当前环境。具体 provider 可以映射到本地目录、容器或远程 Sandbox;上层工具只依赖统一的文件与命令接口。

环境准备好以后,read_file 读取 sales-notes.pdf,把访谈内容作为工具结果交回 ThreadState。模型发现“客户更关注价格”还只是内部判断,于是继续调用搜索工具核对行业数据;稍后写报告时也复用这份工作区。主线现在回到 agent.astream 的循环。

挂在 agent.astream 的模型与工具调用wrap_*_call

Middleware 包在 agent.astream 的调用边界上

主线现在停在一次工具调用。工具真正执行前后,多个 Middleware 会像函数包装器一样套在外面。读取后再写、进度上报和错误处理看似互不相干,但包装顺序会改变谁能看到原始异常、谁能看到规范化结果。

源码在构造 tail 时特意留下顺序说明:ToolProgressMiddleware 要位于 ToolErrorHandlingMiddleware 外面,才能观察后者已经写入 metadata 的结果。

DEERFLOW MAP · MIDDLEWARE一次工具调用外面的包装层

外层记录完整过程,内层规范失败;返回时结果按相反方向穿回去。

before_agent一次 Run 开始
before_model每次模型调用
WRAPPERS · 前注册在外层
wrap_model_call
Provider / AIMessage
wrap_tool_callGuard → Sandbox tool → ToolMessage
after_model按注册顺序逆序回收after_agentRun 结束后捕获与结算
middlewares · tool wrapper order
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 决定某一步怎样被执行和记录,模型仍决定是否调用这一步。插叙结束,工具结果回到 ThreadStateagent.astream 继续下一轮。

挂在 agent.astream 的模型调用SummarizationMiddleware

消息过长时,摘要发生在下一次模型调用前

报告任务经过多轮搜索与文件读取后,messages 会越来越长。模型上下文有长度和成本限制,不能把所有历史原封不动地带进每一次调用。Summarization Middleware 会在达到条件时压缩较早内容,同时保留继续执行所需的事实。

摘要不是 Memory。它服务当前 thread 的上下文窗口,目标是让这条任务能继续;Memory 则用于跨 thread 保留用户或 Agent 的长期经验。混淆两者,会把“这次对话太长”和“下次还要记得”当成同一个问题。

DeerFlowSummarizationMiddleware 在工厂中接收 hooks、模型名称、app_config 和 extensions。当前注册的 memory_flush_hook 会在旧消息从上下文移除前,把即将压缩的消息送进长期 Memory 队列;它不是一个通用的附件或任务保存器。

摘要完成后,新的 summary 与最近消息继续交给模型,主循环没有换成另一条 Run。插叙结束,我们仍停在同一个 agent.astream 中,只是下一次模型输入更短。

挂在 agent.astream 的模型调用get_memory_context

Memory 在模型调用前取回长期经验

假设用户过去明确说过“销售分析优先比较环比,不要把未确认线索写成结论”。这些偏好不一定在当前 thread 消息里,却会影响报告质量。Dynamic Context Middleware 可以在模型调用前按 agent_nameuser_id 取回相关 Memory

检索结果不是把整个记忆库倾倒进提示。配置决定是否注入,Memory manager 负责召回,Middleware 再把结果整理成有限长度的上下文块,与当前日期等动态信息一起进入模型。

DEERFLOW MAP · MEMORY摘要与 Memory 是两条不同路径

摘要压缩当前 thread;Memory 按用户和 Agent 检索跨 thread 经验,并在 Agent 结束时安排更新。

NEW RUN
DynamicContextuser_id + agent_name
Model Context隐藏用户层 Memory 文本
MemoryMiddleware原始 messages
RECALLget_context
add / add_nowaitCAPTURE
LONG-TERM BACKEND
DeerMem默认结构化事实
OpenViking可选远程 Session
Backend contractread / write failure policy
dynamic_context_middleware.py · recall
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_block

injection_enabled 控制召回路径;_get_memory_context 返回 memory_context,随后整理为 memory_block。若关闭注入,模型照常运行,只是不会看到长期经验。

这说明 Memory 是可选增强,不是 ThreadState 的替代品。当前报告的消息、附件和 artifacts 仍在图状态里;长期偏好只在需要时成为额外上下文。

Memory 召回内部manager_class: openviking

OpenViking 是 Memory 后端,不是另一套 Agent

DeerFlow 的 Memory manager 可以有不同实现,OpenViking 是当前支持的一种后端。选择它不会改变 agent.astream 的主循环;变化发生在召回与写入长期经验的接口后面。

配置同时区分总开关、注入开关、manager_classbackend_config。部署者还要决定服务地址、API key、启动失败策略以及检索阈值。

docs/OPENVIKING.md · backend configuration
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 decisions

manager_class: openviking 选择实现,retrieval.top_kscore_threshold 控制候选数量与相关性,max_injection_chars 限制进入模型的字符量。failure_policy.read 设为 fail_open,意味着记忆服务读取失败时任务仍可继续。

因此不要把“启用了 OpenViking”理解成模型自动拥有完美记忆。召回质量取决于写入内容、用户隔离、检索配置和注入预算;Memory 出故障时是否阻断主任务,也是明确的工程选择。

挂在 agent.astream 的工具调用task

task 工具把独立支线交给 Sub-agent

如果主 Agent 同时调查价格、竞品和区域差异,上下文会迅速变乱。它可以调用 task 工具,把“核对华东区竞品降价”作为边界清楚的支线交给 Sub-agent,自己保留整份报告的统筹责任。

这不是偷偷再开一条用户请求。task 本身是一个普通工具调用,有 description、prompt、subagent_typetool_call_id;结果最终仍回到主 Agent 的工具消息。

tools/builtins/task_tool.py · task
@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_toolsubagent_type 选择执行配置,用 prompt 传递明确任务。返回类型可以是字符串或 Command,说明委派结果既可能直接返回文本,也可能带状态更新。

Sub-agent 的价值是上下文隔离与职责分解,而不是数量越多越聪明。任务太小会增加延迟,边界不清则会重复搜索;只有能独立验收的支线才值得委派。

task 工具内部subagent state

Sub-agent 拿到的是新状态,不是整条主对话

委派时最容易误解的一点,是以为 Sub-agent 自动继承主 Agent 的全部消息。实际执行器会组装自己的 system prompt,再把具体 task 作为 HumanMessage 放进新的初始状态。

需要共享的文件、技能或上下文要由运行时明确装配;不相关的主对话不会无条件复制。这样既控制 token,也减少支线被无关讨论干扰。

messages 先放可选的 SystemMessage,再追加表示实际任务的 HumanMessage,最后写入 state。这个新 state 属于支线执行;完成后返回经过选择的结果,而不是把两份消息历史粗暴合并。

当竞品调研结果返回,主 Agent 把它与已有证据一起纳入报告,随后调用写文件工具把成稿放进 Sandbox 的 outputs。Sub-agent 插叙结束;文件此时已经存在,下一步是明确把它交给用户。

主链 06 / 07present_files

present_files 把选定文件登记为 Artifact

Agent 调用 present_files,并传入准备交付的输出路径。工具先规范化路径、检查文件是否位于允许目录,再构造状态更新。只有这一步成功,前端才有稳定、受控的结果列表。

此时文件名第一次成为交付事实:/mnt/user-data/outputs/sales-review.md。此前即使正文提过这个目标名称,也不能说明文件已经生成并可供下载。

present_file_tool.py · artifact update
# 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.updatenormalized_paths 写入 artifacts,并附上一条 ToolMessage。ThreadState 的 reducer 负责合并和去重;页面把这份 artifacts 列表解释为 Artifact,而不是扫描 Sandbox

sales-review.md 因而同时拥有两种事实:它是文件系统里的字节,也是状态里被明确登记的交付物。前一个事实让工具能继续修改,后一个事实让用户能看见、下载和追踪。

挂在 agent.astream 结束之后after_agent

Agent 结束时,Memory 才安排写回

present_files 返回以后,模型确认任务已经满足,Agent 生命周期准备结束。现在才轮到 Memory 的捕获阶段:系统从本次对话中筛选值得长期保存的内容,而不是每产生一个 chunk 就同步修改记忆库。

Memory Middlewareafter_agent 先解析 thread_id、messages、user_idtrace_id;如果当前状态不满足写入条件,就直接返回,不影响主任务结果。注意,此处是 Agent 已经停止选择新工具之后,但 worker 还没有完成 Run 的最终结算。

memory_middleware.py · after_agent
@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 None

after_agent 调用 get_memory_manager().add 安排更新,随后返回 None。源码中的 Queue conversation 表明长期记忆写入与主响应解耦;它不会再回到模型工具循环,也不应该让次要的 Memory 写入阻塞报告交付。

Memory 写回安排完成后,Agent 图产生最后的状态,worker 才离开 agent.astream。主线接下来进入 Run 结算。

主链 07 / 07set_status_if_not_cancelled

worker 验证交付后写入 RunStatus.success

present_files 返回以后,结果先回到 Agent 图。模型确认任务已经满足,agent.astream 才结束;worker 随后计算本次 Run 真正产生的输出路径,并检查交付 receipt 是否完整。

所以“工具调用成功”和“Run 成功”是两个时刻。只有没有 delivery_error,worker 才把候选终态设为 RunStatus.success;若输出校验失败,同一处代码会写入 RunStatus.error

runtime/runs/worker.py · settle Run status
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 中有交付路径。CheckpointRunEventStoreStreamBridge 接下来解释的是这些事实怎样保存或送达,不是新的执行步骤。

挂在 run_agent 的图状态checkpointer

Checkpoint 保存 agent.astream 推进到哪里

主线已经到达 Run 收尾。现在回头解释故障恢复需要的三类事实。第一类是 Checkpoint:LangGraph 在节点边界保存状态快照,使同一 thread 的图能够从已知状态继续,而不必只依赖进程内变量。

Checkpoint 关心 messages、artifacts、todos 等图状态;它不等于事件日志,也不负责把实时 chunk 推给网页。部署可以选择不同持久化实现,因此运行时通过 provider 创建 checkpointer。

make_checkpointer 是 async context manager,返回 AsyncIterator[Checkpointer]。上下文管理器负责建立和关闭后端资源,Agent 工厂只接收统一接口,不必知道它来自内存还是数据库。

有了 Checkpoint,系统知道 Agent 图保存到哪里;但用户还需要查询“这条 Run 发生过什么”,这属于下一类事实。

挂在 worker 的事件记录RunEventStore

RunEventStore 保存可查询的运行历史

RunEventStorethread_idrun_id 保存模型、工具、进度、错误和终态事件。页面刷新后,它可以重新列出历史;运维排查时,也能把同一 thread 中的多条 Run 分开。

持久事件有自己的顺序 seq。查询可以按事件类型、task_idafter_seq 过滤,适合断点续读或只查看某类运行事实。

runtime/events/store/base.py · RunEventStore
@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_idrun_idafter_seq,并返回事件记录列表。这个接口说明持久历史的游标是序列号,不应与实时流后端的消息 ID 随意比较。

Checkpoint 回答“图状态是什么”,RunEventStore 回答“执行过程发生了什么”。两者都能帮助恢复,却保存不同语义;最后还剩实时交付这第三类事实。

挂在 worker 的实时发布bridge.publish

StreamBridge 只负责把正在发生的事件送出去

agent.astream 产生 chunk 后,worker 通过 StreamBridge 发布;SSE consumer 再从 Bridge 消费。Bridge 追求低延迟,不承担全部长期历史。单进程开发可用内存队列,多实例部署则需要共享后端。

选择内存还是 Redis 不改变 worker 的 publish 接口,却会改变跨进程可见性。worker 与 Gateway 若不在同一进程,MemoryStreamBridge 无法让另一个进程读到事件。

DEERFLOW MAP · RUNTIME FACTS同一条 Run 的四种事实

管理状态、图快照、持久事件和实时流各有用途;恢复时不要把它们当成一份数据。

IDENTITYrun_id四份事实的关联键,不是执行本身
图状态Checkpoint

messages、todo、sandbox、artifacts

管理事实RunRecord

status、owner、lease、终止摘要

事件历史RunEventStore

按 seq 查询的 message 与 trace

实时现场StreamBridge

SSE chunk、重放窗口与 StreamGap

stream_bridge/async_provider.py · backend
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_urlqueue_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;刷新后事件丢失,看 RunEventStoreStreamBridge 配置。

到这里,教学命令中的每个部分都有了真实落点:thread 选择持续工作空间,mode 影响运行配置,file 指向线程附件,目标进入图输入;最终文件经过 SandboxArtifact 边界返回界面。这就是 DeerFlow 执行一项任务的完整骨架。

源码复盘9 个真实调用点

沿真实文件再走一遍核心调用栈

CORE CALL STACK这九行才是本章主链

缩进表示调用进入下一层;点击任意一行可回到对应源码。

  1. 01frontend/src/core/threads/hooks.tsthread.submit(...)把消息、附件信息和运行选项交给 Gateway
  2. 02backend/app/gateway/routers/thread_runs.pystart_run(...)流式路由和普通路由共用同一个启动入口
  3. 03backend/app/gateway/services.pycreate_or_reject(...)先得到可查询的 RunRecord
  4. 04backend/app/gateway/services.pyasyncio.create_task(worker)再把实际工作放进后台任务
  5. 05backend/app/gateway/services.pyrun_agent(...)把 RunRecord、配置和运行后端交给 Harness
  6. 06backend/packages/harness/deerflow/runtime/runs/worker.pyagent_factory(...)装配出包含模型、工具和 Middleware 的图
  7. 07backend/packages/harness/deerflow/runtime/runs/worker.pyagent.astream(...)驱动模型与工具循环,并持续产出事件
  8. 08backend/packages/harness/deerflow/tools/builtins/present_file_tool.pypresent_files(...)把 outputs 中选定的文件登记进 artifacts
  9. 09backend/packages/harness/deerflow/runtime/runs/worker.pyset_status_if_not_cancelled(..., RunStatus.success)确认交付无误后写下 Run 的终态

这张调用栈只保留会把控制权交给下一层的代码。缩进表示调用深入,回退表示内部工作完成后回到 worker;文件名则告诉你该去仓库哪里继续读。

MiddlewareMemorySub-agentSandboxCheckpoint 与事件系统没有从调用栈中消失。它们分别挂在 agent_factoryagent.astream、工具调用和 worker 收尾处,但不改变主链从 thread.submit 走到 RunStatus.success 的方向。