- 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/ 噪音排除
373 lines
15 KiB
Markdown
373 lines
15 KiB
Markdown
# 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_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 🚫 外部取消(不走转换校验)
|
||
```
|
||
|
||
```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** | ✅ 真实可用 | 调用 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) |
|