From ab53e87915b155280b5377083eee885955184c92 Mon Sep 17 00:00:00 2001 From: Bryan Helmkamp Date: Sun, 1 Mar 2026 17:15:51 -0500 Subject: [PATCH] Harden git checkpoint and worktree operations for idempotent retry - Add git_cmd() helper disabling maintenance.auto and gc.auto on all host-side git commands; add GIT_REMOTE constant for remote commands - Use --force on branch creation for idempotent retry/resume - Add replace_worktree() that does best-effort remove before add - Add reset_hard() after parallel worktree setup for deterministic state - Add sanitize_ref_component() to clean node IDs in branch names - Add git_replace_worktree_remote() for remote sandbox environments - Add tests for sanitize_ref_component, replace_worktree, reset_hard Co-Authored-By: Claude Opus 4.6 (1M context) --- crates/arc-workflows/src/cli/run.rs | 4 +- crates/arc-workflows/src/engine.rs | 38 ++-- crates/arc-workflows/src/git.rs | 172 ++++++++++++++++--- crates/arc-workflows/src/handler/parallel.rs | 32 +++- crates/arc-workflows/tests/integration.rs | 2 +- 5 files changed, 206 insertions(+), 42 deletions(-) diff --git a/crates/arc-workflows/src/cli/run.rs b/crates/arc-workflows/src/cli/run.rs index 3baeb6e75..849a08932 100644 --- a/crates/arc-workflows/src/cli/run.rs +++ b/crates/arc-workflows/src/cli/run.rs @@ -750,7 +750,7 @@ fn setup_worktree( crate::git::create_branch(original_cwd, &branch_name).map_err(|e| anyhow::anyhow!("{e}"))?; let worktree_path = logs_dir.join("worktree"); - crate::git::add_worktree(original_cwd, &worktree_path, &branch_name) + crate::git::replace_worktree(original_cwd, &worktree_path, &branch_name) .map_err(|e| anyhow::anyhow!("{e}"))?; std::env::set_current_dir(&worktree_path)?; @@ -876,7 +876,7 @@ async fn run_from_branch( // Re-attach worktree to the existing run branch let worktree_path = logs_dir.join("worktree"); - crate::git::add_worktree(&original_cwd, &worktree_path, run_branch) + crate::git::replace_worktree(&original_cwd, &worktree_path, run_branch) .map_err(|e| anyhow::anyhow!("failed to attach worktree to {run_branch}: {e}"))?; std::env::set_current_dir(&worktree_path)?; diff --git a/crates/arc-workflows/src/engine.rs b/crates/arc-workflows/src/engine.rs index 8d46acc11..38e005969 100644 --- a/crates/arc-workflows/src/engine.rs +++ b/crates/arc-workflows/src/engine.rs @@ -555,6 +555,8 @@ async fn git_diff_host(work_dir: PathBuf, base: String) -> Option { } } +pub const GIT_REMOTE: &str = "git -c maintenance.auto=0 -c gc.auto=0"; + /// Run a git checkpoint commit inside a remote execution environment. pub async fn git_checkpoint_remote( exec_env: &dyn ExecutionEnvironment, @@ -565,8 +567,9 @@ pub async fn git_checkpoint_remote( shadow_sha: Option, ) -> Option { // Stage everything + let add_cmd = format!("{GIT_REMOTE} add -A"); let add_result = exec_env - .exec_command("git add -A", 30_000, None, None, None) + .exec_command(&add_cmd, 30_000, None, None, None) .await; if add_result.as_ref().map_or(true, |r| r.exit_code != 0) { return None; @@ -604,18 +607,20 @@ pub async fn git_checkpoint_remote( } // Commit with arc identity using the message file - let commit_cmd = - "git -c user.name=arc -c user.email=arc@local commit --allow-empty -F /tmp/arc-commit-msg"; + let commit_cmd = format!( + "{GIT_REMOTE} -c user.name=arc -c user.email=arc@local commit --allow-empty -F /tmp/arc-commit-msg" + ); let commit_result = exec_env - .exec_command(commit_cmd, 30_000, None, None, None) + .exec_command(&commit_cmd, 30_000, None, None, None) .await; if commit_result.as_ref().map_or(true, |r| r.exit_code != 0) { return None; } // Get the new HEAD SHA + let sha_cmd = format!("{GIT_REMOTE} rev-parse HEAD"); let sha_result = exec_env - .exec_command("git rev-parse HEAD", 10_000, None, None, None) + .exec_command(&sha_cmd, 10_000, None, None, None) .await; match sha_result { Ok(r) if r.exit_code == 0 => Some(r.stdout.trim().to_string()), @@ -625,7 +630,7 @@ pub async fn git_checkpoint_remote( /// Run a git diff inside a remote execution environment. async fn git_diff_remote(exec_env: &dyn ExecutionEnvironment, base: &str) -> Option { - let cmd = format!("git diff {base} HEAD"); + let cmd = format!("{GIT_REMOTE} diff {base} HEAD"); match exec_env.exec_command(&cmd, 30_000, None, None, None).await { Ok(r) if r.exit_code == 0 => Some(r.stdout), _ => None, @@ -640,7 +645,7 @@ pub async fn git_create_branch_at_remote( name: &str, sha: &str, ) -> bool { - let cmd = format!("git branch {name} {sha}"); + let cmd = format!("{GIT_REMOTE} branch --force {name} {sha}"); matches!( exec_env.exec_command(&cmd, 30_000, None, None, None).await, Ok(r) if r.exit_code == 0 @@ -653,7 +658,7 @@ pub async fn git_add_worktree_remote( path: &str, branch: &str, ) -> bool { - let cmd = format!("git worktree add {path} {branch}"); + let cmd = format!("{GIT_REMOTE} worktree add {path} {branch}"); matches!( exec_env.exec_command(&cmd, 30_000, None, None, None).await, Ok(r) if r.exit_code == 0 @@ -665,7 +670,7 @@ pub async fn git_remove_worktree_remote( exec_env: &dyn ExecutionEnvironment, path: &str, ) -> bool { - let cmd = format!("git worktree remove --force {path}"); + let cmd = format!("{GIT_REMOTE} worktree remove --force {path}"); matches!( exec_env.exec_command(&cmd, 30_000, None, None, None).await, Ok(r) if r.exit_code == 0 @@ -677,7 +682,7 @@ pub async fn git_merge_ff_only_remote( exec_env: &dyn ExecutionEnvironment, sha: &str, ) -> bool { - let cmd = format!("git merge --ff-only {sha}"); + let cmd = format!("{GIT_REMOTE} merge --ff-only {sha}"); matches!( exec_env.exec_command(&cmd, 30_000, None, None, None).await, Ok(r) if r.exit_code == 0 @@ -686,8 +691,9 @@ pub async fn git_merge_ff_only_remote( /// Get the current HEAD SHA from a remote execution environment. pub async fn git_head_sha_remote(exec_env: &dyn ExecutionEnvironment) -> Option { + let cmd = format!("{GIT_REMOTE} rev-parse HEAD"); match exec_env - .exec_command("git rev-parse HEAD", 10_000, None, None, None) + .exec_command(&cmd, 10_000, None, None, None) .await { Ok(r) if r.exit_code == 0 => Some(r.stdout.trim().to_string()), @@ -695,6 +701,16 @@ pub async fn git_head_sha_remote(exec_env: &dyn ExecutionEnvironment) -> Option< } } +/// Remove any stale worktree at `path` (best-effort), then add a fresh one. +pub async fn git_replace_worktree_remote( + exec_env: &dyn ExecutionEnvironment, + path: &str, + branch: &str, +) -> bool { + let _ = git_remove_worktree_remote(exec_env, path).await; + git_add_worktree_remote(exec_env, path, branch).await +} + /// Configuration for a pipeline run. pub struct RunConfig { pub logs_root: PathBuf, diff --git a/crates/arc-workflows/src/git.rs b/crates/arc-workflows/src/git.rs index c77a5330e..d03408983 100644 --- a/crates/arc-workflows/src/git.rs +++ b/crates/arc-workflows/src/git.rs @@ -13,11 +13,18 @@ fn git_error(msg: impl Into) -> ArcError { ArcError::Engine(msg.into()) } +/// Return a pre-configured `git` command with auto-maintenance disabled. +fn git_cmd(dir: &Path) -> Command { + let mut cmd = Command::new("git"); + cmd.args(["-c", "maintenance.auto=0", "-c", "gc.auto=0"]) + .current_dir(dir); + cmd +} + /// Assert the working directory is a clean git repo (no uncommitted changes). pub fn ensure_clean(repo: &Path) -> Result<()> { - let output = Command::new("git") + let output = git_cmd(repo) .args(["status", "--porcelain"]) - .current_dir(repo) .output() .map_err(|e| git_error(format!("git status failed: {e}")))?; @@ -35,9 +42,8 @@ pub fn ensure_clean(repo: &Path) -> Result<()> { /// Return the SHA of HEAD. pub fn head_sha(repo: &Path) -> Result { - let output = Command::new("git") + let output = git_cmd(repo) .args(["rev-parse", "HEAD"]) - .current_dir(repo) .output() .map_err(|e| git_error(format!("git rev-parse failed: {e}")))?; @@ -50,9 +56,8 @@ pub fn head_sha(repo: &Path) -> Result { /// Create a new branch at HEAD without checking it out. pub fn create_branch(repo: &Path, name: &str) -> Result<()> { - let output = Command::new("git") - .args(["branch", name, "HEAD"]) - .current_dir(repo) + let output = git_cmd(repo) + .args(["branch", "--force", name, "HEAD"]) .output() .map_err(|e| git_error(format!("git branch failed: {e}")))?; @@ -66,11 +71,10 @@ pub fn create_branch(repo: &Path, name: &str) -> Result<()> { /// Add a git worktree for the given branch at `path`. pub fn add_worktree(repo: &Path, path: &Path, branch: &str) -> Result<()> { - let output = Command::new("git") + let output = git_cmd(repo) .args(["worktree", "add"]) .arg(path) .arg(branch) - .current_dir(repo) .output() .map_err(|e| git_error(format!("git worktree add failed: {e}")))?; @@ -84,10 +88,9 @@ pub fn add_worktree(repo: &Path, path: &Path, branch: &str) -> Result<()> { /// Remove a git worktree. pub fn remove_worktree(repo: &Path, path: &Path) -> Result<()> { - let output = Command::new("git") + let output = git_cmd(repo) .args(["worktree", "remove", "--force"]) .arg(path) - .current_dir(repo) .output() .map_err(|e| git_error(format!("git worktree remove failed: {e}")))?; @@ -101,9 +104,8 @@ pub fn remove_worktree(repo: &Path, path: &Path) -> Result<()> { /// 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) + let output = git_cmd(repo) + .args(["branch", "--force", name, sha]) .output() .map_err(|e| git_error(format!("git branch failed: {e}")))?; @@ -118,9 +120,8 @@ pub fn create_branch_at(repo: &Path, name: &str, sha: &str) -> Result<()> { /// 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") + let output = git_cmd(work_dir) .args(["merge", "--ff-only", sha]) - .current_dir(work_dir) .output() .map_err(|e| git_error(format!("git merge --ff-only failed: {e}")))?; @@ -144,9 +145,8 @@ pub fn checkpoint_commit( shadow_sha: Option<&str>, ) -> Result { // Stage everything - let output = Command::new("git") + let output = git_cmd(work_dir) .args(["add", "-A"]) - .current_dir(work_dir) .output() .map_err(|e| git_error(format!("git add failed: {e}")))?; @@ -177,7 +177,7 @@ pub fn checkpoint_commit( let message = trailerlink::format_message(&subject, "", &trailers); // Commit with arc identity (works even if user.name/email not configured) - let output = Command::new("git") + let output = git_cmd(work_dir) .args([ "-c", "user.name=arc", @@ -188,7 +188,6 @@ pub fn checkpoint_commit( "-m", &message, ]) - .current_dir(work_dir) .output() .map_err(|e| git_error(format!("git commit failed: {e}")))?; @@ -203,9 +202,8 @@ pub fn checkpoint_commit( /// Compute the diff between a base commit and HEAD. /// Returns the patch text (may be empty if no changes). pub fn diff_against(work_dir: &Path, base: &str) -> Result { - let output = Command::new("git") + let output = git_cmd(work_dir) .args(["diff", base, "HEAD"]) - .current_dir(work_dir) .output() .map_err(|e| git_error(format!("git diff failed: {e}")))?; @@ -217,6 +215,44 @@ pub fn diff_against(work_dir: &Path, base: &str) -> Result { Ok(String::from_utf8_lossy(&output.stdout).to_string()) } +/// Remove any stale worktree at `path` (best-effort), then add a fresh one. +pub fn replace_worktree(repo: &Path, path: &Path, branch: &str) -> Result<()> { + let _ = remove_worktree(repo, path); + add_worktree(repo, path, branch) +} + +/// Hard-reset the working directory to a specific SHA. +pub fn reset_hard(work_dir: &Path, sha: &str) -> Result<()> { + let output = git_cmd(work_dir) + .args(["reset", "--hard", sha]) + .output() + .map_err(|e| git_error(format!("git reset --hard failed: {e}")))?; + + if !output.status.success() { + let stderr = String::from_utf8_lossy(&output.stderr); + return Err(git_error(format!("git reset --hard failed: {stderr}"))); + } + + Ok(()) +} + +/// Sanitize a string for use as a git ref component. +/// Lowercases, replaces non-alphanumeric chars with dashes, collapses runs. +pub fn sanitize_ref_component(s: &str) -> String { + let mut result = String::with_capacity(s.len()); + let mut prev_dash = false; + for c in s.chars() { + if c.is_ascii_alphanumeric() { + result.push(c.to_ascii_lowercase()); + prev_dash = false; + } else if !prev_dash { + result.push('-'); + prev_dash = true; + } + } + result.trim_matches('-').to_string() +} + /// Git-native metadata storage for pipeline runs. /// /// Stores checkpoint data, manifests, and graph DOT on an orphan branch @@ -787,4 +823,96 @@ mod tests { .unwrap(); assert_eq!(read_back, artifact_data); } + + #[test] + fn sanitize_ref_component_lowercases() { + assert_eq!(sanitize_ref_component("Hello"), "hello"); + } + + #[test] + fn sanitize_ref_component_replaces_special_chars() { + assert_eq!(sanitize_ref_component("a/b:c d"), "a-b-c-d"); + } + + #[test] + fn sanitize_ref_component_collapses_consecutive_dashes() { + assert_eq!(sanitize_ref_component("a///b"), "a-b"); + } + + #[test] + fn sanitize_ref_component_trims_leading_trailing_dashes() { + assert_eq!(sanitize_ref_component("--abc--"), "abc"); + } + + #[test] + fn sanitize_ref_component_mixed() { + assert_eq!( + sanitize_ref_component("My Node!@#123"), + "my-node-123" + ); + } + + #[test] + fn replace_worktree_on_clean_path() { + let dir = tempfile::tempdir().unwrap(); + init_repo(dir.path()); + create_branch(dir.path(), "rw-branch").unwrap(); + + let wt_path = dir.path().join("rw-worktree"); + replace_worktree(dir.path(), &wt_path, "rw-branch").unwrap(); + assert!(wt_path.join(".git").exists()); + + remove_worktree(dir.path(), &wt_path).unwrap(); + } + + #[test] + fn replace_worktree_replaces_stale() { + let dir = tempfile::tempdir().unwrap(); + init_repo(dir.path()); + create_branch(dir.path(), "stale-branch").unwrap(); + + let wt_path = dir.path().join("stale-wt"); + add_worktree(dir.path(), &wt_path, "stale-branch").unwrap(); + assert!(wt_path.join(".git").exists()); + + // Calling replace_worktree again succeeds (removes stale, re-creates) + replace_worktree(dir.path(), &wt_path, "stale-branch").unwrap(); + assert!(wt_path.join(".git").exists()); + + remove_worktree(dir.path(), &wt_path).unwrap(); + } + + #[test] + fn reset_hard_resets_to_sha() { + let dir = tempfile::tempdir().unwrap(); + init_repo(dir.path()); + let initial_sha = head_sha(dir.path()).unwrap(); + + // Make a commit + fs::write(dir.path().join("file.txt"), "content").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", + "add file", + ]) + .current_dir(dir.path()) + .output() + .unwrap(); + assert_ne!(head_sha(dir.path()).unwrap(), initial_sha); + + // Reset back + reset_hard(dir.path(), &initial_sha).unwrap(); + assert_eq!(head_sha(dir.path()).unwrap(), initial_sha); + assert!(!dir.path().join("file.txt").exists()); + } } diff --git a/crates/arc-workflows/src/handler/parallel.rs b/crates/arc-workflows/src/handler/parallel.rs index 7198e4f21..aba158fbf 100644 --- a/crates/arc-workflows/src/handler/parallel.rs +++ b/crates/arc-workflows/src/handler/parallel.rs @@ -268,7 +268,9 @@ impl Handler for ParallelHandler { let branch_key = &target_id; let branch_name = format!( "arc/run/parallel/{}/{}/{}", - gs.run_id, node.id, branch_key + gs.run_id, + crate::git::sanitize_ref_component(&node.id), + crate::git::sanitize_ref_component(branch_key), ); match &gs.mode { @@ -284,7 +286,8 @@ impl Handler for ParallelHandler { 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) + crate::git::replace_worktree(&wd, &wtp, &bn)?; + crate::git::reset_hard(&wtp, &bs) }) .await .map_err(|e| { @@ -315,7 +318,7 @@ impl Handler for ParallelHandler { "failed to create remote branch {branch_name}" ))); } - let ok = crate::engine::git_add_worktree_remote( + let ok = crate::engine::git_replace_worktree_remote( &*services.execution_env, &wt_path_str, &branch_name, @@ -326,6 +329,20 @@ impl Handler for ParallelHandler { "failed to add remote worktree {wt_path_str}" ))); } + // Reset worktree to the base SHA for a clean start + let reset_cmd = format!( + "{} reset --hard {bsha}", + crate::engine::GIT_REMOTE + ); + let reset_result = services + .execution_env + .exec_command(&reset_cmd, 30_000, Some(&wt_path_str), None, None) + .await; + if !matches!(reset_result, Ok(ref r) if r.exit_code == 0) { + return Err(ArcError::Handler(format!( + "failed to reset remote worktree {wt_path_str}" + ))); + } branch_context.set( "internal.work_dir", serde_json::json!(&wt_path_str), @@ -415,23 +432,26 @@ impl Handler for ParallelHandler { 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 git_r = crate::engine::GIT_REMOTE; + let add_cmd = format!("{git_r} add -A"); let add_result = setup .execution_env - .exec_command("git add -A", 30_000, None, None, None) + .exec_command(&add_cmd, 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}'" + "{git_r} -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_cmd = format!("{git_r} rev-parse HEAD"); let sha_result = setup .execution_env - .exec_command("git rev-parse HEAD", 10_000, None, None, None) + .exec_command(&sha_cmd, 10_000, None, None, None) .await; match sha_result { Ok(r) if r.exit_code == 0 => Some(r.stdout.trim().to_string()), diff --git a/crates/arc-workflows/tests/integration.rs b/crates/arc-workflows/tests/integration.rs index f471cd886..cc6c19c02 100644 --- a/crates/arc-workflows/tests/integration.rs +++ b/crates/arc-workflows/tests/integration.rs @@ -9344,7 +9344,7 @@ async fn parallel_git_branching_host_e2e() { ); // 10. Verify parallel branch refs still exist (for debugging) - let branch_ref_a = format!("arc/run/parallel/{run_id}/fan_out/branch_a"); + 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())