Files
DevFlow/docs/03-模块文档/df-workflow-工作流引擎-2026-06-12.md
绝尘 0e196ee86f 修复: DOC-09+10 df-workflow/df-storage 模块文档过期修正
- DOC-09 df-workflow 4 处:NodeRegistry 删 Default impl(无 panic 铁律)+conditions 默认 true→false(B-260614-02)+try_recv_human_approval 已删+set_waiting/set_skipped 已删(仅 set_cancelled 唯一旁路)+set_* 签名 &mut→&self。
- DOC-10 df-storage:V8→V13 + 表数 13 业务表(+schema_version)+Repo 13(12 impl_repo!+1 手写 SettingsRepo)+MIGRATION_VERSION 常量已删(版本靠 schema_version 表)+迁移表补 V9-V13+knowledges DDL 补 reasoning 列。
主代理抽查:过期项仅历史说明残留。批4
2026-06-15 05:40:37 +08:00

12 KiB
Raw Blame History

df-workflow 工作流引擎

创建: 2026-06-10 | 状态: 初稿 | 最后更新: 2026-06-15


概述

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

当前状态

功能 状态
DAG 数据结构 已实现
拓扑排序 已实现
DagExecutor (顺序执行) 已实现
同层节点并行执行 已实现
Node trait 定义 已实现
状态机 (WorkflowRunStatus) 已实现
EventBus (broadcast) 已实现
条件表达式引擎 仅支持 true/false未识别默认 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:原 Default 注册了一个会 panic 的 script 工厂(unimplemented!),违反项目铁律「无 panic——所有占位代码返回空/默认值」,已删除。所有调用方应显式 new() + 手动 register 真实节点(如 state.rs::build_registry 的做法),实际 ScriptNode/AiNode/HumanNodedf-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> 发送人工审批请求(返回接收者计数)
Default / Clone Defaultnew;Clone 复刻 sender(broadcast sender 可 clone,共享通道)

人工审批当前只发不收:emit_human_approval_request 发出 WorkflowEvent,对应的审批响应(HumanApprovalResponse)由 df-nodes::human_node 通过 EventBus::subscribe() 订阅后自行消费,EventBus 本身不再持审批响应存储/检索接口。

状态机state.rs

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

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

方法 签名 职责
new / get ... -> Self / (&NodeId) -> NodeStatus 创建;取状态(缺失默认 Pending
set_running (&self, NodeId) -> Result<()> Pending → Running
set_completed (&self, NodeId) -> Result<()> Running → Completed
set_failed (&self, NodeId) -> Result<()> Running → Failed
set_cancelled (&self, NodeId) Cancelled唯一受控旁路:不经 transition 合法性校验直接 insert(人工审批取消由 cancel_workflow_node IPC 异步触发,节点可能处于 Running 之外的任意态,走 is_legal 会被拒绝)
is_cancelled (&NodeId) -> bool 是否为 Cancelled
snapshot () -> HashMap<NodeId, NodeStatus> 全量状态快照clone 返回,调用方持独立副本)

set_waiting / set_skipped 已删除:两者同为旁路置位但全仓零调用,已删。Waiting / Skipped 两态当前无对应 setter。Cancelledset_cancelled 作为唯一受控旁路置位,可从任意态跳转。

共享语义:内部 Arc<Mutex<HashMap>>clone() 为浅拷贝Arc 引用计数 +1所有 clone 共享同一底层 HashMap。DagExecutorstate_machine 与下沉到 NodeContext.node_status 的 clone 共享底层;run_workflow 把执行器状态机注册到 AppState 全局表后,cancel_workflow_node IPC 的 set_cancelled 可直达运行中阻塞节点(如 HumanNode)的 is_cancelled 轮询。

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 中存在环")。 FR-D12026-06-14 commit 4b5f096复杂度由 O(V·E) 降为 O(V+E)——一次性遍历边建 adjacency_out(出边表)+ 入度表BFS 分层走索引而非每节点重扫 self.edgesexecutor.rs::run 预建 adjacency_in(入边索引)取前驱,消除逐节点 predecessors() 全表扫描

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(false)tracing::warn!。即条件未识别时阻断而非放行(B-260614-02 反转)——条件分支写错或引擎未实现时不静默放行。

现状: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(false)(保守拒绝),无法据 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 计划 - 决策能力升级

相关文档