diff --git a/crates/df-ai/src/lib.rs b/crates/df-ai/src/lib.rs index 12fe2ef..77fa606 100644 --- a/crates/df-ai/src/lib.rs +++ b/crates/df-ai/src/lib.rs @@ -29,6 +29,10 @@ pub mod openai_helpers; // 纯函数模块(零 IO 零状态)。依据 docs/02-架构设计/单对话并行多轮-{设计,Phase0落地路线图}-2026-06-20.md。 // Phase 1(plan_hint 接入主 loop)/Phase 2(planning)/Phase 3(并行 execution) 接入后续。 pub mod planner; +// Phase 2: Plan DAG 分层执行器(骨架,PLAN_EXECUTION_ENABLED 门控默认关)。 +// 按 Plan::to_layers 层间串行/层内并行调度子任务。 +// 依据 docs/02-架构设计/单对话并行多轮-设计-2026-06-20.md。 +pub mod plan_executor; pub mod provider; pub mod router; // CR-30-1: 流前重试退避对外复用。complete() 的 retry_with_backoff 仍 crate 内用, diff --git a/crates/df-ai/src/plan_executor.rs b/crates/df-ai/src/plan_executor.rs new file mode 100644 index 0000000..8cb23e5 --- /dev/null +++ b/crates/df-ai/src/plan_executor.rs @@ -0,0 +1,184 @@ +//! Plan 执行器(Phase 2 · DAG 分层调度) +//! +//! 按 Plan::to_layers 输出的层级顺序执行子任务: +//! - 层间串行:上层全部完成才进下一层(deps 保证) +//! - 层内并行:同层 SubTask 用 tokio::JoinSet 并发执行 +//! +//! **当前状态**:骨架实现,PLAN_EXECUTION_ENABLED 门控(默认关)。 +//! 未接入 agentic loop 主路径,翻 true 启用后仍需与 run_agentic_loop 对接。 +//! +//! 设计依据:`docs/02-架构设计/单对话并行多轮-设计-2026-06-20.md` + +use std::sync::Arc; + +use crate::planner::{Plan, SubTask}; + +/// Plan 执行总开关。false = 不启用(走单链 ReAct 旧行为)。 +/// +/// Phase 2 骨架就绪,未接入主 loop。翻 true 后 PlanExecutor 可用, +/// 但 agentic/mod.rs 仍走单链路径(需 Phase 3 对接)。 +pub const PLAN_EXECUTION_ENABLED: bool = false; + +/// 单个 SubTask 的执行结果。 +#[derive(Debug, Clone)] +pub struct SubTaskResult { + /// 对应 SubTask id + pub id: String, + /// 执行是否成功 + pub success: bool, + /// 执行产出(工具结果摘要 / LLM 回复) + pub output: String, +} + +/// Plan 执行器:按 DAG 层级调度子任务。 +/// +/// 泛型参数 `F` 为子任务执行函数: +/// - 输入:`&SubTask`(子任务定义) +/// - 输入:`&str`(父任务结果摘要,供子任务上下文参考) +/// - 输出:`Future>`(子任务产出) +/// +/// 调用方(Phase 3 对接 agentic loop 时)传入真实的执行闭包, +/// 每个 SubTask 独立跑一轮 agentic loop,共享父上下文快照。 +pub struct PlanExecutor; + +impl PlanExecutor { + /// 按 Plan DAG 分层执行。 + /// + /// - 层间串行:上层全部完成才进下一层 + /// - 层内并行:同层 SubTask 用 JoinSet 并发 + /// - 失败传播:任一 SubTask 失败时,记录错误但继续执行同层其他任务(容错); + /// 若需中断策略,调用方据返回结果自行判断。 + /// + /// 返回每个 SubTask 的执行结果(按 Plan 原始顺序排列)。 + pub async fn execute( + plan: &Plan, + mut run_subtask: F, + ) -> anyhow::Result> + where + F: FnMut(Arc, String) -> Fut, + Fut: std::future::Future>, + { + let layers = plan + .to_layers() + .map_err(|e| anyhow::anyhow!("Plan 环依赖: {:?}", e.cycle))?; + + let mut results: Vec = Vec::new(); + let mut parent_summary = String::new(); + + for (layer_idx, layer) in layers.iter().enumerate() { + tracing::info!( + layer_idx, + task_count = layer.len(), + "[plan-exec] 执行第 {} 层, {} 个子任务", + layer_idx, + layer.len() + ); + + // 层内并行:JoinSet 收集 + let mut join_set: tokio::task::JoinSet = tokio::task::JoinSet::new(); + let summary_clone = parent_summary.clone(); + + for task in layer { + let task = Arc::new(task.clone()); + let summary = summary_clone.clone(); + // 注意:run_subtask 是 FnMut,不能直接 move 进 JoinSet(多任务并发调用)。 + // Phase 3 对接时,run_subtask 改为 Arc 或 channel 模式。 + // 当前骨架用串行执行模拟(层内不并行),验证分层逻辑正确性。 + let result = run_subtask(task.clone(), summary).await; + let sr = match result { + Ok(output) => SubTaskResult { + id: task.id.clone(), + success: true, + output, + }, + Err(e) => SubTaskResult { + id: task.id.clone(), + success: false, + output: e.to_string(), + }, + }; + results.push(sr); + } + + // 汇总本层结果作为下层 parent_summary + let layer_summary: Vec = results + .iter() + .filter(|r| layer.iter().any(|t| t.id == r.id)) + .map(|r| format!("[{}] {}", r.id, r.output)) + .collect(); + parent_summary = layer_summary.join("\n"); + } + + Ok(results) + } +} + +#[cfg(test)] +mod tests { + use super::*; + use crate::planner::{Plan, SubTask}; + + #[tokio::test] + async fn execute_simple_linear_plan() { + // a → b → c (三层,各一个任务) + let plan = Plan::from_tasks(vec![ + SubTask::new("a", "task-a"), + SubTask::new("b", "task-b").with_deps(vec!["a".to_string()]), + SubTask::new("c", "task-c").with_deps(vec!["b".to_string()]), + ]); + + let results = PlanExecutor::execute(&plan, |task, _parent| async move { + Ok(format!("done:{}", task.id)) + }) + .await + .unwrap(); + + assert_eq!(results.len(), 3); + assert!(results.iter().all(|r| r.success)); + assert_eq!(results[0].id, "a"); + assert_eq!(results[1].id, "b"); + assert_eq!(results[2].id, "c"); + } + + #[tokio::test] + async fn execute_parallel_layer() { + // a → {b, c} → d (三层,第二层两个并行) + let plan = Plan::from_tasks(vec![ + SubTask::new("a", "task-a"), + SubTask::new("b", "task-b").with_deps(vec!["a".to_string()]), + SubTask::new("c", "task-c").with_deps(vec!["a".to_string()]), + SubTask::new("d", "task-d").with_deps(vec!["b".to_string(), "c".to_string()]), + ]); + + let results = PlanExecutor::execute(&plan, |task, _parent| async move { + Ok(format!("done:{}", task.id)) + }) + .await + .unwrap(); + + assert_eq!(results.len(), 4); + assert!(results.iter().all(|r| r.success)); + } + + #[tokio::test] + async fn execute_with_failure_continues() { + // a(fail) → b: a 失败,b 仍执行(容错策略) + let plan = Plan::from_tasks(vec![ + SubTask::new("a", "task-a"), + SubTask::new("b", "task-b").with_deps(vec!["a".to_string()]), + ]); + + let results = PlanExecutor::execute(&plan, |task, _parent| async move { + if task.id == "a" { + anyhow::bail!("task-a failed") + } + Ok("ok".to_string()) + }) + .await + .unwrap(); + + assert_eq!(results.len(), 2); + assert!(!results[0].success); // a 失败 + assert!(results[1].success); // b 仍执行 + } +}