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

373 lines
15 KiB
Markdown
Raw Permalink Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
# DAG 引擎详解
> 来源:基于 `crates/df-workflow/src/` 实际代码核对编写2026-06-14
> 关联模块文档:[df-workflow-工作流引擎-2026-06-12.md](./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`
```rust
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`
```rust
#[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`
```rust
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`
```rust
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 🚫 外部取消(不走转换校验)
```
```rust
fn is_legal(from: &NodeStatus, to: &NodeStatus) -> bool {
matches!((from, to),
(Pending, Running) // 启动
| (Running, Completed) // 完成
| (Running, Failed) // 失败
)
}
```
非法转换直接 `bail!`(如 Pending → Completed 跳过执行),防止状态被意外覆盖。
**共享语义**是关键:
```rust
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`
```rust
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`
```rust
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`
```rust
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 |