diff --git a/crates/arc-workflows/src/engine.rs b/crates/arc-workflows/src/engine.rs index ac111f9ca..8d46acc11 100644 --- a/crates/arc-workflows/src/engine.rs +++ b/crates/arc-workflows/src/engine.rs @@ -501,6 +501,16 @@ fn is_terminal(node: &Node) -> bool { // --- Pipeline engine --- +/// Captured git state for a pipeline run, shared with handlers. +#[derive(Debug, Clone)] +pub struct GitState { + pub mode: GitCheckpointMode, + pub run_id: String, + pub base_sha: String, + pub run_branch: Option, + pub meta_branch: Option, +} + /// How git checkpointing should be performed for a pipeline run. #[derive(Debug, Clone)] pub enum GitCheckpointMode { @@ -512,7 +522,7 @@ pub enum GitCheckpointMode { } /// Run a git checkpoint commit on the host filesystem (local/Docker bind-mount). -async fn git_checkpoint_host( +pub async fn git_checkpoint_host( work_dir: PathBuf, run_id: String, node_id: String, @@ -546,7 +556,7 @@ async fn git_diff_host(work_dir: PathBuf, base: String) -> Option { } /// Run a git checkpoint commit inside a remote execution environment. -async fn git_checkpoint_remote( +pub async fn git_checkpoint_remote( exec_env: &dyn ExecutionEnvironment, run_id: &str, node_id: &str, @@ -622,6 +632,69 @@ async fn git_diff_remote(exec_env: &dyn ExecutionEnvironment, base: &str) -> Opt } } +// --- Remote worktree helpers (for Daytona / sandbox environments) --- + +/// Create a branch at a specific SHA inside a remote execution environment. +pub async fn git_create_branch_at_remote( + exec_env: &dyn ExecutionEnvironment, + name: &str, + sha: &str, +) -> bool { + let cmd = format!("git branch {name} {sha}"); + matches!( + exec_env.exec_command(&cmd, 30_000, None, None, None).await, + Ok(r) if r.exit_code == 0 + ) +} + +/// Add a git worktree inside a remote execution environment. +pub async fn git_add_worktree_remote( + exec_env: &dyn ExecutionEnvironment, + path: &str, + branch: &str, +) -> bool { + let cmd = format!("git worktree add {path} {branch}"); + matches!( + exec_env.exec_command(&cmd, 30_000, None, None, None).await, + Ok(r) if r.exit_code == 0 + ) +} + +/// Remove a git worktree inside a remote execution environment. +pub async fn git_remove_worktree_remote( + exec_env: &dyn ExecutionEnvironment, + path: &str, +) -> bool { + let cmd = format!("git worktree remove --force {path}"); + matches!( + exec_env.exec_command(&cmd, 30_000, None, None, None).await, + Ok(r) if r.exit_code == 0 + ) +} + +/// Fast-forward merge to a given SHA inside a remote execution environment. +pub async fn git_merge_ff_only_remote( + exec_env: &dyn ExecutionEnvironment, + sha: &str, +) -> bool { + let cmd = format!("git merge --ff-only {sha}"); + matches!( + exec_env.exec_command(&cmd, 30_000, None, None, None).await, + Ok(r) if r.exit_code == 0 + ) +} + +/// Get the current HEAD SHA from a remote execution environment. +pub async fn git_head_sha_remote(exec_env: &dyn ExecutionEnvironment) -> Option { + match exec_env + .exec_command("git rev-parse HEAD", 10_000, None, None, None) + .await + { + Ok(r) if r.exit_code == 0 => Some(r.stdout.trim().to_string()), + _ => None, + } +} + /// Configuration for a pipeline run. pub struct RunConfig { pub logs_root: PathBuf, @@ -657,6 +730,7 @@ impl PipelineEngine { registry: Arc::new(registry), emitter, execution_env, + git_state: std::sync::RwLock::new(None), }, interviewer: None, } @@ -675,6 +749,7 @@ impl PipelineEngine { registry: Arc::new(registry), emitter, execution_env, + git_state: std::sync::RwLock::new(None), }, interviewer: Some(interviewer), } @@ -871,6 +946,19 @@ impl PipelineEngine { let run_id = config.run_id.clone(); let artifact_store = ArtifactStore::new(Some(config.logs_root.clone())); + // Populate git_state for handlers (parallel, fan_in) when checkpointing is active + let git_state = match (&config.git_checkpoint, &config.base_sha) { + (Some(mode), Some(base_sha)) => Some(Arc::new(GitState { + mode: mode.clone(), + run_id: run_id.clone(), + base_sha: base_sha.clone(), + run_branch: config.run_branch.clone(), + meta_branch: config.meta_branch.clone(), + })), + _ => None, + }; + self.services.set_git_state(git_state); + self.services.emitter.emit(&PipelineEvent::PipelineStarted { name: graph.name.clone(), run_id: run_id.clone(), diff --git a/crates/arc-workflows/src/git.rs b/crates/arc-workflows/src/git.rs index 81ce13cca..c77a5330e 100644 --- a/crates/arc-workflows/src/git.rs +++ b/crates/arc-workflows/src/git.rs @@ -99,6 +99,39 @@ pub fn remove_worktree(repo: &Path, path: &Path) -> Result<()> { Ok(()) } +/// Create a new branch pointing at a specific SHA (without checking it out). +pub fn create_branch_at(repo: &Path, name: &str, sha: &str) -> Result<()> { + let output = Command::new("git") + .args(["branch", name, sha]) + .current_dir(repo) + .output() + .map_err(|e| git_error(format!("git branch failed: {e}")))?; + + if !output.status.success() { + let stderr = String::from_utf8_lossy(&output.stderr); + return Err(git_error(format!("git branch failed: {stderr}"))); + } + + Ok(()) +} + +/// Fast-forward the current branch to a given SHA. +/// Fails if the merge cannot be done as a fast-forward. +pub fn merge_ff_only(work_dir: &Path, sha: &str) -> Result<()> { + let output = Command::new("git") + .args(["merge", "--ff-only", sha]) + .current_dir(work_dir) + .output() + .map_err(|e| git_error(format!("git merge --ff-only failed: {e}")))?; + + if !output.status.success() { + let stderr = String::from_utf8_lossy(&output.stderr); + return Err(git_error(format!("git merge --ff-only failed: {stderr}"))); + } + + Ok(()) +} + /// Stage all changes and commit in `work_dir` with a structured message /// including trailers for completed node count and shadow commit pointer. /// Returns the new commit SHA. @@ -375,6 +408,99 @@ mod tests { assert!(stdout.contains("test-branch")); } + #[test] + fn create_branch_at_specific_sha() { + let dir = tempfile::tempdir().unwrap(); + init_repo(dir.path()); + let initial_sha = head_sha(dir.path()).unwrap(); + + // Make a second commit so HEAD differs from initial_sha + fs::write(dir.path().join("f.txt"), "x").unwrap(); + Command::new("git") + .args(["add", "-A"]) + .current_dir(dir.path()) + .output() + .unwrap(); + Command::new("git") + .args([ + "-c", + "user.name=test", + "-c", + "user.email=test@test", + "commit", + "-m", + "second", + ]) + .current_dir(dir.path()) + .output() + .unwrap(); + + // Branch at the *initial* SHA (not HEAD) + create_branch_at(dir.path(), "at-initial", &initial_sha).unwrap(); + + let output = Command::new("git") + .args(["rev-parse", "at-initial"]) + .current_dir(dir.path()) + .output() + .unwrap(); + let branch_sha = String::from_utf8_lossy(&output.stdout).trim().to_string(); + assert_eq!(branch_sha, initial_sha); + } + + #[test] + fn merge_ff_only_advances_branch() { + let dir = tempfile::tempdir().unwrap(); + init_repo(dir.path()); + let base_sha = head_sha(dir.path()).unwrap(); + + // Create a branch and worktree, make a commit there + create_branch(dir.path(), "ff-branch").unwrap(); + let wt = dir.path().join("ff-wt"); + add_worktree(dir.path(), &wt, "ff-branch").unwrap(); + fs::write(wt.join("new.txt"), "data").unwrap(); + checkpoint_commit(&wt, "run", "node", "ok", 1, None).unwrap(); + let advanced_sha = head_sha(&wt).unwrap(); + remove_worktree(dir.path(), &wt).unwrap(); + + // Main branch is still at base_sha + assert_eq!(head_sha(dir.path()).unwrap(), base_sha); + + // Fast-forward main to advanced_sha + merge_ff_only(dir.path(), &advanced_sha).unwrap(); + assert_eq!(head_sha(dir.path()).unwrap(), advanced_sha); + } + + #[test] + fn merge_ff_only_fails_on_diverged() { + let dir = tempfile::tempdir().unwrap(); + init_repo(dir.path()); + + // Create divergent history: commit on main + fs::write(dir.path().join("a.txt"), "a").unwrap(); + Command::new("git") + .args(["add", "-A"]) + .current_dir(dir.path()) + .output() + .unwrap(); + Command::new("git") + .args([ + "-c", + "user.name=test", + "-c", + "user.email=test@test", + "commit", + "-m", + "on main", + ]) + .current_dir(dir.path()) + .output() + .unwrap(); + + // A random SHA that isn't an ancestor/descendant + let err = merge_ff_only(dir.path(), "0000000000000000000000000000000000000000"); + assert!(err.is_err()); + } + #[test] fn add_and_remove_worktree() { let dir = tempfile::tempdir().unwrap(); diff --git a/crates/arc-workflows/src/handler/codergen.rs b/crates/arc-workflows/src/handler/codergen.rs index 3fd6f1229..18d081426 100644 --- a/crates/arc-workflows/src/handler/codergen.rs +++ b/crates/arc-workflows/src/handler/codergen.rs @@ -336,6 +336,7 @@ mod tests { execution_env: std::sync::Arc::new(arc_agent::LocalExecutionEnvironment::new( std::env::current_dir().unwrap_or_else(|_| std::path::PathBuf::from(".")), )), + git_state: std::sync::RwLock::new(None), } } diff --git a/crates/arc-workflows/src/handler/conditional.rs b/crates/arc-workflows/src/handler/conditional.rs index 18889cf1d..cea0262b0 100644 --- a/crates/arc-workflows/src/handler/conditional.rs +++ b/crates/arc-workflows/src/handler/conditional.rs @@ -43,6 +43,7 @@ mod tests { execution_env: std::sync::Arc::new(arc_agent::LocalExecutionEnvironment::new( std::env::current_dir().unwrap_or_else(|_| std::path::PathBuf::from(".")), )), + git_state: std::sync::RwLock::new(None), } } diff --git a/crates/arc-workflows/src/handler/exit.rs b/crates/arc-workflows/src/handler/exit.rs index 677b06b11..bdbed8a9d 100644 --- a/crates/arc-workflows/src/handler/exit.rs +++ b/crates/arc-workflows/src/handler/exit.rs @@ -40,6 +40,7 @@ mod tests { execution_env: std::sync::Arc::new(arc_agent::LocalExecutionEnvironment::new( std::env::current_dir().unwrap_or_else(|_| std::path::PathBuf::from(".")), )), + git_state: std::sync::RwLock::new(None), } } diff --git a/crates/arc-workflows/src/handler/fan_in.rs b/crates/arc-workflows/src/handler/fan_in.rs index fced0357d..8555333dd 100644 --- a/crates/arc-workflows/src/handler/fan_in.rs +++ b/crates/arc-workflows/src/handler/fan_in.rs @@ -72,6 +72,31 @@ impl Handler for FanInHandler { return Ok(Outcome::fail("all candidates failed")); } + // --- Fast-forward to winner's HEAD when git isolation is active --- + let best_head_sha = { + let empty_vec = vec![]; + let arr = results.as_array().unwrap_or(&empty_vec); + arr.iter() + .find(|v| v.get("id").and_then(|v| v.as_str()) == Some(&best.id)) + .and_then(|v| v.get("head_sha").and_then(|v| v.as_str()).map(String::from)) + }; + + if let (Some(ref sha), Some(ref gs)) = (&best_head_sha, services.git_state()) { + match &gs.mode { + crate::engine::GitCheckpointMode::Host(work_dir) => { + let wd = work_dir.clone(); + let s = sha.clone(); + let _ = tokio::task::spawn_blocking(move || { + crate::git::merge_ff_only(&wd, &s) + }) + .await; + } + crate::engine::GitCheckpointMode::Remote(_) => { + crate::engine::git_merge_ff_only_remote(&*services.execution_env, sha).await; + } + } + } + let mut outcome = Outcome::success(); outcome.context_updates.insert( "parallel.fan_in.best_id".to_string(), @@ -81,6 +106,12 @@ impl Handler for FanInHandler { "parallel.fan_in.best_outcome".to_string(), serde_json::json!(best.status), ); + if let Some(ref sha) = best_head_sha { + outcome.context_updates.insert( + "parallel.fan_in.best_head_sha".to_string(), + serde_json::json!(sha), + ); + } outcome.notes = Some(format!("Selected best candidate: {}", best.id)); Ok(outcome) @@ -273,6 +304,7 @@ mod tests { execution_env: std::sync::Arc::new(arc_agent::LocalExecutionEnvironment::new( std::env::current_dir().unwrap_or_else(|_| std::path::PathBuf::from(".")), )), + git_state: std::sync::RwLock::new(None), } } diff --git a/crates/arc-workflows/src/handler/manager_loop.rs b/crates/arc-workflows/src/handler/manager_loop.rs index 5b61b0906..d7bf71c88 100644 --- a/crates/arc-workflows/src/handler/manager_loop.rs +++ b/crates/arc-workflows/src/handler/manager_loop.rs @@ -221,6 +221,7 @@ mod tests { execution_env: std::sync::Arc::new(arc_agent::LocalExecutionEnvironment::new( std::env::current_dir().unwrap_or_else(|_| std::path::PathBuf::from(".")), )), + git_state: std::sync::RwLock::new(None), } } diff --git a/crates/arc-workflows/src/handler/mod.rs b/crates/arc-workflows/src/handler/mod.rs index 94d79dd3c..7e9be67c3 100644 --- a/crates/arc-workflows/src/handler/mod.rs +++ b/crates/arc-workflows/src/handler/mod.rs @@ -17,6 +17,7 @@ use arc_agent::ExecutionEnvironment; use async_trait::async_trait; use crate::context::Context; +use crate::engine::GitState; use crate::error::ArcError; use crate::event::EventEmitter; use crate::graph::{shape_to_handler_type, Graph, Node}; @@ -28,6 +29,21 @@ pub struct EngineServices { pub registry: Arc, pub emitter: Arc, pub execution_env: Arc, + /// Git state for the current run. Set via `set_git_state` at the start of + /// `run_internal` and read by parallel/fan-in handlers. + pub(crate) git_state: std::sync::RwLock>>, +} + +impl EngineServices { + /// Read the current git state (if any). + pub fn git_state(&self) -> Option> { + self.git_state.read().unwrap().clone() + } + + /// Set the git state for the current run. + pub fn set_git_state(&self, state: Option>) { + *self.git_state.write().unwrap() = state; + } } /// The handler interface for node execution. diff --git a/crates/arc-workflows/src/handler/parallel.rs b/crates/arc-workflows/src/handler/parallel.rs index 94900d0f3..7198e4f21 100644 --- a/crates/arc-workflows/src/handler/parallel.rs +++ b/crates/arc-workflows/src/handler/parallel.rs @@ -1,11 +1,13 @@ -use std::path::Path; +use std::path::{Path, PathBuf}; use std::sync::Arc; use std::time::Instant; +use arc_agent::ExecutionEnvironment; use async_trait::async_trait; use tokio::sync::Semaphore; use crate::context::Context; +use crate::engine::GitCheckpointMode; use crate::error::ArcError; use crate::event::PipelineEvent; use crate::graph::{Graph, Node}; @@ -13,6 +15,85 @@ use crate::outcome::{Outcome, StageStatus}; use super::{EngineServices, Handler}; +// --------------------------------------------------------------------------- +// WorktreeEnv — decorates an ExecutionEnvironment with a custom working dir +// --------------------------------------------------------------------------- + +/// Wraps an existing `ExecutionEnvironment` so that all operations use a +/// different working directory (the worktree path inside a remote sandbox). +struct WorktreeEnv { + inner: Arc, + worktree_dir: String, +} + +#[async_trait] +impl ExecutionEnvironment for WorktreeEnv { + async fn read_file( + &self, + path: &str, + offset: Option, + limit: Option, + ) -> Result { + self.inner.read_file(path, offset, limit).await + } + async fn write_file(&self, path: &str, content: &str) -> Result<(), String> { + self.inner.write_file(path, content).await + } + async fn delete_file(&self, path: &str) -> Result<(), String> { + self.inner.delete_file(path).await + } + async fn file_exists(&self, path: &str) -> Result { + self.inner.file_exists(path).await + } + async fn list_directory( + &self, + path: &str, + depth: Option, + ) -> Result, String> { + self.inner.list_directory(path, depth).await + } + async fn exec_command( + &self, + command: &str, + timeout_ms: u64, + working_dir: Option<&str>, + env_vars: Option<&std::collections::HashMap>, + cancel_token: Option, + ) -> Result { + // Default to worktree dir when no explicit working_dir is given + let wd = working_dir.unwrap_or(&self.worktree_dir); + self.inner + .exec_command(command, timeout_ms, Some(wd), env_vars, cancel_token) + .await + } + async fn grep( + &self, + pattern: &str, + path: &str, + options: &arc_agent::execution_env::GrepOptions, + ) -> Result, String> { + self.inner.grep(pattern, path, options).await + } + async fn glob(&self, pattern: &str, path: Option<&str>) -> Result, String> { + self.inner.glob(pattern, path).await + } + async fn initialize(&self) -> Result<(), String> { + self.inner.initialize().await + } + async fn cleanup(&self) -> Result<(), String> { + self.inner.cleanup().await + } + fn working_directory(&self) -> &str { + &self.worktree_dir + } + fn platform(&self) -> &str { + self.inner.platform() + } + fn os_version(&self) -> String { + self.inner.os_version() + } +} + /// Convert a Duration's milliseconds to u64, saturating on overflow. fn millis_u64(d: std::time::Duration) -> u64 { u64::try_from(d.as_millis()).unwrap_or(u64::MAX) @@ -94,6 +175,8 @@ fn parse_error_policy(raw: &str) -> ErrorPolicy { struct BranchResult { id: String, outcome: Outcome, + head_sha: Option, + worktree_path: Option, } #[async_trait] @@ -138,18 +221,145 @@ impl Handler for ParallelHandler { let max_parallel = usize::try_from(max_parallel).unwrap_or(4).max(1); let semaphore = Arc::new(Semaphore::new(max_parallel)); + let git_state = services.git_state(); - // Build branch tasks - let mut handles = Vec::new(); + // --- Git isolation: checkpoint "parallel base" before fan-out --- + let base_sha: Option = if let Some(ref gs) = git_state { + match &gs.mode { + GitCheckpointMode::Host(work_dir) => { + let wd = work_dir.clone(); + let rid = gs.run_id.clone(); + let nid = node.id.clone(); + crate::engine::git_checkpoint_host(wd, rid, nid, "parallel_base".into(), 0, None) + .await + } + GitCheckpointMode::Remote(_) => { + crate::engine::git_checkpoint_remote( + &*services.execution_env, + &gs.run_id, + &node.id, + "parallel_base", + 0, + None, + ) + .await + } + } + } else { + None + }; + + // Build per-branch execution environments (sequentially for git setup) + struct BranchSetup { + target_id: String, + branch_index: usize, + branch_context: Context, + execution_env: Arc, + worktree_path: Option, + } + + let mut branch_setups: Vec = Vec::new(); for (branch_index, edge) in branches.iter().enumerate() { let target_id = edge.to.clone(); let branch_context = context.clone_context(); + + let (branch_exec_env, worktree_path): (Arc, Option) = + if let (Some(ref gs), Some(ref bsha)) = (&git_state, &base_sha) { + let branch_key = &target_id; + let branch_name = format!( + "arc/run/parallel/{}/{}/{}", + gs.run_id, node.id, branch_key + ); + + match &gs.mode { + GitCheckpointMode::Host(work_dir) => { + let wt_path = logs_root + .join("parallel") + .join(&node.id) + .join(branch_key) + .join("worktree"); + let wd = work_dir.clone(); + let bn = branch_name.clone(); + let bs = bsha.clone(); + let wtp = wt_path.clone(); + tokio::task::spawn_blocking(move || { + crate::git::create_branch_at(&wd, &bn, &bs)?; + crate::git::add_worktree(&wd, &wtp, &bn) + }) + .await + .map_err(|e| { + ArcError::Handler(format!("worktree setup join error: {e}")) + })??; + branch_context.set( + "internal.work_dir", + serde_json::json!(wt_path.to_string_lossy().as_ref()), + ); + let env: Arc = Arc::new( + arc_agent::LocalExecutionEnvironment::new(wt_path.clone()), + ); + (env, Some(wt_path)) + } + GitCheckpointMode::Remote(_) => { + let wt_path_str = format!( + "/home/daytona/workspace/.arc-parallel/{}/{}", + node.id, branch_key + ); + let ok = crate::engine::git_create_branch_at_remote( + &*services.execution_env, + &branch_name, + bsha, + ) + .await; + if !ok { + return Err(ArcError::Handler(format!( + "failed to create remote branch {branch_name}" + ))); + } + let ok = crate::engine::git_add_worktree_remote( + &*services.execution_env, + &wt_path_str, + &branch_name, + ) + .await; + if !ok { + return Err(ArcError::Handler(format!( + "failed to add remote worktree {wt_path_str}" + ))); + } + branch_context.set( + "internal.work_dir", + serde_json::json!(&wt_path_str), + ); + let env: Arc = Arc::new(WorktreeEnv { + inner: Arc::clone(&services.execution_env), + worktree_dir: wt_path_str.clone(), + }); + (env, Some(PathBuf::from(wt_path_str))) + } + } + } else { + (Arc::clone(&services.execution_env), None) + }; + + branch_setups.push(BranchSetup { + target_id, + branch_index, + branch_context, + execution_env: branch_exec_env, + worktree_path, + }); + } + + // --- Fan out: concurrent execution --- + let mut handles = Vec::new(); + for setup in branch_setups { let registry = Arc::clone(&services.registry); let emitter = Arc::clone(&services.emitter); - let execution_env = Arc::clone(&services.execution_env); let graph = graph.clone(); let logs_root = logs_root.to_path_buf(); let sem = Arc::clone(&semaphore); + let has_git = git_state.is_some(); + let run_id = git_state.as_ref().map(|gs| gs.run_id.clone()); let handle = tokio::spawn(async move { let _permit = sem @@ -158,52 +368,91 @@ impl Handler for ParallelHandler { .map_err(|e| ArcError::Handler(format!("semaphore error: {e}")))?; emitter.emit(&PipelineEvent::ParallelBranchStarted { - branch: target_id.clone(), - index: branch_index, + branch: setup.target_id.clone(), + index: setup.branch_index, }); let branch_start = Instant::now(); - let Some(target_node) = graph.nodes.get(&target_id) else { - let outcome = - Outcome::fail(format!("branch target node not found: {target_id}")); + let Some(target_node) = graph.nodes.get(&setup.target_id) else { + let outcome = Outcome::fail(format!( + "branch target node not found: {}", + setup.target_id + )); emitter.emit(&PipelineEvent::ParallelBranchCompleted { - branch: target_id.clone(), - index: branch_index, + branch: setup.target_id.clone(), + index: setup.branch_index, duration_ms: millis_u64(branch_start.elapsed()), status: "fail".to_string(), }); return Ok(BranchResult { - id: target_id.clone(), + id: setup.target_id.clone(), outcome, + head_sha: None, + worktree_path: setup.worktree_path, }); }; let branch_services = EngineServices { registry: Arc::clone(®istry), emitter: Arc::clone(&emitter), - execution_env: Arc::clone(&execution_env), + execution_env: Arc::clone(&setup.execution_env), + git_state: std::sync::RwLock::new(None), }; let handler = registry.resolve(target_node); let outcome = handler .execute( target_node, - &branch_context, + &setup.branch_context, &graph, &logs_root, &branch_services, ) .await?; + // Checkpoint commit after branch execution (capture head_sha) + let head_sha = if has_git { + let rid = run_id.as_deref().unwrap_or("unknown"); + let nid = &setup.target_id; + let status_str = outcome.status.to_string(); + // Use exec_command to commit and capture HEAD in the branch worktree + let add_result = setup + .execution_env + .exec_command("git add -A", 30_000, None, None, None) + .await; + if add_result.as_ref().is_ok_and(|r| r.exit_code == 0) { + let msg = format!("arc({rid}): {nid} ({status_str})"); + let commit_cmd = format!( + "git -c user.name=arc -c user.email=arc@local commit --allow-empty -m '{msg}'" + ); + let _ = setup + .execution_env + .exec_command(&commit_cmd, 30_000, None, None, None) + .await; + } + let sha_result = setup + .execution_env + .exec_command("git rev-parse HEAD", 10_000, None, None, None) + .await; + match sha_result { + Ok(r) if r.exit_code == 0 => Some(r.stdout.trim().to_string()), + _ => None, + } + } else { + None + }; + emitter.emit(&PipelineEvent::ParallelBranchCompleted { - branch: target_id.clone(), - index: branch_index, + branch: setup.target_id.clone(), + index: setup.branch_index, duration_ms: millis_u64(branch_start.elapsed()), status: outcome.status.to_string(), }); Ok::(BranchResult { - id: target_id, + id: setup.target_id, outcome, + head_sha, + worktree_path: setup.worktree_path, }) }); handles.push(handle); @@ -234,6 +483,8 @@ impl Handler for ParallelHandler { let result = BranchResult { id: String::new(), outcome: Outcome::fail(e.to_string()), + head_sha: None, + worktree_path: None, }; if error_policy == ErrorPolicy::FailFast { results.push(result); @@ -252,6 +503,8 @@ impl Handler for ParallelHandler { let result = BranchResult { id: String::new(), outcome: Outcome::fail(format!("task join error: {join_err}")), + head_sha: None, + worktree_path: None, }; if error_policy == ErrorPolicy::FailFast { results.push(result); @@ -269,6 +522,62 @@ impl Handler for ParallelHandler { } } + // --- Git isolation: clean up worktrees, then ff-merge winner --- + if let Some(ref gs) = git_state { + // Clean up worktrees first + for result in &results { + if let Some(ref wt_path) = result.worktree_path { + match &gs.mode { + GitCheckpointMode::Host(work_dir) => { + let wd = work_dir.clone(); + let wtp = wt_path.clone(); + let _ = tokio::task::spawn_blocking(move || { + crate::git::remove_worktree(&wd, &wtp) + }) + .await; + } + GitCheckpointMode::Remote(_) => { + let wt_str = wt_path.to_string_lossy().to_string(); + crate::engine::git_remove_worktree_remote( + &*services.execution_env, + &wt_str, + ) + .await; + } + } + } + } + + // Fast-forward main branch to first successful branch (lexically sorted). + // This must happen here — before the engine creates its own checkpoint commit + // on the main branch — so that subsequent commits are descendants of the winner. + let mut successful: Vec<_> = results + .iter() + .filter(|r| r.outcome.status == StageStatus::Success && r.head_sha.is_some()) + .collect(); + successful.sort_by(|a, b| a.id.cmp(&b.id)); + if let Some(winner) = successful.first() { + let sha = winner.head_sha.as_ref().unwrap(); + match &gs.mode { + GitCheckpointMode::Host(work_dir) => { + let wd = work_dir.clone(); + let s = sha.clone(); + let _ = tokio::task::spawn_blocking(move || { + crate::git::merge_ff_only(&wd, &s) + }) + .await; + } + GitCheckpointMode::Remote(_) => { + crate::engine::git_merge_ff_only_remote( + &*services.execution_env, + sha, + ) + .await; + } + } + } + } + // Count successes and failures let success_count = results .iter() @@ -284,10 +593,14 @@ impl Handler for ParallelHandler { let results_json: Vec = results .iter() .map(|r| { - serde_json::json!({ + let mut entry = serde_json::json!({ "id": r.id, "status": r.outcome.status.to_string(), - }) + }); + if let Some(ref sha) = r.head_sha { + entry["head_sha"] = serde_json::json!(sha); + } + entry }) .collect(); context.set("parallel.results", serde_json::json!(results_json)); @@ -341,7 +654,7 @@ impl Handler for ParallelHandler { } }; - // Build suggested_next_ids from successful branch targets + // Build suggested_next_ids from branch targets let branch_ids: Vec = results.iter().map(|r| r.id.clone()).collect(); let is_fail = status == StageStatus::Fail; @@ -386,6 +699,7 @@ mod tests { execution_env: Arc::new(arc_agent::LocalExecutionEnvironment::new( std::env::current_dir().unwrap_or_else(|_| std::path::PathBuf::from(".")), )), + git_state: std::sync::RwLock::new(None), } } diff --git a/crates/arc-workflows/src/handler/script.rs b/crates/arc-workflows/src/handler/script.rs index 731034ded..18a5d5a7e 100644 --- a/crates/arc-workflows/src/handler/script.rs +++ b/crates/arc-workflows/src/handler/script.rs @@ -180,6 +180,7 @@ mod tests { execution_env: std::sync::Arc::new(arc_agent::LocalExecutionEnvironment::new( std::env::current_dir().unwrap_or_else(|_| std::path::PathBuf::from(".")), )), + git_state: std::sync::RwLock::new(None), } } diff --git a/crates/arc-workflows/src/handler/start.rs b/crates/arc-workflows/src/handler/start.rs index 9c7c5aa7c..16180bfd8 100644 --- a/crates/arc-workflows/src/handler/start.rs +++ b/crates/arc-workflows/src/handler/start.rs @@ -39,6 +39,7 @@ mod tests { execution_env: std::sync::Arc::new(arc_agent::LocalExecutionEnvironment::new( std::env::current_dir().unwrap_or_else(|_| std::path::PathBuf::from(".")), )), + git_state: std::sync::RwLock::new(None), } } diff --git a/crates/arc-workflows/src/handler/sub_pipeline.rs b/crates/arc-workflows/src/handler/sub_pipeline.rs index d87d869ef..e5cd8c231 100644 --- a/crates/arc-workflows/src/handler/sub_pipeline.rs +++ b/crates/arc-workflows/src/handler/sub_pipeline.rs @@ -172,6 +172,7 @@ mod tests { registry: Arc::new(registry), emitter: Arc::new(EventEmitter::new()), execution_env: local_env(), + git_state: std::sync::RwLock::new(None), } } @@ -180,6 +181,7 @@ mod tests { registry: Arc::new(registry), emitter: Arc::new(EventEmitter::new()), execution_env: local_env(), + git_state: std::sync::RwLock::new(None), } } diff --git a/crates/arc-workflows/src/handler/wait_human.rs b/crates/arc-workflows/src/handler/wait_human.rs index 81c5b6d9c..8a8e54714 100644 --- a/crates/arc-workflows/src/handler/wait_human.rs +++ b/crates/arc-workflows/src/handler/wait_human.rs @@ -279,6 +279,7 @@ mod tests { execution_env: std::sync::Arc::new(arc_agent::LocalExecutionEnvironment::new( std::env::current_dir().unwrap_or_else(|_| std::path::PathBuf::from(".")), )), + git_state: std::sync::RwLock::new(None), } } diff --git a/crates/arc-workflows/tests/daytona_integration.rs b/crates/arc-workflows/tests/daytona_integration.rs index 4f95a4578..e3956376c 100644 --- a/crates/arc-workflows/tests/daytona_integration.rs +++ b/crates/arc-workflows/tests/daytona_integration.rs @@ -537,6 +537,228 @@ async fn daytona_git_checkpoint_remote_emits_events() { env.cleanup().await.unwrap(); } +// --------------------------------------------------------------------------- +// Parallel git branching on Daytona (Remote mode) +// --------------------------------------------------------------------------- + +use arc_workflows::handler::fan_in::FanInHandler; +use arc_workflows::handler::parallel::ParallelHandler; + +/// End-to-end: parallel branches get isolated worktrees in Daytona sandbox, +/// fan-in fast-forwards to winner. +#[tokio::test] +#[ignore] +async fn daytona_parallel_git_branching_e2e() { + let env = create_env().await; + env.initialize().await.unwrap(); + let env: Arc = Arc::new(env); + + // Install git if not available + let git_check = env + .exec_command("git --version", 10_000, None, None, None) + .await; + if git_check.as_ref().map_or(true, |r| r.exit_code != 0) { + let install = env + .exec_command( + "apt-get update -qq && apt-get install -y -qq git >/dev/null 2>&1", + 120_000, + None, + None, + None, + ) + .await + .expect("apt-get install git should not error"); + assert_eq!( + install.exit_code, 0, + "git install failed: {}", + install.stderr + ); + } + + // Set up git in the sandbox (uses existing repo from Daytona project clone) + let (run_id, base_sha, branch_name) = setup_daytona_git(&*env).await; + + // Pipeline: start -> fan_out -> {branch_a, branch_b} -> fan_in -> exit + let mut graph = Graph::new("DaytonaParallelGitBranching"); + graph.attrs.insert( + "goal".to_string(), + AttrValue::String("Test parallel git branching on Daytona".to_string()), + ); + + let mut start = Node::new("start"); + start + .attrs + .insert("shape".to_string(), AttrValue::String("Mdiamond".to_string())); + graph.nodes.insert("start".to_string(), start); + + let mut fan_out = Node::new("fan_out"); + fan_out + .attrs + .insert("shape".to_string(), AttrValue::String("component".to_string())); + graph.nodes.insert("fan_out".to_string(), fan_out); + + let branch_a = Node::new("branch_a"); + graph.nodes.insert("branch_a".to_string(), branch_a); + + let branch_b = Node::new("branch_b"); + graph.nodes.insert("branch_b".to_string(), branch_b); + + let mut fan_in = Node::new("fan_in"); + fan_in.attrs.insert( + "shape".to_string(), + AttrValue::String("tripleoctagon".to_string()), + ); + graph.nodes.insert("fan_in".to_string(), fan_in); + + let mut exit_node = Node::new("exit"); + exit_node + .attrs + .insert("shape".to_string(), AttrValue::String("Msquare".to_string())); + graph.nodes.insert("exit".to_string(), exit_node); + + graph.edges.push(Edge::new("start", "fan_out")); + graph.edges.push(Edge::new("fan_out", "branch_a")); + graph.edges.push(Edge::new("fan_out", "branch_b")); + graph.edges.push(Edge::new("branch_a", "fan_in")); + graph.edges.push(Edge::new("branch_b", "fan_in")); + graph.edges.push(Edge::new("fan_in", "exit")); + + let logs_dir = tempfile::tempdir().unwrap(); + let mut emitter = EventEmitter::new(); + let events = Arc::new(std::sync::Mutex::new(Vec::new())); + { + let events_clone = Arc::clone(&events); + emitter.on_event(move |event| { + events_clone.lock().unwrap().push(event.clone()); + }); + } + + let mut registry = HandlerRegistry::new(Box::new(FileWriterHandler)); + registry.register("start", Box::new(StartHandler)); + registry.register("exit", Box::new(ExitHandler)); + registry.register("parallel", Box::new(ParallelHandler)); + registry.register("parallel.fan_in", Box::new(FanInHandler::new(None))); + + let engine = PipelineEngine::new(registry, Arc::new(emitter), Arc::clone(&env)); + + let config = RunConfig { + logs_root: logs_dir.path().to_path_buf(), + cancel_token: None, + dry_run: false, + run_id: run_id.clone(), + git_checkpoint: Some(GitCheckpointMode::Remote(logs_dir.path().to_path_buf())), + base_sha: Some(base_sha), + run_branch: Some(branch_name), + meta_branch: None, + }; + + let outcome = engine + .run(&graph, &config) + .await + .expect("daytona parallel pipeline should succeed"); + assert_eq!( + outcome.status, + StageStatus::Success, + "pipeline failed: {:?}", + outcome.failure_reason + ); + + // Verify parallel.results has head_sha for each branch + let checkpoint = Checkpoint::load(&logs_dir.path().join("checkpoint.json")) + .expect("checkpoint should load"); + let parallel_results = checkpoint + .context_values + .get("parallel.results") + .expect("parallel.results should be in context"); + let results_arr = parallel_results.as_array().expect("should be an array"); + assert_eq!(results_arr.len(), 2, "should have 2 branch results"); + + // Both branches should have head_sha (40-char hex) + let has_sha = results_arr.iter().all(|v| { + v.get("head_sha") + .and_then(|v| v.as_str()) + .is_some_and(|s| s.len() == 40 && s.chars().all(|c| c.is_ascii_hexdigit())) + }); + assert!(has_sha, "all branches should have 40-char hex head_sha"); + + // Branch SHAs should differ (each branch made unique changes) + let sha_a = results_arr + .iter() + .find(|v| v.get("id").and_then(|v| v.as_str()) == Some("branch_a")) + .and_then(|v| v.get("head_sha").and_then(|v| v.as_str())) + .unwrap(); + let sha_b = results_arr + .iter() + .find(|v| v.get("id").and_then(|v| v.as_str()) == Some("branch_b")) + .and_then(|v| v.get("head_sha").and_then(|v| v.as_str())) + .unwrap(); + assert_ne!(sha_a, sha_b, "branch SHAs should differ"); + + // Verify fan_in selected a winner and set best_head_sha + let best_id = checkpoint + .context_values + .get("parallel.fan_in.best_id") + .and_then(|v| v.as_str().map(String::from)) + .expect("fan_in should have selected a best_id"); + assert_eq!(best_id, "branch_a", "heuristic should pick branch_a (lexical)"); + + let best_head_sha = checkpoint + .context_values + .get("parallel.fan_in.best_head_sha") + .and_then(|v| v.as_str().map(String::from)); + assert!( + best_head_sha.is_some(), + "fan_in should have set best_head_sha" + ); + + // Verify winner's file exists in sandbox + let winner_check = env + .exec_command("cat branch_a.txt", 10_000, None, None, None) + .await + .expect("cat should succeed"); + assert_eq!(winner_check.exit_code, 0, "winner's file should exist"); + assert!( + winner_check.stdout.contains("branch_a"), + "winner's file should have correct content, got: {}", + winner_check.stdout + ); + + // Verify events + { + let events = events.lock().unwrap(); + let parallel_started: Vec<_> = events + .iter() + .filter(|e| { + matches!( + e, + arc_workflows::event::PipelineEvent::ParallelStarted { .. } + ) + }) + .collect(); + assert_eq!( + parallel_started.len(), + 1, + "should have exactly one ParallelStarted event" + ); + let parallel_completed: Vec<_> = events + .iter() + .filter(|e| { + matches!( + e, + arc_workflows::event::PipelineEvent::ParallelCompleted { .. } + ) + }) + .collect(); + assert_eq!( + parallel_completed.len(), + 1, + "should have exactly one ParallelCompleted event" + ); + } + + env.cleanup().await.expect("Daytona cleanup should succeed"); +} + // --------------------------------------------------------------------------- // CLI Backend on Daytona — real CLI tools via exec_command // --------------------------------------------------------------------------- diff --git a/crates/arc-workflows/tests/integration.rs b/crates/arc-workflows/tests/integration.rs index 262f80867..f471cd886 100644 --- a/crates/arc-workflows/tests/integration.rs +++ b/crates/arc-workflows/tests/integration.rs @@ -8743,6 +8743,33 @@ fn parse_real_gemini_json() { // --------------------------------------------------------------------------- use arc_workflows::engine::GitCheckpointMode; +use arc_workflows::handler::fan_in::FanInHandler; +use arc_workflows::handler::parallel::ParallelHandler; + +/// A handler that writes a file named `{node_id}.txt` into the execution environment's +/// working directory. Used to verify git worktree isolation in parallel branches. +struct FileWriterHandler; + +#[async_trait::async_trait] +impl Handler for FileWriterHandler { + async fn execute( + &self, + node: &Node, + _context: &Context, + _graph: &Graph, + _logs_root: &Path, + services: &arc_workflows::handler::EngineServices, + ) -> Result { + let work_dir = services.execution_env.working_directory().to_string(); + let file_path = format!("{}/{}.txt", work_dir, node.id); + services + .execution_env + .write_file(&file_path, &format!("written by {}", node.id)) + .await + .map_err(|e| ArcError::Handler(format!("write_file failed: {e}")))?; + Ok(Outcome::success()) + } +} /// End-to-end test: pipeline with `GitCheckpointMode::Host` emits `GitCheckpoint` /// events with valid commit SHAs and writes `diff.patch` per stage. @@ -9078,3 +9105,300 @@ async fn git_checkpoint_host_writes_shadow_branch() { .current_dir(repo.path()) .output(); } + +// --------------------------------------------------------------------------- +// Host e2e: parallel git branching with worktree isolation +// --------------------------------------------------------------------------- + +/// End-to-end: parallel branches get isolated worktrees, fan-in fast-forwards to winner. +/// +/// Pipeline: start -> fan_out -> {branch_a, branch_b} -> fan_in -> exit +/// +/// Each branch writes a unique file. After fan-in, only the winner's file should +/// be present in the main worktree. +#[tokio::test] +async fn parallel_git_branching_host_e2e() { + // 1. Create a temporary git repo with an initial commit + let repo = tempfile::tempdir().unwrap(); + std::process::Command::new("git") + .args(["init"]) + .current_dir(repo.path()) + .output() + .unwrap(); + std::process::Command::new("git") + .args([ + "-c", + "user.name=test", + "-c", + "user.email=test@test", + "commit", + "--allow-empty", + "-m", + "init", + ]) + .current_dir(repo.path()) + .output() + .unwrap(); + + // 2. Set up run branch and worktree (same as cli/run.rs) + let base_sha = { + let out = std::process::Command::new("git") + .args(["rev-parse", "HEAD"]) + .current_dir(repo.path()) + .output() + .unwrap(); + String::from_utf8_lossy(&out.stdout).trim().to_string() + }; + let run_id = "par-git-test"; + let run_branch = format!("arc/run/{run_id}"); + std::process::Command::new("git") + .args(["branch", &run_branch, "HEAD"]) + .current_dir(repo.path()) + .output() + .unwrap(); + let worktree_path = repo.path().join("worktree"); + std::process::Command::new("git") + .args(["worktree", "add"]) + .arg(&worktree_path) + .arg(&run_branch) + .current_dir(repo.path()) + .output() + .unwrap(); + + // 3. Build pipeline: start -> fan_out -> {branch_a, branch_b} -> fan_in -> exit + let mut graph = Graph::new("ParallelGitBranching"); + graph.attrs.insert( + "goal".to_string(), + AttrValue::String("Test parallel git branching".to_string()), + ); + + let mut start = Node::new("start"); + start + .attrs + .insert("shape".to_string(), AttrValue::String("Mdiamond".to_string())); + graph.nodes.insert("start".to_string(), start); + + let mut fan_out = Node::new("fan_out"); + fan_out + .attrs + .insert("shape".to_string(), AttrValue::String("component".to_string())); + graph.nodes.insert("fan_out".to_string(), fan_out); + + let branch_a = Node::new("branch_a"); + graph.nodes.insert("branch_a".to_string(), branch_a); + + let branch_b = Node::new("branch_b"); + graph.nodes.insert("branch_b".to_string(), branch_b); + + let mut fan_in = Node::new("fan_in"); + fan_in.attrs.insert( + "shape".to_string(), + AttrValue::String("tripleoctagon".to_string()), + ); + graph.nodes.insert("fan_in".to_string(), fan_in); + + let mut exit = Node::new("exit"); + exit.attrs + .insert("shape".to_string(), AttrValue::String("Msquare".to_string())); + graph.nodes.insert("exit".to_string(), exit); + + graph.edges.push(Edge::new("start", "fan_out")); + graph.edges.push(Edge::new("fan_out", "branch_a")); + graph.edges.push(Edge::new("fan_out", "branch_b")); + graph.edges.push(Edge::new("branch_a", "fan_in")); + graph.edges.push(Edge::new("branch_b", "fan_in")); + graph.edges.push(Edge::new("fan_in", "exit")); + + // 4. Set up engine with FileWriterHandler for branches + let logs_dir = tempfile::tempdir().unwrap(); + let mut emitter = EventEmitter::new(); + let events = collect_events(&mut emitter); + + let env: Arc = + Arc::new(arc_agent::LocalExecutionEnvironment::new(worktree_path.clone())); + + let mut registry = HandlerRegistry::new(Box::new(FileWriterHandler)); + registry.register("start", Box::new(StartHandler)); + registry.register("exit", Box::new(ExitHandler)); + registry.register("parallel", Box::new(ParallelHandler)); + registry.register( + "parallel.fan_in", + Box::new(FanInHandler::new(None)), // heuristic select — picks branch_a (lexical tiebreak) + ); + + let engine = PipelineEngine::new(registry, Arc::new(emitter), env); + + let config = RunConfig { + logs_root: logs_dir.path().to_path_buf(), + cancel_token: None, + dry_run: false, + run_id: run_id.into(), + git_checkpoint: Some(GitCheckpointMode::Host(worktree_path.clone())), + base_sha: Some(base_sha.clone()), + run_branch: Some(run_branch.clone()), + meta_branch: None, + }; + + // 5. Run pipeline + let outcome = engine + .run(&graph, &config) + .await + .expect("parallel pipeline should succeed"); + assert_eq!( + outcome.status, + StageStatus::Success, + "pipeline failed: {:?}", + outcome.failure_reason + ); + + // 6. Verify parallel.results has head_sha for each branch + let checkpoint = Checkpoint::load(&logs_dir.path().join("checkpoint.json")) + .expect("checkpoint should load"); + let parallel_results = checkpoint + .context_values + .get("parallel.results") + .expect("parallel.results should be in context"); + let results_arr = parallel_results.as_array().expect("should be an array"); + assert_eq!(results_arr.len(), 2, "should have 2 branch results"); + + // Both branches should have head_sha + let branch_a_result = results_arr + .iter() + .find(|v| v.get("id").and_then(|v| v.as_str()) == Some("branch_a")) + .expect("branch_a result should exist"); + let branch_b_result = results_arr + .iter() + .find(|v| v.get("id").and_then(|v| v.as_str()) == Some("branch_b")) + .expect("branch_b result should exist"); + + let sha_a = branch_a_result + .get("head_sha") + .and_then(|v| v.as_str()) + .expect("branch_a should have head_sha"); + let sha_b = branch_b_result + .get("head_sha") + .and_then(|v| v.as_str()) + .expect("branch_b should have head_sha"); + + assert_eq!(sha_a.len(), 40, "SHA should be 40 hex chars"); + assert_eq!(sha_b.len(), 40, "SHA should be 40 hex chars"); + assert_ne!(sha_a, sha_b, "branch SHAs should differ"); + + // 7. Verify fan_in selected a winner and set best_head_sha + let best_id = checkpoint + .context_values + .get("parallel.fan_in.best_id") + .and_then(|v| v.as_str().map(String::from)) + .expect("fan_in should have selected a best_id"); + let best_head_sha = checkpoint + .context_values + .get("parallel.fan_in.best_head_sha") + .and_then(|v| v.as_str().map(String::from)) + .expect("fan_in should have set best_head_sha"); + + // Heuristic select with both success: lexical tiebreak picks "branch_a" + assert_eq!(best_id, "branch_a", "heuristic should pick branch_a (lexical)"); + + // 8. Verify winner's file is in the main worktree, loser's is NOT + let winner_file = worktree_path.join(format!("{best_id}.txt")); + assert!( + winner_file.exists(), + "winner's file ({best_id}.txt) should exist in main worktree after ff-merge" + ); + let winner_content = std::fs::read_to_string(&winner_file).unwrap(); + assert!( + winner_content.contains(&format!("written by {best_id}")), + "winner's file should have correct content" + ); + + let loser_id = if best_id == "branch_a" { + "branch_b" + } else { + "branch_a" + }; + let loser_file = worktree_path.join(format!("{loser_id}.txt")); + assert!( + !loser_file.exists(), + "loser's file ({loser_id}.txt) should NOT exist in main worktree" + ); + + // 9. Verify the main worktree HEAD matches the winner's head_sha + let main_head = { + let out = std::process::Command::new("git") + .args(["rev-parse", "HEAD"]) + .current_dir(&worktree_path) + .output() + .unwrap(); + String::from_utf8_lossy(&out.stdout).trim().to_string() + }; + // After fan-in ff-only + engine's own checkpoint commits, HEAD should be a + // descendant of best_head_sha. + let is_ancestor = std::process::Command::new("git") + .args(["merge-base", "--is-ancestor", &best_head_sha, &main_head]) + .current_dir(&worktree_path) + .output() + .unwrap(); + assert!( + is_ancestor.status.success(), + "best_head_sha ({best_head_sha}) should be an ancestor of current HEAD ({main_head})" + ); + + // 10. Verify parallel branch refs still exist (for debugging) + let branch_ref_a = format!("arc/run/parallel/{run_id}/fan_out/branch_a"); + let ref_check = std::process::Command::new("git") + .args(["rev-parse", "--verify", &branch_ref_a]) + .current_dir(repo.path()) + .output() + .unwrap(); + assert!( + ref_check.status.success(), + "parallel branch ref should still exist for debugging" + ); + + // 11. Verify final.patch contains the winner's changes + let final_patch = logs_dir.path().join("final.patch"); + assert!( + final_patch.exists(), + "final.patch should exist in logs_root" + ); + let patch_content = std::fs::read_to_string(&final_patch).unwrap(); + assert!( + patch_content.contains(&format!("{best_id}.txt")), + "final.patch should contain winner's file" + ); + assert!( + !patch_content.contains(&format!("{loser_id}.txt")), + "final.patch should NOT contain loser's file" + ); + + // 12. Verify events + let events = events.lock().unwrap(); + let parallel_started: Vec<_> = events + .iter() + .filter(|e| matches!(e, PipelineEvent::ParallelStarted { .. })) + .collect(); + assert_eq!( + parallel_started.len(), + 1, + "should have exactly one ParallelStarted event" + ); + + let parallel_completed: Vec<_> = events + .iter() + .filter(|e| matches!(e, PipelineEvent::ParallelCompleted { .. })) + .collect(); + assert_eq!( + parallel_completed.len(), + 1, + "should have exactly one ParallelCompleted event" + ); + + // Cleanup + let _ = std::process::Command::new("git") + .args(["worktree", "remove", "--force"]) + .arg(&worktree_path) + .current_dir(repo.path()) + .output(); +} + +// Daytona parallel git branching test is in daytona_integration.rs