Files
DevFlow/docs/03-模块文档/DAG引擎详解-2026-06-14.md
绝尘 04032a2a8d 重构: 文档汇总+进度看板+孤儿任务清理脚本+gitignore 噪音排除
- 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/ 噪音排除
2026-06-15 05:14:21 +08:00

15 KiB
Raw Blame History

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_statusStateMachine 的 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_runningPending→Running       │
│  ③ 收集上游输出作为 inputs                             │
│  ④ 构建 NodeContext                                   │
│  ⑤ 创建 async future不立即执行                     │
└──────────────────────────────────────────────────────┘
                    │
                    ▼
┌─ 阶段二:并发执行 ───────────────────────────────────┐
│  futures::future::join_all(node_futures).await       │
│                                                      │
│  同层所有节点真正并发跑tokio 异步运行时)              │
│  等待全部完成(或某个失败)                             │
└──────────────────────────────────────────────────────┘
                    │
                    ▼
┌─ 阶段三:收尾 ───────────────────────────────────────┐
│  遍历结果:                                            │
│  ✅ 成功 → set_completed + emit NodeCompleted + 存输出 │
│  ❌ 失败 → set_failed + emit NodeFailed + 记录错误     │
│  🚫 取消 → 跳过 set_failedCancelled 是终态)         │
│                                                      │
│  任一失败 → 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 真实可用 调用 LLMOpenAI/Anthropicconfig 驱动 providerbase_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-R8P0