squash合并: - 意图识别层论证(8维度+10业界佐证) - 多主题上下文管理愿景+并存论证+补充论证(多轮agentic) - 架构设计文档物理分类(四子目录+INDEX+命名规范+引用同步+边界清晰化) - 前端架构技术债清单归档
12 KiB
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 / NodeOutput(node.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) |
NodeOutputderiveDebug/Clone/Serialize/Deserialize;NodeContext仅Debug/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/HumanNode由df-nodescrate 注册。
DAG 执行流程
1. 接收 WorkflowDef (DAG 定义)
2. 拓扑排序 → 得到执行层 (layers)
3. 逐层执行:
- 同层节点并行(`futures::future::join_all`,已实现)
- 阻塞/非阻塞节点当前同等异步执行(`is_blocking` 未被 executor 消费,人工等待逻辑规划中、当前未实现)
4. 状态变更通过 EventBus 广播
5. 节点完成后仅更新内存态(`StateMachine` + `outputs` HashMap),无持久化快照(规划中)
DAG 定义序列化(dag_def.rs)
区分两层表示:运行时 Dag 持 Box<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_type,config 一律填 Value::Null(无法从 trait object 反推) |
事件类型
事件枚举定义在 df-core::events::WorkflowEvent,DagExecutor 通过 EventBus::send 广播。
执行器实际广播的(executor.rs):
NodeStarted/NodeCompleted/NodeFailedWorkflowCompleted
枚举已定义但执行器当前未触发(df-core::events 中存在,DagExecutor 不发):
NodeProgress(节点进度)NodeOutput(节点输出流)WorkflowPaused(暂停等待外部输入)WorkflowFailed(工作流失败,带 failed_node)HumanApprovalRequest/HumanApprovalResponse(人工审批)
注:枚举中无
WorkflowStarted/WorkflowResumed,旧文档所述为误。
EventBus(eventbus.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 |
— | Default 走 new;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_legal):Pending → Running,Running → Completed,Running → 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。Cancelled由set_cancelled作为唯一受控旁路置位,可从任意态跳转。共享语义:内部
Arc<Mutex<HashMap>>,clone()为浅拷贝(Arc 引用计数 +1),所有 clone 共享同一底层 HashMap。DagExecutor的state_machine与下沉到NodeContext.node_status的 clone 共享底层;run_workflow把执行器状态机注册到AppState全局表后,cancel_workflow_nodeIPC 的set_cancelled可直达运行中阻塞节点(如HumanNode)的is_cancelled轮询。
Dag 对外 API(dag.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-D1(2026-06-14 commit 4b5f096):复杂度由 O(V·E) 降为 O(V+E)——一次性遍历边建 adjacency_out(出边表)+ 入度表,BFS 分层走索引而非每节点重扫 self.edges。executor.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 — EventBus(broadcast)
├── registry.rs — NodeRegistry(节点类型注册表)
└── conditions.rs — 条件表达式引擎
🔮 决策能力接入需求(B 路线)
2026-06-12 记录。executor 分层并行已具骨架,conditions 为最大空壳。
AI Chat(B 路线)从单链 ReAct 升级为规划式协作时,本引擎需补两处:
-
condition.rs::ConditionEngine::evaluate()—— 当前只认"true"/"false"字面量,未识别表达式默认Ok(false)(保守拒绝),无法据 AI 输出做条件分支。需补:- JSON Path 取值(从上游节点输出读字段)
- 比较运算(
==!=><>=<=) contains/ 字符串匹配and/or/not组合
-
executor 分层并行接入 agentic loop ——
topological_layers+join_all已实现(测试test_same_layer_runs_in_parallel),但仅喂静态 DAG。需让run_agentic_loop内动态生成的 DAG 走这条并行通道,而非单轮串行process_tool_calls。