DeerFlow 2.0 深度解析:从 Deep Research 到 Super Agent Harness 的架构进化
DeerFlow 2.0 深度解析:从 Deep Research 到 Super Agent Harness 的架构进化
🦌 Super Agent Harness
🔗 LangGraph Orchestration
⚡ Parallel Sub-Agents
🧠 Persistent Memory
DeerFlow(Deep Exploration and Efficient Research Flow)是字节跳动开源的超级代理框架(Super Agent Harness)。2.0 版本是一次彻底的重写,它不再是简单的研究工具,而是一个完整的 Agent 运行时:
- 中间件管道 — 40+ 可组合中间件,从上下文压缩到循环检测
- 子代理并行引擎 — 自动任务分解,多代理并行执行,结构化结果汇聚
- 渐进式技能加载 — 按需加载 SKILL.md,保持上下文窗口精简
- 沙箱执行环境 — 每个任务独立文件系统,支持 Docker/K8s 隔离
- 持久化长期记忆 — 跨会话记忆,Markdown 事实存储,去重与迁移
- 上下文工程 — 激进摘要、工具结果卸载、严格 Tool-Call 恢复
- 生产级运行时 — Redis Stream Bridge、RunJournal 审计、Read-Before-Write 门控
| 维度 | DeerFlow 1.x | DeerFlow 2.0 | AutoGen | CrewAI |
|---|---|---|---|---|
| 定位 | Deep Research | Super Agent Harness | 多代理对话 | 角色协作 |
| 架构 | 固定管道 | 中间件管道 | 对话图 | 顺序/层级 |
| 子代理 | ❌ | ✅ 并行执行 | ✅ 对话式 | ✅ 委派式 |
| 技能系统 | ❌ | ✅ 渐进加载 | ❌ | ❌ |
| 沙箱 | ❌ | ✅ Docker/K8s | ⚠️ 有限 | ❌ |
| 长期记忆 | ❌ | ✅ Markdown 事实 | ⚠️ 插件 | ⚠️ 插件 |
| 上下文工程 | 基础 | 激进压缩 | 基础 | 基础 |
🎯 引言:为什么需要 Super Agent Harness?
想象一个场景:你让 AI “研究量子计算的最新进展,生成一份带图表的报告,同时制作一个演示幻灯片”。
传统 Agent 框架会怎样?要么在单一上下文中挣扎(token 爆炸),要么在固定管道中迷失(无法动态分解)。
DeerFlow 2.0 的回答:不是给你一个框架去组装,而是给你一个完整的运行时——电池全含,即插即用。
- 上下文爆炸 — 长任务中 token 消耗指数增长
- 单点故障 — 一个 Agent 做所有事,失败即全盘崩溃
- 能力固化 — 工具和能力在启动时全部加载,浪费上下文
- 无状态 — 每次对话从零开始,无法积累知识
- 安全真空 — 代码执行没有隔离边界
DeerFlow 的设计哲学:一个 Harness(线束/框架),而非一个 Framework(框架)。区别在于——Framework 让你组装,Harness 让你使用。它自带文件系统、记忆、技能、沙箱感知执行,以及为复杂多步任务进行规划和生成子代理的能力。
🏗️ 整体架构
🦌 DeerFlow 2.0 核心组件
🔬 Lead Agent:中间件管道的艺术
DeerFlow 2.0 最核心的架构创新是中间件管道(Middleware Pipeline)。Lead Agent 不是一个简单的 LLM 调用循环,而是一个由 40+ 中间件组成的精密处理链。
Agent 工厂模式
Lead Agent 通过 make_lead_agent 工厂函数创建,每次请求都会根据运行时配置动态组装:
def make_lead_agent(config: RunnableConfig): """LangGraph graph factory; keep the signature compatible with LangGraph Server.""" runtime_config = _get_runtime_config(config) # ... return _make_lead_agent(config, app_config=runtime_app_config)
def _make_lead_agent(config: RunnableConfig, *, app_config: AppConfig): cfg = _get_runtime_config(config)
# 🔑 动态模型解析:request → agent config → global default model_name = _resolve_model_name( requested_model_name or agent_model_name, app_config=resolved_app_config )
# 🔑 思考模式自适应 thinking_enabled = bool(_resolve_runtime_option( cfg, "thinking_enabled", agent_thinking, True )) if thinking_enabled and not model_config.supports_thinking: thinking_enabled = False # 优雅降级
# 组装最终 Agent return create_agent( model=create_chat_model(name=model_name, thinking_enabled=thinking_enabled, ...), tools=final_tools, middleware=build_middlewares(config, model_name=model_name, ...), system_prompt=apply_prompt_template(...), state_schema=get_thread_state_schema(mode), )模型解析遵循严格的优先级:请求参数 > 自定义 Agent 配置 > 全局默认。
_resolve_runtime_option 使用 key in cfg(而非 cfg.get(key))区分”请求省略了字段”和”请求设置为 falsy 值”——这确保了 thinking_enabled: false 不会被错误地回退到默认值。
中间件链:精密的洋葱模型
# backend/packages/harness/deerflow/agents/lead_agent/agent.py - build_middlewares
def build_middlewares(config, model_name, agent_name=None, ...): middlewares = build_lead_runtime_middlewares(...) # 基础运行时中间件
# 1️⃣ 动态上下文注入(日期 + 记忆) middlewares.append(DynamicContextMiddleware(agent_name=agent_name))
# 2️⃣ 技能激活(/skill-name 斜杠命令) middlewares.append(SkillActivationMiddleware(available_skills=available_skills))
# 3️⃣ 技能工具策略(allowed-tools 过滤) middlewares.append(SkillToolPolicyMiddleware(available_skills=available_skills))
# 4️⃣ 持久上下文(摘要 + 账本 + 技能注入) middlewares.append(DurableContextMiddleware(...))
# 5️⃣ 上下文摘要压缩 if summarization_middleware is not None: middlewares.append(summarization_middleware)
# 6️⃣ 计划模式 TodoList if is_plan_mode: middlewares.append(TodoMiddleware(...))
# 7️⃣ Token 用量追踪 middlewares.append(TokenUsageMiddleware())
# 8️⃣ 标题生成 middlewares.append(TitleMiddleware())
# 9️⃣ 记忆中间件 middlewares.append(MemoryMiddleware(agent_name=agent_name))
# 🔟 图像查看(仅视觉模型) if model_config.supports_vision: middlewares.append(ViewImageMiddleware())
# 1️⃣1️⃣ MCP 路由提示 if mcp_routing_middleware is not None: middlewares.append(mcp_routing_middleware)
# 1️⃣2️⃣ 延迟工具过滤 if deferred_setup.deferred_names: middlewares.append(DeferredToolFilterMiddleware(...))
# 1️⃣3️⃣ SystemMessage 合并(兼容严格后端) middlewares.append(SystemMessageCoalescingMiddleware())
# 1️⃣4️⃣ 子代理并发限制 if subagent_enabled: middlewares.append(SubagentLimitMiddleware(max_concurrent=3, max_total=N))
# 1️⃣5️⃣ 循环检测 middlewares.append(LoopDetectionMiddleware.from_config(...))
# 1️⃣6️⃣ Token 预算硬停 middlewares.append(TokenBudgetMiddleware.from_config(...))
# 1️⃣7️⃣ 自定义扩展中间件 middlewares.extend(load_configured_extension_middlewares(...))
# 1️⃣8️⃣ 终端响应保护 middlewares.append(TerminalResponseMiddleware())
# 1️⃣9️⃣ 安全终止检测 middlewares.append(SafetyFinishReasonMiddleware.from_config(...))
# 2️⃣0️⃣ 澄清拦截(必须最后) middlewares.append(ClarificationMiddleware()) return middlewares💡 中间件顺序的深意
中间件的注册顺序决定了执行顺序。DeerFlow 的设计遵循几个关键约束:
Summarization 在前 — 先压缩上下文,后续中间件处理更少的 token
Clarification 在最后 — 所有处理完成后才拦截澄清请求
Safety 在 Terminal 之后 — LangChain 反序 after_model 调度,Safety 先执行
SystemMessage 合并在工具过滤后 — 确保所有注入完成后再合并
🔀 子代理执行引擎:并行任务的精密编排
子代理系统是 DeerFlow 2.0 最具工程深度的模块。它解决了一个核心问题:如何让多个 Agent 安全、高效、可取消地并行工作?
执行架构
SubagentExecutor 核心源码
class SubagentExecutor: """Executor for running subagents."""
def __init__(self, config, tools, app_config=None, parent_model=None, sandbox_state=None, thread_data=None, thread_id=None, ...): self.config = config self.trace_id = trace_id or str(uuid.uuid4())[:8]
# 🔑 工具过滤:allowlist + denylist 双重机制 self._base_tools = _filter_tools(tools, config.tools, config.disallowed_tools)
# 🔑 守卫中间件收集:用于检测 token/loop cap self._stop_reason_middlewares: list[Any] = []
async def _aexecute(self, task: str, result_holder=None) -> SubagentResult: """异步执行子代理任务""" # 构建初始状态(含技能加载 + 工具授权) state, final_tools, deferred_setup = await self._build_initial_state(task) agent = self._create_agent(final_tools, deferred_setup=deferred_setup)
# Token 收集器 collector = SubagentTokenCollector(caller=f"subagent:{self.config.name}")
run_config: RunnableConfig = { "recursion_limit": self.config.max_turns, # 🔑 turns 预算 "callbacks": [collector], }
# 🔑 流式执行 + 协作式取消 final_state = None async for chunk in agent.astream(state, config=run_config, context=context, stream_mode="values"): # 协作式取消检测 if result.cancel_event.is_set(): result.try_set_terminal(SubagentStatus.CANCELLED, error="Cancelled by user") return result
final_state = chunk # 实时捕获步骤消息 capture_new_step_messages(messages, ai_messages, seen_message_ids, ...)
# 🔑 守卫停止原因检测 stop_reason = self._consume_guard_stop_reason() result.try_set_terminal(SubagentStatus.COMPLETED, result=final_result, stop_reason=stop_reason) return result子代理的终止不是简单的”完成/失败”二元状态,而是一个精密的多层保护:
- Turn 预算 (
recursion_limit) — 防止无限循环 - Token 预算 (
TokenBudgetMiddleware) — 防止 token 爆炸 - 循环检测 (
LoopDetectionMiddleware) — 检测重复工具调用
每种保护触发后,stop_reason 会标记为 turn_capped / token_capped / loop_capped,让 Lead Agent 能区分”正常完成”和”被截断”。
隔离事件循环:解决 async 嵌套难题
子代理执行面临一个经典的 Python 异步难题:父代理已经在事件循环中运行,子代理如何创建自己的异步执行?
# backend/packages/harness/deerflow/subagents/executor.py - 隔离事件循环
# 🔑 全局持久化事件循环(避免每次执行创建/销毁)_isolated_subagent_loop: asyncio.AbstractEventLoop | None = None_isolated_subagent_loop_thread: threading.Thread | None = None
def _get_isolated_subagent_loop() -> asyncio.AbstractEventLoop: """返回用于隔离子代理执行的持久化事件循环""" global _isolated_subagent_loop, _isolated_subagent_loop_thread
with _isolated_subagent_loop_lock: # 健康检查:线程存活 + 循环未关闭 + 循环运行中 loop_is_usable = ( _isolated_subagent_loop is not None and not _isolated_subagent_loop.is_closed() and _isolated_subagent_loop.is_running() and thread_is_alive )
if not loop_is_usable: loop = asyncio.new_event_loop() started_event = threading.Event() thread = threading.Thread( target=_run_isolated_subagent_loop, args=(loop, started_event), name="subagent-persistent-loop", daemon=True, ) thread.start() # 🔑 等待循环就绪(5秒超时) if not started_event.wait(timeout=5): raise RuntimeError("Timed out starting isolated subagent event loop") _isolated_subagent_loop = loop
return _isolated_subagent_loop🔄 为什么需要隔离事件循环?
Python 的 asyncio 不允许在已运行的事件循环中调用 asyncio.run()。DeerFlow 的解决方案:
关键细节:copy_context() 确保子代理继承父代理的 ContextVar 状态(如 user_id、trace_id),而不会污染父循环。
线程安全的终端状态:try_set_terminal
# backend/packages/harness/deerflow/subagents/executor.py - SubagentResult
@dataclassclass SubagentResult: task_id: str trace_id: str status: SubagentStatus result: str | None = None error: str | None = None stop_reason: str | None = None # 🔑 token_capped / turn_capped / loop_capped cancel_event: threading.Event = field(default_factory=threading.Event) _state_lock: threading.Lock = field(default_factory=threading.Lock)
def try_set_terminal(self, status, *, result=None, error=None, stop_reason=None, ...) -> bool: """设置终端状态 — 恰好一次(exactly-once 语义)""" if not status.is_terminal: raise ValueError(f"Status {status} is not terminal")
with self._state_lock: # 🔑 第一个终端转换获胜;后续写入被忽略 if self.status.is_terminal: return False
self.result = result self.error = error self.stop_reason = stop_reason self.status = status return True后台超时/取消和执行工作线程可能在同一个 result holder 上竞争。try_set_terminal 通过 _state_lock + “先到先得”语义确保:
- 超时线程设置
TIMED_OUT→ 执行线程的后续COMPLETED被忽略 - 执行线程设置
COMPLETED→ 超时线程的后续TIMED_OUT被忽略
这避免了”已完成的任务被错误标记为超时”的 bug。
📦 技能系统:渐进式能力加载
DeerFlow 的技能系统是其区别于其他 Agent 框架的标志性设计。技能不是代码插件,而是结构化的 Markdown 能力模块。
技能目录结构
/mnt/skills/├── public/ ← 内置技能(只读)│ ├── research/SKILL.md│ ├── report-generation/SKILL.md│ ├── slide-creation/SKILL.md│ ├── web-page/SKILL.md│ └── image-generation/SKILL.md├── custom/ ← 用户自定义技能│ └── your-skill/SKILL.md└── integrations/ ← 托管集成(只读) └── lark-cli/lark-doc/SKILL.md渐进式加载:Deferred Discovery
传统做法是将所有工具 schema 塞入系统提示。DeerFlow 2.0 引入了**延迟发现(Deferred Discovery)**机制:
# backend/packages/harness/deerflow/skills/catalog.py - SkillCatalog
@dataclass(frozen=True)class SkillCatalog: """不可变技能目录 — 纯搜索,无变更""" skills: tuple[Skill, ...]
def search(self, query: str) -> list[Skill]: """三种查询形式(镜像 DeferredToolCatalog)"""
# 1️⃣ 精确选择:"select:data-analysis,deep-research" if query.startswith("select:"): wanted = {n.strip() for n in query[7:].split(",")} return [s for s in self.skills if s.name in wanted]
# 2️⃣ 必选前缀:"+podcast gen" → 名称必含 podcast,按 gen 排序 if query.startswith("+"): parts = query[1:].split(None, 1) required = parts[0].lower() candidates = [s for s in self.skills if required in s.name.lower()] # ...排序逻辑 return candidates[:MAX_RESULTS]
# 3️⃣ 自由文本正则:"chart visualization" regex = _compile_catalog_regex(query) scored = [] for s in self.skills: searchable = f"{s.name} {s.description or ''}" if regex.search(searchable): # 🔑 名称匹配得分高于描述匹配 scored.append((2 if regex.search(s.name) else 1, s)) scored.sort(key=lambda x: x[0], reverse=True) return [s for _, s in scored][:MAX_RESULTS]🎯 渐进式加载的三层架构
<skill_index>),不含完整描述describe_skill 工具获取元数据/skill-name 或 read_file 读取完整 SKILL.md效果:系统提示保持精简(prefix-cache 友好),同时不牺牲能力发现。对于 token 敏感的模型尤为重要。
技能安全:SkillScan
安装外部技能存在安全风险。DeerFlow 引入了 SkillScan — 一个原生确定性安全扫描器:
# 技能安装流程中的安全检查# Phase 1: 离线确定性扫描(无 Semgrep 依赖)# - 阻止高置信度 CRITICAL 发现(私钥、shell 执行)# - 警告级发现传递给 LLM 扫描器进行上下文审查
# Phase 2: LLM 上下文审查# - 对 Phase 1 的警告进行语义分析# - 判断是否为真正的安全威胁SkillScan 是最佳努力的行为范围界定,不是硬安全边界:
- 通过其他工具加载技能指令不会被捕获
- 活跃技能条目可能从有界上下文中被驱逐
- 注册表失败时,fail-closed 到框架安全工具
🧠 上下文工程:四层防线对抗 Token 爆炸
长任务中,上下文管理是 Agent 框架的生死线。DeerFlow 2.0 构建了四层递进防线:
DeerFlowSummarizationMiddleware:摘要模型分离
class DeerFlowSummarizationMiddleware(AgentMiddleware): """上下文压缩 — 摘要模型与运行模型分离"""
def __init__(self, model, *, summarize_model=None, ...): # 🔑 关键设计:摘要用便宜模型,运行用强模型 self._summarize_model = summarize_model or model self._trigger_tokens = trigger_tokens # 触发压缩的阈值 self._keep_recent_tokens = keep_recent # 保留最近 N token
async def abefore_model(self, state, runtime): messages = state["messages"] total_tokens = estimate_token_count(messages)
if total_tokens < self._trigger_tokens: return None # 未触发,跳过
# 🔑 before_summarization hooks(技能账本、委派记录保存) context = ContextCompactionResult() for hook in self._before_summarization_hooks: hook(messages, context)
# 压缩:保留 system + 最近消息,中间部分生成摘要 summary = await self._summarize_model.ainvoke([ SystemMessage("Summarize the conversation concisely..."), *messages_to_summarize ])
# 🔑 摘要结果注入为 SystemMessage(name="summary") return {"messages": [RemoveMessage(id=m.id) for m in old] + [summary_msg]}摘要模型和运行模型分离是一个精妙的成本优化:
- 运行模型(如 Claude Sonnet)负责推理和工具调用
- 摘要模型(如 GPT-4o-mini)负责压缩历史——成本降低 10-50x
context_compaction.py 中的 compact_thread_context 还实现了模型优先级链:请求指定 > 配置摘要模型 > 运行模型回退。
ToolOutputBudget:超大结果磁盘卸载
class ToolOutputBudgetMiddleware(AgentMiddleware): """工具输出预算 — 防止单次工具调用耗尽上下文"""
# 策略: # 1. 结果 > budget → 写入沙箱磁盘文件 # 2. 替换为 head + tail 截断预览 + 文件路径引用 # 3. 回退:若磁盘写入失败,纯内存截断 # 4. 生成 ToolOutputSynopsis(确定性,无 LLM)ToolOutputSynopsis:确定性结构摘要
_MAX_SYNOPSIS_INPUT_BYTES = 5_000_000 # 🔑 5MB DoS 防护上限
def build_tool_output_synopsis(content: str, *, tool_name: str = "") -> ToolOutputSynopsis: """无 LLM 的确定性摘要 — 按内容类型差异化处理""" # 自动检测:JSON / CSV / YAML / XML / Code / Text # JSON → 结构形状 + 键统计 + 采样 # CSV → 列名 + 行数 + 前 50 行采样 # Code → import 列表 + 符号列表 # XML → defusedxml 防实体扩展攻击 # 超过 5MB → 仅 head/tail 原始采样(防 DoS)🔄 循环检测:双层防护
Agent 最常见的失败模式是工具调用循环——反复调用相同工具却不改变策略。DeerFlow 实现了精密的双层检测:
class LoopDetectionMiddleware(AgentMiddleware): """双层循环检测:Hash 精确匹配 + 频率统计"""
# 🔑 Layer 1: 精确哈希检测 # 连续 N 次完全相同的工具调用序列 → 判定循环 def _hash_tool_calls(self, tool_calls) -> str: """MD5 排序哈希 — 对参数排序后取哈希""" normalized = sorted( json.dumps(tc, sort_keys=True) for tc in tool_calls ) return hashlib.md5("|".join(normalized).encode()).hexdigest()
# 🔑 Layer 2: 频率检测 # 同一工具在窗口内调用频率过高 → 判定循环 # 即使参数不完全相同(如每次改一个字符)_stable_tool_key:按工具类型差异化
def _stable_tool_key(tool_call: dict) -> str: """生成稳定的工具调用标识 — 不同工具不同策略""" name = tool_call.get("name", "") args = tool_call.get("args", {})
if name == "read_file": # 🔑 按 200 行桶分组 — 读同一区域视为相同 start = args.get("start_line", 0) bucket = start // 200 return f"read_file:{args.get('path')}:{bucket}"
if name in ("write_file", "str_replace"): # 写操作:路径 + 内容哈希 content_hash = hashlib.md5(str(args.get("content", "")).encode()).hexdigest()[:8] return f"{name}:{args.get('path')}:{content_hash}"
# 默认:名称 + 完整参数哈希 return f"{name}:{hashlib.md5(json.dumps(args, sort_keys=True).encode()).hexdigest()[:12]}"循环检测的警告消息通过 wrap_model_call 注入(而非 after_model),原因是:
after_model追加的消息会被add_messagesreducer 放到末尾- 这会打断 AIMessage → ToolMessage 的严格配对
- 严格 Provider(如 OpenAI)会因配对断裂返回 HTTP 400
wrap_model_call 可以在正确位置(紧跟 AIMessage 之后)插入警告。
💾 长期记忆:跨会话知识积累
class MemoryMiddleware(AgentMiddleware): """per-agent 隔离的持久化记忆"""
def __init__(self, agent_name: str | None = None): # 🔑 每个 agent 有独立的记忆命名空间 self._namespace = f"memory:{agent_name or 'lead'}"
async def aafter_agent(self, state, runtime): """Agent 完成后异步提取记忆(不阻塞响应)""" # 1. 从最终对话中提取关键事实 # 2. 去重:与已有记忆比对 # 3. 写入 Markdown 事实存储 # 4. Debouncing:短时间内多次完成只触发一次- Markdown 事实存储 — 人类可读、可审计、可手动编辑
- per-agent 隔离 — Lead Agent 和子代理的记忆互不污染
- Debouncing — 避免高频写入(Goal 续跑时多轮完成只提取一次)
- DynamicContextMiddleware 注入 — 下次对话时自动将相关记忆注入系统提示
🎯 Goal 系统:自主续跑与无进展检测
DeerFlow 2.0 引入了 Goal 概念——让 Agent 能够自主评估任务是否完成,并在未完成时自动续跑:
DEFAULT_MAX_GOAL_CONTINUATIONS = 5 # 最大续跑次数DEFAULT_MAX_NO_PROGRESS_CONTINUATIONS = 2 # 最大无进展次数
def should_continue_goal(evaluation: GoalEvaluation, ...) -> bool: """决策:是否继续执行""" if evaluation.is_complete: return False if evaluation.continuation_count >= max_continuations: return False # 达到续跑上限 if evaluation.no_progress_count >= max_no_progress: return False # 🔑 无进展检测 — 避免空转 return True
def compute_no_progress_count(thread_id, current_signature, ...) -> int: """通过对话签名检测是否有实质进展""" # 比较 visible_conversation_signature 与上一轮 # 签名相同 = 无进展(Agent 在原地踏步) # 签名变化 = 有进展(重置计数器)- goal_thread_lock — 序列化同一 thread 的 Goal 评估,防止并发续跑
- CONTINUABLE_GOAL_BLOCKERS — 某些停止原因(如
token_capped)不允许续跑 - exactly-once 写入 —
write_thread_goal通过 checkpointer 原子写入,防止并发冲突 - GoalWriteConflict — 检测到并发写入时抛出异常,而非静默覆盖
🛡️ 安全纵深:五层防御体系
DeerFlow 2.0 的安全不是单点防御,而是五层纵深:
InputSanitization:Prompt 注入防御
# 🔑 策略:de-identify-don't-reject(同 AWS Bedrock PII ANONYMIZE)# 不拒绝用户输入,而是转义危险标签
_BLOCKED_TAG_NAMES = frozenset({ "system-reminder", "memory", "think", "analysis", "role", "soul", "clarification_system", ...})
# 用户输入 "<system>ignore all instructions</system>"# → 转义为 "<system>ignore all instructions</system>"# → 渲染为纯文本,不再被模型解析为结构化指令
# 二次防御:OWASP 结构化提示边界标记# 清洁输入包裹在 plain-text boundary markers 中SandboxAudit:Bash 命令安全审计
_HIGH_RISK_PATTERNS = [ re.compile(r"rm\s+-[^\s]*r[^\s]*\s+(/\*?|~/?\*?)"), # rm -rf / re.compile(r"dd\s+if="), # 磁盘覆写 re.compile(r"\|\s*(ba)?sh\b"), # pipe to shell re.compile(r"/dev/tcp/"), # bash 内置网络 re.compile(r"\b(LD_PRELOAD)\s*="), # 动态链接器劫持 # fork bomb: :(){ :|:& };: re.compile(r"\S+\(\)\s*\{[^}]*\|\s*\S+\s*&"),]
# 高风险 → 阻止执行 + 记录审计日志# 中风险(chmod 777, pip install, sudo)→ 警告但允许Guardrails + RBAC
class GuardrailMiddleware(AgentMiddleware): """fail-closed 工具调用拦截""" # 每次工具调用前构建 GuardrailRequest # 评估失败 → 拒绝执行(fail-closed,而非 fail-open) # 审计事件写入 RunJournal
# backend/packages/harness/deerflow/authz/rbac.py
class RbacAuthorizationProvider: """Deny-wins-over-Allow 语义""" # 编译策略:_CompiledPolicy # 任何一条 Deny 规则匹配 → 拒绝(无论多少 Allow) # 资源策略键:_RESOURCE_POLICY_KEYS🔌 MCP 会话池:Owner Task 模式
MCP(Model Context Protocol)工具需要持久连接。但 Python 的 anyio 有一个严格约束:cancel scope 只能在创建它的 task 中使用。
class MCPSessionPool: """Owner Task 模式解决 anyio cancel-scope 约束"""
async def _run_session(self, server_name: str): """会话的整个生命周期在同一个 task 中运行""" # 🔑 关键:连接、使用、关闭都在这个 task 内 # 外部通过 queue 发送请求,而非直接调用 session async with mcp_client_session(...) as session: while True: request = await self._request_queue.get() result = await session.call_tool(...) await self._response_queue.put(result)
# 🔑 LRU 驱逐:空闲会话超时后优雅关闭 # 🔑 inflight 去重:同一 server 的并发请求复用同一会话anyio 的 cancel scope 是 task-local 的。如果 session 在 task A 中创建,但在 task B 中使用,会触发 RuntimeError: cancel scope accessed from different task。
Owner Task 模式的解法:每个 session 有一个专属 task,外部通过消息队列与其通信。这样 cancel scope 始终在同一个 task 内闭合。
⚡ 生产级优化:从 Demo 到 Production 的鸿沟
以下模块是 DeerFlow 从“能跑”到“能上线”的关键差距填充。每一个都对应一个真实的生产事故场景。
StreamBridge:Redis 跨进程 SSE
class RedisStreamBridge(StreamBridge): """跨进程事件流 — 多 worker 部署的必备"""
supports_cross_process = True # 🔑 与 MemoryStreamBridge 的关键区别 _XREAD_COUNT = 64 # 每次读取批量 _MAX_SUBSCRIBE_RETRIES = 3 # 重连上限
async def subscribe(self, run_id: str, *, last_event_id: str | None = None): """支持 Last-Event-ID 重连""" # 🔑 客户端断线重连时,从上次位置继续读取 # 不丢失事件,不重复事件 # HEARTBEAT_SENTINEL 保持连接活跃 # END_SENTINEL 标记流结束单进程 MemoryStreamBridge 在 uvicorn --workers 4 时会失败:客户端 SSE 连接到 worker A,但 Agent 运行在 worker B。Redis Streams 提供跨进程的事件总线,并通过 XREAD 的 last_event_id 语义实现断线重连。
RunJournal:全链路审计
# backend/packages/harness/deerflow/runtime/journal.py (883 行)
class RunJournal: """LangChain callback → RunEvent 标准化审计"""
# 🔑 作为 LangChain CallbackHandler 挂载 # 自动捕获: # - LLM 调用开始/结束 + token 用量 # - 工具调用开始/结束 + 结果 # - 中间件事件(guardrail 拦截、循环检测等) # - 错误和异常
# 🔑 累积 token 用量(按模型分组) # token_usage_by_model: {"claude-sonnet": {"input": 12000, "output": 3000}}
# 🔑 批量写入 RunEventStore(避免 per-step I/O) async def flush(self): await self._event_store.put_batch(self._buffer)ReadBeforeWrite:SHA-256 版本门控
class ReadBeforeWriteMiddleware(AgentMiddleware): """防止 'append-only, never read back' 的重复输出"""
# 问题:Agent 向文件追加同一内容 5 次(因为它从不回读) # 解法:修改已有文件前必须先 read_file
# 🔑 设计不变量: # 1. read_file 时在 ToolMessage.additional_kwargs 盖上 SHA-256 mark # 2. 摘要删除 read 结果 → mark 也消失 → 门控不通过 # 3. 写操作永远不刷新 mark(写改变哈希,强制重新读取) # 4. per-(scope, path) 序列化:防止同轮并发写都通过同一旧 mark # 5. Fail-open:沙箱故障时不阻止工具执行
_BLOCK_MESSAGE = ( "Error: {tool_name} blocked — {path} already exists and you have not " "read its current version. Call read_file on it, then retry." )ToolProgress:结果质量状态机
class ToolProgressMiddleware(AgentMiddleware): """per-(thread, tool) 结果质量状态机"""
# 状态转换:ACTIVE → WARNED → BLOCKED # # ACTIVE: 工具正常返回结果 # WARNED: 工具连续返回空/错误结果 → 注入警告 # BLOCKED: 警告后仍无改善 → 阻止调用 # # 🔑 与 LoopDetection 的分工: # LoopDetection = 检测“重复相同调用” # ToolProgress = 检测“调用不同参数但结果始终无效”LLMErrorHandling:进程级并发控制
# backend/packages/harness/deerflow/agents/middlewares/llm_error_handling_middleware.py (962 行)
class _ProcessWideLimiter: """跨事件循环的 LLM 调用并发限制器"""
# 🔑 为什么不用 asyncio.Semaphore? # Semaphore 绑定到第一个使用它的事件循环 # Lead Agent 在主循环,子代理在隔离循环 # → 必须用 threading 原语实现跨循环限制
# 🔑 关键设计: # - Cap 不可变:启动时解析一次,永不修改 # - Lossless waiter handoff:取消时许可不丢失 # - call_soon_threadsafe 跨循环唤醒
# 智能重试分类:_BUSY_PATTERNS = ("server busy", "服务繁忙", "稍后重试", ...)_QUOTA_PATTERNS = ("insufficient_quota", "余额不足", ...)_BURST_PATTERNS = ("limit_burst_rate", "请求速率增长过快", ...)
# 🔑 Per-exception 重试预算:_RETRY_BUDGET_OVERRIDES = { "StreamChunkTimeoutError": 2, # 已等待 120s,只重试 1 次}DanglingToolCall:消息完整性修复
class DanglingToolCallMiddleware(AgentMiddleware): """修复悬挂工具调用 + 孤儿 ToolMessage"""
# 问题场景: # 1. 用户中断 → AIMessage 有 tool_calls 但无对应 ToolMessage # 2. 摘要删除了 AIMessage → ToolMessage 成为孤儿 # 3. 严格 Provider 拒绝:HTTP 400 "tool_call_id not found"
# 解法: # - 为每个悬挂 tool_call 插入合成 ToolMessage (status="error") # - 删除无匹配 tool_call 的孤儿 ToolMessage # - 清理畸形的 tool-call name/args # 🔑 使用 wrap_model_call 在正确位置插入(而非末尾追加)其他生产级优化
🏭 生产级优化矩阵
批量 32 chunks 发射,避免浏览器对增长 JSON 的二次方解析。write_file 参数流从每 token 一次变为每 32 tokens 一次。
FLUSH_THRESHOLD=25,批量 put_batch 持久化。避免每个子代理步骤都触发一次数据库写入锁。
请求级密钥通过 ContextVar 传播,永不进入 prompt。工具通过 owner_token 访问,而非字符串拼接。
Gateway 重启后,租约过期的运行被标记为 orphan_recovered。SQLite 写入有 5 次指数退避重试。
确定性委派追踪。状态指导:“已完成,不要重新委派” / “已失败,可修改计划后重试”。防止重复委派同一任务。
检测 Provider 的 finish_reason=‘length’ 截断。保留内容用于审计,但标记 stop_reason 为 model_length_capped。
🌐 Gateway 请求生命周期
一个完整的用户请求如何流经 DeerFlow:
💡 设计哲学总结
🦌 DeerFlow 2.0 设计哲学
DeerFlow 2.0 的核心洞察:Agent 框架的难点不在 LLM 调用,而在 LLM 调用之间的一切——上下文管理、错误恢复、并发控制、安全边界、状态一致性。这些“不性感”的工程问题,才是 Demo 和 Production 的鸿沟。
🚀 快速上手
# 克隆项目git clone https://github.com/bytedance/deer-flow.gitcd deer-flow
# 安装依赖cd backend && uv sync && cd ..cd frontend && pnpm install && cd ..
# 配置模型cp backend/.env.example backend/.env# 编辑 .env 设置 LLM_API_KEY 和 LLM_MODEL
# 启动后端cd backend && uv run uvicorn app.main:app --reload
# 启动前端cd frontend && pnpm devbackend/packages/harness/deerflow/agents/— Agent 核心(中间件、Lead Agent、子代理)backend/packages/harness/deerflow/runtime/— 运行时(StreamBridge、Goal、Journal)backend/packages/harness/deerflow/mcp/— MCP 会话池backend/packages/harness/deerflow/sandbox/— 沙箱抽象backend/app/— FastAPI Gateway + IM 通道frontend/— Next.js 前端
本文基于 DeerFlow 2.0 源码分析,项目仍在快速迭代中。建议结合源码阅读,感受每个设计决策背后的工程权衡。
文章分享
如果这篇文章对你有帮助,欢迎分享给更多人!



