首页
/ DeerFlow Task Tool 改进解析:消除 LLM 轮询浪费,让子代理任务"一次调用、后台等到底"

DeerFlow Task Tool 改进解析:消除 LLM 轮询浪费,让子代理任务"一次调用、后台等到底"

2026-09-04 20:07:45作者:裴麒琰

在 DeerFlow 这类长时程(long-horizon)SuperAgent 框架中,主代理(lead agent)通过 task 工具把有界任务委派给子代理(subagent)执行。本篇基于仓库内的 task_tool_improvements.md 展开,完整讲清这次改进的核心动机——消除 LLM 反复调用 task_status 轮询造成的无效 API 请求——并逐项拆解四项变更:移除 run_in_background 参数、引入后端轮询、将 task_status 从 LLM 可见工具中移除、以及提示词文档的同步更新。读完本文,你将掌握该工具的完整调用链路(execute_async → 后台任务注册表 → 后端轮询 → 结果回填),以及双层超时保护与线程/事件循环架构的源码级原理,能够独立分析、验证和配置 DeerFlow 的子代理执行机制。

改进背景:LLM 轮询为什么是浪费

在改进之前,如果子代理任务以"后台模式"启动,LLM 必须自己管理完成状态:先调用 task(..., run_in_background=True) 拿到 task_id,然后反复调用 task_status(task_id),直到任务完成。文档给出的"改进前"模式如下:

# LLM 必须自己管理轮询
task_id = task(
    subagent_type="bash",
    prompt="Run tests",
    description="Run tests",
    run_in_background=True
)
# 然后 LLM 必须反复轮询:
while True:
    status = task_status(task_id)
    if completed:
        break

这条路径的问题在于:每一次 task_status 调用都会触发一次完整的 LLM 请求(工具调用 → 模型推理 → 再决策),对于耗时几十秒到几分钟的子代理任务,LLM 可能产生十几次甚至几十次纯状态查询的往返。这些请求不产生任何新信息,却消耗 token、API 配额和墙钟时间。改进的核心思路一句话概括:把"何时任务结束"的判断从 LLM 侧移到后端侧——工具调用本身阻塞直到完成,轮询由后端线程完成,LLM 只做一次工具调用。

改进后的调用形态(来自原文档):

# 工具会阻塞直到完成,轮询发生在后端
result = task(
    subagent_type="bash",
    prompt="Run tests",
    description="Run tests"
)
# 调用返回后结果立即可用

变更一:移除 run_in_background 参数

文档明确:run_in_background 参数已从 task 工具中删除。所有子代理任务默认异步执行,但工具本身会处理"等待完成",对 LLM 表现为同步阻塞语义。这一设计的关键在于"异步执行"与"同步等待"的分离:

  • 异步执行:子代理任务在独立于父运行的执行面(当前实现中是持久隔离事件循环,见下文)中启动,不占用父运行主循环的资源;
  • 同步等待task 工具协程内部循环读取后台任务状态,直到终态才把 ToolMessage 交还主代理。

对照当前仓库的 task 工具实现task_tool 的签名只保留了 promptsubagent_type、可选的 acceptance_criteriadescription,确实不存在任何 run_in_backgroundtask_id 出参——LLM 侧的接口已经被收敛到"描述任务、拿回结果"这一步。

变更二:后端轮询(Backend Polling)

文档给出的后端轮询逻辑骨架:

# 启动后台执行
task_id = executor.execute_async(prompt)

# 在后端轮询任务完成
while True:
    result = get_background_task_result(task_id)

    # 检查任务是否完成或失败
    if result.status == SubagentStatus.COMPLETED:
        return f"[Subagent: {subagent_type}]\n\n{result.result}"
    elif result.status == SubagentStatus.FAILED:
        return f"[Subagent: {subagent_type}] Task failed: {result.error}"

    # 下次轮询前等待
    time.sleep(2)

    # 超时保护(5 分钟)
    if poll_count > 150:
        return "Task timed out after 5 minutes"

文档总结的效果:LLM 只发起一次工具调用、不再有无谓的 LLM 轮询请求、后端统一负责状态检查、并有超时保护兜底。

当前实现中的轮询循环

当前仓库 task_tool.py 中的真实轮询实现比文档骨架更完整,核心流程如下:

# 保留 provider 的 tool-call ID 用于流/消息关联,
# 但用服务端生成的 execution ID 做进程级后台任务控制
execution_id = executor.execute_async(prompt, task_id=tool_call_id)

# 在后端轮询任务完成(消除 LLM 轮询的需要)
poll_count = 0
# 轮询超时:执行超时 + 60s 缓冲,每 5 秒检查一次
max_poll_count = (config.timeout_seconds + 60) // 5

while True:
    result = get_background_task_result(execution_id)
    # ... 状态变化日志、token 用量汇总、task_running 事件推送 ...
    if result.status == SubagentStatus.COMPLETED:
        # 报告 token 用量、发 task_completed 事件、清理注册表、返回结果
        ...
    # 仍在运行:等 5 秒后继续
    await asyncio.sleep(5)
    poll_count += 1
    if poll_count > max_poll_count:
        # 轮询超时兜底:请求协作取消 + 延迟清理,返回 polling_timed_out
        ...

从源码结构看,当前实相对原文档描述有几处演进(以当前仓库为准):

  1. 轮询节奏:文档中的示例是每 2 秒一次、poll_count > 150(约 5 分钟)截止;当前代码改为每 5 秒一次(task_tool.py#L1101),上限 max_poll_count 不再写死为 150,而是由配置的 timeout_seconds 推导:(timeout_seconds + 60) // 5,即轮询预算始终比执行超时多留 60 秒缓冲(task_tool.py#L904-L905)。
  2. 终态覆盖更全:轮询循环处理的不只是 COMPLETED/FAILED,还包括 CANCELLED(被取消)与 TIMED_OUT(执行超时),每个分支都会发送对应的自定义事件(task_completed / task_failed / task_cancelled / task_timed_out)供前端消费,并统一上报子代理的 token 用量快照。
  3. 任务"消失"的防御:若注册表中查不到 execution_id,工具会以 task_failed 返回 "Task disappeared from background tasks",而不是无限等待。

轮询期间工具还会把子代理产生的增量 AI 消息作为 task_running 事件推给前端流(含每条消息的序号与累计 token 用量),也就是说"后端轮询"同时承担了进度透传的职责,父 LLM 无需感知即可让用户侧看到任务进展。

变更三:task_status 不再暴露给 LLM

文档说明:task_status_tool 不再暴露给 LLM,仅保留在代码库中供潜在的内部/调试用途。当前仓库的工具注册表印证了这一点,tools.py#L35-L38

SUBAGENT_TOOLS = [
    task_tool,
    # task_status_tool is no longer exposed to LLM (backend handles polling internally)
]

从源码结构看,还有一层防止递归套娃的设计:子代理自身的工具组装会传入 subagent_enabled=False(见 task_tool.py#L836-L848),且 SubagentConfigdisallowed_tools 默认为 ["task"]——子代理既不会递归调用 task 再开子代理,也拿不到任何状态查询类工具,其可见世界被收敛为"完成分配给我的任务并写出报告"。

变更四:同步更新提示词与文档

文档提到已更新 prompt.py 中的 SUBAGENT_SECTION,删除所有关于后台任务与轮询的措辞,简化使用示例,并明确说明工具会自动等待完成。与之配套的是 task 工具自身的 docstring(它由 @tool("task", parse_docstring=True) 解析后成为 LLM 可见的工具描述),当前版本在 task_tool.py#L643-L721 中给出了完整的委派决策准则:

  • 该委派的情形:可真正并行且相互独立的任务(实质性缩短墙钟时间)、子代理具备主代理路径没有的专家工具/技能/模型、需要上下文隔离的高上下文量调查;
  • 不该委派的情形:仅因任务复杂/多步骤/仓库大、把相互依赖的步骤拆给并行子代理、并行写重叠文件或有共享可变状态、需要用户交互澄清;
  • 成本提醒:重复的仓库探索、结果协调与验证、主代理用直接工具更便宜的活。

这套"何时委派"的提示词约束与后端轮询改进是同一枚硬币的两面:既然一次工具调用就能拿到完整结果,提示词就把 LLM 的注意力从"怎么管理异步"转移到"值不值得委派"上。

实现纵深:后台任务注册表与执行生命周期

状态机

子代理执行的状态由 SubagentStatus 定义:PENDING → RUNNING → {COMPLETED, FAILED, CANCELLED, TIMED_OUT},其中后四种为终态(is_terminal)。每个执行对应一个 SubagentResult 数据对象,除状态与结果文本外,还携带 trace_id(关联父子日志)、ai_messages(增量消息流)、token_usage_records(累计 token 用量)、stop_reason(被护栏截断的原因,如 token_capped/turn_capped)等字段。

注册表

进程级的后台任务注册表定义在 executor.py#L350-L354

# 后台任务结果的全局存储
_background_tasks: dict[str, SubagentResult] = {}
_background_tasks_lock = threading.Lock()
_background_futures: dict[str, Future[SubagentResult]] = {}

配套的操作函数围绕"并发安全 + 不泄漏内存"设计:

函数 位置 行为要点
get_background_task_result(execution_id) executor.py#L1861-L1871 加锁读取注册表条目;工具轮询循环每次迭代调用它
request_cancel_background_task(execution_id) executor.py#L1838-L1858 设置 cancel_event 并尝试 Future.cancel();子代理在 agent.astream() 迭代中协作式检查该事件,在下个迭代边界停止
cleanup_background_task(execution_id) executor.py#L1884-L1914 仅移除终态条目,避免与后台执行器仍在写入的条目发生竞态;非终态条目只打日志不删
force_cleanup_background_task(execution_id) executor.py#L1917-L1934 无条件移除的最后手段,仅在"条目存在但状态对象已不可读"的中断兜底路径使用

值得注意的细节:execute_async 返回的是服务端生成的 UUIDexecutor.py#L1756),而非 provider 的 tool-call ID——后者可能在不同父运行中重复,不能用作进程级注册表键;tool-call ID 仅作为 external_task_id 保留用于日志关联。

启动与超时:execute_async

executor.execute_async 的完整逻辑:

def execute_async(self, task: str, task_id: str | None = None) -> str:
    execution_id = str(uuid.uuid4())
    result = SubagentResult(task_id=execution_id, external_task_id=task_id,
                            trace_id=self.trace_id, status=SubagentStatus.PENDING)
    ...
    with _background_tasks_lock:
        _background_tasks[execution_id] = result

    async def run_with_timeout() -> SubagentResult:
        try:
            return await asyncio.wait_for(
                self._aexecute(task, result),
                timeout=self.config.timeout_seconds,   # ← 执行级超时
            )
        except TimeoutError:
            result.cancel_event.set()
            result.try_set_terminal(SubagentStatus.TIMED_OUT,
                error=f"Execution timed out after {self.config.timeout_seconds} seconds", ...)
            return result
        except asyncio.CancelledError:
            ...  # 置为 CANCELLED
        except Exception as exc:
            ...  # 置为 FAILED

    execution_future = _submit_to_isolated_loop_in_context(parent_context, run_with_timeout)
    ...
    return execution_id

这里体现了文档所说的"双层超时保护":

  1. 执行超时(第一层):asyncio.wait_forconfig.timeout_seconds 为上限包裹真正的执行协程;超时后主动置 cancel_event 并落 TIMED_OUT 终态;
  2. 轮询超时(第二层):工具侧的 max_poll_count = (timeout_seconds + 60) // 5 兜底。源码注释直言它是"以防线程池超时不生效"的安全网(task_tool.py#L1104-L1106)。触发时会请求协作取消、调度延迟清理,并返回 polling_timed_out 状态让主代理知道后台任务可能卡死。

两层保护保证:即使子代理执行完全挂起,系统也绝不会无限等待。

超时配置:从文档值到当前仓库值

原文档给出的配置快照(改进时的默认值):

# packages/harness/deerflow/subagents/config.py
@dataclass
class SubagentConfig:
    # ...
    timeout_seconds: int = 300  # 默认 5 分钟

当前仓库的 SubagentConfig 已演进:

@dataclass
class SubagentConfig:
    name: str
    description: str
    system_prompt: str | None = None
    tools: list[str] | None = None
    disallowed_tools: list[str] | None = field(default_factory=lambda: ["task"])
    skills: list[str] | None = None
    model: str = "inherit"          # "inherit" 表示沿用父代理模型
    max_turns: int = 50
    timeout_seconds: int = 900     # 15 分钟(裸兜底值)

结合其 docstring 与 config.example.yaml#L1540-L1546 的说明,当前的生效规则是:内置子代理的有效执行上限取全局 subagents.timeout_seconds默认 1800 秒 = 30 分钟),由注册表层叠加;SubagentConfig.timeout_seconds 的 900 仅在不存在不同的全局值时作为裸兜底。自定义子代理则用各自的 timeout_seconds(默认 900),可按需在配置中覆盖,例如示例中的深研型代理设为 2700 秒、命令执行型设为 300 秒(见 config.example.yaml#L1572-L1582)。

# config.example.yaml 中 subagents 段(节选)
# subagents:
#   # 内置子代理的默认超时(秒)(默认: 1800 = 30 分钟)
#   # 自定义代理使用各自的 timeout_seconds(默认 900),除非显式覆盖。
#   timeout_seconds: 1800
#   # 可选:全局 max-turns 覆盖,作用于所有子代理
#   # max_turns: 120

由于轮询上限由该值推导((timeout_seconds + 60) // 5),调大执行超时会同步放宽轮询预算,两层保护始终对齐。

执行面架构:从"双线程池"到"持久隔离事件循环"

文档描述改进时期采用的是双线程池方案以避免嵌套线程池:

  • 调度池_scheduler_pool,max 4 workers):编排后台任务执行,运行 run_task() 管理任务生命周期;
  • 执行池_execution_pool,max 8 workers):真正运行子代理并执行超时控制。

工作方式:

# execute_async() 中
_scheduler_pool.submit(run_task)          # 提交编排任务

# run_task() 中
future = _execution_pool.submit(self.execute, task)  # 提交执行
exec_result = future.result(timeout=timeout_seconds)  # 带超时地等待

文档列出的收益:调度与执行关注点分离、无嵌套线程池、超时施加在正确的层级、资源利用更好。

从当前源码结构看,这一"编排/执行分离"的设计意图被保留,但实现载体换成了持久隔离事件循环(persistent isolated subagent event loop),定义在 executor.py#L560-L643:一个专用守护线程上运行着一个进程级长生命周期的 asyncio 事件循环,所有子代理执行(无论同步 execute() 还是异步 execute_async())都通过 _submit_to_isolated_loop_in_contextrun_coroutine_threadsafe 提交到该循环,并保留 ContextVar 状态(trace、认证上下文等)。这样做的好处与文档所述的线程池收益一脉相承:

  • 共享的异步客户端(HTTP client 等)不会被绑定到每次调用即销毁的短命循环上;
  • 超时同样施加在正确的层级:异步路径用 asyncio.wait_for,同步路径用 future.result(timeout=...)executor.py#L1678-L1700);
  • 循环在 atexit 中注册了优雅关闭(executor.py#L582-L616)。

可以推断:文档中的"两个池"是当时解决嵌套线程池问题的形态;当前形态把"编排面"(父运行循环上的工具协程轮询)与"执行面"(隔离循环上的 _aexecute)分离得更彻底,延迟清理等跨循环生命周期管理也因此有了稳定的宿主(见下节)。

中断、取消与延迟清理:文档之外的健壮性补充

当前实现在文档骨架之外还补齐了完整的异常兜底,这部分对理解"工具调用阻塞直到完成"的边界条件很有价值(均位于 task_tool.py#L191-L483):

  • 父运行被取消asyncio.CancelledError):工具先 request_cancel_background_task 协作取消子代理,再经 _finalize_interrupted_subagent 有界地等待终态、把最终 token 用量快照上报给父运行日志(RunJournal),然后移除注册表条目;
  • 轮询器自身意外出错:走同样的收尾路径,但只等一个短的宽限期(_UNEXPECTED_EXIT_GRACE_SECONDS = 5.0 秒,task_tool.py#L56-L63),避免父运行被一个子代理的长超时拖住;
  • 延迟清理:若等待结束时子代理仍非终态,清理任务被调度到进程级持久子代理循环上执行(run_on_isolated_subagent_loop),使其能存活过同步调用路径中 asyncio.run() 的循环拆除;对持续不可读的状态对象还有 force_cleanup_background_task 兜底。

这些机制共同保证:无论父运行是正常完成、被用户取消还是崩溃,后台子代理条目都不会在注册表中永久泄漏,且子代理已消耗的 token 尽可能完整入账。

测试与验证

原文档给出的验证方式(可直接照做):

  1. 启动一个耗时数秒的子代理任务;
  2. 验证工具调用阻塞直到完成;
  3. 验证结果被直接返回;
  4. 验证全程没有 task_status 调用。

示例测试场景:

# 这应该阻塞约 10 秒然后返回结果
result = task(
    subagent_type="bash",
    prompt="sleep 10 && echo 'Done'",
    description="Test task"
)
# result 应包含 "Done"

当前仓库中有对应的自动化测试资产可以进一步印证行为契约:test_task_tool_core_logic.py 覆盖任务工具的核心逻辑,test_task_tool_usage_recorder.py 覆盖子代理用量上报,test_subagent_executor.pytest_subagent_timeout_config.py 覆盖执行器与超时配置。读者可以在仓库中按这些文件名检索,查看轮询、超时与清理路径各自的断言细节。

迁移说明

对于此前使用 run_in_background=True 的用户或代码,文档给出的迁移步骤:

  • 删除该参数;
  • 删除所有自研的轮询逻辑;
  • 工具会自动等待完成,无需其他改动——API 保持向后兼容(仅移除了该参数)。

小结

这次改进的本质是一次职责下放:把"任务结束了吗"这一高频、低信息量的判断从 LLM 推理循环中剥离,交给确定性的后端轮询循环。收益具体而可验证——LLM 对每个子代理任务恰好一次工具调用;token 与 API 成本下降;状态检查行为一致;执行与轮询双层超时保证不无限等待。而当前仓库在该骨架上又向前演进:轮询节奏与上限改为由 subagents.timeout_seconds(默认 30 分钟)统一推导、终态扩展为四态、执行面迁移到持久隔离事件循环、并补齐了取消/崩溃/泄漏三类兜底路径。理解这条演进线,对阅读 DeerFlow 任何与子代理生命周期相关的代码(task 工具执行器配置定义)都会大有裨益。

登录后查看全文
热门项目推荐
相关项目推荐