后端: - 工作流推进链(D-03):advance_task/状态机/闸门走 df-nodes Node trait,conditions 条件引擎扩展 - 想法评估闭环:启发式评分+对抗评估,df-ideas/scoring + df-storage/idea_eval_repo + idea 前端打通 - 全局事件数据总线:df-ai/context+context_helpers+augmentation 跨模块解耦 - AI planner/plan_hint/intent:aichat B 路线并行多轮基础 - patch_file 加固(TD-03/04):读改写整体锁防 lost update,expected_hash 合约闭环 - 压缩超时兜底(F-15 卡死根治) - F-09 多会话并发:LlmConcurrency per-conv + streamingGuard 前端守护 + verify 脚本 - 知识注入 DRY/skills/audit 扩展 清理: - aichat 技术债(误报 allow/死导入/过时注释 30 项) - URGENT.md 删除(11 项加急全解决/迁 todo) - 文档整理(todo/待决策/待审查/ARCHITECTURE/INDEX + 总线/技术债审查新文档)
520 lines
18 KiB
Rust
520 lines
18 KiB
Rust
//! DagExecutor 测试套件 — 从 executor.rs 抽离(纯搬迁,零行为变更)
|
||
//!
|
||
//! 抽离说明:executor.rs 主体 run 逻辑之外曾带大量内联测试,占比偏高(抽离前的历史
|
||
//! 行数分布不再维护),独立成文件降低阅读 run 主体的滚动噪音。测试逻辑、断言、节点 mock
|
||
//! 全部原样搬迁。
|
||
//! 通过 `#[path = "executor_helpers.rs"]` 内联回 executor.rs 的 #[cfg(test)] mod,
|
||
//! 不进 lib.rs、不改 pub 路径、不影响外部调用方。
|
||
//!
|
||
//! SW-01 TOCTOU 相关测试(test_cancelled_node_skips_set_failed /
|
||
//! test_cancelled_node_skips_set_completed / test_cancelled_node_emits_node_cancelled_event)
|
||
//! 原样保留,锁定 executor.rs 节点执行完毕后 Ok/Err 两分支的 is_cancelled 短路行为
|
||
//! (已取消节点跳过 set_completed/set_failed,保持 Cancelled 终态)。
|
||
|
||
use super::*;
|
||
use crate::node::{Node, NodeResult, NodeSchema};
|
||
use async_trait::async_trait;
|
||
use std::time::Duration;
|
||
|
||
/// 测试节点:sleep 指定毫秒后返回空输出
|
||
struct SleepNode {
|
||
sleep_ms: u64,
|
||
}
|
||
|
||
#[async_trait]
|
||
impl Node for SleepNode {
|
||
async fn execute(&self, _ctx: NodeContext) -> NodeResult {
|
||
tokio::time::sleep(Duration::from_millis(self.sleep_ms)).await;
|
||
Ok(NodeOutput::empty())
|
||
}
|
||
|
||
fn schema(&self) -> NodeSchema {
|
||
NodeSchema {
|
||
params: serde_json::Value::Null,
|
||
output: serde_json::Value::Null,
|
||
}
|
||
}
|
||
|
||
fn node_type(&self) -> &str {
|
||
"sleep"
|
||
}
|
||
}
|
||
|
||
/// 测试节点:直接返回错误
|
||
struct FailNode;
|
||
|
||
/// 测试节点:把 ctx.config 的指定 key 字符串原样回显到 output(验证节点级 config 下沉)
|
||
struct EchoConfigNode {
|
||
key: String,
|
||
}
|
||
|
||
#[async_trait]
|
||
impl Node for EchoConfigNode {
|
||
async fn execute(&self, ctx: NodeContext) -> NodeResult {
|
||
let v = ctx
|
||
.config
|
||
.get(&self.key)
|
||
.and_then(|v| v.as_str())
|
||
.unwrap_or("")
|
||
.to_string();
|
||
Ok(NodeOutput::from_value(serde_json::json!({ "echo": v })))
|
||
}
|
||
|
||
fn schema(&self) -> NodeSchema {
|
||
NodeSchema {
|
||
params: serde_json::Value::Null,
|
||
output: serde_json::Value::Null,
|
||
}
|
||
}
|
||
|
||
fn node_type(&self) -> &str {
|
||
"echo_config"
|
||
}
|
||
}
|
||
|
||
#[async_trait]
|
||
impl Node for FailNode {
|
||
async fn execute(&self, _ctx: NodeContext) -> NodeResult {
|
||
Err(anyhow::anyhow!("故意失败"))
|
||
}
|
||
|
||
fn schema(&self) -> NodeSchema {
|
||
NodeSchema {
|
||
params: serde_json::Value::Null,
|
||
output: serde_json::Value::Null,
|
||
}
|
||
}
|
||
|
||
fn node_type(&self) -> &str {
|
||
"fail"
|
||
}
|
||
}
|
||
|
||
/// ④-1:节点级 config 下沉验证。
|
||
/// DagDef 节点 config={foo:"bar"} + 全局 config={foo:"GLOBAL"} → NodeContext.config.foo=="bar"
|
||
/// (节点级覆盖全局级)。同时验证缺节点 config 的节点回退全局 config(零行为破坏现有调用方)。
|
||
#[tokio::test]
|
||
async fn node_config_overrides_global_in_node_context() {
|
||
use crate::registry::NodeRegistry;
|
||
|
||
// 构造 DagDef:n1 节点写 config={foo:"bar"};n2 节点 config 空(验证回退全局)
|
||
let mut def = crate::dag_def::DagDef::new();
|
||
def.add_node("n1".to_string(), "echo".to_string(), serde_json::json!({ "foo": "bar" }));
|
||
def.add_node("n2".to_string(), "echo".to_string(), serde_json::json!({}));
|
||
|
||
// 注册 echo 工厂 → EchoConfigNode(读 config.foo)
|
||
let mut registry = NodeRegistry::new();
|
||
registry.register("echo", |_cfg| {
|
||
Box::new(EchoConfigNode { key: "foo".to_string() })
|
||
});
|
||
|
||
let dag = registry.build_dag(&def).expect("build_dag");
|
||
|
||
let mut executor = DagExecutor::new(EventBus::new(), "test-nodecfg".to_string());
|
||
// 全局 config 写 foo=GLOBAL(验证被 n1 节点级覆盖,n2 回退用全局)
|
||
let outputs = executor
|
||
.run(&dag, serde_json::json!({ "foo": "GLOBAL" }))
|
||
.await
|
||
.expect("run");
|
||
|
||
// n1:节点级 foo="bar" 覆盖全局 "GLOBAL"
|
||
assert_eq!(
|
||
outputs["n1"].data["echo"],
|
||
serde_json::json!("bar"),
|
||
"节点级 config 应覆盖全局"
|
||
);
|
||
// n2:节点 config 空 → 回退全局 foo="GLOBAL"(零行为破坏现有调用方)
|
||
assert_eq!(
|
||
outputs["n2"].data["echo"],
|
||
serde_json::json!("GLOBAL"),
|
||
"空节点 config 应回退全局"
|
||
);
|
||
}
|
||
|
||
#[tokio::test]
|
||
async fn test_same_layer_runs_in_parallel() {
|
||
// 两个无依赖节点位于同一层,各 sleep 100ms
|
||
let mut dag = Dag::new();
|
||
dag.add_node("a".to_string(), Box::new(SleepNode { sleep_ms: 100 }));
|
||
dag.add_node("b".to_string(), Box::new(SleepNode { sleep_ms: 100 }));
|
||
|
||
let mut executor = DagExecutor::new(EventBus::new(), "test-exec".to_string());
|
||
let start = std::time::Instant::now();
|
||
let outputs = executor.run(&dag, serde_json::Value::Null).await.unwrap();
|
||
let elapsed = start.elapsed();
|
||
|
||
assert_eq!(outputs.len(), 2);
|
||
// 串行需要约 200ms,并行应明显小于 180ms
|
||
assert!(
|
||
elapsed < Duration::from_millis(180),
|
||
"同层节点应并行执行,实际耗时 {:?}",
|
||
elapsed
|
||
);
|
||
}
|
||
|
||
#[tokio::test]
|
||
async fn test_layer_failure_aborts_following_layers() {
|
||
// a(失败) 与 b(成功) 同层,c 依赖 a,失败后 c 不应执行
|
||
let mut dag = Dag::new();
|
||
dag.add_node("a".to_string(), Box::new(FailNode));
|
||
dag.add_node("b".to_string(), Box::new(SleepNode { sleep_ms: 10 }));
|
||
dag.add_node("c".to_string(), Box::new(SleepNode { sleep_ms: 10 }));
|
||
dag.add_edge("a".to_string(), "c".to_string());
|
||
|
||
let mut executor = DagExecutor::new(EventBus::new(), "test-exec".to_string());
|
||
let err = executor
|
||
.run(&dag, serde_json::Value::Null)
|
||
.await
|
||
.unwrap_err();
|
||
assert!(err.to_string().contains("节点 a 执行失败"));
|
||
|
||
// 同层成功节点状态正常更新,下游节点保持 Pending
|
||
use df_types::types::NodeStatus;
|
||
assert_eq!(executor.state_machine.get(&"a".to_string()), NodeStatus::Failed);
|
||
assert_eq!(executor.state_machine.get(&"b".to_string()), NodeStatus::Completed);
|
||
assert_eq!(executor.state_machine.get(&"c".to_string()), NodeStatus::Pending);
|
||
}
|
||
|
||
/// 测试节点:执行中通过共享 node_status 自取消(模拟外部 IPC set_cancelled),返回 Err
|
||
struct CancelSelfNode;
|
||
|
||
#[async_trait]
|
||
impl Node for CancelSelfNode {
|
||
async fn execute(&self, ctx: NodeContext) -> NodeResult {
|
||
// 通过共享 node_status 置 Cancelled(写共享 HashMap,executor.state_machine 可见)
|
||
ctx.node_status.set_cancelled(ctx.node_id.clone());
|
||
Err(anyhow::anyhow!("人工审批被取消"))
|
||
}
|
||
|
||
fn schema(&self) -> NodeSchema {
|
||
NodeSchema {
|
||
params: serde_json::Value::Null,
|
||
output: serde_json::Value::Null,
|
||
}
|
||
}
|
||
|
||
fn node_type(&self) -> &str {
|
||
"cancel_self"
|
||
}
|
||
}
|
||
|
||
/// B-03b-R1:取消的节点返回 Err 时,executor 跳过 set_failed(Cancelled→Failed transition 非法会 bail),
|
||
/// 状态保持 Cancelled,run 返回取消相关 Err 而非状态转换错误。
|
||
#[tokio::test]
|
||
async fn test_cancelled_node_skips_set_failed() {
|
||
use df_types::types::NodeStatus;
|
||
let mut dag = Dag::new();
|
||
dag.add_node("x".to_string(), Box::new(CancelSelfNode));
|
||
|
||
let mut executor = DagExecutor::new(EventBus::new(), "test-cancel".to_string());
|
||
let result = executor.run(&dag, serde_json::Value::Null).await;
|
||
|
||
// run 返回 Err(取消致中止后续层),但不 panic/bail transition
|
||
assert!(result.is_err());
|
||
let err = result.unwrap_err().to_string();
|
||
assert!(
|
||
err.contains("取消") || err.contains("执行失败"),
|
||
"应返回取消相关错误,实际: {}",
|
||
err
|
||
);
|
||
// 状态保持 Cancelled(未被 set_failed 覆盖为 Failed)—— R1 修法生效的核心验证
|
||
assert_eq!(
|
||
executor.state_machine.get(&"x".to_string()),
|
||
NodeStatus::Cancelled
|
||
);
|
||
}
|
||
|
||
/// 测试节点:执行中通过共享 node_status 自取消(模拟外部 IPC set_cancelled),返回 Ok
|
||
/// 模拟 HumanNode select! 收合法 HumanApprovalResponse 后 return Ok,但 join_all 让出窗口内
|
||
/// cancel_workflow_node IPC 把节点改 Cancelled 的 TOCTOU 场景。
|
||
struct CancelSelfThenOkNode;
|
||
|
||
#[async_trait]
|
||
impl Node for CancelSelfThenOkNode {
|
||
async fn execute(&self, ctx: NodeContext) -> NodeResult {
|
||
ctx.node_status.set_cancelled(ctx.node_id.clone());
|
||
Ok(NodeOutput::empty())
|
||
}
|
||
|
||
fn schema(&self) -> NodeSchema {
|
||
NodeSchema {
|
||
params: serde_json::Value::Null,
|
||
output: serde_json::Value::Null,
|
||
}
|
||
}
|
||
|
||
fn node_type(&self) -> &str {
|
||
"cancel_self_then_ok"
|
||
}
|
||
}
|
||
|
||
/// R-P1-3:Ok 后取消的 TOCTOU 防护。节点返回 Ok 但执行期间已被 set_cancelled,
|
||
/// executor Ok 分支必须跳过 set_completed(transition Cancelled→Completed 非法会 bail),
|
||
/// 状态保持 Cancelled,run 返回 Ok(已批准审批不应误报失败)。
|
||
#[tokio::test]
|
||
async fn test_cancelled_node_skips_set_completed() {
|
||
use df_types::types::NodeStatus;
|
||
let mut dag = Dag::new();
|
||
dag.add_node("y".to_string(), Box::new(CancelSelfThenOkNode));
|
||
|
||
let mut executor = DagExecutor::new(EventBus::new(), "test-cancel-ok".to_string());
|
||
let result = executor.run(&dag, serde_json::Value::Null).await;
|
||
|
||
// run 返回 Ok —— Ok 节点不应因 TOCTOU 取消被误判为失败
|
||
assert!(
|
||
result.is_ok(),
|
||
"Ok 后取消应短路 set_completed 而非 bail,实际 err: {:?}",
|
||
result.err()
|
||
);
|
||
// 状态保持 Cancelled(未被 set_completed 覆盖为 Completed,也不 bail)—— R-P1-3 修法核心验证
|
||
assert_eq!(
|
||
executor.state_machine.get(&"y".to_string()),
|
||
NodeStatus::Cancelled
|
||
);
|
||
}
|
||
|
||
/// 复核-新③:取消节点(Err 路径)emit 的应是 NodeCancelled(非 NodeFailed),
|
||
/// 前端按 type 区分「取消」与「失败」。订阅事件总线收集所有事件断言。
|
||
#[tokio::test]
|
||
async fn test_cancelled_node_emits_node_cancelled_event() {
|
||
use df_types::events::WorkflowEvent;
|
||
use df_types::types::NodeStatus;
|
||
|
||
let mut dag = Dag::new();
|
||
dag.add_node("c".to_string(), Box::new(CancelSelfNode));
|
||
|
||
let bus = EventBus::new();
|
||
let mut rx = bus.subscribe();
|
||
let mut executor = DagExecutor::new(bus, "test-cancel-event".to_string());
|
||
let _ = executor.run(&dag, serde_json::Value::Null).await;
|
||
|
||
// 收集所有事件(run 已结束,事件总线无新事件)
|
||
let mut events: Vec<WorkflowEvent> = Vec::new();
|
||
while let Ok(ev) = rx.try_recv() {
|
||
events.push(ev);
|
||
}
|
||
|
||
// 必含 NodeCancelled { node_id: "c" }
|
||
let cancelled = events.iter().any(|e| matches!(
|
||
e,
|
||
WorkflowEvent::NodeCancelled { node_id } if node_id == "c"
|
||
));
|
||
assert!(cancelled, "取消节点应 emit NodeCancelled, 实际事件: {:?}", events);
|
||
|
||
// 不应含 node "c" 的 NodeFailed(失败事件)—— 语义双标修复的核心验证
|
||
let failed = events.iter().any(|e| matches!(
|
||
e,
|
||
WorkflowEvent::NodeFailed { node_id, .. } if node_id == "c"
|
||
));
|
||
assert!(!failed, "取消节点不应 emit NodeFailed, 实际事件: {:?}", events);
|
||
|
||
// 状态保持 Cancelled
|
||
assert_eq!(
|
||
executor.state_machine.get(&"c".to_string()),
|
||
NodeStatus::Cancelled
|
||
);
|
||
}
|
||
|
||
// ============================================================================
|
||
// conditions-eval feature:边条件路由测试(feature flag 开启时才编译)
|
||
// ============================================================================
|
||
|
||
/// 测试节点:把 config.ok 作为 bool 输出到 { "ok": <bool> }。
|
||
/// 边 condition "$.ok == true" 据此判定路由。
|
||
#[cfg(feature = "conditions-eval")]
|
||
struct BoolFlagNode {
|
||
ok: bool,
|
||
}
|
||
|
||
#[cfg(feature = "conditions-eval")]
|
||
#[async_trait]
|
||
impl Node for BoolFlagNode {
|
||
async fn execute(&self, _ctx: NodeContext) -> NodeResult {
|
||
Ok(NodeOutput::from_value(serde_json::json!({ "ok": self.ok })))
|
||
}
|
||
|
||
fn schema(&self) -> NodeSchema {
|
||
NodeSchema {
|
||
params: serde_json::Value::Null,
|
||
output: serde_json::Value::Null,
|
||
}
|
||
}
|
||
|
||
fn node_type(&self) -> &str {
|
||
"bool_flag"
|
||
}
|
||
}
|
||
|
||
/// 测试节点:执行时记录到共享 Mutex,用于断言节点是否被执行。
|
||
#[cfg(feature = "conditions-eval")]
|
||
struct RecordingNode {
|
||
flag: std::sync::Arc<std::sync::Mutex<Vec<String>>>,
|
||
name: String,
|
||
}
|
||
|
||
#[cfg(feature = "conditions-eval")]
|
||
#[async_trait]
|
||
impl Node for RecordingNode {
|
||
async fn execute(&self, _ctx: NodeContext) -> NodeResult {
|
||
self.flag.lock().unwrap().push(self.name.clone());
|
||
Ok(NodeOutput::empty())
|
||
}
|
||
|
||
fn schema(&self) -> NodeSchema {
|
||
NodeSchema {
|
||
params: serde_json::Value::Null,
|
||
output: serde_json::Value::Null,
|
||
}
|
||
}
|
||
|
||
fn node_type(&self) -> &str {
|
||
"recording"
|
||
}
|
||
}
|
||
|
||
/// conditions-eval 开启时:边 condition 求值 true → 下游执行。
|
||
///
|
||
/// 结构:branch(true) --condition "$.ok == true"--> target_run
|
||
/// branch(true) --condition "$.ok == false"--> target_skip
|
||
/// branch.ok==true:第一条边满足 → target_run 执行;第二条边 false → target_skip 跳过。
|
||
#[cfg(feature = "conditions-eval")]
|
||
#[tokio::test]
|
||
async fn conditions_eval_routes_by_edge_condition() {
|
||
let mut dag = Dag::new();
|
||
dag.add_node("branch".to_string(), Box::new(BoolFlagNode { ok: true }));
|
||
|
||
let ran = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
|
||
let ran_run = ran.clone();
|
||
let ran_skip = ran.clone();
|
||
|
||
dag.add_node(
|
||
"target_run".to_string(),
|
||
Box::new(RecordingNode { flag: ran_run, name: "target_run".to_string() }),
|
||
);
|
||
dag.add_node(
|
||
"target_skip".to_string(),
|
||
Box::new(RecordingNode { flag: ran_skip, name: "target_skip".to_string() }),
|
||
);
|
||
|
||
// 两条带条件的出边:branch.ok==true 命中第一条
|
||
dag.add_edge_with_condition(
|
||
"branch".to_string(),
|
||
"target_run".to_string(),
|
||
"$.ok == true".to_string(),
|
||
);
|
||
dag.add_edge_with_condition(
|
||
"branch".to_string(),
|
||
"target_skip".to_string(),
|
||
"$.ok == false".to_string(),
|
||
);
|
||
|
||
let mut executor = DagExecutor::new(EventBus::new(), "test-cond-route".to_string());
|
||
let outputs = executor
|
||
.run(&dag, serde_json::Value::Null)
|
||
.await
|
||
.expect("run");
|
||
|
||
// target_run 被执行(条件命中);target_skip 被跳过(条件不满足,不入 outputs)
|
||
assert!(
|
||
outputs.contains_key("target_run"),
|
||
"条件命中的下游应执行,outputs: {:?}",
|
||
outputs.keys().collect::<Vec<_>>()
|
||
);
|
||
assert!(
|
||
!outputs.contains_key("target_skip"),
|
||
"条件不满足的下游应跳过(不入 outputs),outputs: {:?}",
|
||
outputs.keys().collect::<Vec<_>>()
|
||
);
|
||
|
||
let ran = ran.lock().unwrap();
|
||
assert!(
|
||
ran.contains(&"target_run".to_string()),
|
||
"target_run 应被执行,ran: {:?}",
|
||
ran
|
||
);
|
||
assert!(
|
||
!ran.contains(&"target_skip".to_string()),
|
||
"target_skip 不应被执行,ran: {:?}",
|
||
ran
|
||
);
|
||
}
|
||
|
||
/// conditions-eval 开启时:节点所有入边均带 condition 且全部 false → 跳过执行。
|
||
/// 验证状态机保持 Pending(跳过节点不 set_running/Completed/Failed)。
|
||
#[cfg(feature = "conditions-eval")]
|
||
#[tokio::test]
|
||
async fn conditions_eval_skip_keeps_pending() {
|
||
use df_types::types::NodeStatus;
|
||
|
||
let mut dag = Dag::new();
|
||
dag.add_node("src".to_string(), Box::new(BoolFlagNode { ok: true }));
|
||
dag.add_node(
|
||
"skip_target".to_string(),
|
||
Box::new(RecordingNode {
|
||
flag: std::sync::Arc::new(std::sync::Mutex::new(Vec::new())),
|
||
name: "skip_target".to_string(),
|
||
}),
|
||
);
|
||
// src.ok==true,但边要求 == false → 不满足 → 跳过
|
||
dag.add_edge_with_condition(
|
||
"src".to_string(),
|
||
"skip_target".to_string(),
|
||
"$.ok == false".to_string(),
|
||
);
|
||
|
||
let mut executor = DagExecutor::new(EventBus::new(), "test-cond-skip".to_string());
|
||
let outputs = executor
|
||
.run(&dag, serde_json::Value::Null)
|
||
.await
|
||
.expect("run");
|
||
|
||
assert!(!outputs.contains_key("skip_target"), "skip_target 应被跳过不入 outputs");
|
||
// 跳过的节点保持 Pending(未被状态机转换)
|
||
assert_eq!(
|
||
executor.state_machine.get(&"skip_target".to_string()),
|
||
NodeStatus::Pending,
|
||
"跳过的节点应保持 Pending"
|
||
);
|
||
assert_eq!(
|
||
executor.state_machine.get(&"src".to_string()),
|
||
NodeStatus::Completed,
|
||
"源节点应正常完成"
|
||
);
|
||
}
|
||
|
||
/// conditions-eval 开启时:无条件入边 + 条件入边混合,无条件的总放行(不被跳过)。
|
||
#[cfg(feature = "conditions-eval")]
|
||
#[tokio::test]
|
||
async fn conditions_eval_unconditional_edge_always_passes() {
|
||
let mut dag = Dag::new();
|
||
dag.add_node("a".to_string(), Box::new(BoolFlagNode { ok: true }));
|
||
dag.add_node("b".to_string(), Box::new(BoolFlagNode { ok: false }));
|
||
dag.add_node(
|
||
"mix".to_string(),
|
||
Box::new(RecordingNode {
|
||
flag: std::sync::Arc::new(std::sync::Mutex::new(Vec::new())),
|
||
name: "mix".to_string(),
|
||
}),
|
||
);
|
||
// 一条无条件边(b → mix)+ 一条条件边(a → mix, 要求 a.ok==false 不满足)
|
||
dag.add_edge("b".to_string(), "mix".to_string());
|
||
dag.add_edge_with_condition(
|
||
"a".to_string(),
|
||
"mix".to_string(),
|
||
"$.ok == false".to_string(),
|
||
);
|
||
|
||
let mut executor = DagExecutor::new(EventBus::new(), "test-cond-mix".to_string());
|
||
let outputs = executor
|
||
.run(&dag, serde_json::Value::Null)
|
||
.await
|
||
.expect("run");
|
||
|
||
// 存在无条件入边 → mix 必执行(无条件边不被跳过逻辑计入 has_any_cond_edge 的「全部 false」)
|
||
assert!(
|
||
outputs.contains_key("mix"),
|
||
"存在无条件入边时节点应执行,outputs: {:?}",
|
||
outputs.keys().collect::<Vec<_>>()
|
||
);
|
||
}
|