From 7c95f6fa8eee400e24bdaba34c9377c25870aa24 Mon Sep 17 00:00:00 2001 From: Bryan Helmkamp Date: Thu, 23 Apr 2026 22:37:13 -0400 Subject: [PATCH] refactor(workflow): use RunServices builders for child engines MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Reverses the #[cfg(test)] gating on RunServices::with_run_store / with_emitter / with_sandbox / with_cancel_requested — manager_loop and parallel handlers have production callers that were unpacking 7 fields into locals just to reconstruct RunServices::new(...). manager_loop builds its child via parent_run.with_run_store(...).with_cancel_requested(None). parallel builds each branch via parent_run.with_sandbox(...). Co-Authored-By: Claude Opus 4.7 (1M context) --- .../src/handler/manager_loop.rs | 20 ++++---------- .../fabro-workflow/src/handler/parallel.rs | 26 +++++-------------- lib/crates/fabro-workflow/src/services.rs | 4 --- 3 files changed, 11 insertions(+), 39 deletions(-) diff --git a/lib/crates/fabro-workflow/src/handler/manager_loop.rs b/lib/crates/fabro-workflow/src/handler/manager_loop.rs index 13acd9a10..32ef6fe39 100644 --- a/lib/crates/fabro-workflow/src/handler/manager_loop.rs +++ b/lib/crates/fabro-workflow/src/handler/manager_loop.rs @@ -23,7 +23,6 @@ use crate::pipeline; use crate::pipeline::types::Initialized; use crate::run_dir::visit_from_context; use crate::run_options::RunOptions; -use crate::services::RunServices; /// Orchestrates a child workflow engine, polling for completion or stop /// conditions. @@ -228,12 +227,8 @@ impl Handler for SubWorkflowHandler { } let before_snapshot = context.snapshot(); - let emitter = Arc::clone(&services.run.emitter); - let sandbox = Arc::clone(&services.run.sandbox); + let parent_run = Arc::clone(&services.run); let registry = Arc::clone(&services.registry); - let hook_runner = services.run.hook_runner.clone(); - let provider = services.run.provider; - let llm_source = Arc::clone(&services.run.llm_source); let env = services.env.clone(); let inputs = services.inputs.clone(); let dry_run = services.dry_run; @@ -253,6 +248,9 @@ impl Handler for SubWorkflowHandler { // Spawn child engine let mut child_handle = tokio::spawn(async move { + let child_run = parent_run + .with_run_store(run_store.into()) + .with_cancel_requested(None); let initialized = Initialized { graph: child_graph, source: String::new(), @@ -263,15 +261,7 @@ impl Handler for SubWorkflowHandler { artifact_sink: Some(ArtifactSink::Store(artifact_store)), run_control: None, engine: Arc::new(EngineServices { - run: RunServices::new( - run_store.into(), - emitter, - sandbox, - hook_runner, - None, - provider, - llm_source, - ), + run: child_run, registry, git_state: std::sync::RwLock::new(None), env, diff --git a/lib/crates/fabro-workflow/src/handler/parallel.rs b/lib/crates/fabro-workflow/src/handler/parallel.rs index d6c63d9eb..035db5e0c 100644 --- a/lib/crates/fabro-workflow/src/handler/parallel.rs +++ b/lib/crates/fabro-workflow/src/handler/parallel.rs @@ -19,7 +19,6 @@ use crate::millis_u64; use crate::outcome::{FailureCategory, FailureDetail, Outcome, OutcomeExt, StageStatus}; use crate::run_dir::visit_from_context; use crate::sandbox_git::{GIT_REMOTE, git_checkpoint, git_merge_ff_only, git_remove_worktree}; -use crate::services::RunServices; /// Fans out execution to multiple branches concurrently. /// Each branch gets an isolated context clone and runs independently. @@ -286,16 +285,11 @@ impl Handler for ParallelHandler { // --- Fan out: concurrent execution --- let mut handles = Vec::new(); for setup in branch_setups { + let parent_run = Arc::clone(&services.run); let registry = Arc::clone(&services.registry); - let emitter = Arc::clone(&services.run.emitter); - let hook_runner = services.run.hook_runner.clone(); - let run_store = services.run.run_store.clone(); let env = services.env.clone(); let inputs = services.inputs.clone(); let dry_run = services.dry_run; - let cancel_requested = services.run.cancel_requested.clone(); - let provider = services.run.provider; - let llm_source = Arc::clone(&services.run.llm_source); let workflow_path = services.workflow_path.clone(); let workflow_bundle = services.workflow_bundle.clone(); let graph = graph.clone(); @@ -321,7 +315,7 @@ impl Handler for ParallelHandler { .await .map_err(|e| Error::handler(format!("semaphore error: {e}")))?; - emitter.emit_scoped( + parent_run.emitter.emit_scoped( &Event::ParallelBranchStarted { parallel_group_id: group_id.clone(), parallel_branch_id: setup.parallel_branch_id.clone(), @@ -337,7 +331,7 @@ impl Handler for ParallelHandler { "branch target node not found: {}", setup.target_id )); - emitter.emit_scoped( + parent_run.emitter.emit_scoped( &Event::ParallelBranchCompleted { parallel_group_id: group_id.clone(), parallel_branch_id: setup.parallel_branch_id.clone(), @@ -358,15 +352,7 @@ impl Handler for ParallelHandler { }; let branch_services = EngineServices { - run: RunServices::new( - run_store.clone(), - Arc::clone(&emitter), - Arc::clone(&setup.sandbox), - hook_runner.clone(), - cancel_requested, - provider, - llm_source, - ), + run: parent_run.with_sandbox(Arc::clone(&setup.sandbox)), registry: Arc::clone(®istry), git_state: std::sync::RwLock::new(None), env: env.clone(), @@ -419,7 +405,7 @@ impl Handler for ParallelHandler { match sha_result { Ok(r) if r.exit_code == 0 => { let sha = r.stdout.trim().to_string(); - emitter.emit_scoped( + parent_run.emitter.emit_scoped( &Event::GitCommit { node_id: Some(setup.target_id.clone()), sha: sha.clone(), @@ -434,7 +420,7 @@ impl Handler for ParallelHandler { None }; - emitter.emit_scoped( + parent_run.emitter.emit_scoped( &Event::ParallelBranchCompleted { parallel_group_id: group_id.clone(), parallel_branch_id: setup.parallel_branch_id.clone(), diff --git a/lib/crates/fabro-workflow/src/services.rs b/lib/crates/fabro-workflow/src/services.rs index 01323418a..1751212ad 100644 --- a/lib/crates/fabro-workflow/src/services.rs +++ b/lib/crates/fabro-workflow/src/services.rs @@ -112,7 +112,6 @@ impl RunServices { ) } - #[cfg(test)] #[must_use] pub fn with_run_store(self: &Arc, run_store: RunStoreHandle) -> Arc { Arc::new(Self { @@ -121,7 +120,6 @@ impl RunServices { }) } - #[cfg(test)] #[must_use] pub fn with_emitter(self: &Arc, emitter: Arc) -> Arc { Arc::new(Self { @@ -130,7 +128,6 @@ impl RunServices { }) } - #[cfg(test)] #[must_use] pub fn with_sandbox(self: &Arc, sandbox: Arc) -> Arc { Arc::new(Self { @@ -139,7 +136,6 @@ impl RunServices { }) } - #[cfg(test)] #[must_use] pub fn with_cancel_requested( self: &Arc,