Files
DevFlow/docs/03-模块文档/df-workflow-工作流引擎.md
绝尘 cf017f81e2 新增: Phase2 阶段收尾(Sprint 1-20)
重构:删 5 零引用 crate(df-evolve/plugin/stages/task/traceability)+ 清死模块、ai.rs 拆 11 子 module、ai.ts 拆 6 composable、i18n 拆目录
功能:知识库全栈(df-project/scan + CRUD + 时间线 + 前端)、Settings 拆分、appSettings KV 迁移、模型池、LLM 并发 Semaphore
修复:审批持久化根治、ConditionEngine 默认拒绝、NodeRegistry unimplemented 清除、promote 补偿删除、工具结果截断 50KB、路径校验防 symlink 逃逸
文档:B-03 人工审批设计、决策记录三分档、规格契约自检、经验记录、todo 看板、PROGRESS 更新

详见 PROGRESS.md。src-tauri/儿童每日打卡应用/ 与本项目无关,已排除。
2026-06-14 14:08:20 +08:00

11 KiB
Raw Blame History

df-workflow 工作流引擎

创建: 2026-06-10 | 状态: 初稿


概述

df-workflow 是 DevFlow 的核心引擎,负责 DAG 定义、拓扑排序、节点执行、状态流转和事件广播。引擎不感知具体业务逻辑。

当前状态

功能 状态
DAG 数据结构 已实现
拓扑排序 已实现
DagExecutor (顺序执行) 已实现
同层节点并行执行 已实现
Node trait 定义 已实现
状态机 (WorkflowRunStatus) 已实现
EventBus (broadcast) 已实现
条件表达式引擎 仅支持 true/false
断点续跑 待实施(引擎整体无暂停/恢复/快照机制,缺乏底层基础设施支撑)

核心设计

Node trait

#[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;
}

// NodeResult = anyhow::Result<NodeOutput> 的别名
// ctx 按值传递(非引用),由 executor 在每层构建

NodeContext / NodeOutputnode.rs

执行上下文与输出结构,由 executor 在每层为每个节点构建。

NodeContext 字段:

字段 类型 职责
node_id NodeId 当前节点 ID
inputs HashMap<String, NodeOutput> 上游节点输出,key 为上游节点 ID
config serde_json::Value 节点配置参数
execution_id String 工作流执行 ID
event_bus EventBus 事件总线(可 Clone,广播状态变更)
node_status StateMachine 节点状态机(用于检查取消状态)

NodeOutput 字段与构造器:

成员 签名 说明
data serde_json::Value 输出数据
metadata HashMap<String, String> 输出元数据
empty() () -> Self 空输出(data = Null,空 metadata)
from_value(data) (Value) -> Self 从 JSON 值构造(空 metadata)

NodeOutput derive Debug/Clone/Serialize/Deserialize;NodeContextDebug/Clone(StateMachine 非 Serialize)。

节点注册机制NodeRegistry

节点工厂注册表,据 DagDef.node_type 字符串创建 Box<dyn Node> 实例,桥接可序列化定义与运行时 trait object。

方法 签名 职责
new () -> Self 创建空注册表
register (&mut self, type_name: &str, factory: F) 注册一个节点工厂(F: Fn(&Value) -> Box<dyn Node>
create (&self, type_name: &str, config: &Value) -> Result<Box<dyn Node>> 按类型名 + 配置创建节点实例,未注册则报错
build_dag (&self, def: &DagDef) -> Result<Dag> DagDef 构建完整运行时 Dag(建节点 + 加边,自动分发条件边)
is_registered (&self, type_name: &str) -> bool 检查类型是否已注册
registered_types (&self) -> Vec<&str> 列出所有已注册类型名

Default 实现仅注册占位 script 工厂(unimplemented!),实际 ScriptNodedf-nodes crate 注册。

DAG 执行流程

1. 接收 WorkflowDef (DAG 定义)
2. 拓扑排序 → 得到执行层 (layers)
3. 逐层执行:
   - 同层节点并行(`futures::future::join_all`,已实现)
   - 阻塞/非阻塞节点当前同等异步执行(`is_blocking` 未被 executor 消费,人工等待逻辑规划中、当前未实现)
4. 状态变更通过 EventBus 广播
5. 节点完成后仅更新内存态(`StateMachine` + `outputs` HashMap无持久化快照规划中

DAG 定义序列化dag_def.rs

区分两层表示:运行时 DagBox<dyn Node>trait object不可序列化可持久化的 DagDef / NodeDef / EdgeDef#[derive(Serialize, Deserialize)],用于模板与存盘。NodeRegistry::build_dag 负责从 DagDef 还原运行时 Dag

方法 签名 职责
DagDef::new () -> Self 空定义
DagDef::add_node (&mut self, id, node_type, config: Value) 加节点定义label 默认 None
DagDef::add_edge (&mut self, source, target) 加普通边condition=None
add_edge_with_condition (&mut self, source, target, condition) 加带条件表达式的边
from_dag_edges (&dag: &Dag) -> Self 从运行时 Dag 反推定义;注意只能还原边的 condition 与节点的 node_typeconfig 一律填 Value::Null(无法从 trait object 反推)

事件类型

事件枚举定义在 df-core::events::WorkflowEventDagExecutor 通过 EventBus::send 广播。

执行器实际广播的executor.rs

  • NodeStarted / NodeCompleted / NodeFailed
  • WorkflowCompleted

枚举已定义但执行器当前未触发df-core::events 中存在DagExecutor 不发):

  • NodeProgress(节点进度)
  • NodeOutput(节点输出流)
  • WorkflowPaused(暂停等待外部输入)
  • WorkflowFailed(工作流失败,带 failed_node
  • HumanApprovalRequest / HumanApprovalResponse(人工审批)

注:枚举中无 WorkflowStarted / WorkflowResumed,旧文档所述为误。

EventBuseventbus.rs

基于 tokio::sync::broadcast 的发布/订阅,内部持 broadcast::Sender<WorkflowEvent>

方法/成员 签名 职责
DEFAULT_CAPACITY const usize = 256 默认通道容量
new () -> Self 用默认容量建总线
with_capacity (usize) -> Self 指定容量建总线
send (&self, WorkflowEvent) -> () 广播事件(忽略接收者已关闭错误,异步)
subscribe (&self) -> broadcast::Receiver<WorkflowEvent> 订阅事件流
emit_human_approval_request (&self, WorkflowEvent) -> Result<usize, SendError> 发送人工审批请求(返回接收者计数)
try_recv_human_approval (&self, execution_id, node_id) -> Option<HumanApprovalResponse> TODO 占位:当前恒返回 None,审批响应存储/检索未实现,需配合前端
Default / Clone Defaultnew;Clone 复刻 sender(broadcast sender 可 clone,共享通道)

emit_human_approval_requesttry_recv_human_approval 为人工审批占位接口,后者未实现,执行器当前不消费审批响应(对应 is_blocking 等待逻辑亦未落地)。

状态机state.rs

StateMachine 维护 HashMap<NodeId, NodeStatus>,校验节点状态转换,非法转换返回错误。

合法转换链is_legalPending → RunningRunning → CompletedRunning → Failed。其余均拒绝。

方法 签名 职责
new / get ... -> Self / (&NodeId) -> NodeStatus 创建;取状态(缺失默认 Pending
set_running (&mut self, NodeId) -> Result<()> Pending → Running
set_completed (&mut self, NodeId) -> Result<()> Running → Completed
set_failed (&mut self, NodeId) -> Result<()> Running → Failed
set_waiting (&mut self, NodeId) 设为 Waiting绕过转换校验(直接 set
set_skipped (&mut self, NodeId) 设为 Skipped绕过转换校验
is_cancelled (&NodeId) -> bool 是否为 Cancelled(同样无对应 setter外部直接 set
snapshot () -> &HashMap<NodeId, NodeStatus> 全量状态快照引用

Waiting / Skipped / Cancelled 三态暂未纳入 is_legal 校验链,对应的 set_* 直接 insert,可从任意态跳转。

Dag 对外 APIdag.rs

方法 签名 职责
new () -> Self 空 DAG
add_node (&mut self, id: NodeId, node: Box<dyn Node>) 加节点
add_edge (&mut self, source, target) 加普通边condition=None
add_edge_with_condition (&mut self, source, target, condition: String) 加带条件边
predecessors (&NodeId) -> Vec<NodeId> 上游节点 ID
successors (&NodeId) -> Vec<NodeId> 下游节点 ID
topological_layers () -> Result<Vec<Vec<NodeId>>> BFS 分层拓扑排序,同层可并行;检测到环时报 Workflow 错误"DAG 中存在环"

Edge { source, target, condition: Option<String> } — condition 为可选条件表达式,由 conditions.rs 求值。

条件表达式引擎conditions.rs

ConditionEngine::evaluate(expr: &str, context: &Value) -> Result<bool> — 仅支持 "true" / "false" 字面量(区分大小写,先 trim() 去空白)。

默认放行(安全风险):空串、"True"/"FALSE" 等大小写不匹配字面量、以及任意非 true/false 字面量(如 "yes""1""$.status == 'completed'")均回退为 Ok(true)tracing::warn!。即条件不匹配时放行而非阻断,DAG 边全通 —— 无法据上游输出做条件分支。

现状:JSON Path、比较运算、contains、逻辑组合均未实现(context 参数当前未被使用,仅占位对齐签名)。详见 B 路线决策接入需求

依赖关系

df-core (类型、事件、错误)
  ← df-workflow
    ← df-nodes (节点实现)
    ← df-stages (阶段节点)

文件结构

crates/df-workflow/src/
├── lib.rs        — 模块入口
├── dag.rs        — DAG 数据结构与拓扑排序
├── dag_def.rs    — 可序列化 DAG 定义DagDef/NodeDef/EdgeDef
├── executor.rs   — DagExecutor
├── node.rs       — Node trait + NodeContext/Output
├── state.rs      — 状态机
├── eventbus.rs   — EventBusbroadcast
├── registry.rs   — NodeRegistry节点类型注册表
└── conditions.rs — 条件表达式引擎

🔮 决策能力接入需求B 路线)

2026-06-12 记录。executor 分层并行已具骨架conditions 为最大空壳。

AI ChatB 路线)从单链 ReAct 升级为规划式协作时,本引擎需补两处:

  1. condition.rs::ConditionEngine::evaluate() —— 当前只认 "true" / "false" 字面量,默认 Ok(true)DAG 边全通,无法据 AI 输出做条件分支。需补:

    • JSON Path 取值(从上游节点输出读字段)
    • 比较运算(== != > < >= <=
    • contains / 字符串匹配
    • and / or / not 组合
  2. executor 分层并行接入 agentic loop —— topological_layers + join_all 已实现(测试 test_same_layer_runs_in_parallel),但仅喂静态 DAG。需让 run_agentic_loop 内动态生成的 DAG 走这条并行通道,而非单轮串行 process_tool_calls

详见 Phase 2 计划 - 决策能力升级

相关文档