From e6e913eabc8e62a82ef406d03e0915c2acda687c Mon Sep 17 00:00:00 2001 From: Bryan Helmkamp Date: Thu, 26 Mar 2026 09:24:11 -0400 Subject: [PATCH] refactor: rename local variables/fields to align with Options suffix Follow-up to fafc0a3c. Renames local variables, function parameters, and struct fields that hold renamed types (ExecutorOptions, RunCreateOptions, RunOptions) from config/settings to options/run_options for consistency. Co-Authored-By: Claude Opus 4.6 (1M context) --- lib/crates/fabro-api/src/server.rs | 14 +- lib/crates/fabro-cli/src/commands/resume.rs | 18 +- lib/crates/fabro-cli/src/commands/run.rs | 6 +- lib/crates/fabro-core/src/executor.rs | 20 +- .../src/handler/manager_loop.rs | 4 +- .../fabro-workflows/src/lifecycle/disk.rs | 4 +- .../fabro-workflows/src/lifecycle/git.rs | 55 +- .../fabro-workflows/src/lifecycle/mod.rs | 24 +- .../fabro-workflows/src/operations/create.rs | 26 +- .../fabro-workflows/src/operations/start.rs | 14 +- .../fabro-workflows/src/pipeline/execute.rs | 28 +- .../src/pipeline/execute/tests.rs | 70 ++- .../fabro-workflows/src/pipeline/finalize.rs | 36 +- .../src/pipeline/initialize.rs | 6 +- .../src/pipeline/pull_request.rs | 4 +- .../fabro-workflows/src/pipeline/retro.rs | 8 +- .../fabro-workflows/src/pipeline/types.rs | 8 +- .../fabro-workflows/src/test_support.rs | 28 +- .../tests/daytona_integration.rs | 24 +- .../fabro-workflows/tests/integration.rs | 586 +++++++++--------- 20 files changed, 505 insertions(+), 478 deletions(-) diff --git a/lib/crates/fabro-api/src/server.rs b/lib/crates/fabro-api/src/server.rs index e4804db96..fd7e409a6 100644 --- a/lib/crates/fabro-api/src/server.rs +++ b/lib/crates/fabro-api/src/server.rs @@ -653,7 +653,7 @@ async fn execute_run(state: Arc, run_id: String) { } }; let run_record = persisted.run_record().clone(); - let config = RunOptions { + let run_options = RunOptions { config: run_record.config, run_dir: run_dir.clone(), cancel_token: Some(cancel_token), @@ -673,7 +673,7 @@ async fn execute_run(state: Arc, run_id: String) { let sandbox = Arc::clone(&sandbox); let registry = Arc::new(registry); let run_id = run_id.clone(); - let config = config.clone(); + let run_options = run_options.clone(); let hooks = state.hooks.clone(); let dry_run = state.dry_run; async move { @@ -690,7 +690,7 @@ async fn execute_run(state: Arc, run_id: String) { setup_command_timeout_ms: 300_000, devcontainer_phases: Vec::new(), }, - run_options: config, + run_options, hooks: fabro_hooks::HookConfig { hooks }, sandbox_env: HashMap::new(), checkpoint: None, @@ -719,13 +719,13 @@ async fn execute_run(state: Arc, run_id: String) { }; // Save final checkpoint - let checkpoint = Checkpoint::load(&config.run_dir.join("checkpoint.json")).ok(); + let checkpoint = Checkpoint::load(&run_options.run_dir.join("checkpoint.json")).ok(); // Auto-derive retro and accumulate aggregate usage if let Some(ref cp) = checkpoint { let failed = result.is_err(); let completed_stages = fabro_workflows::build_completed_stages(cp, failed); - let stage_durations = fabro_retro::retro::extract_stage_durations(&config.run_dir); + let stage_durations = fabro_retro::retro::extract_stage_durations(&run_options.run_dir); let retro = fabro_retro::retro::derive_retro( &run_id, "workflow", @@ -734,7 +734,7 @@ async fn execute_run(state: Arc, run_id: String) { 0, &stage_durations, ); - let _ = retro.save(&config.run_dir); + let _ = retro.save(&run_options.run_dir); // Accumulate aggregate usage let mut agg = state @@ -778,7 +778,7 @@ async fn execute_run(state: Arc, run_id: String) { if let Some(ctx) = final_context { managed_run.context = Some(ctx); } - managed_run.run_dir = Some(config.run_dir.clone()); + managed_run.run_dir = Some(run_options.run_dir.clone()); managed_run.event_tx = None; } drop(runs); diff --git a/lib/crates/fabro-cli/src/commands/resume.rs b/lib/crates/fabro-cli/src/commands/resume.rs index 5182017ea..8a3ac03af 100644 --- a/lib/crates/fabro-cli/src/commands/resume.rs +++ b/lib/crates/fabro-cli/src/commands/resume.rs @@ -107,7 +107,7 @@ struct ResumeContext { /// Kept as Arc so the sandbox event callbacks can emit through it. Listeners /// that need to be added later (e.g. ProgressUI) are registered separately. emitter: Arc, - settings: RunOptions, + run_options: RunOptions, setup_commands: Vec, /// Devcontainer lifecycle phases (on_create, post_create, post_start) resolved from config. devcontainer_phases: Vec<(String, Vec)>, @@ -459,7 +459,7 @@ async fn prepare_from_checkpoint( }; let sandbox: Arc = Arc::new(fabro_agent::ReadBeforeWriteSandbox::new(sandbox)); - let settings = RunOptions { + let run_options = RunOptions { config: settings_config, run_dir: run_dir.clone(), cancel_token: None, @@ -497,7 +497,7 @@ async fn prepare_from_checkpoint( run_cfg, sandbox, emitter, - settings, + run_options, setup_commands, devcontainer_phases, devcontainer_env, @@ -940,7 +940,7 @@ async fn prepare_from_branch( .unwrap_or_default(); setup_commands.extend(sandbox.resume_setup_commands(&run_branch)); - let settings = RunOptions { + let run_options = RunOptions { config: settings_config, run_dir: run_dir.clone(), cancel_token: None, @@ -981,7 +981,7 @@ async fn prepare_from_branch( run_cfg, sandbox, emitter, - settings, + run_options, setup_commands, devcontainer_phases, devcontainer_env, @@ -1009,7 +1009,7 @@ async fn run_resumed( mut run_cfg, sandbox, emitter, - mut settings, + mut run_options, setup_commands, devcontainer_phases, devcontainer_env, @@ -1195,7 +1195,7 @@ async fn run_resumed( } } }; - settings.dry_run = dry_run_mode; + run_options.dry_run = dry_run_mode; if let Some(ref mut cfg) = run_cfg { run_config::resolve_sandbox_env(cfg)?; @@ -1314,7 +1314,7 @@ async fn run_resumed( let pr_config = if dry_run_mode { None } else { - settings.pull_request().cloned() + run_options.pull_request().cloned() }; let started = start( persisted, @@ -1326,7 +1326,7 @@ async fn run_resumed( sandbox: Arc::clone(&sandbox), registry: Arc::new(registry), lifecycle, - run_options: settings, + run_options, hooks: fabro_hooks::HookConfig { hooks: run_cfg .as_ref() diff --git a/lib/crates/fabro-cli/src/commands/run.rs b/lib/crates/fabro-cli/src/commands/run.rs index 7c58f215e..988b2ae85 100644 --- a/lib/crates/fabro-cli/src/commands/run.rs +++ b/lib/crates/fabro-cli/src/commands/run.rs @@ -1645,7 +1645,7 @@ async fn run_command_impl( None }; - let config = RunOptions { + let run_options = RunOptions { config: settings_config, run_dir: run_dir.clone(), cancel_token: None, @@ -1691,7 +1691,7 @@ async fn run_command_impl( let pr_config = if dry_run_mode { None } else { - config.pull_request().cloned() + run_options.pull_request().cloned() }; let started = start( persisted, @@ -1703,7 +1703,7 @@ async fn run_command_impl( sandbox: Arc::clone(&sandbox), registry: Arc::new(registry), lifecycle, - run_options: config, + run_options, hooks: fabro_hooks::HookConfig { hooks: run_cfg .as_ref() diff --git a/lib/crates/fabro-core/src/executor.rs b/lib/crates/fabro-core/src/executor.rs index 5c4193db7..a9b1e1979 100644 --- a/lib/crates/fabro-core/src/executor.rs +++ b/lib/crates/fabro-core/src/executor.rs @@ -25,7 +25,7 @@ pub struct ExecutorOptions { pub struct Executor { handler: Arc>, lifecycle: Box>, - settings: ExecutorOptions, + options: ExecutorOptions, } enum NextStep { @@ -38,7 +38,7 @@ enum NextStep { pub struct ExecutorBuilder { handler: Arc>, lifecycle: Option>>, - settings: ExecutorOptions, + options: ExecutorOptions, } impl ExecutorBuilder { @@ -46,7 +46,7 @@ impl ExecutorBuilder { Self { handler, lifecycle: None, - settings: ExecutorOptions::default(), + options: ExecutorOptions::default(), } } @@ -56,17 +56,17 @@ impl ExecutorBuilder { } pub fn cancel_token(mut self, token: Arc) -> Self { - self.settings.cancel_token = Some(token); + self.options.cancel_token = Some(token); self } pub fn stall_token(mut self, token: CancellationToken) -> Self { - self.settings.stall_token = Some(token); + self.options.stall_token = Some(token); self } pub fn max_node_visits(mut self, limit: usize) -> Self { - self.settings.max_node_visits = Some(limit); + self.options.max_node_visits = Some(limit); self } @@ -74,7 +74,7 @@ impl ExecutorBuilder { Executor { handler: self.handler, lifecycle: self.lifecycle.unwrap_or_else(|| Box::new(NoopLifecycle)), - settings: self.settings, + options: self.options, } } } @@ -89,7 +89,7 @@ impl Executor { loop { // Check cancellation - if let Some(ref token) = self.settings.cancel_token { + if let Some(ref token) = self.options.cancel_token { if token.load(Ordering::Relaxed) { state.cancelled = true; let outcome = Outcome::fail("run cancelled"); @@ -151,7 +151,7 @@ impl Executor { }); } } - if let Some(global_max) = self.settings.max_node_visits { + if let Some(global_max) = self.options.max_node_visits { if visits >= global_max { return Err(CoreError::VisitLimitExceeded { node_id: node.id().to_string(), @@ -176,7 +176,7 @@ impl Executor { } NodeDecision::Continue => { // Execute with retry, racing against stall token - let mut result = if let Some(ref stall) = self.settings.stall_token { + let mut result = if let Some(ref stall) = self.options.stall_token { tokio::select! { r = self.execute_with_retry(&node, &state, graph) => r?, () = stall.cancelled() => { diff --git a/lib/crates/fabro-workflows/src/handler/manager_loop.rs b/lib/crates/fabro-workflows/src/handler/manager_loop.rs index e1e27f167..3150d7056 100644 --- a/lib/crates/fabro-workflows/src/handler/manager_loop.rs +++ b/lib/crates/fabro-workflows/src/handler/manager_loop.rs @@ -143,7 +143,7 @@ impl Handler for SubWorkflowHandler { let child_cancel = Arc::clone(&cancel_token); let git_state = services.git_state(); - let child_config = RunOptions { + let child_run_options = RunOptions { config: fabro_config::FabroConfig::default(), run_dir: child_logs, cancel_token: Some(cancel_token), @@ -184,7 +184,7 @@ impl Handler for SubWorkflowHandler { let initialized = Initialized { graph: child_graph, source: String::new(), - settings: child_config, + run_options: child_run_options, checkpoint: None, seed_context: Some(child_context), emitter, diff --git a/lib/crates/fabro-workflows/src/lifecycle/disk.rs b/lib/crates/fabro-workflows/src/lifecycle/disk.rs index e079c8eb5..08bba5c4c 100644 --- a/lib/crates/fabro-workflows/src/lifecycle/disk.rs +++ b/lib/crates/fabro-workflows/src/lifecycle/disk.rs @@ -25,7 +25,7 @@ pub struct DiskLifecycle { pub run_dir: PathBuf, pub run_id: String, pub graph: Arc, - pub config: Arc, + pub run_options: Arc, pub emitter: Arc, pub circuit_breaker: Arc, pub checkpoint_enabled: bool, @@ -39,7 +39,7 @@ impl RunLifecycle for DiskLifecycle { _state: &WfRunState, ) -> fabro_core::error::Result<()> { // Write start.json - write_start_record(&self.run_dir, &self.config); + write_start_record(&self.run_dir, &self.run_options); // Write run status as Running crate::run_status::write_run_status( &self.run_dir, diff --git a/lib/crates/fabro-workflows/src/lifecycle/git.rs b/lib/crates/fabro-workflows/src/lifecycle/git.rs index 31ab10073..b01d5d15e 100644 --- a/lib/crates/fabro-workflows/src/lifecycle/git.rs +++ b/lib/crates/fabro-workflows/src/lifecycle/git.rs @@ -35,7 +35,7 @@ pub struct GitLifecycle { pub emitter: Arc, pub run_dir: PathBuf, pub run_id: String, - pub config: Arc, + pub run_options: Arc, pub start_node_id: Option, // Cross-lifecycle data (shared with EventLifecycle) pub checkpoint_git_result: Arc>>, @@ -55,13 +55,13 @@ impl RunLifecycle for GitLifecycle { // Init metadata branch (best-effort) if let (Some(_), Some(repo_path)) = ( - self.config + self.run_options .git .as_ref() .and_then(|g| g.meta_branch.as_ref()), - self.config.host_repo_path.as_ref(), + self.run_options.host_repo_path.as_ref(), ) { - let store = crate::git::MetadataStore::new(repo_path, &self.config.git_author); + let store = crate::git::MetadataStore::new(repo_path, &self.run_options.git_author); let run_json = std::fs::read(self.run_dir.join("run.json")).ok(); let start_json = std::fs::read(self.run_dir.join("start.json")).ok(); let sandbox_json = std::fs::read(self.run_dir.join("sandbox.json")).ok(); @@ -97,20 +97,20 @@ impl RunLifecycle for GitLifecycle { let node_id = node.id(); // Skip git checkpoint for the start node (always empty) or if git disabled - if self.start_node_id.as_deref() == Some(node_id) || self.config.git.is_none() { + if self.start_node_id.as_deref() == Some(node_id) || self.run_options.git.is_none() { *self.checkpoint_git_result.lock().unwrap() = None; return Ok(()); } // Shadow commit (best-effort, metadata branch) let shadow_sha: Option = if let (Some(_), Some(repo_path)) = ( - self.config + self.run_options .git .as_ref() .and_then(|g| g.meta_branch.as_ref()), - self.config.host_repo_path.as_ref(), + self.run_options.host_repo_path.as_ref(), ) { - let store = crate::git::MetadataStore::new(repo_path, &self.config.git_author); + let store = crate::git::MetadataStore::new(repo_path, &self.run_options.git_author); // Build checkpoint JSON for shadow branch let checkpoint_path = self.run_dir.join("checkpoint.json"); std::fs::read(&checkpoint_path).ok().and_then(|cp_json| { @@ -158,8 +158,8 @@ impl RunLifecycle for GitLifecycle { &result.outcome.status.to_string(), completed_count, shadow_sha, - self.config.checkpoint_exclude_globs(), - &self.config.git_author, + self.run_options.checkpoint_exclude_globs(), + &self.run_options.git_author, ) .await; @@ -186,18 +186,21 @@ impl RunLifecycle for GitLifecycle { } // Push run branch (skip in dry-run mode) - if !self.config.dry_run { - if let Some(branch) = - self.config.git.as_ref().and_then(|g| g.run_branch.as_ref()) + if !self.run_options.dry_run { + if let Some(branch) = self + .run_options + .git + .as_ref() + .and_then(|g| g.run_branch.as_ref()) { let push_ok = if self.sandbox.git_push_branch(branch).await { true - } else if let Some(repo_path) = self.config.host_repo_path.as_ref() { + } else if let Some(repo_path) = self.run_options.host_repo_path.as_ref() { let refspec = format!("refs/heads/{branch}"); git_push_host( repo_path, &refspec, - &self.config.github_app, + &self.run_options.github_app, "run branch", ) .await @@ -208,17 +211,17 @@ impl RunLifecycle for GitLifecycle { } // Push metadata branch (always from host) if let (Some(meta_branch), Some(repo_path)) = ( - self.config + self.run_options .git .as_ref() .and_then(|g| g.meta_branch.as_ref()), - self.config.host_repo_path.as_ref(), + self.run_options.host_repo_path.as_ref(), ) { let refspec = format!("refs/heads/{meta_branch}"); let meta_push_ok = git_push_host( repo_path, &refspec, - &self.config.github_app, + &self.run_options.github_app, "metadata branch", ) .await; @@ -235,7 +238,12 @@ impl RunLifecycle for GitLifecycle { .lock() .unwrap() .clone() - .or_else(|| self.config.git.as_ref().and_then(|g| g.base_sha.clone())) + .or_else(|| { + self.run_options + .git + .as_ref() + .and_then(|g| g.base_sha.clone()) + }) .unwrap_or_else(|| sha.clone()); let diff_dest = node_dir(&self.run_dir, node_id, visit).join("diff.patch"); @@ -275,9 +283,14 @@ impl RunLifecycle for GitLifecycle { async fn on_run_end(&self, outcome: &Outcome, _state: &WfRunState) { // Write final.patch on success if (outcome.status == StageStatus::Success || outcome.status == StageStatus::PartialSuccess) - && self.config.git.is_some() + && self.run_options.git.is_some() { - if let Some(base_sha) = self.config.git.as_ref().and_then(|g| g.base_sha.clone()) { + if let Some(base_sha) = self + .run_options + .git + .as_ref() + .and_then(|g| g.base_sha.clone()) + { let diff_dest = self.run_dir.join("final.patch"); match git_diff(&*self.sandbox, &base_sha).await { Ok(patch) if !patch.is_empty() => { diff --git a/lib/crates/fabro-workflows/src/lifecycle/mod.rs b/lib/crates/fabro-workflows/src/lifecycle/mod.rs index b2b9a0650..ba4d6c169 100644 --- a/lib/crates/fabro-workflows/src/lifecycle/mod.rs +++ b/lib/crates/fabro-workflows/src/lifecycle/mod.rs @@ -79,7 +79,7 @@ impl WorkflowLifecycle { sandbox: Arc, graph: Arc, run_dir: PathBuf, - config: Arc, + run_options: Arc, is_resume: bool, ) -> Self { let restarted_from: Arc>> = Arc::new(Mutex::new(None)); @@ -91,7 +91,7 @@ impl WorkflowLifecycle { let circuit_breaker = Arc::new(CircuitBreakerLifecycle::new(loop_restart_signature_limit)); - let has_run_branch = config + let has_run_branch = run_options .git .as_ref() .and_then(|g| g.run_branch.as_ref()) @@ -106,11 +106,11 @@ impl WorkflowLifecycle { let event = EventLifecycle { emitter: Arc::clone(&emitter), graph_name: graph.name.clone(), - run_id: config.run_id.clone(), + run_id: run_options.run_id.clone(), run_start: Mutex::new(Instant::now()), restarted_from: Arc::clone(&restarted_from), - base_sha: config.git.as_ref().and_then(|g| g.base_sha.clone()), - run_branch: config.git.as_ref().and_then(|g| g.run_branch.clone()), + base_sha: run_options.git.as_ref().and_then(|g| g.base_sha.clone()), + run_branch: run_options.git.as_ref().and_then(|g| g.run_branch.clone()), worktree_dir: working_directory.clone(), goal: (!graph.goal().is_empty()).then(|| graph.goal().to_string()), artifact_store: Arc::clone(&artifact_store), @@ -122,7 +122,7 @@ impl WorkflowLifecycle { hook_runner, sandbox: Arc::clone(&sandbox), hook_work_dir: working_directory.clone().map(PathBuf::from), - run_id: config.run_id.clone(), + run_id: run_options.run_id.clone(), graph_name: graph.name.clone(), }; @@ -130,9 +130,9 @@ impl WorkflowLifecycle { let disk = DiskLifecycle { run_dir: run_dir.clone(), - run_id: config.run_id.clone(), + run_id: run_options.run_id.clone(), graph: Arc::clone(&graph), - config: Arc::clone(&config), + run_options: Arc::clone(&run_options), emitter: Arc::clone(&emitter), circuit_breaker: Arc::clone(&circuit_breaker), checkpoint_enabled: true, @@ -145,8 +145,8 @@ impl WorkflowLifecycle { artifact_store: Arc::clone(&artifact_store), emitter: Arc::clone(&emitter), run_dir: run_dir.clone(), - run_id: config.run_id.clone(), - config: Arc::clone(&config), + run_id: run_options.run_id.clone(), + run_options: Arc::clone(&run_options), start_node_id, checkpoint_git_result: Arc::clone(&checkpoint_git_result), last_git_sha: Arc::clone(&last_git_sha), @@ -158,7 +158,7 @@ impl WorkflowLifecycle { Some(run_dir.clone()), Arc::clone(&emitter), run_dir, - config.asset_globs().to_vec(), + run_options.asset_globs().to_vec(), ); Self { @@ -174,7 +174,7 @@ impl WorkflowLifecycle { checkpoint_git_result, is_initial_resume: AtomicBool::new(is_resume), graph, - run_id: config.run_id.clone(), + run_id: run_options.run_id.clone(), working_directory, } } diff --git a/lib/crates/fabro-workflows/src/operations/create.rs b/lib/crates/fabro-workflows/src/operations/create.rs index 2f9d6302b..5b8ff959e 100644 --- a/lib/crates/fabro-workflows/src/operations/create.rs +++ b/lib/crates/fabro-workflows/src/operations/create.rs @@ -62,13 +62,13 @@ pub fn validate_from_file(path: &Path) -> Result { } /// Parse, transform, validate, normalize config, and persist a run. -pub fn create(dot_source: &str, settings: RunCreateOptions) -> Result { +pub fn create(dot_source: &str, options: RunCreateOptions) -> Result { let validated = preprocess_and_validate( dot_source, - settings.base_dir.clone(), + options.base_dir.clone(), Vec::new(), - Some(&settings.config), - settings.goal_override.as_deref(), + Some(&options.config), + options.goal_override.as_deref(), )?; if validated.has_errors() { @@ -77,19 +77,19 @@ pub fn create(dot_source: &str, settings: RunCreateOptions) -> Result Result { let source = std::fs::read_to_string(path) .map_err(|e| FabroError::Parse(format!("Failed to read {}: {e}", path.display())))?; let base_dir = path.parent().unwrap_or(Path::new(".")); - settings.base_dir = Some(base_dir.to_path_buf()); - create(&source, settings) + options.base_dir = Some(base_dir.to_path_buf()); + create(&source, options) } /// Build a persisted workflow from an already-materialized graph. @@ -99,13 +99,13 @@ pub fn create_from_file( #[doc(hidden)] pub fn create_from_graph( mut graph: Graph, - settings: RunCreateOptions, + options: RunCreateOptions, ) -> Result { - if let Some(goal_override) = settings.goal_override.as_deref() { + if let Some(goal_override) = options.goal_override.as_deref() { apply_goal_override(&mut graph, Some(goal_override)); } let validated = Validated::new(graph, String::new(), vec![]); - persist_validated(validated, settings) + persist_validated(validated, options) } fn preprocess_and_validate( @@ -150,7 +150,7 @@ fn apply_goal_override(graph: &mut Graph, goal_override: Option<&str>) { fn persist_validated( validated: Validated, - settings: RunCreateOptions, + options: RunCreateOptions, ) -> Result { let RunCreateOptions { mut config, @@ -163,7 +163,7 @@ fn persist_validated( host_repo_path, goal_override: _, base_dir: _, - } = settings; + } = options; finalize_config(&mut config, validated.graph()); diff --git a/lib/crates/fabro-workflows/src/operations/start.rs b/lib/crates/fabro-workflows/src/operations/start.rs index 53ecaa10b..04d8de4ec 100644 --- a/lib/crates/fabro-workflows/src/operations/start.rs +++ b/lib/crates/fabro-workflows/src/operations/start.rs @@ -78,10 +78,10 @@ pub async fn start(persisted: Persisted, options: StartOptions) -> Result Result RunOptions { + fn test_run_options(run_dir: &std::path::Path) -> RunOptions { RunOptions { config: FabroConfig::default(), run_dir: run_dir.to_path_buf(), @@ -410,7 +410,7 @@ mod tests { sandbox, registry, lifecycle, - run_options: test_settings(run_dir), + run_options: test_run_options(run_dir), hooks: fabro_hooks::HookConfig { hooks: vec![] }, sandbox_env: HashMap::new(), checkpoint: None, diff --git a/lib/crates/fabro-workflows/src/pipeline/execute.rs b/lib/crates/fabro-workflows/src/pipeline/execute.rs index 82c896727..eb17ef5e5 100644 --- a/lib/crates/fabro-workflows/src/pipeline/execute.rs +++ b/lib/crates/fabro-workflows/src/pipeline/execute.rs @@ -33,7 +33,7 @@ pub async fn execute(init: Initialized) -> Executed { let Initialized { graph, source: _, - settings, + run_options, checkpoint, seed_context, emitter, @@ -48,15 +48,15 @@ pub async fn execute(init: Initialized) -> Executed { let graph_arc = Arc::new(graph.clone()); let wf_graph = WorkflowGraph(Arc::clone(&graph_arc)); - let git_state = settings.git.as_ref().and_then(|git| { + let git_state = run_options.git.as_ref().and_then(|git| { let base_sha = git.base_sha.clone()?; Some(Arc::new(GitState { - run_id: settings.run_id.clone(), + run_id: run_options.run_id.clone(), base_sha, run_branch: git.run_branch.clone(), meta_branch: git.meta_branch.clone(), - checkpoint_exclude_globs: settings.checkpoint_exclude_globs().to_vec(), - git_author: settings.git_author.clone(), + checkpoint_exclude_globs: run_options.checkpoint_exclude_globs().to_vec(), + git_author: run_options.git_author.clone(), })) }); @@ -72,17 +72,17 @@ pub async fn execute(init: Initialized) -> Executed { let handler = Arc::new(WorkflowNodeHandler { services: shared_services, - run_dir: settings.run_dir.clone(), + run_dir: run_options.run_dir.clone(), graph: Arc::clone(&graph_arc), }); - let settings_arc = Arc::new(settings.clone()); + let settings_arc = Arc::new(run_options.clone()); let lifecycle = WorkflowLifecycle::new( Arc::clone(&emitter), hook_runner.clone(), Arc::clone(&sandbox), graph_arc, - settings.run_dir.clone(), + run_options.run_dir.clone(), settings_arc, checkpoint.is_some(), ); @@ -134,7 +134,7 @@ pub async fn execute(init: Initialized) -> Executed { return Executed { graph, outcome: Err(err), - settings, + run_options, hook_runner, emitter, sandbox, @@ -155,7 +155,7 @@ pub async fn execute(init: Initialized) -> Executed { return Executed { graph, outcome: Err(err), - settings, + run_options, hook_runner, emitter, sandbox, @@ -171,7 +171,7 @@ pub async fn execute(init: Initialized) -> Executed { return Executed { graph, outcome: Err(err), - settings, + run_options, hook_runner, emitter, sandbox, @@ -187,7 +187,7 @@ pub async fn execute(init: Initialized) -> Executed { let graph_max = graph.max_node_visits(); let max_node_visits = if graph_max > 0 { Some(graph_max as usize) - } else if settings.dry_run { + } else if run_options.dry_run { Some(10) } else { None @@ -235,7 +235,7 @@ pub async fn execute(init: Initialized) -> Executed { ExecutorBuilder::new(handler as Arc>) .lifecycle(Box::new(lifecycle)); - if let Some(ref cancel) = settings.cancel_token { + if let Some(ref cancel) = run_options.cancel_token { builder = builder.cancel_token(cancel.clone()); } if let Some(token) = stall_token.clone() { @@ -290,7 +290,7 @@ pub async fn execute(init: Initialized) -> Executed { Executed { graph, outcome, - settings, + run_options, hook_runner, emitter, sandbox, diff --git a/lib/crates/fabro-workflows/src/pipeline/execute/tests.rs b/lib/crates/fabro-workflows/src/pipeline/execute/tests.rs index 29f8e6e42..1e7ae39a4 100644 --- a/lib/crates/fabro-workflows/src/pipeline/execute/tests.rs +++ b/lib/crates/fabro-workflows/src/pipeline/execute/tests.rs @@ -66,7 +66,7 @@ fn make_registry() -> HandlerRegistry { registry } -fn test_settings(run_dir: &Path, run_id: &str) -> RunOptions { +fn test_run_options(run_dir: &Path, run_id: &str) -> RunOptions { RunOptions { run_dir: run_dir.to_path_buf(), cancel_token: None, @@ -160,7 +160,7 @@ async fn execute_runs_start_to_exit_and_returns_final_context() { setup_command_timeout_ms: 1_000, devcontainer_phases: vec![], }, - run_options: test_settings(&run_dir, "run-test"), + run_options: test_run_options(&run_dir, "run-test"), hooks: HookConfig { hooks: vec![] }, sandbox_env: HashMap::new(), checkpoint: None, @@ -189,22 +189,22 @@ async fn run_with_lifecycle( emitter: Arc, sandbox: Arc, graph: &Graph, - settings: RunOptions, + run_options: RunOptions, lifecycle: LifecycleOptions, ) -> Result { - let run_dir = settings.run_dir.clone(); - let run_id = settings.run_id.clone(); + let run_dir = run_options.run_dir.clone(); + let run_id = run_options.run_id.clone(); std::fs::create_dir_all(&run_dir).unwrap(); let initialized = initialize( persisted_workflow(graph.clone(), String::new(), &run_dir, &run_id), InitOptions { run_id, - dry_run: settings.dry_run, + dry_run: run_options.dry_run, emitter, sandbox, registry: Arc::new(registry), lifecycle, - run_options: settings, + run_options, hooks: HookConfig { hooks: vec![] }, sandbox_env: HashMap::new(), checkpoint: None, @@ -375,7 +375,7 @@ async fn execute_runs_simple_workflow() { Arc::new(EventEmitter::new()), local_env(), &simple_graph(), - &test_settings(dir.path(), "test-run"), + &test_run_options(dir.path(), "test-run"), ) .await .unwrap(); @@ -390,7 +390,7 @@ async fn execute_saves_checkpoint() { Arc::new(EventEmitter::new()), local_env(), &simple_graph(), - &test_settings(dir.path(), "test-run"), + &test_run_options(dir.path(), "test-run"), ) .await .unwrap(); @@ -412,7 +412,7 @@ async fn execute_emits_events() { Arc::new(emitter), local_env(), &simple_graph(), - &test_settings(dir.path(), "test-run"), + &test_run_options(dir.path(), "test-run"), ) .await .unwrap(); @@ -428,7 +428,7 @@ async fn execute_error_when_no_start_node() { Arc::new(EventEmitter::new()), local_env(), &Graph::new("empty"), - &test_settings(dir.path(), "test-run"), + &test_run_options(dir.path(), "test-run"), ) .await; assert!(result.is_err()); @@ -442,7 +442,7 @@ async fn execute_mirrors_graph_goal_to_context() { Arc::new(EventEmitter::new()), local_env(), &simple_graph(), - &test_settings(dir.path(), "test-run"), + &test_run_options(dir.path(), "test-run"), ) .await .unwrap(); @@ -491,7 +491,7 @@ async fn execute_conditional_routing_uses_unconditional_success_path() { Arc::new(EventEmitter::new()), local_env(), &g, - &test_settings(dir.path(), "test-run"), + &test_run_options(dir.path(), "test-run"), ) .await .unwrap(); @@ -504,8 +504,8 @@ async fn execute_conditional_routing_uses_unconditional_success_path() { #[tokio::test] async fn execute_writes_start_json_and_node_status() { let dir = tempfile::tempdir().unwrap(); - let mut settings = test_settings(dir.path(), "test-run"); - settings.git = Some(GitCheckpointOptions { + let mut run_options = test_run_options(dir.path(), "test-run"); + run_options.git = Some(GitCheckpointOptions { base_sha: Some("abc123".into()), run_branch: Some("fabro/run/test-run".into()), meta_branch: None, @@ -516,7 +516,7 @@ async fn execute_writes_start_json_and_node_status() { Arc::new(EventEmitter::new()), local_env(), &simple_graph(), - &settings, + &run_options, ) .await .unwrap(); @@ -575,7 +575,7 @@ async fn timeout_causes_fail_status_json() { Arc::new(EventEmitter::new()), local_env(), &g, - &test_settings(dir.path(), "test-run"), + &test_run_options(dir.path(), "test-run"), ) .await .unwrap(); @@ -604,8 +604,8 @@ async fn execute_cancelled_mid_run() { let cancel_token_clone = Arc::clone(&cancel_token); let mut registry = make_registry(); registry.register("slow", Box::new(SlowHandler { sleep_ms: 200 })); - let mut settings = test_settings(dir.path(), "test-run"); - settings.cancel_token = Some(cancel_token); + let mut run_options = test_run_options(dir.path(), "test-run"); + run_options.cancel_token = Some(cancel_token); tokio::spawn(async move { tokio::time::sleep(Duration::from_millis(50)).await; @@ -617,7 +617,7 @@ async fn execute_cancelled_mid_run() { Arc::new(EventEmitter::new()), local_env(), &g, - &settings, + &run_options, ) .await; assert!(matches!(result, Err(FabroError::Cancelled))); @@ -635,7 +635,7 @@ async fn max_node_visits_errors_on_cycle() { Arc::new(EventEmitter::new()), local_env(), &g, - &test_settings(dir.path(), "test-run"), + &test_run_options(dir.path(), "test-run"), ) .await; let err = result.unwrap_err().to_string(); @@ -670,7 +670,7 @@ async fn panic_handler_writes_panic_txt() { Arc::new(EventEmitter::new()), local_env(), &g, - &test_settings(dir.path(), "test-run"), + &test_run_options(dir.path(), "test-run"), ) .await; @@ -691,7 +691,7 @@ async fn loop_circuit_breaker_aborts_on_repeated_failure() { Arc::new(EventEmitter::new()), local_env(), &looping_fail_graph(), - &test_settings(dir.path(), "test-run"), + &test_run_options(dir.path(), "test-run"), ) .await; let err = result.unwrap_err().to_string(); @@ -740,7 +740,7 @@ async fn stall_watchdog_triggers_on_hung_handler() { Arc::new(EventEmitter::new()), local_env(), &g, - &test_settings(dir.path(), "test-run"), + &test_run_options(dir.path(), "test-run"), ) .await; let err = result.unwrap_err().to_string(); @@ -804,7 +804,7 @@ async fn retry_emits_stage_started_per_attempt() { Arc::new(emitter), local_env(), &g, - &test_settings(dir.path(), "retry-events-test"), + &test_run_options(dir.path(), "retry-events-test"), ) .await .unwrap(); @@ -845,7 +845,7 @@ async fn run_with_lifecycle_emits_initialize_and_setup_events() { Arc::new(emitter), local_env(), &simple_graph(), - test_settings(dir.path(), "order-test"), + test_run_options(dir.path(), "order-test"), test_lifecycle(vec!["echo ok".to_string()]), ) .await @@ -916,17 +916,23 @@ async fn git_checkpoint_skips_start_node() { }); let sandbox: Arc = Arc::new(fabro_agent::LocalSandbox::new(repo.to_path_buf())); - let mut settings = test_settings(run_tmp.path(), "git-cp-test"); - settings.git = Some(GitCheckpointOptions { + let mut run_options = test_run_options(run_tmp.path(), "git-cp-test"); + run_options.git = Some(GitCheckpointOptions { base_sha: Some(base_sha), run_branch: None, meta_branch: Some(crate::git::MetadataStore::branch_name("git-cp-test")), }); - settings.host_repo_path = Some(repo.to_path_buf()); + run_options.host_repo_path = Some(repo.to_path_buf()); - run_graph(make_registry(), Arc::new(emitter), sandbox, &g, &settings) - .await - .unwrap(); + run_graph( + make_registry(), + Arc::new(emitter), + sandbox, + &g, + &run_options, + ) + .await + .unwrap(); let collected = events.lock().unwrap(); let checkpoint_node_ids: Vec<&str> = collected diff --git a/lib/crates/fabro-workflows/src/pipeline/finalize.rs b/lib/crates/fabro-workflows/src/pipeline/finalize.rs index 4b02d3444..236c44758 100644 --- a/lib/crates/fabro-workflows/src/pipeline/finalize.rs +++ b/lib/crates/fabro-workflows/src/pipeline/finalize.rs @@ -151,15 +151,18 @@ pub fn persist_terminal_outcome( /// /// This captures the last diff.patch (written after the final checkpoint) and retro.json. /// Best-effort: errors are logged as warnings. -pub async fn write_finalize_commit(config: &RunOptions, run_dir: &Path) { +pub async fn write_finalize_commit(run_options: &RunOptions, run_dir: &Path) { let (Some(meta_branch), Some(repo_path)) = ( - config.git.as_ref().and_then(|g| g.meta_branch.as_ref()), - config.host_repo_path.as_ref(), + run_options + .git + .as_ref() + .and_then(|g| g.meta_branch.as_ref()), + run_options.host_repo_path.as_ref(), ) else { return; }; - let store = crate::git::MetadataStore::new(repo_path, &config.git_author); + let store = crate::git::MetadataStore::new(repo_path, &run_options.git_author); let mut entries = crate::git::scan_node_files(run_dir); if let Ok(retro_bytes) = std::fs::read(run_dir.join("retro.json")) { entries.push(("retro.json".to_string(), retro_bytes)); @@ -168,14 +171,19 @@ pub async fn write_finalize_commit(config: &RunOptions, run_dir: &Path) { .iter() .map(|(k, v)| (k.as_str(), v.as_slice())) .collect(); - if let Err(e) = store.write_files(&config.run_id, &refs, "finalize run") { + if let Err(e) = store.write_files(&run_options.run_id, &refs, "finalize run") { tracing::warn!(error = %e, "Failed to write finalize commit to metadata branch"); return; } let refspec = format!("refs/heads/{meta_branch}"); - crate::sandbox_git::git_push_host(repo_path, &refspec, &config.github_app, "finalize metadata") - .await; + crate::sandbox_git::git_push_host( + repo_path, + &refspec, + &run_options.github_app, + "finalize metadata", + ) + .await; } async fn run_hooks( @@ -220,7 +228,7 @@ pub async fn finalize( let Retroed { graph, outcome, - settings, + run_options, hook_runner, emitter, sandbox, @@ -238,7 +246,7 @@ pub async fn finalize( options.last_git_sha.clone(), ); - write_finalize_commit(&settings, &options.run_dir).await; + write_finalize_commit(&run_options, &options.run_dir).await; if options.preserve_sandbox { let info = sandbox.sandbox_info(); @@ -279,12 +287,12 @@ pub async fn finalize( persist_terminal_outcome(&options.run_dir, &conclusion, run_status, status_reason); Ok(Concluded { - run_id: settings.run_id.clone(), + run_id: run_options.run_id.clone(), outcome, conclusion, - pushed_branch: settings.git.as_ref().and_then(|g| g.run_branch.clone()), + pushed_branch: run_options.git.as_ref().and_then(|g| g.run_branch.clone()), graph, - settings, + run_options, emitter, }) } @@ -301,7 +309,7 @@ mod tests { use crate::pipeline::types::Retroed; use crate::run_options::RunOptions; - fn test_settings(run_dir: &std::path::Path) -> RunOptions { + fn test_run_options(run_dir: &std::path::Path) -> RunOptions { RunOptions { config: FabroConfig::default(), run_dir: run_dir.to_path_buf(), @@ -326,7 +334,7 @@ mod tests { let retroed = Retroed { graph: Graph::new("test"), outcome: Ok(Outcome::success()), - settings: test_settings(&run_dir), + run_options: test_run_options(&run_dir), hook_runner: None, emitter: Arc::new(EventEmitter::new()), sandbox: Arc::new(fabro_agent::LocalSandbox::new( diff --git a/lib/crates/fabro-workflows/src/pipeline/initialize.rs b/lib/crates/fabro-workflows/src/pipeline/initialize.rs index 5a18664e3..7be3c0876 100644 --- a/lib/crates/fabro-workflows/src/pipeline/initialize.rs +++ b/lib/crates/fabro-workflows/src/pipeline/initialize.rs @@ -173,7 +173,7 @@ pub async fn initialize( Ok(Initialized { graph, source, - settings: options.run_options, + run_options: options.run_options, checkpoint: options.checkpoint, seed_context: options.seed_context, emitter: options.emitter, @@ -298,7 +298,7 @@ mod tests { .await .unwrap(); - assert_eq!(initialized.settings.run_dir, run_dir); + assert_eq!(initialized.run_options.run_dir, run_dir); assert_eq!(initialized.source, source); assert!(initialized.hook_runner.is_none()); assert_eq!( @@ -402,7 +402,7 @@ mod tests { .await .unwrap(); - assert_eq!(initialized.settings.run_dir, run_dir); + assert_eq!(initialized.run_options.run_dir, run_dir); assert_eq!(initialized.source, source); } } diff --git a/lib/crates/fabro-workflows/src/pipeline/pull_request.rs b/lib/crates/fabro-workflows/src/pipeline/pull_request.rs index 502fb4883..4d94c1bbe 100644 --- a/lib/crates/fabro-workflows/src/pipeline/pull_request.rs +++ b/lib/crates/fabro-workflows/src/pipeline/pull_request.rs @@ -482,7 +482,7 @@ pub async fn pull_request(concluded: Concluded, options: &PullRequestOptions) -> conclusion, pushed_branch, graph, - settings, + run_options, emitter, } = concluded; @@ -499,7 +499,7 @@ pub async fn pull_request(concluded: Concluded, options: &PullRequestOptions) -> .await .unwrap_or_default(); if let (Some(base_branch), Some(run_branch), Some(creds), Some(origin)) = ( - &settings.base_branch, + &run_options.base_branch, pushed_branch.as_deref(), &options.github_app, &options.origin_url, diff --git a/lib/crates/fabro-workflows/src/pipeline/retro.rs b/lib/crates/fabro-workflows/src/pipeline/retro.rs index 24a106d92..6f9a022c4 100644 --- a/lib/crates/fabro-workflows/src/pipeline/retro.rs +++ b/lib/crates/fabro-workflows/src/pipeline/retro.rs @@ -116,7 +116,7 @@ pub async fn retro(executed: Executed, options: &RetroOptions) -> Retroed { let Executed { graph, outcome, - settings, + run_options, hook_runner, emitter, sandbox, @@ -133,7 +133,7 @@ pub async fn retro(executed: Executed, options: &RetroOptions) -> Retroed { Retroed { graph, outcome, - settings, + run_options, hook_runner, emitter, sandbox, @@ -176,7 +176,7 @@ mod tests { checkpoint.save(&run_dir.join("checkpoint.json")).unwrap(); } - fn test_settings(run_dir: &std::path::Path) -> RunOptions { + fn test_run_options(run_dir: &std::path::Path) -> RunOptions { RunOptions { config: FabroConfig::default(), run_dir: run_dir.to_path_buf(), @@ -207,7 +207,7 @@ mod tests { let executed = Executed { graph: Graph::new("test"), outcome: Ok(crate::outcome::Outcome::success()), - settings: test_settings(&run_dir), + run_options: test_run_options(&run_dir), hook_runner: None, emitter: Arc::clone(&emitter), sandbox: Arc::clone(&sandbox), diff --git a/lib/crates/fabro-workflows/src/pipeline/types.rs b/lib/crates/fabro-workflows/src/pipeline/types.rs index e1468603e..59ef6886c 100644 --- a/lib/crates/fabro-workflows/src/pipeline/types.rs +++ b/lib/crates/fabro-workflows/src/pipeline/types.rs @@ -208,7 +208,7 @@ pub struct InitOptions { pub struct Initialized { pub graph: Graph, pub source: String, - pub settings: RunOptions, + pub run_options: RunOptions, pub(crate) checkpoint: Option, pub(crate) seed_context: Option, pub emitter: Arc, @@ -224,7 +224,7 @@ pub struct Initialized { pub struct Executed { pub graph: Graph, pub outcome: Result, - pub settings: RunOptions, + pub run_options: RunOptions, pub hook_runner: Option>, pub emitter: Arc, pub sandbox: Arc, @@ -237,7 +237,7 @@ pub struct Executed { pub struct Retroed { pub graph: Graph, pub outcome: Result, - pub settings: RunOptions, + pub run_options: RunOptions, pub hook_runner: Option>, pub emitter: Arc, pub sandbox: Arc, @@ -253,7 +253,7 @@ pub struct Concluded { pub conclusion: Conclusion, pub pushed_branch: Option, pub graph: Graph, - pub settings: RunOptions, + pub run_options: RunOptions, pub emitter: Arc, } diff --git a/lib/crates/fabro-workflows/src/test_support.rs b/lib/crates/fabro-workflows/src/test_support.rs index 8766d3cf5..747c4db76 100644 --- a/lib/crates/fabro-workflows/src/test_support.rs +++ b/lib/crates/fabro-workflows/src/test_support.rs @@ -23,14 +23,14 @@ fn initialized( emitter: Arc, sandbox: Arc, graph: &fabro_graphviz::graph::Graph, - settings: &RunOptions, + run_options: &RunOptions, options: InitializedOptions, ) -> Initialized { - std::fs::create_dir_all(&settings.run_dir).expect("failed to create run dir"); + std::fs::create_dir_all(&run_options.run_dir).expect("failed to create run dir"); Initialized { graph: graph.clone(), source: String::new(), - settings: settings.clone(), + run_options: run_options.clone(), checkpoint: options.checkpoint, seed_context: None, emitter, @@ -38,7 +38,7 @@ fn initialized( registry: Arc::new(registry), hook_runner: options.hook_runner, env: options.env, - dry_run: settings.dry_run, + dry_run: run_options.dry_run, } } @@ -47,14 +47,14 @@ pub async fn run_graph( emitter: Arc, sandbox: Arc, graph: &fabro_graphviz::graph::Graph, - settings: &RunOptions, + run_options: &RunOptions, ) -> Result { let executed = pipeline::execute(initialized( registry, emitter, sandbox, graph, - settings, + run_options, InitializedOptions { hook_runner: None, env: HashMap::new(), @@ -70,7 +70,7 @@ pub async fn run_graph_with_hooks( emitter: Arc, sandbox: Arc, graph: &fabro_graphviz::graph::Graph, - settings: &RunOptions, + run_options: &RunOptions, hook_runner: Arc, env: Option>, ) -> Result { @@ -79,7 +79,7 @@ pub async fn run_graph_with_hooks( emitter, sandbox, graph, - settings, + run_options, InitializedOptions { hook_runner: Some(hook_runner), env: env.unwrap_or_default(), @@ -95,7 +95,7 @@ pub async fn run_graph_from_checkpoint( emitter: Arc, sandbox: Arc, graph: &fabro_graphviz::graph::Graph, - settings: &RunOptions, + run_options: &RunOptions, checkpoint: &Checkpoint, ) -> Result { let executed = pipeline::execute(initialized( @@ -103,7 +103,7 @@ pub async fn run_graph_from_checkpoint( emitter, sandbox, graph, - settings, + run_options, InitializedOptions { hook_runner: None, env: HashMap::new(), @@ -137,7 +137,7 @@ impl WorkflowRunner { pub async fn run( &self, graph: &fabro_graphviz::graph::Graph, - settings: &RunOptions, + run_options: &RunOptions, ) -> Result { let registry = self .registry @@ -150,7 +150,7 @@ impl WorkflowRunner { Arc::clone(&self.emitter), Arc::clone(&self.sandbox), graph, - settings, + run_options, ) .await } @@ -158,7 +158,7 @@ impl WorkflowRunner { pub async fn run_from_checkpoint( &self, graph: &fabro_graphviz::graph::Graph, - settings: &RunOptions, + run_options: &RunOptions, checkpoint: &Checkpoint, ) -> Result { let registry = self @@ -172,7 +172,7 @@ impl WorkflowRunner { Arc::clone(&self.emitter), Arc::clone(&self.sandbox), graph, - settings, + run_options, checkpoint, ) .await diff --git a/lib/crates/fabro-workflows/tests/daytona_integration.rs b/lib/crates/fabro-workflows/tests/daytona_integration.rs index 8f40aeb0d..23ac81cfd 100644 --- a/lib/crates/fabro-workflows/tests/daytona_integration.rs +++ b/lib/crates/fabro-workflows/tests/daytona_integration.rs @@ -388,7 +388,7 @@ async fn daytona_pipeline_artifact_offload_and_sync() { registry.register("exit", Box::new(ExitHandler)); let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::new()), env.clone()); - let config = RunOptions { + let run_options = RunOptions { config: FabroConfig::default(), run_dir: dir.path().to_path_buf(), cancel_token: None, @@ -403,7 +403,7 @@ async fn daytona_pipeline_artifact_offload_and_sync() { git: None, }; let outcome = engine - .run(&graph, &config) + .run(&graph, &run_options) .await .expect("pipeline should succeed"); assert_eq!(outcome.status, StageStatus::Success); @@ -579,7 +579,7 @@ async fn daytona_git_checkpoint_remote_emits_events() { registry.register("exit", Box::new(ExitHandler)); let engine = WorkflowRunner::new(registry, Arc::new(emitter), env.clone()); - let config = RunOptions { + let run_options = RunOptions { config: FabroConfig::default(), run_dir: dir.path().to_path_buf(), cancel_token: None, @@ -598,7 +598,7 @@ async fn daytona_git_checkpoint_remote_emits_events() { }), }; let outcome = engine - .run(&graph, &config) + .run(&graph, &run_options) .await .expect("pipeline should succeed"); assert_eq!(outcome.status, StageStatus::Success); @@ -765,7 +765,7 @@ async fn daytona_parallel_git_branching_e2e() { let engine = WorkflowRunner::new(registry, Arc::new(emitter), Arc::clone(&env)); - let config = RunOptions { + let run_options = RunOptions { config: FabroConfig::default(), run_dir: run_tmp.path().to_path_buf(), cancel_token: None, @@ -784,7 +784,7 @@ async fn daytona_parallel_git_branching_e2e() { }), }; let outcome = engine - .run(&graph, &config) + .run(&graph, &run_options) .await .expect("daytona parallel pipeline should succeed"); assert_eq!( @@ -1141,7 +1141,7 @@ async fn daytona_git_checkpoint_with_shadow_branch() { let meta_branch = MetadataStore::branch_name(&run_id); let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::new()), env.clone()); - let config = RunOptions { + let run_options = RunOptions { config: FabroConfig::default(), run_dir: dir.path().to_path_buf(), cancel_token: None, @@ -1160,7 +1160,7 @@ async fn daytona_git_checkpoint_with_shadow_branch() { }), }; let outcome = engine - .run(&graph, &config) + .run(&graph, &run_options) .await .expect("pipeline should succeed"); assert_eq!(outcome.status, StageStatus::Success); @@ -1281,7 +1281,7 @@ async fn daytona_asset_collection() { graph.edges.push(Edge::new("start", "create_assets")); graph.edges.push(Edge::new("create_assets", "exit")); - let config = RunOptions { + let run_options = RunOptions { config: FabroConfig { assets: Some(fabro_config::run::AssetsConfig { include: vec!["test-results/**".to_string()], @@ -1301,7 +1301,7 @@ async fn daytona_asset_collection() { git: None, }; let outcome = engine - .run(&graph, &config) + .run(&graph, &run_options) .await .expect("pipeline should succeed"); assert_eq!(outcome.status, StageStatus::Success); @@ -1537,7 +1537,7 @@ async fn daytona_git_push_run_branch_to_origin() { registry.register("exit", Box::new(ExitHandler)); let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::new()), env.clone()); - let config = RunOptions { + let run_options = RunOptions { config: FabroConfig::default(), run_dir: dir.path().to_path_buf(), cancel_token: None, @@ -1556,7 +1556,7 @@ async fn daytona_git_push_run_branch_to_origin() { }), }; let outcome = engine - .run(&graph, &config) + .run(&graph, &run_options) .await .expect("pipeline should succeed"); assert_eq!(outcome.status, StageStatus::Success); diff --git a/lib/crates/fabro-workflows/tests/integration.rs b/lib/crates/fabro-workflows/tests/integration.rs index 3fc2842dd..a17b60e1e 100644 --- a/lib/crates/fabro-workflows/tests/integration.rs +++ b/lib/crates/fabro-workflows/tests/integration.rs @@ -193,7 +193,7 @@ async fn end_to_end_linear_pipeline() { Arc::new(EventEmitter::new()), local_env(), ); - let config = RunOptions { + let run_options = RunOptions { config: FabroConfig::default(), run_dir: dir.path().to_path_buf(), cancel_token: None, @@ -208,7 +208,7 @@ async fn end_to_end_linear_pipeline() { git: None, }; let outcome = engine - .run(&graph, &config) + .run(&graph, &run_options) .await .expect("run should succeed"); assert_eq!(outcome.status, StageStatus::Success); @@ -332,7 +332,7 @@ async fn end_to_end_branching_pipeline() { registry.register("conditional", Box::new(ConditionalHandler)); let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::new()), local_env()); - let config = RunOptions { + let run_options = RunOptions { config: FabroConfig::default(), run_dir: dir.path().to_path_buf(), cancel_token: None, @@ -347,7 +347,7 @@ async fn end_to_end_branching_pipeline() { git: None, }; let outcome = engine - .run(&graph, &config) + .run(&graph, &run_options) .await .expect("run should succeed"); assert_eq!(outcome.status, StageStatus::Success); @@ -452,7 +452,7 @@ async fn end_to_end_human_gate_pipeline() { registry.register("human", Box::new(HumanHandler::new(interviewer))); let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::new()), local_env()); - let config = RunOptions { + let run_options = RunOptions { config: FabroConfig::default(), run_dir: dir.path().to_path_buf(), cancel_token: None, @@ -467,7 +467,7 @@ async fn end_to_end_human_gate_pipeline() { git: None, }; let outcome = engine - .run(&graph, &config) + .run(&graph, &run_options) .await .expect("run should succeed"); assert_eq!(outcome.status, StageStatus::Success); @@ -548,7 +548,7 @@ async fn human_gate_aborted_input_fails_closed_without_fail_route() { registry.register("human", Box::new(HumanHandler::new(interviewer))); let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::new()), local_env()); - let config = RunOptions { + let run_options = RunOptions { config: FabroConfig::default(), run_dir: dir.path().to_path_buf(), cancel_token: None, @@ -563,7 +563,7 @@ async fn human_gate_aborted_input_fails_closed_without_fail_route() { git: None, }; let outcome = engine - .run(&graph, &config) + .run(&graph, &run_options) .await .expect("engine should return Ok with fail outcome"); assert_eq!( @@ -659,7 +659,7 @@ async fn human_gate_aborted_input_routes_via_outcome_fail_condition() { registry.register("human", Box::new(HumanHandler::new(interviewer))); let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::new()), local_env()); - let config = RunOptions { + let run_options = RunOptions { config: FabroConfig::default(), run_dir: dir.path().to_path_buf(), cancel_token: None, @@ -674,7 +674,7 @@ async fn human_gate_aborted_input_routes_via_outcome_fail_condition() { git: None, }; let outcome = engine - .run(&graph, &config) + .run(&graph, &run_options) .await .expect("aborted human gate should follow explicit fail route"); assert_eq!(outcome.status, StageStatus::Success); @@ -772,7 +772,7 @@ async fn goal_gate_routes_to_retry_target_on_failure() { registry.register("always_fail", Box::new(AlwaysFailHandler)); let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::new()), local_env()); - let config = RunOptions { + let run_options = RunOptions { config: FabroConfig::default(), run_dir: dir.path().to_path_buf(), cancel_token: None, @@ -786,7 +786,7 @@ async fn goal_gate_routes_to_retry_target_on_failure() { host_repo_path: None, git: None, }; - let result = engine.run(&graph, &config).await; + let result = engine.run(&graph, &run_options).await; assert!( result.is_ok(), "goal gate unsatisfied with no retry_target should return Ok(fail outcome)" @@ -893,7 +893,7 @@ async fn goal_gate_routes_to_retry_target_when_present() { ); let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::new()), local_env()); - let config = RunOptions { + let run_options = RunOptions { config: FabroConfig::default(), run_dir: dir.path().to_path_buf(), cancel_token: None, @@ -908,7 +908,7 @@ async fn goal_gate_routes_to_retry_target_when_present() { git: None, }; let outcome = engine - .run(&graph, &config) + .run(&graph, &run_options) .await .expect("run should eventually succeed after retry"); assert_eq!(outcome.status, StageStatus::Success); @@ -1205,7 +1205,7 @@ async fn retry_on_failure_then_succeed() { ); let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::new()), local_env()); - let config = RunOptions { + let run_options = RunOptions { config: FabroConfig::default(), run_dir: dir.path().to_path_buf(), cancel_token: None, @@ -1220,7 +1220,7 @@ async fn retry_on_failure_then_succeed() { git: None, }; let outcome = engine - .run(&graph, &config) + .run(&graph, &run_options) .await .expect("should succeed after retry"); assert_eq!(outcome.status, StageStatus::Success); @@ -1280,7 +1280,7 @@ async fn pipeline_with_many_nodes() { Arc::new(EventEmitter::new()), local_env(), ); - let config = RunOptions { + let run_options = RunOptions { config: FabroConfig::default(), run_dir: dir.path().to_path_buf(), cancel_token: None, @@ -1295,7 +1295,7 @@ async fn pipeline_with_many_nodes() { git: None, }; let outcome = engine - .run(&graph, &config) + .run(&graph, &run_options) .await .expect("large pipeline should succeed"); assert_eq!(outcome.status, StageStatus::Success); @@ -1602,7 +1602,7 @@ async fn smoke_test_with_mock_codergen_backend() { registry.register("conditional", Box::new(ConditionalHandler)); let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::new()), local_env()); - let config = RunOptions { + let run_options = RunOptions { config: FabroConfig::default(), run_dir: dir.path().to_path_buf(), cancel_token: None, @@ -1617,7 +1617,7 @@ async fn smoke_test_with_mock_codergen_backend() { git: None, }; let outcome = engine - .run(&graph, &config) + .run(&graph, &run_options) .await .expect("smoke test should succeed"); assert_eq!(outcome.status, StageStatus::Success); @@ -1703,7 +1703,7 @@ async fn end_to_end_parallel_fan_out_fan_in() { ); let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::new()), local_env()); - let config = RunOptions { + let run_options = RunOptions { config: FabroConfig::default(), run_dir: dir.path().to_path_buf(), cancel_token: None, @@ -1718,7 +1718,7 @@ async fn end_to_end_parallel_fan_out_fan_in() { git: None, }; let outcome = engine - .run(&graph, &config) + .run(&graph, &run_options) .await .expect("parallel pipeline should succeed"); assert_eq!(outcome.status, StageStatus::Success); @@ -1816,7 +1816,7 @@ async fn resume_from_checkpoint_completes_pipeline() { registry.register("exit", Box::new(ExitHandler)); let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::new()), local_env()); - let config = RunOptions { + let run_options = RunOptions { config: FabroConfig::default(), run_dir: dir.path().to_path_buf(), cancel_token: None, @@ -1831,7 +1831,7 @@ async fn resume_from_checkpoint_completes_pipeline() { git: None, }; let outcome = engine - .run_from_checkpoint(&graph, &config, &checkpoint) + .run_from_checkpoint(&graph, &run_options, &checkpoint) .await .expect("resume should succeed"); assert_eq!(outcome.status, StageStatus::Success); @@ -1915,7 +1915,7 @@ async fn resume_from_checkpoint_preserves_goal_gate_outcomes() { registry.register("exit", Box::new(ExitHandler)); let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::new()), local_env()); - let config = RunOptions { + let run_options = RunOptions { config: FabroConfig::default(), run_dir: dir.path().to_path_buf(), cancel_token: None, @@ -1932,7 +1932,7 @@ async fn resume_from_checkpoint_preserves_goal_gate_outcomes() { // This should succeed because goal gate for gated_work is satisfied // via restored outcomes let outcome = engine - .run_from_checkpoint(&graph, &config, &checkpoint) + .run_from_checkpoint(&graph, &run_options, &checkpoint) .await .expect("resume with goal gate should succeed"); assert_eq!(outcome.status, StageStatus::Success); @@ -1958,7 +1958,7 @@ async fn graph_goal_in_context() { Arc::new(EventEmitter::new()), local_env(), ); - let config = RunOptions { + let run_options = RunOptions { config: FabroConfig::default(), run_dir: dir.path().to_path_buf(), cancel_token: None, @@ -1972,7 +1972,7 @@ async fn graph_goal_in_context() { host_repo_path: None, git: None, }; - engine.run(&graph, &config).await.expect("run"); + engine.run(&graph, &run_options).await.expect("run"); let cp = Checkpoint::load(&dir.path().join("checkpoint.json")).unwrap(); assert_eq!( @@ -1994,7 +1994,7 @@ async fn event_streaming_lifecycle() { let emitter = EventEmitter::new(); let events = collect_events(&emitter); let engine = WorkflowRunner::new(make_linear_registry(), Arc::new(emitter), local_env()); - let config = RunOptions { + let run_options = RunOptions { config: FabroConfig::default(), run_dir: dir.path().to_path_buf(), cancel_token: None, @@ -2008,7 +2008,7 @@ async fn event_streaming_lifecycle() { host_repo_path: None, git: None, }; - engine.run(&graph, &config).await.expect("run"); + engine.run(&graph, &run_options).await.expect("run"); let collected = events.lock().unwrap(); assert!(collected @@ -2074,7 +2074,7 @@ async fn context_flow_between_stages() { Arc::new(EventEmitter::new()), local_env(), ); - let config = RunOptions { + let run_options = RunOptions { config: FabroConfig::default(), run_dir: dir.path().to_path_buf(), cancel_token: None, @@ -2088,7 +2088,7 @@ async fn context_flow_between_stages() { host_repo_path: None, git: None, }; - engine.run(&graph, &config).await.expect("run"); + engine.run(&graph, &run_options).await.expect("run"); let cp = Checkpoint::load(&dir.path().join("checkpoint.json")).unwrap(); assert_eq!( @@ -2127,7 +2127,7 @@ async fn tool_handler_e2e() { Arc::new(EventEmitter::new()), local_env(), ); - let config = RunOptions { + let run_options = RunOptions { config: FabroConfig::default(), run_dir: dir.path().to_path_buf(), cancel_token: None, @@ -2141,7 +2141,7 @@ async fn tool_handler_e2e() { host_repo_path: None, git: None, }; - let outcome = engine.run(&graph, &config).await.expect("run"); + let outcome = engine.run(&graph, &run_options).await.expect("run"); assert_eq!(outcome.status, StageStatus::Success); let cp = Checkpoint::load(&dir.path().join("checkpoint.json")).unwrap(); @@ -2197,7 +2197,7 @@ async fn auto_approve_interviewer_e2e() { Arc::new(EventEmitter::new()), local_env(), ); - let config = RunOptions { + let run_options = RunOptions { config: FabroConfig::default(), run_dir: dir.path().to_path_buf(), cancel_token: None, @@ -2211,7 +2211,7 @@ async fn auto_approve_interviewer_e2e() { host_repo_path: None, git: None, }; - let outcome = engine.run(&graph, &config).await.expect("run"); + let outcome = engine.run(&graph, &run_options).await.expect("run"); assert_eq!(outcome.status, StageStatus::Success); let cp = Checkpoint::load(&dir.path().join("checkpoint.json")).unwrap(); @@ -2234,7 +2234,7 @@ async fn codergen_without_backend_simulated() { Arc::new(EventEmitter::new()), local_env(), ); - let config = RunOptions { + let run_options = RunOptions { config: FabroConfig::default(), run_dir: dir.path().to_path_buf(), cancel_token: None, @@ -2248,7 +2248,7 @@ async fn codergen_without_backend_simulated() { host_repo_path: None, git: None, }; - engine.run(&graph, &config).await.expect("run"); + engine.run(&graph, &run_options).await.expect("run"); let response = std::fs::read_to_string(dir.path().join("nodes").join("code").join("response.md")).unwrap(); @@ -2339,7 +2339,7 @@ async fn branching_loop_back_on_failure() { }), ); let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::new()), local_env()); - let config = RunOptions { + let run_options = RunOptions { config: FabroConfig::default(), run_dir: dir.path().to_path_buf(), cancel_token: None, @@ -2353,7 +2353,7 @@ async fn branching_loop_back_on_failure() { host_repo_path: None, git: None, }; - let outcome = engine.run(&graph, &config).await.expect("run"); + let outcome = engine.run(&graph, &run_options).await.expect("run"); assert_eq!(outcome.status, StageStatus::Success); let cp = Checkpoint::load(&dir.path().join("checkpoint.json")).unwrap(); @@ -2422,7 +2422,7 @@ async fn human_gate_loops_back() { registry.register("exit", Box::new(ExitHandler)); registry.register("human", Box::new(HumanHandler::new(interviewer))); let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::new()), local_env()); - let config = RunOptions { + let run_options = RunOptions { config: FabroConfig::default(), run_dir: dir.path().to_path_buf(), cancel_token: None, @@ -2436,7 +2436,7 @@ async fn human_gate_loops_back() { host_repo_path: None, git: None, }; - let outcome = engine.run(&graph, &config).await.expect("run"); + let outcome = engine.run(&graph, &run_options).await.expect("run"); assert_eq!(outcome.status, StageStatus::Success); let cp = Checkpoint::load(&dir.path().join("checkpoint.json")).unwrap(); @@ -2480,7 +2480,7 @@ async fn scenario_ship_a_feature() { Arc::new(emitter), local_env(), ); - let config = RunOptions { + let run_options = RunOptions { config: FabroConfig::default(), run_dir: dir.path().to_path_buf(), cancel_token: None, @@ -2494,7 +2494,7 @@ async fn scenario_ship_a_feature() { host_repo_path: None, git: None, }; - let outcome = engine.run(&graph, &config).await.expect("run"); + let outcome = engine.run(&graph, &run_options).await.expect("run"); assert_eq!(outcome.status, StageStatus::Success); let cp = Checkpoint::load(&dir.path().join("checkpoint.json")).unwrap(); @@ -2566,7 +2566,7 @@ async fn scenario_parallel_expert_review() { registry.register("human", Box::new(HumanHandler::new(interviewer))); let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::new()), local_env()); - let config = RunOptions { + let run_options = RunOptions { config: FabroConfig::default(), run_dir: dir.path().to_path_buf(), cancel_token: None, @@ -2580,7 +2580,7 @@ async fn scenario_parallel_expert_review() { host_repo_path: None, git: None, }; - let outcome = engine.run(&graph, &config).await.expect("run"); + let outcome = engine.run(&graph, &run_options).await.expect("run"); assert_eq!(outcome.status, StageStatus::Success); let cp = Checkpoint::load(&dir.path().join("checkpoint.json")).unwrap(); @@ -2650,7 +2650,7 @@ async fn scenario_node_retries_on_retry_status() { }), ); let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::new()), local_env()); - let config = RunOptions { + let run_options = RunOptions { config: FabroConfig::default(), run_dir: dir.path().to_path_buf(), cancel_token: None, @@ -2664,7 +2664,7 @@ async fn scenario_node_retries_on_retry_status() { host_repo_path: None, git: None, }; - let outcome = engine.run(&graph, &config).await.expect("run"); + let outcome = engine.run(&graph, &run_options).await.expect("run"); assert_eq!(outcome.status, StageStatus::Success); let cp = Checkpoint::load(&dir.path().join("checkpoint.json")).unwrap(); @@ -2712,7 +2712,7 @@ async fn scenario_loop_restart_resets_context() { }), ); let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::new()), local_env()); - let config = RunOptions { + let run_options = RunOptions { config: FabroConfig::default(), run_dir: dir.path().to_path_buf(), cancel_token: None, @@ -2726,7 +2726,7 @@ async fn scenario_loop_restart_resets_context() { host_repo_path: None, git: None, }; - let outcome = engine.run(&graph, &config).await.expect("run"); + let outcome = engine.run(&graph, &run_options).await.expect("run"); assert_eq!(outcome.status, StageStatus::Success); assert!(call_count.load(std::sync::atomic::Ordering::SeqCst) >= 2); } @@ -2780,7 +2780,7 @@ async fn scenario_bug_triage_router() { registry.register("exit", Box::new(ExitHandler)); registry.register("conditional", Box::new(ConditionalHandler)); let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::new()), local_env()); - let config = RunOptions { + let run_options = RunOptions { config: FabroConfig::default(), run_dir: dir.path().to_path_buf(), cancel_token: None, @@ -2794,7 +2794,7 @@ async fn scenario_bug_triage_router() { host_repo_path: None, git: None, }; - let outcome = engine.run(&graph, &config).await.expect("run"); + let outcome = engine.run(&graph, &run_options).await.expect("run"); assert_eq!(outcome.status, StageStatus::Success); let cp = Checkpoint::load(&dir.path().join("checkpoint.json")).unwrap(); @@ -2839,7 +2839,7 @@ async fn scenario_crash_recovery() { registry.register("start", Box::new(StartHandler)); registry.register("exit", Box::new(ExitHandler)); let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::new()), local_env()); - let config = RunOptions { + let run_options = RunOptions { config: FabroConfig::default(), run_dir: dir.path().to_path_buf(), cancel_token: None, @@ -2854,7 +2854,7 @@ async fn scenario_crash_recovery() { git: None, }; let outcome = engine - .run_from_checkpoint(&graph, &config, &checkpoint) + .run_from_checkpoint(&graph, &run_options, &checkpoint) .await .expect("run"); assert_eq!(outcome.status, StageStatus::Success); @@ -2948,7 +2948,7 @@ async fn manager_loop_stop_condition_satisfied_e2e() { registry.register("done_setter", Box::new(DoneSetterHandler)); registry.register("stack.manager_loop", Box::new(SubWorkflowHandler)); let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::new()), local_env()); - let config = RunOptions { + let run_options = RunOptions { config: FabroConfig::default(), run_dir: dir.path().to_path_buf(), cancel_token: None, @@ -2962,7 +2962,7 @@ async fn manager_loop_stop_condition_satisfied_e2e() { host_repo_path: None, git: None, }; - let outcome = engine.run(&graph, &config).await.expect("run"); + let outcome = engine.run(&graph, &run_options).await.expect("run"); let cp = Checkpoint::load(&dir.path().join("checkpoint.json")).unwrap(); let manager_outcome = cp.node_outcomes.get("manager").expect("manager outcome"); @@ -3025,7 +3025,7 @@ async fn manager_loop_max_cycles_exceeded_e2e() { registry.register("exit", Box::new(ExitHandler)); registry.register("stack.manager_loop", Box::new(SubWorkflowHandler)); let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::new()), local_env()); - let config = RunOptions { + let run_options = RunOptions { config: FabroConfig::default(), run_dir: dir.path().to_path_buf(), cancel_token: None, @@ -3039,7 +3039,7 @@ async fn manager_loop_max_cycles_exceeded_e2e() { host_repo_path: None, git: None, }; - let outcome = engine.run(&graph, &config).await.expect("run"); + let outcome = engine.run(&graph, &run_options).await.expect("run"); let cp = Checkpoint::load(&dir.path().join("checkpoint.json")).unwrap(); let manager_outcome = cp.node_outcomes.get("manager").expect("manager outcome"); @@ -3161,7 +3161,7 @@ async fn conditional_branching_success_fail_paths() { registry.register("exit", Box::new(ExitHandler)); registry.register("always_fail", Box::new(AlwaysFailHandler)); let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::new()), local_env()); - let config = RunOptions { + let run_options = RunOptions { config: FabroConfig::default(), run_dir: dir.path().to_path_buf(), cancel_token: None, @@ -3175,7 +3175,7 @@ async fn conditional_branching_success_fail_paths() { host_repo_path: None, git: None, }; - let outcome = engine.run(&graph, &config).await.expect("run"); + let outcome = engine.run(&graph, &run_options).await.expect("run"); assert_eq!(outcome.status, StageStatus::Success); let cp = Checkpoint::load(&dir.path().join("checkpoint.json")).unwrap(); @@ -3214,7 +3214,7 @@ async fn edge_selection_condition_match_wins_over_weight() { registry.register("start", Box::new(StartHandler)); registry.register("exit", Box::new(ExitHandler)); let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::new()), local_env()); - let config = RunOptions { + let run_options = RunOptions { config: FabroConfig::default(), run_dir: dir.path().to_path_buf(), cancel_token: None, @@ -3228,7 +3228,7 @@ async fn edge_selection_condition_match_wins_over_weight() { host_repo_path: None, git: None, }; - engine.run(&graph, &config).await.expect("run"); + engine.run(&graph, &run_options).await.expect("run"); let cp = Checkpoint::load(&dir.path().join("checkpoint.json")).unwrap(); assert!(cp.completed_nodes.contains(&"cond_target".to_string())); @@ -3261,7 +3261,7 @@ async fn edge_selection_weight_breaks_ties() { registry.register("start", Box::new(StartHandler)); registry.register("exit", Box::new(ExitHandler)); let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::new()), local_env()); - let config = RunOptions { + let run_options = RunOptions { config: FabroConfig::default(), run_dir: dir.path().to_path_buf(), cancel_token: None, @@ -3275,7 +3275,7 @@ async fn edge_selection_weight_breaks_ties() { host_repo_path: None, git: None, }; - engine.run(&graph, &config).await.expect("run"); + engine.run(&graph, &run_options).await.expect("run"); let cp = Checkpoint::load(&dir.path().join("checkpoint.json")).unwrap(); assert!(cp.completed_nodes.contains(&"high".to_string())); @@ -3300,7 +3300,7 @@ async fn edge_selection_lexical_tiebreak() { registry.register("start", Box::new(StartHandler)); registry.register("exit", Box::new(ExitHandler)); let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::new()), local_env()); - let config = RunOptions { + let run_options = RunOptions { config: FabroConfig::default(), run_dir: dir.path().to_path_buf(), cancel_token: None, @@ -3314,7 +3314,7 @@ async fn edge_selection_lexical_tiebreak() { host_repo_path: None, git: None, }; - engine.run(&graph, &config).await.expect("run"); + engine.run(&graph, &run_options).await.expect("run"); let cp = Checkpoint::load(&dir.path().join("checkpoint.json")).unwrap(); assert!(cp.completed_nodes.contains(&"alpha".to_string())); @@ -3358,7 +3358,7 @@ async fn context_updates_visible_across_nodes() { registry.register("conditional", Box::new(ConditionalHandler)); registry.register("context_setter", Box::new(ContextSetterHandler)); let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::new()), local_env()); - let config = RunOptions { + let run_options = RunOptions { config: FabroConfig::default(), run_dir: dir.path().to_path_buf(), cancel_token: None, @@ -3372,7 +3372,7 @@ async fn context_updates_visible_across_nodes() { host_repo_path: None, git: None, }; - engine.run(&graph, &config).await.expect("run"); + engine.run(&graph, &run_options).await.expect("run"); let cp = Checkpoint::load(&dir.path().join("checkpoint.json")).unwrap(); assert!(cp.completed_nodes.contains(&"yes".to_string())); @@ -3402,7 +3402,7 @@ async fn stylesheet_applies_model_override() { Arc::new(EventEmitter::new()), local_env(), ); - let config = RunOptions { + let run_options = RunOptions { config: FabroConfig::default(), run_dir: dir.path().to_path_buf(), cancel_token: None, @@ -3416,7 +3416,7 @@ async fn stylesheet_applies_model_override() { host_repo_path: None, git: None, }; - let outcome = engine.run(&graph, &config).await.expect("run"); + let outcome = engine.run(&graph, &run_options).await.expect("run"); assert_eq!(outcome.status, StageStatus::Success); } @@ -3458,7 +3458,7 @@ async fn custom_handler_registration_and_execution() { registry.register("exit", Box::new(ExitHandler)); registry.register("my_custom", Box::new(CustomHandler)); let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::new()), local_env()); - let config = RunOptions { + let run_options = RunOptions { config: FabroConfig::default(), run_dir: dir.path().to_path_buf(), cancel_token: None, @@ -3472,7 +3472,7 @@ async fn custom_handler_registration_and_execution() { host_repo_path: None, git: None, }; - engine.run(&graph, &config).await.expect("run"); + engine.run(&graph, &run_options).await.expect("run"); let cp = Checkpoint::load(&dir.path().join("checkpoint.json")).unwrap(); assert_eq!( @@ -3529,7 +3529,7 @@ async fn integration_smoke_plan_implement_review_done() { Arc::new(emitter), local_env(), ); - let config = RunOptions { + let run_options = RunOptions { config: FabroConfig::default(), run_dir: dir.path().to_path_buf(), cancel_token: None, @@ -3543,7 +3543,7 @@ async fn integration_smoke_plan_implement_review_done() { host_repo_path: None, git: None, }; - let outcome = engine.run(&graph, &config).await.expect("run"); + let outcome = engine.run(&graph, &run_options).await.expect("run"); assert_eq!(outcome.status, StageStatus::Success); // Verify all nodes completed @@ -3633,7 +3633,7 @@ async fn manager_loop_runs_child_engine_e2e() { registry.register("stack.manager_loop", Box::new(SubWorkflowHandler)); let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::new()), local_env()); - let config = RunOptions { + let run_options = RunOptions { config: FabroConfig::default(), run_dir: dir.path().to_path_buf(), cancel_token: None, @@ -3648,7 +3648,7 @@ async fn manager_loop_runs_child_engine_e2e() { git: None, }; let outcome = engine - .run(&graph, &config) + .run(&graph, &run_options) .await .expect("manager loop E2E should succeed"); assert_eq!(outcome.status, StageStatus::Success); @@ -3767,7 +3767,7 @@ async fn manager_loop_context_flows_e2e() { registry.register("stack.manager_loop", Box::new(SubWorkflowHandler)); let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::new()), local_env()); - let config = RunOptions { + let run_options = RunOptions { config: FabroConfig::default(), run_dir: dir.path().to_path_buf(), cancel_token: None, @@ -3781,7 +3781,7 @@ async fn manager_loop_context_flows_e2e() { host_repo_path: None, git: None, }; - let outcome = engine.run(&graph, &config).await.expect("run"); + let outcome = engine.run(&graph, &run_options).await.expect("run"); assert_eq!(outcome.status, StageStatus::Success); // Check that child's context updates were propagated through the manager @@ -3840,7 +3840,7 @@ async fn manager_loop_child_dotfile_e2e() { registry.register("stack.manager_loop", Box::new(SubWorkflowHandler)); let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::new()), local_env()); - let config = RunOptions { + let run_options = RunOptions { config: FabroConfig::default(), run_dir: dir.path().to_path_buf(), cancel_token: None, @@ -3854,7 +3854,7 @@ async fn manager_loop_child_dotfile_e2e() { host_repo_path: None, git: None, }; - let outcome = engine.run(&graph, &config).await.expect("run"); + let outcome = engine.run(&graph, &run_options).await.expect("run"); assert_eq!(outcome.status, StageStatus::Success); } @@ -3953,7 +3953,7 @@ async fn graph_merge_e2e_through_engine() { Arc::new(EventEmitter::new()), local_env(), ); - let config = RunOptions { + let run_options = RunOptions { config: FabroConfig::default(), run_dir: dir.path().to_path_buf(), cancel_token: None, @@ -3968,7 +3968,7 @@ async fn graph_merge_e2e_through_engine() { git: None, }; let outcome = engine - .run(&main_graph, &config) + .run(&main_graph, &run_options) .await .expect("graph merge E2E should succeed"); assert_eq!(outcome.status, StageStatus::Success); @@ -4103,7 +4103,7 @@ async fn fidelity_default_is_compact() { ); let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::new()), local_env()); - let config = RunOptions { + let run_options = RunOptions { config: FabroConfig::default(), run_dir: dir.path().to_path_buf(), cancel_token: None, @@ -4117,7 +4117,7 @@ async fn fidelity_default_is_compact() { host_repo_path: None, git: None, }; - engine.run(&graph, &config).await.expect("run"); + engine.run(&graph, &run_options).await.expect("run"); let fidelities = captures.fidelities.lock().unwrap(); assert_eq!(fidelities.len(), 1); @@ -4160,7 +4160,7 @@ async fn fidelity_graph_default_applied() { ); let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::new()), local_env()); - let config = RunOptions { + let run_options = RunOptions { config: FabroConfig::default(), run_dir: dir.path().to_path_buf(), cancel_token: None, @@ -4174,7 +4174,7 @@ async fn fidelity_graph_default_applied() { host_repo_path: None, git: None, }; - engine.run(&graph, &config).await.expect("run"); + engine.run(&graph, &run_options).await.expect("run"); let fidelities = captures.fidelities.lock().unwrap(); assert_eq!(fidelities[0].1, "truncate"); @@ -4213,7 +4213,7 @@ async fn fidelity_node_overrides_graph_default() { ); let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::new()), local_env()); - let config = RunOptions { + let run_options = RunOptions { config: FabroConfig::default(), run_dir: dir.path().to_path_buf(), cancel_token: None, @@ -4227,7 +4227,7 @@ async fn fidelity_node_overrides_graph_default() { host_repo_path: None, git: None, }; - engine.run(&graph, &config).await.expect("run"); + engine.run(&graph, &run_options).await.expect("run"); let fidelities = captures.fidelities.lock().unwrap(); assert_eq!(fidelities[0].1, "summary:medium"); @@ -4272,7 +4272,7 @@ async fn fidelity_edge_overrides_node_and_graph() { ); let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::new()), local_env()); - let config = RunOptions { + let run_options = RunOptions { config: FabroConfig::default(), run_dir: dir.path().to_path_buf(), cancel_token: None, @@ -4286,7 +4286,7 @@ async fn fidelity_edge_overrides_node_and_graph() { host_repo_path: None, git: None, }; - engine.run(&graph, &config).await.expect("run"); + engine.run(&graph, &run_options).await.expect("run"); let fidelities = captures.fidelities.lock().unwrap(); assert_eq!(fidelities[0].1, "summary:high"); @@ -4321,7 +4321,7 @@ async fn fidelity_full_produces_empty_preamble() { ); let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::new()), local_env()); - let config = RunOptions { + let run_options = RunOptions { config: FabroConfig::default(), run_dir: dir.path().to_path_buf(), cancel_token: None, @@ -4335,7 +4335,7 @@ async fn fidelity_full_produces_empty_preamble() { host_repo_path: None, git: None, }; - engine.run(&graph, &config).await.expect("run"); + engine.run(&graph, &run_options).await.expect("run"); let fidelities = captures.fidelities.lock().unwrap(); assert_eq!(fidelities[0].1, "full"); @@ -4380,7 +4380,7 @@ async fn fidelity_truncate_preamble_minimal() { ); let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::new()), local_env()); - let config = RunOptions { + let run_options = RunOptions { config: FabroConfig::default(), run_dir: dir.path().to_path_buf(), cancel_token: None, @@ -4394,7 +4394,7 @@ async fn fidelity_truncate_preamble_minimal() { host_repo_path: None, git: None, }; - engine.run(&graph, &config).await.expect("run"); + engine.run(&graph, &run_options).await.expect("run"); let preambles = captures.preambles.lock().unwrap(); let preamble = &preambles[0].1; @@ -4452,7 +4452,7 @@ async fn fidelity_summary_low_mode() { ); let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::new()), local_env()); - let config = RunOptions { + let run_options = RunOptions { config: FabroConfig::default(), run_dir: dir.path().to_path_buf(), cancel_token: None, @@ -4466,7 +4466,7 @@ async fn fidelity_summary_low_mode() { host_repo_path: None, git: None, }; - engine.run(&graph, &config).await.expect("run"); + engine.run(&graph, &run_options).await.expect("run"); let fidelities = captures.fidelities.lock().unwrap(); assert_eq!(fidelities[0].1, "summary:low"); @@ -4519,7 +4519,7 @@ async fn fidelity_summary_medium_mode() { ); let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::new()), local_env()); - let config = RunOptions { + let run_options = RunOptions { config: FabroConfig::default(), run_dir: dir.path().to_path_buf(), cancel_token: None, @@ -4533,7 +4533,7 @@ async fn fidelity_summary_medium_mode() { host_repo_path: None, git: None, }; - engine.run(&graph, &config).await.expect("run"); + engine.run(&graph, &run_options).await.expect("run"); let fidelities = captures.fidelities.lock().unwrap(); assert_eq!(fidelities[0].1, "summary:medium"); @@ -4586,7 +4586,7 @@ async fn fidelity_summary_high_mode() { ); let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::new()), local_env()); - let config = RunOptions { + let run_options = RunOptions { config: FabroConfig::default(), run_dir: dir.path().to_path_buf(), cancel_token: None, @@ -4600,7 +4600,7 @@ async fn fidelity_summary_high_mode() { host_repo_path: None, git: None, }; - engine.run(&graph, &config).await.expect("run"); + engine.run(&graph, &run_options).await.expect("run"); let fidelities = captures.fidelities.lock().unwrap(); assert_eq!(fidelities[0].1, "summary:high"); @@ -4646,7 +4646,7 @@ async fn fidelity_full_sets_thread_id_in_context() { ); let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::new()), local_env()); - let config = RunOptions { + let run_options = RunOptions { config: FabroConfig::default(), run_dir: dir.path().to_path_buf(), cancel_token: None, @@ -4660,7 +4660,7 @@ async fn fidelity_full_sets_thread_id_in_context() { host_repo_path: None, git: None, }; - engine.run(&graph, &config).await.expect("run"); + engine.run(&graph, &run_options).await.expect("run"); let thread_ids = captures.thread_ids.lock().unwrap(); assert_eq!(thread_ids[0].0, "work"); @@ -4717,7 +4717,7 @@ async fn fidelity_full_nodes_share_thread_id() { ); let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::new()), local_env()); - let config = RunOptions { + let run_options = RunOptions { config: FabroConfig::default(), run_dir: dir.path().to_path_buf(), cancel_token: None, @@ -4731,7 +4731,7 @@ async fn fidelity_full_nodes_share_thread_id() { host_repo_path: None, git: None, }; - engine.run(&graph, &config).await.expect("run"); + engine.run(&graph, &run_options).await.expect("run"); let thread_ids = captures.thread_ids.lock().unwrap(); assert_eq!(thread_ids[0].0, "step_a"); @@ -4798,7 +4798,7 @@ async fn fidelity_resume_degrades_full_to_summary_high() { ); let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::new()), local_env()); - let config = RunOptions { + let run_options = RunOptions { config: FabroConfig::default(), run_dir: dir.path().to_path_buf(), cancel_token: None, @@ -4813,7 +4813,7 @@ async fn fidelity_resume_degrades_full_to_summary_high() { git: None, }; engine - .run_from_checkpoint(&graph, &config, &checkpoint) + .run_from_checkpoint(&graph, &run_options, &checkpoint) .await .expect("resume should succeed"); @@ -4895,7 +4895,7 @@ async fn fidelity_resume_degrade_only_affects_first_hop() { ); let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::new()), local_env()); - let config = RunOptions { + let run_options = RunOptions { config: FabroConfig::default(), run_dir: dir.path().to_path_buf(), cancel_token: None, @@ -4910,7 +4910,7 @@ async fn fidelity_resume_degrade_only_affects_first_hop() { git: None, }; engine - .run_from_checkpoint(&graph, &config, &checkpoint) + .run_from_checkpoint(&graph, &run_options, &checkpoint) .await .expect("resume should succeed"); @@ -4979,7 +4979,7 @@ async fn fidelity_resume_no_degrade_when_not_full() { ); let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::new()), local_env()); - let config = RunOptions { + let run_options = RunOptions { config: FabroConfig::default(), run_dir: dir.path().to_path_buf(), cancel_token: None, @@ -4994,7 +4994,7 @@ async fn fidelity_resume_no_degrade_when_not_full() { git: None, }; engine - .run_from_checkpoint(&graph, &config, &checkpoint) + .run_from_checkpoint(&graph, &run_options, &checkpoint) .await .expect("resume should succeed"); @@ -5021,7 +5021,7 @@ async fn fidelity_stored_in_checkpoint_context() { registry.register("exit", Box::new(ExitHandler)); let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::new()), local_env()); - let config = RunOptions { + let run_options = RunOptions { config: FabroConfig::default(), run_dir: dir.path().to_path_buf(), cancel_token: None, @@ -5035,7 +5035,7 @@ async fn fidelity_stored_in_checkpoint_context() { host_repo_path: None, git: None, }; - engine.run(&graph, &config).await.expect("run"); + engine.run(&graph, &run_options).await.expect("run"); let cp = Checkpoint::load(&dir.path().join("checkpoint.json")).unwrap(); assert_eq!( @@ -5107,7 +5107,7 @@ async fn fidelity_precedence_multi_node_pipeline() { ); let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::new()), local_env()); - let config = RunOptions { + let run_options = RunOptions { config: FabroConfig::default(), run_dir: dir.path().to_path_buf(), cancel_token: None, @@ -5121,7 +5121,7 @@ async fn fidelity_precedence_multi_node_pipeline() { host_repo_path: None, git: None, }; - engine.run(&graph, &config).await.expect("run"); + engine.run(&graph, &run_options).await.expect("run"); let fidelities = captures.fidelities.lock().unwrap(); assert_eq!(fidelities[0].0, "step_a"); @@ -5175,7 +5175,7 @@ async fn fidelity_compact_preamble_includes_completed_stages_and_context() { ); let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::new()), local_env()); - let config = RunOptions { + let run_options = RunOptions { config: FabroConfig::default(), run_dir: dir.path().to_path_buf(), cancel_token: None, @@ -5189,7 +5189,7 @@ async fn fidelity_compact_preamble_includes_completed_stages_and_context() { host_repo_path: None, git: None, }; - engine.run(&graph, &config).await.expect("run"); + engine.run(&graph, &run_options).await.expect("run"); let preambles = captures.preambles.lock().unwrap(); // step_b's preamble should contain structured summary of completed work @@ -5250,7 +5250,7 @@ async fn fidelity_summary_low_excludes_context_values_in_pipeline() { }), ); let engine_low = WorkflowRunner::new(registry_low, Arc::new(EventEmitter::new()), local_env()); - let config_low = RunOptions { + let run_options_low = RunOptions { config: FabroConfig::default(), run_dir: dir_low.path().to_path_buf(), cancel_token: None, @@ -5265,7 +5265,7 @@ async fn fidelity_summary_low_excludes_context_values_in_pipeline() { git: None, }; engine_low - .run(&graph_low, &config_low) + .run(&graph_low, &run_options_low) .await .expect("run low"); @@ -5317,7 +5317,7 @@ async fn fidelity_summary_low_excludes_context_values_in_pipeline() { }), ); let engine_med = WorkflowRunner::new(registry_med, Arc::new(EventEmitter::new()), local_env()); - let config_med = RunOptions { + let run_options_med = RunOptions { config: FabroConfig::default(), run_dir: dir_med.path().to_path_buf(), cancel_token: None, @@ -5332,7 +5332,7 @@ async fn fidelity_summary_low_excludes_context_values_in_pipeline() { git: None, }; engine_med - .run(&graph_med, &config_med) + .run(&graph_med, &run_options_med) .await .expect("run med"); @@ -5388,7 +5388,7 @@ async fn fidelity_thread_id_fallback_to_previous_node_in_pipeline() { ); let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::new()), local_env()); - let config = RunOptions { + let run_options = RunOptions { config: FabroConfig::default(), run_dir: dir.path().to_path_buf(), cancel_token: None, @@ -5402,7 +5402,7 @@ async fn fidelity_thread_id_fallback_to_previous_node_in_pipeline() { host_repo_path: None, git: None, }; - engine.run(&graph, &config).await.expect("run"); + engine.run(&graph, &run_options).await.expect("run"); let thread_ids = captures.thread_ids.lock().unwrap(); // step_a should have previous node = start @@ -5442,7 +5442,7 @@ async fn fidelity_thread_id_from_node_class_in_pipeline() { ); let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::new()), local_env()); - let config = RunOptions { + let run_options = RunOptions { config: FabroConfig::default(), run_dir: dir.path().to_path_buf(), cancel_token: None, @@ -5456,7 +5456,7 @@ async fn fidelity_thread_id_from_node_class_in_pipeline() { host_repo_path: None, git: None, }; - engine.run(&graph, &config).await.expect("run"); + engine.run(&graph, &run_options).await.expect("run"); let thread_ids = captures.thread_ids.lock().unwrap(); assert_eq!(thread_ids[0].0, "work"); @@ -5499,7 +5499,7 @@ async fn fidelity_edge_thread_id_override_in_pipeline() { ); let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::new()), local_env()); - let config = RunOptions { + let run_options = RunOptions { config: FabroConfig::default(), run_dir: dir.path().to_path_buf(), cancel_token: None, @@ -5513,7 +5513,7 @@ async fn fidelity_edge_thread_id_override_in_pipeline() { host_repo_path: None, git: None, }; - engine.run(&graph, &config).await.expect("run"); + engine.run(&graph, &run_options).await.expect("run"); let thread_ids = captures.thread_ids.lock().unwrap(); assert_eq!(thread_ids[0].0, "work"); @@ -5557,7 +5557,7 @@ async fn fidelity_full_without_explicit_thread_id_uses_previous_node() { ); let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::new()), local_env()); - let config = RunOptions { + let run_options = RunOptions { config: FabroConfig::default(), run_dir: dir.path().to_path_buf(), cancel_token: None, @@ -5571,7 +5571,7 @@ async fn fidelity_full_without_explicit_thread_id_uses_previous_node() { host_repo_path: None, git: None, }; - engine.run(&graph, &config).await.expect("run"); + engine.run(&graph, &run_options).await.expect("run"); let fidelities = captures.fidelities.lock().unwrap(); assert_eq!(fidelities[0].1, "full"); @@ -5625,7 +5625,7 @@ async fn fidelity_from_parsed_dot_pipeline() { ); let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::new()), local_env()); - let config = RunOptions { + let run_options = RunOptions { config: FabroConfig::default(), run_dir: dir.path().to_path_buf(), cancel_token: None, @@ -5639,7 +5639,7 @@ async fn fidelity_from_parsed_dot_pipeline() { host_repo_path: None, git: None, }; - engine.run(&graph, &config).await.expect("run"); + engine.run(&graph, &run_options).await.expect("run"); let fidelities = captures.fidelities.lock().unwrap(); // step_a: no node fidelity, no edge fidelity -> graph default "truncate" @@ -5673,7 +5673,7 @@ async fn fidelity_checkpoint_roundtrip_preserves_fidelity() { registry.register("exit", Box::new(ExitHandler)); let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::new()), local_env()); - let config = RunOptions { + let run_options = RunOptions { config: FabroConfig::default(), run_dir: dir.path().to_path_buf(), cancel_token: None, @@ -5687,7 +5687,7 @@ async fn fidelity_checkpoint_roundtrip_preserves_fidelity() { host_repo_path: None, git: None, }; - engine.run(&graph, &config).await.expect("run"); + engine.run(&graph, &run_options).await.expect("run"); // Load, save, load again to verify roundtrip let checkpoint_path = dir.path().join("checkpoint.json"); @@ -5743,7 +5743,7 @@ async fn fidelity_node_thread_id_overrides_edge_thread_id_in_pipeline() { ); let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::new()), local_env()); - let config = RunOptions { + let run_options = RunOptions { config: FabroConfig::default(), run_dir: dir.path().to_path_buf(), cancel_token: None, @@ -5757,7 +5757,7 @@ async fn fidelity_node_thread_id_overrides_edge_thread_id_in_pipeline() { host_repo_path: None, git: None, }; - engine.run(&graph, &config).await.expect("run"); + engine.run(&graph, &run_options).await.expect("run"); let thread_ids = captures.thread_ids.lock().unwrap(); assert_eq!(thread_ids[0].0, "work"); @@ -5830,7 +5830,7 @@ async fn fidelity_resume_preserves_context_values_across_checkpoint() { ); let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::new()), local_env()); - let config = RunOptions { + let run_options = RunOptions { config: FabroConfig::default(), run_dir: dir.path().to_path_buf(), cancel_token: None, @@ -5845,7 +5845,7 @@ async fn fidelity_resume_preserves_context_values_across_checkpoint() { git: None, }; engine - .run_from_checkpoint(&graph, &config, &checkpoint) + .run_from_checkpoint(&graph, &run_options, &checkpoint) .await .expect("resume should succeed"); @@ -6045,7 +6045,7 @@ mod real_llm { ); let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::new()), local_env()); - let config = RunOptions { + let run_options = RunOptions { config: FabroConfig::default(), run_dir: dir.path().to_path_buf(), cancel_token: None, @@ -6061,7 +6061,7 @@ mod real_llm { }; let outcome = tokio::time::timeout( std::time::Duration::from_secs(120), - engine.run(&graph, &config), + engine.run(&graph, &run_options), ) .await .expect("should not timeout") @@ -6159,7 +6159,7 @@ mod real_llm { ); let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::new()), local_env()); - let config = RunOptions { + let run_options = RunOptions { config: FabroConfig::default(), run_dir: dir.path().to_path_buf(), cancel_token: None, @@ -6175,7 +6175,7 @@ mod real_llm { }; let outcome = tokio::time::timeout( std::time::Duration::from_secs(120), - engine.run(&graph, &config), + engine.run(&graph, &run_options), ) .await .expect("should not timeout") @@ -6298,7 +6298,7 @@ mod real_llm { registry.register("human", Box::new(HumanHandler::new(interviewer))); let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::new()), local_env()); - let config = RunOptions { + let run_options = RunOptions { config: FabroConfig::default(), run_dir: dir.path().to_path_buf(), cancel_token: None, @@ -6314,7 +6314,7 @@ mod real_llm { }; let outcome = tokio::time::timeout( std::time::Duration::from_secs(120), - engine.run(&graph, &config), + engine.run(&graph, &run_options), ) .await .expect("should not timeout") @@ -6405,7 +6405,7 @@ mod real_llm { ); let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::new()), local_env()); - let config = RunOptions { + let run_options = RunOptions { config: FabroConfig::default(), run_dir: dir.path().to_path_buf(), cancel_token: None, @@ -6421,7 +6421,7 @@ mod real_llm { }; let outcome = tokio::time::timeout( std::time::Duration::from_secs(30), - engine.run(&graph, &config), + engine.run(&graph, &run_options), ) .await .expect("should not timeout") @@ -6501,7 +6501,7 @@ async fn human_gate_freeform_only_routes_text() { registry.register("human", Box::new(HumanHandler::new(interviewer))); let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::new()), local_env()); - let config = RunOptions { + let run_options = RunOptions { config: FabroConfig::default(), run_dir: dir.path().to_path_buf(), cancel_token: None, @@ -6516,7 +6516,7 @@ async fn human_gate_freeform_only_routes_text() { git: None, }; let outcome = engine - .run(&graph, &config) + .run(&graph, &run_options) .await .expect("run should succeed"); assert_eq!(outcome.status, StageStatus::Success); @@ -6631,7 +6631,7 @@ async fn human_gate_freeform_with_fixed_choice_match() { registry.register("human", Box::new(HumanHandler::new(interviewer))); let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::new()), local_env()); - let config = RunOptions { + let run_options = RunOptions { config: FabroConfig::default(), run_dir: dir.path().to_path_buf(), cancel_token: None, @@ -6646,7 +6646,7 @@ async fn human_gate_freeform_with_fixed_choice_match() { git: None, }; let outcome = engine - .run(&graph, &config) + .run(&graph, &run_options) .await .expect("run should succeed"); assert_eq!(outcome.status, StageStatus::Success); @@ -6746,7 +6746,7 @@ async fn human_gate_freeform_fallback_on_unmatched_text() { registry.register("human", Box::new(HumanHandler::new(interviewer))); let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::new()), local_env()); - let config = RunOptions { + let run_options = RunOptions { config: FabroConfig::default(), run_dir: dir.path().to_path_buf(), cancel_token: None, @@ -6761,7 +6761,7 @@ async fn human_gate_freeform_fallback_on_unmatched_text() { git: None, }; let outcome = engine - .run(&graph, &config) + .run(&graph, &run_options) .await .expect("run should succeed"); assert_eq!(outcome.status, StageStatus::Success); @@ -6874,7 +6874,7 @@ async fn human_gate_freeform_sets_allow_freeform_on_question() { registry.register("human", Box::new(HumanHandler::new(interviewer))); let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::new()), local_env()); - let config = RunOptions { + let run_options = RunOptions { config: FabroConfig::default(), run_dir: dir.path().to_path_buf(), cancel_token: None, @@ -6889,7 +6889,7 @@ async fn human_gate_freeform_sets_allow_freeform_on_question() { git: None, }; let outcome = engine - .run(&graph, &config) + .run(&graph, &run_options) .await .expect("run should succeed"); assert_eq!(outcome.status, StageStatus::Success); @@ -6982,7 +6982,7 @@ async fn human_gate_without_freeform_sets_allow_freeform_false() { registry.register("human", Box::new(HumanHandler::new(interviewer))); let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::new()), local_env()); - let config = RunOptions { + let run_options = RunOptions { config: FabroConfig::default(), run_dir: dir.path().to_path_buf(), cancel_token: None, @@ -6997,7 +6997,7 @@ async fn human_gate_without_freeform_sets_allow_freeform_false() { git: None, }; let outcome = engine - .run(&graph, &config) + .run(&graph, &run_options) .await .expect("run should succeed"); assert_eq!(outcome.status, StageStatus::Success); @@ -7219,13 +7219,13 @@ struct HookTestRunner { } impl HookTestRunner { - async fn run(&self, graph: &Graph, config: &RunOptions) -> Result { + async fn run(&self, graph: &Graph, run_options: &RunOptions) -> Result { run_graph_with_hooks( make_linear_registry(), Arc::clone(&self.emitter), local_env(), graph, - config, + run_options, Arc::clone(&self.hook_runner), None, ) @@ -7262,7 +7262,7 @@ fn engine_with_hooks_and_events( ) } -fn make_run_config(dir: &std::path::Path) -> RunOptions { +fn make_run_options(dir: &std::path::Path) -> RunOptions { RunOptions { config: FabroConfig::default(), run_dir: dir.to_path_buf(), @@ -7337,9 +7337,9 @@ async fn hook_run_start_proceed_allows_run() { let engine = engine_with_hooks(hooks); let graph = parse(simple_linear_dot()).unwrap(); let dir = tempfile::tempdir().unwrap(); - let config = make_run_config(dir.path()); + let run_options = make_run_options(dir.path()); - let outcome = engine.run(&graph, &config).await.unwrap(); + let outcome = engine.run(&graph, &run_options).await.unwrap(); assert_eq!(outcome.status, StageStatus::Success); } @@ -7349,9 +7349,9 @@ async fn hook_run_start_block_prevents_run() { let (engine, events) = engine_with_hooks_and_events(hooks); let graph = parse(simple_linear_dot()).unwrap(); let dir = tempfile::tempdir().unwrap(); - let config = make_run_config(dir.path()); + let run_options = make_run_options(dir.path()); - let result = engine.run(&graph, &config).await; + let result = engine.run(&graph, &run_options).await; assert!(result.is_err(), "RunStart block should cause error"); let err = result.unwrap_err(); assert!( @@ -7387,9 +7387,9 @@ async fn hook_run_start_block_with_json_reason() { let engine = engine_with_hooks(hooks); let graph = parse(simple_linear_dot()).unwrap(); let dir = tempfile::tempdir().unwrap(); - let config = make_run_config(dir.path()); + let run_options = make_run_options(dir.path()); - let result = engine.run(&graph, &config).await; + let result = engine.run(&graph, &run_options).await; assert!(result.is_err()); let err = result.unwrap_err(); assert!( @@ -7406,9 +7406,9 @@ async fn hook_stage_start_proceed_allows_execution() { let engine = engine_with_hooks(hooks); let graph = parse(simple_linear_dot()).unwrap(); let dir = tempfile::tempdir().unwrap(); - let config = make_run_config(dir.path()); + let run_options = make_run_options(dir.path()); - let outcome = engine.run(&graph, &config).await.unwrap(); + let outcome = engine.run(&graph, &run_options).await.unwrap(); assert_eq!(outcome.status, StageStatus::Success); // Work node should have executed (response.md exists) @@ -7432,9 +7432,9 @@ async fn hook_stage_start_skip_bypasses_node() { let (engine, events) = engine_with_hooks_and_events(hooks); let graph = parse(simple_linear_dot()).unwrap(); let dir = tempfile::tempdir().unwrap(); - let config = make_run_config(dir.path()); + let run_options = make_run_options(dir.path()); - let outcome = engine.run(&graph, &config).await.unwrap(); + let outcome = engine.run(&graph, &run_options).await.unwrap(); // Pipeline reached exit with goal gates satisfied — per spec, SUCCESS. assert_eq!(outcome.status, StageStatus::Success); @@ -7469,9 +7469,9 @@ async fn hook_stage_start_block_aborts_run() { let engine = engine_with_hooks(hooks); let graph = parse(simple_linear_dot()).unwrap(); let dir = tempfile::tempdir().unwrap(); - let config = make_run_config(dir.path()); + let run_options = make_run_options(dir.path()); - let result = engine.run(&graph, &config).await; + let result = engine.run(&graph, &run_options).await; assert!(result.is_err(), "StageStart block should abort the run"); } @@ -7488,9 +7488,9 @@ async fn hook_stage_start_matcher_filters_by_node_id() { let engine = engine_with_hooks(hooks); let graph = parse(two_step_dot()).unwrap(); let dir = tempfile::tempdir().unwrap(); - let config = make_run_config(dir.path()); + let run_options = make_run_options(dir.path()); - let outcome = engine.run(&graph, &config).await.unwrap(); + let outcome = engine.run(&graph, &run_options).await.unwrap(); // Pipeline reached exit with goal gates satisfied — per spec, SUCCESS. assert_eq!(outcome.status, StageStatus::Success); @@ -7525,9 +7525,9 @@ async fn hook_stage_start_matcher_no_match_proceeds() { let engine = engine_with_hooks(hooks); let graph = parse(simple_linear_dot()).unwrap(); let dir = tempfile::tempdir().unwrap(); - let config = make_run_config(dir.path()); + let run_options = make_run_options(dir.path()); - let outcome = engine.run(&graph, &config).await.unwrap(); + let outcome = engine.run(&graph, &run_options).await.unwrap(); assert_eq!(outcome.status, StageStatus::Success); } @@ -7544,9 +7544,9 @@ async fn hook_stage_complete_fires_after_success() { )]; let engine = engine_with_hooks(hooks); let graph = parse(two_step_dot()).unwrap(); - let config = make_run_config(dir.path()); + let run_options = make_run_options(dir.path()); - let outcome = engine.run(&graph, &config).await.unwrap(); + let outcome = engine.run(&graph, &run_options).await.unwrap(); assert_eq!(outcome.status, StageStatus::Success); // Marker file should exist and contain node IDs @@ -7573,9 +7573,9 @@ async fn hook_stage_complete_failure_does_not_block_pipeline() { let engine = engine_with_hooks(hooks); let graph = parse(simple_linear_dot()).unwrap(); let dir = tempfile::tempdir().unwrap(); - let config = make_run_config(dir.path()); + let run_options = make_run_options(dir.path()); - let outcome = engine.run(&graph, &config).await.unwrap(); + let outcome = engine.run(&graph, &run_options).await.unwrap(); assert_eq!( outcome.status, StageStatus::Success, @@ -7596,9 +7596,9 @@ async fn hook_run_complete_fires_on_success() { )]; let engine = engine_with_hooks(hooks); let graph = parse(simple_linear_dot()).unwrap(); - let config = make_run_config(dir.path()); + let run_options = make_run_options(dir.path()); - let outcome = engine.run(&graph, &config).await.unwrap(); + let outcome = engine.run(&graph, &run_options).await.unwrap(); assert_eq!(outcome.status, StageStatus::Success); assert!( @@ -7626,9 +7626,9 @@ async fn hook_run_complete_does_not_fire_on_blocked_run() { ]; let engine = engine_with_hooks(hooks); let graph = parse(simple_linear_dot()).unwrap(); - let config = make_run_config(dir.path()); + let run_options = make_run_options(dir.path()); - let _ = engine.run(&graph, &config).await; + let _ = engine.run(&graph, &run_options).await; assert!( !marker.exists(), @@ -7655,9 +7655,9 @@ async fn hook_run_failed_fires_on_stage_block() { ]; let engine = engine_with_hooks(hooks); let graph = parse(simple_linear_dot()).unwrap(); - let config = make_run_config(dir.path()); + let run_options = make_run_options(dir.path()); - let _ = engine.run(&graph, &config).await; + let _ = engine.run(&graph, &run_options).await; // RunFailed may or may not fire depending on the error path — a StageStart // block causes an engine error, which doesn't go through the normal @@ -7680,9 +7680,9 @@ async fn hook_receives_env_vars() { )]; let engine = engine_with_hooks(hooks); let graph = parse(simple_linear_dot()).unwrap(); - let config = make_run_config(dir.path()); + let run_options = make_run_options(dir.path()); - engine.run(&graph, &config).await.unwrap(); + engine.run(&graph, &run_options).await.unwrap(); assert!(env_file.exists(), "Env file should be written by hook"); let content = std::fs::read_to_string(&env_file).unwrap(); @@ -7729,9 +7729,9 @@ async fn multiple_hooks_same_event_all_fire() { ]; let engine = engine_with_hooks(hooks); let graph = parse(simple_linear_dot()).unwrap(); - let config = make_run_config(dir.path()); + let run_options = make_run_options(dir.path()); - engine.run(&graph, &config).await.unwrap(); + engine.run(&graph, &run_options).await.unwrap(); assert!(marker1.exists(), "First hook should have fired"); assert!(marker2.exists(), "Second hook should have fired"); @@ -7744,9 +7744,9 @@ async fn no_hooks_configured_runs_normally() { let engine = engine_with_hooks(vec![]); let graph = parse(simple_linear_dot()).unwrap(); let dir = tempfile::tempdir().unwrap(); - let config = make_run_config(dir.path()); + let run_options = make_run_options(dir.path()); - let outcome = engine.run(&graph, &config).await.unwrap(); + let outcome = engine.run(&graph, &run_options).await.unwrap(); assert_eq!(outcome.status, StageStatus::Success); } @@ -7767,9 +7767,9 @@ async fn hook_edge_selected_override_redirects_routing() { let (engine, events) = engine_with_hooks_and_events(hooks); let graph = parse(branching_dot()).unwrap(); let dir = tempfile::tempdir().unwrap(); - let config = make_run_config(dir.path()); + let run_options = make_run_options(dir.path()); - let outcome = engine.run(&graph, &config).await.unwrap(); + let outcome = engine.run(&graph, &run_options).await.unwrap(); assert_eq!(outcome.status, StageStatus::Success); // Verify pathB was executed (override worked) @@ -7796,9 +7796,9 @@ async fn hook_edge_selected_block_aborts_run() { let engine = engine_with_hooks(hooks); let graph = parse(branching_dot()).unwrap(); let dir = tempfile::tempdir().unwrap(); - let config = make_run_config(dir.path()); + let run_options = make_run_options(dir.path()); - let result = engine.run(&graph, &config).await; + let result = engine.run(&graph, &run_options).await; assert!(result.is_err(), "EdgeSelected block should abort the run"); } @@ -7815,9 +7815,9 @@ async fn hook_checkpoint_saved_fires() { )]; let engine = engine_with_hooks(hooks); let graph = parse(simple_linear_dot()).unwrap(); - let config = make_run_config(dir.path()); + let run_options = make_run_options(dir.path()); - let outcome = engine.run(&graph, &config).await.unwrap(); + let outcome = engine.run(&graph, &run_options).await.unwrap(); assert_eq!(outcome.status, StageStatus::Success); // Checkpoint is saved after each node @@ -7837,10 +7837,10 @@ async fn hook_stage_start_exit_2_blocks() { let engine = engine_with_hooks(hooks); let graph = parse(simple_linear_dot()).unwrap(); let dir = tempfile::tempdir().unwrap(); - let config = make_run_config(dir.path()); + let run_options = make_run_options(dir.path()); // exit 2 without JSON defaults to Block - let result = engine.run(&graph, &config).await; + let result = engine.run(&graph, &run_options).await; assert!(result.is_err(), "exit 2 should block"); } @@ -7919,9 +7919,9 @@ async fn hook_config_merge_run_overrides_by_name() { let engine = engine_with_hooks(merged.hooks); let graph = parse(simple_linear_dot()).unwrap(); let dir = tempfile::tempdir().unwrap(); - let config = make_run_config(dir.path()); + let run_options = make_run_options(dir.path()); - let outcome = engine.run(&graph, &config).await.unwrap(); + let outcome = engine.run(&graph, &run_options).await.unwrap(); assert_eq!(outcome.status, StageStatus::Success); } @@ -7973,7 +7973,7 @@ async fn hook_blocking_override_makes_non_blocking_event_blocking() { let engine = engine_with_hooks(hooks); let graph = parse(simple_linear_dot()).unwrap(); let dir = tempfile::tempdir().unwrap(); - let config = make_run_config(dir.path()); + let run_options = make_run_options(dir.path()); // This test verifies that the blocking override is respected // Note: StageComplete hooks run AFTER execution, so they use the @@ -7981,7 +7981,7 @@ async fn hook_blocking_override_makes_non_blocking_event_blocking() { // for StageComplete since it's always after the fact). This is correct // behavior — the blocking flag only affects the runner's execution // strategy (sequential vs parallel), not the engine's decision handling. - let outcome = engine.run(&graph, &config).await.unwrap(); + let outcome = engine.run(&graph, &run_options).await.unwrap(); assert_eq!(outcome.status, StageStatus::Success); } @@ -7995,11 +7995,11 @@ async fn hook_non_blocking_override_on_blocking_event() { let engine = engine_with_hooks(hooks); let graph = parse(simple_linear_dot()).unwrap(); let dir = tempfile::tempdir().unwrap(); - let config = make_run_config(dir.path()); + let run_options = make_run_options(dir.path()); // With blocking=false, the RunStart hook failure should NOT block the run // because the runner treats it as non-blocking (doesn't merge decisions) - let outcome = engine.run(&graph, &config).await.unwrap(); + let outcome = engine.run(&graph, &run_options).await.unwrap(); assert_eq!(outcome.status, StageStatus::Success); } @@ -8018,9 +8018,9 @@ async fn hook_matcher_regex_pattern() { let engine = engine_with_hooks(hooks); let graph = parse(two_step_dot()).unwrap(); let dir = tempfile::tempdir().unwrap(); - let config = make_run_config(dir.path()); + let run_options = make_run_options(dir.path()); - let outcome = engine.run(&graph, &config).await.unwrap(); + let outcome = engine.run(&graph, &run_options).await.unwrap(); // Pipeline reached exit with goal gates satisfied — per spec, SUCCESS. assert_eq!(outcome.status, StageStatus::Success); @@ -8054,9 +8054,9 @@ async fn hook_json_proceed_explicit() { let engine = engine_with_hooks(hooks); let graph = parse(simple_linear_dot()).unwrap(); let dir = tempfile::tempdir().unwrap(); - let config = make_run_config(dir.path()); + let run_options = make_run_options(dir.path()); - let outcome = engine.run(&graph, &config).await.unwrap(); + let outcome = engine.run(&graph, &run_options).await.unwrap(); assert_eq!(outcome.status, StageStatus::Success); } @@ -8069,9 +8069,9 @@ async fn hook_json_block_with_reason() { let engine = engine_with_hooks(hooks); let graph = parse(simple_linear_dot()).unwrap(); let dir = tempfile::tempdir().unwrap(); - let config = make_run_config(dir.path()); + let run_options = make_run_options(dir.path()); - let result = engine.run(&graph, &config).await; + let result = engine.run(&graph, &run_options).await; assert!(result.is_err()); assert!(result .unwrap_err() @@ -8095,9 +8095,9 @@ async fn hook_sandbox_false_runs_on_host() { let engine = engine_with_hooks(hooks); let graph = parse(simple_linear_dot()).unwrap(); - let config = make_run_config(dir.path()); + let run_options = make_run_options(dir.path()); - engine.run(&graph, &config).await.unwrap(); + engine.run(&graph, &run_options).await.unwrap(); assert!(marker.exists(), "Host hook should write marker file"); assert_eq!(std::fs::read_to_string(&marker).unwrap().trim(), "host"); @@ -8180,9 +8180,9 @@ async fn hook_prompt_proceed_allows_run() { let engine = engine_with_hooks(hooks); let graph = parse(simple_linear_dot()).unwrap(); let dir = tempfile::tempdir().unwrap(); - let config = make_run_config(dir.path()); + let run_options = make_run_options(dir.path()); - let outcome = engine.run(&graph, &config).await.unwrap(); + let outcome = engine.run(&graph, &run_options).await.unwrap(); assert_eq!(outcome.status, StageStatus::Success); } @@ -8208,9 +8208,9 @@ async fn hook_prompt_block_prevents_run() { let engine = engine_with_hooks(hooks); let graph = parse(simple_linear_dot()).unwrap(); let dir = tempfile::tempdir().unwrap(); - let config = make_run_config(dir.path()); + let run_options = make_run_options(dir.path()); - let result = engine.run(&graph, &config).await; + let result = engine.run(&graph, &run_options).await; assert!( result.is_err(), "Prompt hook block should cause error, got: {result:?}" @@ -8239,9 +8239,9 @@ async fn hook_agent_proceed_allows_run() { let engine = engine_with_hooks(hooks); let graph = parse(simple_linear_dot()).unwrap(); let dir = tempfile::tempdir().unwrap(); - let config = make_run_config(dir.path()); + let run_options = make_run_options(dir.path()); - let outcome = engine.run(&graph, &config).await.unwrap(); + let outcome = engine.run(&graph, &run_options).await.unwrap(); assert_eq!(outcome.status, StageStatus::Success); } @@ -8273,9 +8273,9 @@ async fn hook_agent_with_tool_use() { }]; let engine = engine_with_hooks(hooks); let graph = parse(simple_linear_dot()).unwrap(); - let config = make_run_config(dir.path()); + let run_options = make_run_options(dir.path()); - let outcome = engine.run(&graph, &config).await.unwrap(); + let outcome = engine.run(&graph, &run_options).await.unwrap(); assert_eq!(outcome.status, StageStatus::Success); } @@ -8292,9 +8292,9 @@ async fn hooks_do_not_duplicate_workflow_events() { let (engine, events) = engine_with_hooks_and_events(hooks); let graph = parse(simple_linear_dot()).unwrap(); let dir = tempfile::tempdir().unwrap(); - let config = make_run_config(dir.path()); + let run_options = make_run_options(dir.path()); - engine.run(&graph, &config).await.unwrap(); + engine.run(&graph, &run_options).await.unwrap(); let captured = events.lock().unwrap(); @@ -8367,7 +8367,7 @@ async fn arc_e2e_with_real_llm() { let run_dir = tempfile::tempdir().unwrap(); let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::new()), local_env()); - let config = RunOptions { + let run_options = RunOptions { config: FabroConfig::default(), run_dir: run_dir.path().to_path_buf(), cancel_token: None, @@ -8382,7 +8382,7 @@ async fn arc_e2e_with_real_llm() { git: None, }; let outcome = engine - .run(&graph, &config) + .run(&graph, &run_options) .await .expect("run should succeed"); @@ -8495,7 +8495,7 @@ async fn run_fidelity_prompt_pipeline(fidelity: &str) -> String { ); let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::new()), local_env()); - let config = RunOptions { + let run_options = RunOptions { config: FabroConfig::default(), run_dir: dir.path().to_path_buf(), cancel_token: None, @@ -8510,7 +8510,7 @@ async fn run_fidelity_prompt_pipeline(fidelity: &str) -> String { git: None, }; engine - .run(&graph, &config) + .run(&graph, &run_options) .await .expect("pipeline should succeed"); @@ -8694,7 +8694,7 @@ async fn large_context_values_are_offloaded_to_artifact_store() { let emitter = EventEmitter::new(); let events = collect_events(&emitter); let engine = WorkflowRunner::new(registry, Arc::new(emitter), local_env()); - let config = RunOptions { + let run_options = RunOptions { config: FabroConfig::default(), run_dir: dir.path().to_path_buf(), cancel_token: None, @@ -8709,7 +8709,7 @@ async fn large_context_values_are_offloaded_to_artifact_store() { git: None, }; let outcome = engine - .run(&graph, &config) + .run(&graph, &run_options) .await .expect("pipeline should succeed"); assert_eq!(outcome.status, StageStatus::Success); @@ -8913,7 +8913,7 @@ async fn artifact_pointers_rewritten_for_remote_sandbox() { let remote_env = Arc::new(RemoteMockEnv::new("/sandbox")); let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::new()), remote_env.clone()); - let config = RunOptions { + let run_options = RunOptions { config: FabroConfig::default(), run_dir: dir.path().to_path_buf(), cancel_token: None, @@ -8928,7 +8928,7 @@ async fn artifact_pointers_rewritten_for_remote_sandbox() { git: None, }; let outcome = engine - .run(&graph, &config) + .run(&graph, &run_options) .await .expect("pipeline should succeed"); assert_eq!(outcome.status, StageStatus::Success); @@ -9043,7 +9043,7 @@ async fn node_dir_uses_visit_count_on_revisit() { ); let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::new()), local_env()); - let config = RunOptions { + let run_options = RunOptions { config: FabroConfig::default(), run_dir: dir.path().to_path_buf(), cancel_token: None, @@ -9058,7 +9058,7 @@ async fn node_dir_uses_visit_count_on_revisit() { git: None, }; let outcome = engine - .run(&graph, &config) + .run(&graph, &run_options) .await .expect("pipeline should succeed"); assert_eq!(outcome.status, StageStatus::Success); @@ -10015,7 +10015,7 @@ async fn full_pipeline_with_cli_backend_node() { let dir = tempfile::tempdir().unwrap(); let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::new()), env); - let config = RunOptions { + let run_options = RunOptions { config: FabroConfig::default(), run_dir: dir.path().to_path_buf(), cancel_token: None, @@ -10030,7 +10030,7 @@ async fn full_pipeline_with_cli_backend_node() { git: None, }; let outcome = engine - .run(&graph, &config) + .run(&graph, &run_options) .await .expect("pipeline should succeed"); assert_eq!(outcome.status, StageStatus::Success); @@ -10146,7 +10146,7 @@ async fn stylesheet_backend_property_routes_to_cli() { let dir = tempfile::tempdir().unwrap(); let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::new()), env); - let config = RunOptions { + let run_options = RunOptions { config: FabroConfig::default(), run_dir: dir.path().to_path_buf(), cancel_token: None, @@ -10161,7 +10161,7 @@ async fn stylesheet_backend_property_routes_to_cli() { git: None, }; let outcome = engine - .run(&graph, &config) + .run(&graph, &run_options) .await .expect("pipeline should succeed"); assert_eq!(outcome.status, StageStatus::Success); @@ -10425,7 +10425,7 @@ async fn git_checkpoint_host_emits_events_and_diff_patch() { registry.register("exit", Box::new(ExitHandler)); let engine = WorkflowRunner::new(registry, Arc::new(emitter), env); - let config = RunOptions { + let run_options = RunOptions { config: FabroConfig::default(), run_dir: run_dir.path().to_path_buf(), cancel_token: None, @@ -10445,7 +10445,7 @@ async fn git_checkpoint_host_emits_events_and_diff_patch() { }; // 5. Run pipeline let outcome = engine - .run(&graph, &config) + .run(&graph, &run_options) .await .expect("pipeline should succeed"); assert_eq!(outcome.status, StageStatus::Success); @@ -10628,7 +10628,7 @@ async fn git_checkpoint_host_writes_shadow_branch() { let engine = WorkflowRunner::new(registry, Arc::new(emitter), env); let meta_branch = MetadataStore::branch_name(run_id); - let config = RunOptions { + let run_options = RunOptions { config: FabroConfig::default(), run_dir: run_dir.path().to_path_buf(), cancel_token: None, @@ -10648,7 +10648,7 @@ async fn git_checkpoint_host_writes_shadow_branch() { }; // 5. Run pipeline let outcome = engine - .run(&graph, &config) + .run(&graph, &run_options) .await .expect("pipeline should succeed"); assert_eq!(outcome.status, StageStatus::Success); @@ -10826,7 +10826,7 @@ async fn parallel_git_branching_host_e2e() { let engine = WorkflowRunner::new(registry, Arc::new(emitter), env); - let config = RunOptions { + let run_options = RunOptions { config: FabroConfig::default(), run_dir: run_dir.path().to_path_buf(), cancel_token: None, @@ -10846,7 +10846,7 @@ async fn parallel_git_branching_host_e2e() { }; // 5. Run pipeline let outcome = engine - .run(&graph, &config) + .run(&graph, &run_options) .await .expect("parallel pipeline should succeed"); assert_eq!( @@ -11089,7 +11089,7 @@ async fn git_checkpoint_host_skips_empty_diff_patch() { registry.register("exit", Box::new(ExitHandler)); let engine = WorkflowRunner::new(registry, Arc::new(emitter), env); - let config = RunOptions { + let run_options = RunOptions { config: FabroConfig::default(), run_dir: run_dir.path().to_path_buf(), cancel_token: None, @@ -11108,7 +11108,7 @@ async fn git_checkpoint_host_skips_empty_diff_patch() { }), }; let outcome = engine - .run(&graph, &config) + .run(&graph, &run_options) .await .expect("pipeline should succeed"); assert_eq!(outcome.status, StageStatus::Success); @@ -11471,7 +11471,7 @@ async fn e2e_circuit_breaker_deterministic_self_loop() { ); let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::new()), local_env()); - let config = RunOptions { + let run_options = RunOptions { config: FabroConfig::default(), run_dir: dir.path().to_path_buf(), cancel_token: None, @@ -11485,7 +11485,7 @@ async fn e2e_circuit_breaker_deterministic_self_loop() { host_repo_path: None, git: None, }; - let result = engine.run(&graph, &config).await; + let result = engine.run(&graph, &run_options).await; assert!(result.is_err(), "pipeline should abort, not loop forever"); let err = result.unwrap_err().to_string(); assert!( @@ -11518,7 +11518,7 @@ async fn e2e_circuit_breaker_custom_limit() { ); let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::new()), local_env()); - let config = RunOptions { + let run_options = RunOptions { config: FabroConfig::default(), run_dir: dir.path().to_path_buf(), cancel_token: None, @@ -11532,7 +11532,7 @@ async fn e2e_circuit_breaker_custom_limit() { host_repo_path: None, git: None, }; - let result = engine.run(&graph, &config).await; + let result = engine.run(&graph, &run_options).await; assert!(result.is_err()); let err = result.unwrap_err().to_string(); assert!( @@ -11558,7 +11558,7 @@ async fn e2e_circuit_breaker_ignores_transient_failures() { registry.register("test_handler", Box::new(TransientInfraFailHandler)); let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::new()), local_env()); - let config = RunOptions { + let run_options = RunOptions { config: FabroConfig::default(), run_dir: dir.path().to_path_buf(), cancel_token: None, @@ -11572,7 +11572,7 @@ async fn e2e_circuit_breaker_ignores_transient_failures() { host_repo_path: None, git: None, }; - let result = engine.run(&graph, &config).await; + let result = engine.run(&graph, &run_options).await; assert!(result.is_err()); let err = result.unwrap_err().to_string(); // Should hit visit limit, NOT circuit breaker @@ -11605,7 +11605,7 @@ async fn e2e_circuit_breaker_different_reasons_separate_counters() { ); let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::new()), local_env()); - let config = RunOptions { + let run_options = RunOptions { config: FabroConfig::default(), run_dir: dir.path().to_path_buf(), cancel_token: None, @@ -11619,7 +11619,7 @@ async fn e2e_circuit_breaker_different_reasons_separate_counters() { host_repo_path: None, git: None, }; - let result = engine.run(&graph, &config).await; + let result = engine.run(&graph, &run_options).await; assert!(result.is_err()); let err = result.unwrap_err().to_string(); // Should hit visit limit because each failure has a unique signature @@ -11645,7 +11645,7 @@ async fn e2e_circuit_breaker_loop_restart() { ); let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::new()), local_env()); - let config = RunOptions { + let run_options = RunOptions { config: FabroConfig::default(), run_dir: dir.path().to_path_buf(), cancel_token: None, @@ -11659,7 +11659,7 @@ async fn e2e_circuit_breaker_loop_restart() { host_repo_path: None, git: None, }; - let result = engine.run(&graph, &config).await; + let result = engine.run(&graph, &run_options).await; assert!( result.is_err(), "pipeline should abort, not restart forever" @@ -11707,7 +11707,7 @@ async fn e2e_failure_signature_persisted_in_context() { ); let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::new()), local_env()); - let config = RunOptions { + let run_options = RunOptions { config: FabroConfig::default(), run_dir: dir.path().to_path_buf(), cancel_token: None, @@ -11721,7 +11721,7 @@ async fn e2e_failure_signature_persisted_in_context() { host_repo_path: None, git: None, }; - let outcome = engine.run(&graph, &config).await.unwrap(); + let outcome = engine.run(&graph, &run_options).await.unwrap(); // Pipeline reaches exit (terminal) with goal gates satisfied. // Per spec, reaching exit with satisfied goal gates returns SUCCESS. assert_eq!(outcome.status, StageStatus::Success); @@ -11771,7 +11771,7 @@ async fn e2e_failure_signature_hint_overrides_reason_in_context() { registry.register("hint_handler", Box::new(SignatureHintHandler)); let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::new()), local_env()); - let config = RunOptions { + let run_options = RunOptions { config: FabroConfig::default(), run_dir: dir.path().to_path_buf(), cancel_token: None, @@ -11785,7 +11785,7 @@ async fn e2e_failure_signature_hint_overrides_reason_in_context() { host_repo_path: None, git: None, }; - let _outcome = engine.run(&graph, &config).await.unwrap(); + let _outcome = engine.run(&graph, &run_options).await.unwrap(); let cp = Checkpoint::load(&dir.path().join("checkpoint.json")).unwrap(); let sig_str = cp @@ -11827,7 +11827,7 @@ async fn e2e_signature_maps_persist_in_checkpoint() { ); let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::new()), local_env()); - let config = RunOptions { + let run_options = RunOptions { config: FabroConfig::default(), run_dir: dir.path().to_path_buf(), cancel_token: None, @@ -11841,7 +11841,7 @@ async fn e2e_signature_maps_persist_in_checkpoint() { host_repo_path: None, git: None, }; - let outcome = engine.run(&graph, &config).await.unwrap(); + let outcome = engine.run(&graph, &run_options).await.unwrap(); assert_eq!(outcome.status, StageStatus::Success); // Load checkpoint and verify signature maps @@ -11954,7 +11954,7 @@ async fn e2e_circuit_breaker_emits_events_before_abort() { ); let engine = WorkflowRunner::new(registry, Arc::new(emitter), local_env()); - let config = RunOptions { + let run_options = RunOptions { config: FabroConfig::default(), run_dir: dir.path().to_path_buf(), cancel_token: None, @@ -11968,7 +11968,7 @@ async fn e2e_circuit_breaker_emits_events_before_abort() { host_repo_path: None, git: None, }; - let result = engine.run(&graph, &config).await; + let result = engine.run(&graph, &run_options).await; assert!(result.is_err()); let events = events.lock().unwrap(); @@ -12021,7 +12021,7 @@ async fn e2e_circuit_breaker_does_not_fire_below_limit() { ); let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::new()), local_env()); - let config = RunOptions { + let run_options = RunOptions { config: FabroConfig::default(), run_dir: dir.path().to_path_buf(), cancel_token: None, @@ -12035,7 +12035,7 @@ async fn e2e_circuit_breaker_does_not_fire_below_limit() { host_repo_path: None, git: None, }; - let outcome = engine.run(&graph, &config).await.unwrap(); + let outcome = engine.run(&graph, &run_options).await.unwrap(); assert_eq!( outcome.status, StageStatus::Success, @@ -12117,7 +12117,7 @@ async fn e2e_circuit_breaker_multi_stage_impl_verify_cycle() { ); let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::new()), local_env()); - let config = RunOptions { + let run_options = RunOptions { config: FabroConfig::default(), run_dir: dir.path().to_path_buf(), cancel_token: None, @@ -12131,7 +12131,7 @@ async fn e2e_circuit_breaker_multi_stage_impl_verify_cycle() { host_repo_path: None, git: None, }; - let result = engine.run(&graph, &config).await; + let result = engine.run(&graph, &run_options).await; assert!( result.is_err(), "should detect impl/verify cycle, not loop forever" @@ -12214,7 +12214,7 @@ async fn e2e_loop_restart_blocked_for_deterministic_failure() { ); let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::new()), local_env()); - let config = RunOptions { + let run_options = RunOptions { config: FabroConfig::default(), run_dir: dir.path().to_path_buf(), cancel_token: None, @@ -12228,7 +12228,7 @@ async fn e2e_loop_restart_blocked_for_deterministic_failure() { host_repo_path: None, git: None, }; - let result = engine.run(&graph, &config).await; + let result = engine.run(&graph, &run_options).await; assert!( result.is_err(), "deterministic failure should not loop_restart" @@ -12254,7 +12254,7 @@ async fn e2e_loop_restart_blocked_for_structural_failure() { ); let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::new()), local_env()); - let config = RunOptions { + let run_options = RunOptions { config: FabroConfig::default(), run_dir: dir.path().to_path_buf(), cancel_token: None, @@ -12268,7 +12268,7 @@ async fn e2e_loop_restart_blocked_for_structural_failure() { host_repo_path: None, git: None, }; - let result = engine.run(&graph, &config).await; + let result = engine.run(&graph, &run_options).await; assert!( result.is_err(), "structural failure should not loop_restart" @@ -12294,7 +12294,7 @@ async fn e2e_loop_restart_blocked_for_budget_exhausted_failure() { ); let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::new()), local_env()); - let config = RunOptions { + let run_options = RunOptions { config: FabroConfig::default(), run_dir: dir.path().to_path_buf(), cancel_token: None, @@ -12308,7 +12308,7 @@ async fn e2e_loop_restart_blocked_for_budget_exhausted_failure() { host_repo_path: None, git: None, }; - let result = engine.run(&graph, &config).await; + let result = engine.run(&graph, &run_options).await; assert!( result.is_err(), "budget_exhausted failure should not loop_restart" @@ -12334,7 +12334,7 @@ async fn e2e_loop_restart_blocked_for_canceled_failure() { ); let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::new()), local_env()); - let config = RunOptions { + let run_options = RunOptions { config: FabroConfig::default(), run_dir: dir.path().to_path_buf(), cancel_token: None, @@ -12348,7 +12348,7 @@ async fn e2e_loop_restart_blocked_for_canceled_failure() { host_repo_path: None, git: None, }; - let result = engine.run(&graph, &config).await; + let result = engine.run(&graph, &run_options).await; assert!(result.is_err(), "canceled failure should not loop_restart"); let err = result.unwrap_err().to_string(); assert!( @@ -12371,7 +12371,7 @@ async fn e2e_loop_restart_blocked_for_compilation_loop_failure() { ); let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::new()), local_env()); - let config = RunOptions { + let run_options = RunOptions { config: FabroConfig::default(), run_dir: dir.path().to_path_buf(), cancel_token: None, @@ -12385,7 +12385,7 @@ async fn e2e_loop_restart_blocked_for_compilation_loop_failure() { host_repo_path: None, git: None, }; - let result = engine.run(&graph, &config).await; + let result = engine.run(&graph, &run_options).await; assert!( result.is_err(), "compilation_loop failure should not loop_restart" @@ -12412,7 +12412,7 @@ async fn e2e_loop_restart_allowed_for_transient_infra() { ); let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::new()), local_env()); - let config = RunOptions { + let run_options = RunOptions { config: FabroConfig::default(), run_dir: dir.path().to_path_buf(), cancel_token: None, @@ -12426,7 +12426,7 @@ async fn e2e_loop_restart_allowed_for_transient_infra() { host_repo_path: None, git: None, }; - let result = engine.run(&graph, &config).await; + let result = engine.run(&graph, &run_options).await; assert!( result.is_ok(), "transient_infra failure should be allowed to loop_restart, got: {:?}", @@ -12516,7 +12516,7 @@ async fn e2e_stall_watchdog_triggers_from_dot_parsed_pipeline() { }); let engine = WorkflowRunner::new(registry, Arc::new(emitter), local_env()); - let config = RunOptions { + let run_options = RunOptions { config: FabroConfig::default(), run_dir: dir.path().to_path_buf(), cancel_token: None, @@ -12530,7 +12530,7 @@ async fn e2e_stall_watchdog_triggers_from_dot_parsed_pipeline() { host_repo_path: None, git: None, }; - let result = engine.run(&graph, &config).await; + let result = engine.run(&graph, &run_options).await; assert!(result.is_err(), "expected stall watchdog error"); let err = result.unwrap_err().to_string(); assert!( @@ -12572,7 +12572,7 @@ async fn e2e_stall_watchdog_kept_alive_by_handler_events() { ); let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::new()), local_env()); - let config = RunOptions { + let run_options = RunOptions { config: FabroConfig::default(), run_dir: dir.path().to_path_buf(), cancel_token: None, @@ -12587,7 +12587,7 @@ async fn e2e_stall_watchdog_kept_alive_by_handler_events() { git: None, }; let outcome = engine - .run(&graph, &config) + .run(&graph, &run_options) .await .expect("pipeline should succeed"); assert_eq!(outcome.status, StageStatus::Success); @@ -12618,7 +12618,7 @@ async fn e2e_stall_watchdog_disabled_with_zero_timeout() { registry.register("slow", Box::new(SlowTestHandler { sleep_ms: 50 })); let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::new()), local_env()); - let config = RunOptions { + let run_options = RunOptions { config: FabroConfig::default(), run_dir: dir.path().to_path_buf(), cancel_token: None, @@ -12633,7 +12633,7 @@ async fn e2e_stall_watchdog_disabled_with_zero_timeout() { git: None, }; let outcome = engine - .run(&graph, &config) + .run(&graph, &run_options) .await .expect("pipeline should succeed"); assert_eq!(outcome.status, StageStatus::Success); @@ -12683,7 +12683,7 @@ async fn e2e_stall_watchdog_with_explicit_timeout_override() { registry.register("hanging", Box::new(HangingHandler)); let engine = WorkflowRunner::new(registry, Arc::new(EventEmitter::new()), local_env()); - let config = RunOptions { + let run_options = RunOptions { config: FabroConfig::default(), run_dir: dir.path().to_path_buf(), cancel_token: None, @@ -12698,7 +12698,7 @@ async fn e2e_stall_watchdog_with_explicit_timeout_override() { git: None, }; let start = std::time::Instant::now(); - let result = engine.run(&graph, &config).await; + let result = engine.run(&graph, &run_options).await; let elapsed = start.elapsed(); assert!(result.is_err(), "expected stall watchdog error"); @@ -12814,7 +12814,7 @@ async fn asset_collection_local_sandbox_success() { graph.edges.push(Edge::new("start", "create_assets")); graph.edges.push(Edge::new("create_assets", "exit")); - let config = RunOptions { + let run_options = RunOptions { config: FabroConfig { assets: Some(fabro_config::run::AssetsConfig { include: vec!["test-results/**".to_string()], @@ -12834,7 +12834,7 @@ async fn asset_collection_local_sandbox_success() { git: None, }; let outcome = engine - .run(&graph, &config) + .run(&graph, &run_options) .await .expect("run should succeed"); assert_eq!(outcome.status, StageStatus::Success); @@ -12928,7 +12928,7 @@ async fn asset_collection_local_sandbox_on_failure() { graph.edges.push(Edge::new("start", "create_assets")); graph.edges.push(Edge::new("create_assets", "exit")); - let config = RunOptions { + let run_options = RunOptions { config: FabroConfig { assets: Some(fabro_config::run::AssetsConfig { include: vec!["test-results/**".to_string()], @@ -12948,7 +12948,7 @@ async fn asset_collection_local_sandbox_on_failure() { git: None, }; let outcome = engine - .run(&graph, &config) + .run(&graph, &run_options) .await .expect("run should succeed"); // The pipeline completes with goal gates satisfied — per spec, SUCCESS at exit node. @@ -13025,7 +13025,7 @@ async fn asset_collection_docker_sandbox() { graph.edges.push(Edge::new("start", "create_assets")); graph.edges.push(Edge::new("create_assets", "exit")); - let run_config = RunOptions { + let run_options = RunOptions { config: FabroConfig { assets: Some(fabro_config::run::AssetsConfig { include: vec!["test-results/**".to_string()], @@ -13045,7 +13045,7 @@ async fn asset_collection_docker_sandbox() { git: None, }; let outcome = engine - .run(&graph, &run_config) + .run(&graph, &run_options) .await .expect("pipeline should succeed"); assert_eq!(outcome.status, StageStatus::Success); @@ -13099,7 +13099,7 @@ async fn wait_timer_e2e() { Arc::new(EventEmitter::new()), local_env(), ); - let config = RunOptions { + let run_options = RunOptions { config: FabroConfig::default(), run_dir: dir.path().to_path_buf(), cancel_token: None, @@ -13113,6 +13113,6 @@ async fn wait_timer_e2e() { host_repo_path: None, git: None, }; - let outcome = engine.run(&graph, &config).await.expect("run"); + let outcome = engine.run(&graph, &run_options).await.expect("run"); assert_eq!(outcome.status, StageStatus::Success); }