Add parallel git branching with per-branch worktree isolation

Parallel branches now get isolated git worktrees so concurrent file
writes don't collide. Works across Local, Docker (bind-mount), and
Daytona (remote exec_command) environments.

Key changes:
- git.rs: add create_branch_at() and merge_ff_only() helpers
- engine.rs: add GitState struct, remote worktree helpers
  (git_create_branch_at_remote, git_add_worktree_remote, etc.)
- handler/mod.rs: add git_state field to EngineServices (RwLock)
- handler/parallel.rs: WorktreeEnv wrapper, per-branch worktree
  setup/teardown, checkpoint commits per branch, ff-merge winner
  before returning to engine
- handler/fan_in.rs: ff-merge to winner's HEAD, set best_head_sha
- E2E tests for Host (local) and Daytona (remote) modes

When git_state is None, behavior is unchanged.

Co-Authored-By: Claude Opus 4.6 (1M context) <noreply@anthropic.com>
This commit is contained in:
Bryan Helmkamp 2026-03-01 14:24:00 -05:00
parent c01dae2077
commit 4dd2980a56
15 changed files with 1153 additions and 22 deletions

View file

@ -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<String>,
pub meta_branch: Option<String>,
}
/// 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<String> {
}
/// 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<String> {
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(),

View file

@ -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();

View file

@ -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),
}
}

View file

@ -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),
}
}

View file

@ -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),
}
}

View file

@ -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),
}
}

View file

@ -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),
}
}

View file

@ -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<HandlerRegistry>,
pub emitter: Arc<EventEmitter>,
pub execution_env: Arc<dyn ExecutionEnvironment>,
/// 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<Option<Arc<GitState>>>,
}
impl EngineServices {
/// Read the current git state (if any).
pub fn git_state(&self) -> Option<Arc<GitState>> {
self.git_state.read().unwrap().clone()
}
/// Set the git state for the current run.
pub fn set_git_state(&self, state: Option<Arc<GitState>>) {
*self.git_state.write().unwrap() = state;
}
}
/// The handler interface for node execution.

View file

@ -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<dyn ExecutionEnvironment>,
worktree_dir: String,
}
#[async_trait]
impl ExecutionEnvironment for WorktreeEnv {
async fn read_file(
&self,
path: &str,
offset: Option<usize>,
limit: Option<usize>,
) -> Result<String, String> {
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<bool, String> {
self.inner.file_exists(path).await
}
async fn list_directory(
&self,
path: &str,
depth: Option<usize>,
) -> Result<Vec<arc_agent::execution_env::DirEntry>, 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<String, String>>,
cancel_token: Option<tokio_util::sync::CancellationToken>,
) -> Result<arc_agent::execution_env::ExecResult, String> {
// 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<Vec<String>, String> {
self.inner.grep(pattern, path, options).await
}
async fn glob(&self, pattern: &str, path: Option<&str>) -> Result<Vec<String>, 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<String>,
worktree_path: Option<PathBuf>,
}
#[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<String> = 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<dyn ExecutionEnvironment>,
worktree_path: Option<PathBuf>,
}
let mut branch_setups: Vec<BranchSetup> = 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<dyn ExecutionEnvironment>, Option<PathBuf>) =
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<dyn ExecutionEnvironment> = 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<dyn ExecutionEnvironment> = 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(&registry),
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, ArcError>(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<serde_json::Value> = 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<String> = 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),
}
}

View file

@ -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),
}
}

View file

@ -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),
}
}

View file

@ -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),
}
}

View file

@ -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),
}
}

View file

@ -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<dyn ExecutionEnvironment> = 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
// ---------------------------------------------------------------------------

View file

@ -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<Outcome, ArcError> {
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<dyn arc_agent::ExecutionEnvironment> =
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