mirror of
https://github.com/fabro-sh/fabro.git
synced 2026-10-07 03:00:29 +00:00
refactor(workflow): use RunServices builders for child engines
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) <noreply@anthropic.com>
This commit is contained in:
parent
1ea72ca2e8
commit
7c95f6fa8e
3 changed files with 11 additions and 39 deletions
|
|
@ -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,
|
||||
|
|
|
|||
|
|
@ -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(),
|
||||
|
|
|
|||
|
|
@ -112,7 +112,6 @@ impl RunServices {
|
|||
)
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
#[must_use]
|
||||
pub fn with_run_store(self: &Arc<Self>, run_store: RunStoreHandle) -> Arc<Self> {
|
||||
Arc::new(Self {
|
||||
|
|
@ -121,7 +120,6 @@ impl RunServices {
|
|||
})
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
#[must_use]
|
||||
pub fn with_emitter(self: &Arc<Self>, emitter: Arc<Emitter>) -> Arc<Self> {
|
||||
Arc::new(Self {
|
||||
|
|
@ -130,7 +128,6 @@ impl RunServices {
|
|||
})
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
#[must_use]
|
||||
pub fn with_sandbox(self: &Arc<Self>, sandbox: Arc<dyn Sandbox>) -> Arc<Self> {
|
||||
Arc::new(Self {
|
||||
|
|
@ -139,7 +136,6 @@ impl RunServices {
|
|||
})
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
#[must_use]
|
||||
pub fn with_cancel_requested(
|
||||
self: &Arc<Self>,
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue