- docs/02 架构设计: 新增 aichat审查/异步审批构想/流式渲染调研/generating状态机/密钥迁移健壮性/工作流脚本执行边界/条件表达式引擎/F-07 trait下沉/Agent架构说明/任务推进构想/功能创意池;更新功能决策记录+归档/对抗论证/文档记录规范/经验记录
- docs/03 模块文档: 新增 AI对话引擎/DAG引擎详解;更新 df-knowledge/df-nodes/df-storage/df-workflow/df-ai
- docs/05 代码审查: 新增 全栈审查/全局review/架构审查/近期改动审查/工作区多角度走查/自研memo流式渲染审查
- docs/09 问题排查: 新增 aichat-apikey-401
- docs/INDEX+README 索引同步;docs/todo 待办看板(2026-06-15 汇总)
- PROGRESS.md Sprint 22-25;URGENT.md 加急清单快照(5 项 P0 已全修)
- scripts/cleanup_orphan_tasks.{py,sh} 孤儿任务清理工具
- .gitignore 补 *.broken.bak + tmp/ 噪音排除
15 KiB
DAG 引擎详解
来源:基于
crates/df-workflow/src/实际代码核对编写,2026-06-14 关联模块文档:df-workflow-工作流引擎-2026-06-12.md(接口级参考) 代码版本:Sprint 22 收尾(StateMachine Arc 共享 + 审批取消闭环已落地)
一、什么是 DAG
DAG = Directed Acyclic Graph(有向无环图)。DevFlow 用它来编排工作流——把一个复杂任务拆成多个节点,按依赖关系连接,引擎自动算出执行顺序。
┌─────────┐
│ 检查环境 │ ← 第 1 层(无依赖,先执行)
└────┬────┘
│
┌────┴────┐
│ 运行测试 │ ← 第 2 层(依赖「检查环境」完成)
└────┬────┘
│
┌────┴────┐
│ 构建产物 │ ← 第 3 层(依赖「运行测试」完成)
└─────────┘
也可以有分支和并行:
┌──────────┐
│ 代码扫描 │
└──┬───┬───┘ ← 第 1 层
│ │
┌────┘ └────┐
▼ ▼
┌──────┐ ┌──────────┐
│ 单元测试 │ │ AI 代码审查 │ ← 第 2 层(两个并行执行)
└──┬───┘ └────┬─────┘
│ │
└──────┬───────┘
▼
┌────────────┐
│ 汇总报告 │ ← 第 3 层(等上面两个都完成)
└────────────┘
二、核心组件(6 个文件,各司其职)
用户定义工作流(JSON)
│
▼
┌──────────┐ 序列化/反序列化 ┌──────────┐
│ DagDef │ ←──────────────────→ │ 数据库持久化 │
│ (dag_def) │ └──────────┘
└─────┬─────┘
│ NodeRegistry.build_dag()
▼
┌──────────┐ 拓扑排序分层 ┌──────────────┐
│ Dag │ ──────────────────→ │ Vec<Vec<节点>> │
│ (dag) │ │ 按层组织 │
└──────────┘ └──────┬───────┘
│
┌──────┴───────┐
▼ │
┌────────────┐ │
│ DagExecutor │ ←─── StateMachine
│ (executor) │ (状态转换校验)
└──────┬─────┘
│
┌──────┴─────┐
│ EventBus │ → 事件广播到前端
│ (eventbus) │
└────────────┘
2.1 DagDef(定义层)— 工作流的"蓝图"
源文件:
dag_def.rs
pub struct DagDef {
pub nodes: HashMap<String, NodeDef>, // 节点定义
pub edges: Vec<EdgeDef>, // 依赖关系
}
pub struct NodeDef {
pub id: String, // 如 "check-env"
pub node_type: String, // 如 "script"、"ai"、"human"
pub config: serde_json::Value, // 节点参数(命令、prompt 等)
}
pub struct EdgeDef {
pub source: String, // 如 "check-env"
pub target: String, // 如 "run-tests"
pub condition: Option<String>, // 可选条件(如 "true")
}
这是可序列化的,存到 SQLite 的 workflow_defs.dag 字段(JSON),也可以从模板加载。
2.2 Node(节点抽象)— 所有节点的统一接口
源文件:
node.rs
#[async_trait]
pub trait Node: Send + Sync {
async fn execute(&self, ctx: NodeContext) -> NodeResult;
fn schema(&self) -> NodeSchema;
fn is_blocking(&self) -> bool { false } // 阻塞节点(如人工审批)
fn node_type(&self) -> &str;
}
每个节点执行时拿到一个 NodeContext:
pub struct NodeContext {
pub node_id: String,
pub inputs: HashMap<String, NodeOutput>, // 上游节点的输出
pub config: serde_json::Value, // 节点配置
pub execution_id: String, // 本次执行 ID
pub event_bus: EventBus, // 事件总线(发事件给前端)
pub node_status: StateMachine, // 共享状态机(检查是否被取消)
}
关键设计:node_status 是 StateMachine 的 clone,但内部是 Arc<Mutex<HashMap>>,所以 clone 共享同一底层数据——外部 IPC 调 set_cancelled 写入后,运行中的节点能立即读到。
2.3 Dag(运行时图)+ 拓扑排序
源文件:
dag.rs
DagDef 是数据蓝图,Dag 是装入真实节点实例的运行时图。核心方法是 topological_layers():
输入:节点 + 边的依赖关系
输出:Vec<Vec<NodeId>> —— 分层的节点 ID
算法:Kahn 算法(BFS 分层)
1. 算每个节点的入度(有多少上游)
2. 入度=0 的节点入队(第一层)
3. 取出一层节点,将它们的下游入度 -1
4. 入度归零的下游入队(下一层)
5. 重复直到所有节点处理完
6. 如果处理数 ≠ 节点总数 → 有环,报错
复杂度 O(V+E):一次性遍历边构建邻接索引(出边表)与入度表,BFS 分层只走索引查询(每条边仅被访问一次)。
这就是引擎的核心智能:用户只需画依赖关系(谁依赖谁),引擎自动算出执行顺序和并行机会。
2.4 DagExecutor(执行器)— 三阶段调度
源文件:
executor.rs
pub async fn run(&mut self, dag: &Dag, config: Value) -> Result<HashMap<NodeId, NodeOutput>>
对每一层执行三阶段模式:
┌─ 阶段一:准备 ──────────────────────────────────────┐
│ 遍历层内每个节点: │
│ ① emit NodeStarted 事件 │
│ ② state_machine.set_running(Pending→Running) │
│ ③ 收集上游输出作为 inputs │
│ ④ 构建 NodeContext │
│ ⑤ 创建 async future(不立即执行) │
└──────────────────────────────────────────────────────┘
│
▼
┌─ 阶段二:并发执行 ───────────────────────────────────┐
│ futures::future::join_all(node_futures).await │
│ │
│ 同层所有节点真正并发跑(tokio 异步运行时) │
│ 等待全部完成(或某个失败) │
└──────────────────────────────────────────────────────┘
│
▼
┌─ 阶段三:收尾 ───────────────────────────────────────┐
│ 遍历结果: │
│ ✅ 成功 → set_completed + emit NodeCompleted + 存输出 │
│ ❌ 失败 → set_failed + emit NodeFailed + 记录错误 │
│ 🚫 取消 → 跳过 set_failed(Cancelled 是终态) │
│ │
│ 任一失败 → return Err,后续层不执行 │
│ 全部成功 → 进入下一层 │
└──────────────────────────────────────────────────────┘
为什么三阶段不合并? 避免 Rust 的借用冲突——阶段一需要 &self(读状态机),阶段二的 future 捕获节点引用独立运行,阶段三回来再更新 &mut self。
入边索引优化:执行前预建 adjacency_in: HashMap<NodeId, Vec<NodeId>>(O(E) 一次构建),避免内层循环中调用 dag.predecessors()(每次 O(E) 全表扫描,整体退化 O(V·E))。
2.5 StateMachine(状态机)— 防止非法状态转换
源文件:
state.rs
合法转换图:
Pending ──→ Running ──→ Completed ✅ 正常完成
│
└──→ Failed ❌ 执行失败
任何节点 ←── set_cancelled 🚫 外部取消(不走转换校验)
fn is_legal(from: &NodeStatus, to: &NodeStatus) -> bool {
matches!((from, to),
(Pending, Running) // 启动
| (Running, Completed) // 完成
| (Running, Failed) // 失败
)
}
非法转换直接 bail!(如 Pending → Completed 跳过执行),防止状态被意外覆盖。
共享语义是关键:
pub struct StateMachine {
states: Arc<Mutex<HashMap<NodeId, NodeStatus>>>,
}
clone() 只是 Arc 引用计数 +1,所有副本共享同一 HashMap。这让取消机制可以跨边界工作:
前端点「取消审批」
│
▼
cancel_workflow_node IPC
│
▼
AppState.workflow_state_registry[execution_id].set_cancelled(node_id)
│ (写入共享 HashMap)
▼
HumanNode.execute 内部 ctx.node_status.is_cancelled() → true
│ (读到写入,因为是同一份 HashMap)
▼
返回 Err("人工审批被取消")
│
▼
Executor 检测到 is_cancelled → 跳过 set_failed(避免 Cancelled→Failed 非法转换)
2.6 EventBus(事件总线)— 对外广播
源文件:
eventbus.rs
pub struct EventBus {
sender: broadcast::Sender<WorkflowEvent>, // tokio broadcast
}
容量 256,所有事件通过 broadcast 广播。Tauri IPC 层订阅后转发给前端:
Executor 发事件 → EventBus broadcast → IPC 层 listen → app.emit("workflow-event") → 前端实时展示
事件类型:NodeStarted / NodeCompleted / NodeFailed / WorkflowCompleted / HumanApprovalRequest / HumanApprovalResponse
2.7 ConditionEngine(条件分支)— 当前最简实现
源文件:
conditions.rs
pub fn evaluate(expr: &str, _context: &Value) -> Result<bool> {
if expr.trim() == "true" { return Ok(true); }
if expr.trim() == "false" { return Ok(false); }
// 其他一律 false(保守拒绝,不静默放行)
Ok(false)
}
边上的 condition 字段目前只认 "true" / "false" 字面量。升级为真表达式求值是待办(T-260614-11)。
2.8 NodeRegistry(节点注册表)— 工厂模式创建节点
源文件:
registry.rs
pub struct NodeRegistry {
factories: HashMap<String, NodeFactory>,
}
核心方法 build_dag(def: &DagDef) -> Result<Dag>:遍历 DagDef 的节点定义,按 node_type 查工厂函数创建真实节点实例,组装成运行时 Dag。
不实现 Default trait:原 Default 注册了一个会 panic 的 script 工厂(unimplemented!),违反项目铁律「无 panic」。所有调用方必须显式 new() + 手动 register 真实节点。
三、一次完整的执行流程
以 ProjectDetail.vue 中的 3 节点工作流为例:
用户点「运行工作流」
│
▼
IPC: run_workflow(dag_json)
│
▼
① NodeRegistry.build_dag(dag_def)
→ 从 JSON 创建 3 个真实节点实例 + 边
→ 运行时 Dag
│
▼
② Dag.topological_layers()
→ 算出分层:[["check-env"], ["run-tests"], ["build"]]
│
▼
③ DagExecutor.new(event_bus, execution_id).run(dag, config)
│
├─ 第 1 层 ["check-env"]
│ ├─ set_running("check-env")
│ ├─ ScriptNode.execute(ctx) → tokio::process::Command("cmd /C node -v")
│ └─ set_completed("check-env") + emit NodeCompleted
│
├─ 第 2 层 ["run-tests"]
│ ├─ inputs = { "check-env": 上一步输出 }
│ ├─ set_running("run-tests")
│ ├─ ScriptNode.execute(ctx) → cmd /C npm test
│ └─ set_completed("run-tests")
│
└─ 第 3 层 ["build"]
├─ set_running("build")
├─ ScriptNode.execute(ctx) → cmd /C npm run build
└─ set_completed("build") + emit WorkflowCompleted
│
▼
④ 返回 HashMap<节点ID, 输出> + 前端实时看到日志
四、现有节点类型
| 节点 | 实现状态 | 能力 |
|---|---|---|
| ScriptNode | ✅ 真实可用 | 执行 Shell 命令(Windows: cmd /C),非零退出码判断为失败 |
| AiNode | ✅ 真实可用 | 调用 LLM(OpenAI/Anthropic),config 驱动 provider(base_url/api_key/model/prompt),支持上游输入优先 |
| HumanNode | ✅ 审批闭环 | 阻塞等待人工审批:subscribe→send→select!(响应/超时/取消),execution_id+node_id 双键过滤 |
| DockerNode | ❌ 未实现 | 设计意图:Docker 容器内执行命令(隔离构建/测试环境) |
| GitNode | ❌ 未实现 | 设计意图:Git 操作(commit/push/merge/分支管理) |
| HTTPNode | ❌ 未实现 | 设计意图:发起 HTTP 请求(调外部 API/Webhook) |
| NotifyNode | ❌ 未实现 | 设计意图:多渠道通知(桌面/飞书/钉钉/Webhook) |
| SubflowNode | ❌ 未实现 | 设计意图:嵌套子工作流(复杂流程拆分复用) |
五、当前局限
| 局限 | 说明 | 对应待办 |
|---|---|---|
| 条件引擎只有 true/false | 无法做 $.status == 'ok' 这种真条件分支 |
T-260614-11 |
| 无断点续跑 | 执行中崩溃后无法从失败节点恢复 | Phase 4 |
| 无暂停/恢复 | 只能取消,不能暂停后继续 | Phase 4 |
| 无 DAG 可视化编辑器 | 用户只能写 JSON 定义 | Phase 5 |
| 持久化仅存定义 | 执行状态不落盘(纯内存),重启丢失 | 待需求驱动 |
| 缺 human 节点端到端测试 | 单测绿但前端无 human DAG 入口,审批闭环未端到端验证 | B-03b-R8(P0) |