From e97c16f03f28e1083d16f19ad67b3c9d5a6d09f2 Mon Sep 17 00:00:00 2001 From: Bryan Helmkamp Date: Sat, 28 Feb 2026 10:51:37 -0500 Subject: [PATCH] Git worktree isolation and per-node checkpoint commits Create a dedicated git branch + worktree per pipeline run (Local env only) and commit after every node checkpoint. This gives each run an isolated working directory and a full git trail of changes per stage. New module: git.rs with ensure_clean, head_sha, create_branch, add/remove_worktree, and checkpoint_commit (using arc identity). Engine changes: RunConfig gains run_id and work_dir fields; after each checkpoint save, a git commit is created in the worktree and the SHA is stored in checkpoint.git_commit_sha. CLI changes: for Local execution, the repo cleanliness is verified before any log files are written, then a worktree is created on branch arc/run/{uuid}, cwd is switched into it, and cleanup runs after the engine completes. Handler changes: run_hook() accepts work_dir so hooks execute in the worktree. Co-Authored-By: Claude Opus 4.6 --- crates/arc-attractor/src/checkpoint.rs | 4 + crates/arc-attractor/src/cli/run.rs | 67 ++++- crates/arc-attractor/src/engine.rs | 96 ++++++- crates/arc-attractor/src/git.rs | 271 ++++++++++++++++++ crates/arc-attractor/src/handler/codergen.rs | 23 +- crates/arc-attractor/src/lib.rs | 1 + crates/arc-attractor/src/server.rs | 2 +- .../tests/daytona_integration.rs | 2 + crates/arc-attractor/tests/integration.rs | 176 +++++++++++- 9 files changed, 620 insertions(+), 22 deletions(-) create mode 100644 crates/arc-attractor/src/git.rs diff --git a/crates/arc-attractor/src/checkpoint.rs b/crates/arc-attractor/src/checkpoint.rs index 8e0e6eb53..79f4a8572 100644 --- a/crates/arc-attractor/src/checkpoint.rs +++ b/crates/arc-attractor/src/checkpoint.rs @@ -24,6 +24,9 @@ pub struct Checkpoint { /// The node to resume execution at (the next node after the checkpoint's `current_node`). #[serde(default, skip_serializing_if = "Option::is_none")] pub next_node_id: Option, + /// SHA of the git commit created at this checkpoint (when running in a worktree). + #[serde(default, skip_serializing_if = "Option::is_none")] + pub git_commit_sha: Option, } impl Checkpoint { @@ -45,6 +48,7 @@ impl Checkpoint { logs: context.logs_snapshot(), node_outcomes, next_node_id, + git_commit_sha: None, } } diff --git a/crates/arc-attractor/src/cli/run.rs b/crates/arc-attractor/src/cli/run.rs index 10805ab22..bbf62b094 100644 --- a/crates/arc-attractor/src/cli/run.rs +++ b/crates/arc-attractor/src/cli/run.rs @@ -93,7 +93,27 @@ pub async fn run_command(args: RunArgs, styles: &'static Styles) -> anyhow::Resu bail!("Validation failed"); } - // 2. Create logs directory + // 2. Pre-flight: check git cleanliness before creating any files + // (must happen before logs dir is created, which may be inside the repo) + let execution_env_kind_preview = { + let toml_exec = task_cfg + .as_ref() + .and_then(|c| c.execution.as_ref()) + .and_then(|e| e.environment.as_deref()) + .map(|s| s.parse::()) + .transpose() + .ok() + .flatten(); + args.execution_env.or(toml_exec).unwrap_or_default() + }; + let original_cwd = std::env::current_dir()?; + let git_clean = if execution_env_kind_preview == ExecutionEnvKind::Local { + crate::git::ensure_clean(&original_cwd).is_ok() + } else { + false + }; + + // 3. Create logs directory let logs_dir = args.logs_dir.unwrap_or_else(|| { let base = dirs::home_dir() .expect("could not determine home directory") @@ -233,6 +253,22 @@ pub async fn run_command(args: RunArgs, styles: &'static Styles) -> anyhow::Resu .map_err(|e| anyhow::anyhow!("Invalid execution environment in TOML: {e}"))?; let execution_env_kind = args.execution_env.or(toml_execution_env).unwrap_or_default(); + // Set up git worktree for local execution (must happen before cwd is captured) + let (worktree_run_id, worktree_work_dir, worktree_path) = if git_clean { + match setup_worktree(&original_cwd, &logs_dir) { + Ok((rid, wd, wt)) => (Some(rid), Some(wd), Some(wt)), + Err(e) => { + eprintln!( + "{yellow}Warning:{reset} Git worktree setup failed ({e}), running without worktree.", + yellow = styles.yellow, reset = styles.reset, + ); + (None, None, None) + } + } + } else { + (None, None, None) + }; + let cwd = std::env::current_dir().unwrap_or_else(|_| PathBuf::from(".")); let daytona_config = task_cfg .as_ref() @@ -437,6 +473,8 @@ pub async fn run_command(args: RunArgs, styles: &'static Styles) -> anyhow::Resu logs_root: logs_dir.clone(), cancel_token: None, dry_run: dry_run_mode, + run_id: worktree_run_id, + work_dir: worktree_work_dir, }; let run_start = Instant::now(); @@ -450,6 +488,12 @@ pub async fn run_command(args: RunArgs, styles: &'static Styles) -> anyhow::Resu }; let run_duration_ms = run_start.elapsed().as_millis() as u64; + // Restore cwd and clean up worktree (best-effort) + let _ = std::env::set_current_dir(&original_cwd); + if let Some(ref wt) = worktree_path { + let _ = crate::git::remove_worktree(&original_cwd, wt); + } + { let (status, failure_reason) = match &engine_result { Ok(o) => (o.status.to_string(), o.failure_reason.clone()), @@ -532,6 +576,27 @@ pub async fn run_command(args: RunArgs, styles: &'static Styles) -> anyhow::Resu } } +/// Set up a git worktree for an isolated pipeline run. +/// Caller must have already verified the repo is clean via `git::ensure_clean`. +/// Returns (run_id, work_dir, worktree_path) on success. +fn setup_worktree( + original_cwd: &std::path::Path, + logs_dir: &std::path::Path, +) -> anyhow::Result<(String, PathBuf, PathBuf)> { + let run_id = uuid::Uuid::new_v4().to_string(); + let branch_name = format!("arc/run/{run_id}"); + 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) + .map_err(|e| anyhow::anyhow!("{e}"))?; + + std::env::set_current_dir(&worktree_path)?; + + Ok((run_id, worktree_path.clone(), worktree_path)) +} + #[cfg(test)] mod tests { #[test] diff --git a/crates/arc-attractor/src/engine.rs b/crates/arc-attractor/src/engine.rs index c16d66816..65fe9a7bd 100644 --- a/crates/arc-attractor/src/engine.rs +++ b/crates/arc-attractor/src/engine.rs @@ -467,6 +467,10 @@ pub struct RunConfig { pub logs_root: PathBuf, pub cancel_token: Option>, pub dry_run: bool, + /// Pre-assigned run ID. Generated if `None`. + pub run_id: Option, + /// Git worktree path for checkpoint commits. + pub work_dir: Option, } /// The pipeline execution engine. @@ -695,12 +699,12 @@ impl PipelineEngine { mut node_visits: HashMap, ) -> Result { let run_start = Instant::now(); - let run_id = uuid::Uuid::new_v4().to_string(); + let run_id = config.run_id.clone().unwrap_or_else(|| uuid::Uuid::new_v4().to_string()); let artifact_store = ArtifactStore::new(Some(config.logs_root.clone())); self.services.emitter.emit(&PipelineEvent::PipelineStarted { name: graph.name.clone(), - id: run_id, + id: run_id.clone(), }); self.inform(&format!("Pipeline started: {}", graph.name), "pipeline"); @@ -784,6 +788,12 @@ impl PipelineEngine { current_node_id = start_node.id.clone(); } + // Store run_id and work_dir in context for handlers + context.set("internal.run_id", serde_json::json!(run_id)); + if let Some(ref wd) = config.work_dir { + context.set("internal.work_dir", serde_json::json!(wd.to_string_lossy().as_ref())); + } + loop { // Check for cancellation before processing each node if let Some(ref token) = config.cancel_token { @@ -990,7 +1000,7 @@ impl PipelineEngine { let next_node_id_for_checkpoint = next_edge.map(|e| e.to.clone()); // Step 6: Save checkpoint with all state - let checkpoint = Checkpoint::from_context( + let mut checkpoint = Checkpoint::from_context( &context, &node.id, completed_nodes.clone(), @@ -1007,6 +1017,32 @@ impl PipelineEngine { }); } + // Step 6b: Git checkpoint commit (when running in a worktree) + if let Some(ref work_dir) = config.work_dir { + let wd = work_dir.clone(); + let rid = run_id.clone(); + let nid = node.id.clone(); + let status_str = outcome.status.to_string(); + match tokio::task::spawn_blocking(move || { + crate::git::checkpoint_commit(&wd, &rid, &nid, &status_str) + }) + .await + { + Ok(Ok(sha)) => { + checkpoint.git_commit_sha = Some(sha); + if let Err(e) = checkpoint.save(&checkpoint_path) { + context.append_log(format!("checkpoint re-save with SHA failed: {e}")); + } + } + Ok(Err(e)) => { + context.append_log(format!("git checkpoint commit failed: {e}")); + } + Err(e) => { + context.append_log(format!("git checkpoint commit task panicked: {e}")); + } + } + } + // Step 7: Follow selected edge match next_edge { None => { @@ -1736,6 +1772,8 @@ mod tests { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false, + run_id: None, + work_dir: None, }; let outcome = engine.run(&g, &config).await.unwrap(); assert_eq!(outcome.status, StageStatus::Success); @@ -1750,6 +1788,8 @@ mod tests { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false, + run_id: None, + work_dir: None, }; engine.run(&g, &config).await.unwrap(); let checkpoint_path = dir.path().join("checkpoint.json"); @@ -1773,6 +1813,8 @@ mod tests { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false, + run_id: None, + work_dir: None, }; engine.run(&g, &config).await.unwrap(); @@ -1791,6 +1833,8 @@ mod tests { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false, + run_id: None, + work_dir: None, }; let result = engine.run(&g, &config).await; assert!(result.is_err()); @@ -1805,6 +1849,8 @@ mod tests { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false, + run_id: None, + work_dir: None, }; engine.run(&g, &config).await.unwrap(); @@ -1832,6 +1878,8 @@ mod tests { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false, + run_id: None, + work_dir: None, }; let outcome = engine.run(&g, &config).await.unwrap(); assert_eq!(outcome.status, StageStatus::Success); @@ -1885,6 +1933,8 @@ mod tests { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false, + run_id: None, + work_dir: None, }; engine.run(&g, &config).await.unwrap(); @@ -1956,6 +2006,8 @@ mod tests { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false, + run_id: None, + work_dir: None, }; engine.run(&g, &config).await.unwrap(); @@ -1979,6 +2031,8 @@ mod tests { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false, + run_id: None, + work_dir: None, }; engine.run(&g, &config).await.unwrap(); @@ -1999,6 +2053,8 @@ mod tests { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false, + run_id: None, + work_dir: None, }; engine.run(&g, &config).await.unwrap(); @@ -2143,7 +2199,7 @@ mod tests { let dir = tempfile::tempdir().unwrap(); let g = simple_graph(); let engine = PipelineEngine::new(make_registry(), Arc::new(EventEmitter::new()), local_env()); - let config = RunConfig { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false }; + let config = RunConfig { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false, run_id: None, work_dir: None }; engine.run(&g, &config).await.unwrap(); let manifest_path = dir.path().join("manifest.json"); @@ -2165,7 +2221,7 @@ mod tests { g.edges.push(Edge::new("start", "exit")); let engine = PipelineEngine::new(make_registry(), Arc::new(EventEmitter::new()), local_env()); - let config = RunConfig { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false }; + let config = RunConfig { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false, run_id: None, work_dir: None }; engine.run(&g, &config).await.unwrap(); let manifest_path = dir.path().join("manifest.json"); @@ -2201,7 +2257,7 @@ mod tests { let mut registry = make_registry(); registry.register("always_fail", Box::new(AlwaysFailHandler)); let engine = PipelineEngine::new(registry, Arc::new(EventEmitter::new()), local_env()); - let config = RunConfig { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false }; + let config = RunConfig { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false, run_id: None, work_dir: None }; let outcome = engine.run(&g, &config).await.unwrap(); assert_eq!(outcome.status, StageStatus::Success); @@ -2237,7 +2293,7 @@ mod tests { let mut registry = make_registry(); registry.register("always_fail", Box::new(AlwaysFailHandler)); let engine = PipelineEngine::new(registry, Arc::new(EventEmitter::new()), local_env()); - let config = RunConfig { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false }; + let config = RunConfig { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false, run_id: None, work_dir: None }; let result = engine.run(&g, &config).await; assert!(result.is_ok()); @@ -2276,7 +2332,7 @@ mod tests { let mut registry = make_registry(); registry.register("slow", Box::new(SlowHandler { sleep_ms: 500 })); let engine = PipelineEngine::new(registry, Arc::new(EventEmitter::new()), local_env()); - let config = RunConfig { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false }; + let config = RunConfig { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false, run_id: None, work_dir: None }; let result = engine.run(&g, &config).await; assert!(result.is_ok()); @@ -2310,7 +2366,7 @@ mod tests { let mut registry = make_registry(); registry.register("slow", Box::new(SlowHandler { sleep_ms: 10 })); let engine = PipelineEngine::new(registry, Arc::new(EventEmitter::new()), local_env()); - let config = RunConfig { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false }; + let config = RunConfig { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false, run_id: None, work_dir: None }; let outcome = engine.run(&g, &config).await.unwrap(); assert_eq!(outcome.status, StageStatus::Success); } @@ -2341,7 +2397,7 @@ mod tests { let mut registry = make_registry(); registry.register("slow", Box::new(SlowHandler { sleep_ms: 500 })); let engine = PipelineEngine::new(registry, Arc::new(EventEmitter::new()), local_env()); - let config = RunConfig { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false }; + let config = RunConfig { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false, run_id: None, work_dir: None }; let outcome = engine.run(&g, &config).await.unwrap(); assert_eq!(outcome.status, StageStatus::Success); @@ -2395,6 +2451,8 @@ mod tests { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false, + run_id: None, + work_dir: None, }; engine.run(&g, &config).await.unwrap(); // Give spawned inform tasks time to complete @@ -2422,6 +2480,8 @@ mod tests { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false, + run_id: None, + work_dir: None, }; engine.run(&g, &config).await.unwrap(); // Give spawned inform tasks time to complete @@ -2447,6 +2507,8 @@ mod tests { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false, + run_id: None, + work_dir: None, }; let outcome = engine.run(&g, &config).await.unwrap(); assert_eq!(outcome.status, StageStatus::Success); @@ -2464,6 +2526,8 @@ mod tests { logs_root: dir.path().to_path_buf(), cancel_token: Some(cancel_token), dry_run: false, + run_id: None, + work_dir: None, }; let result = engine.run(&g, &config).await; assert!(result.is_err()); @@ -2480,6 +2544,8 @@ mod tests { logs_root: dir.path().to_path_buf(), cancel_token: Some(cancel_token), dry_run: false, + run_id: None, + work_dir: None, }; let outcome = engine.run(&g, &config).await.unwrap(); assert_eq!(outcome.status, StageStatus::Success); @@ -2508,6 +2574,8 @@ mod tests { logs_root: dir.path().to_path_buf(), cancel_token: Some(cancel_token), dry_run: false, + run_id: None, + work_dir: None, }; // Set cancel after a short delay (while the slow handler is running) @@ -2580,6 +2648,8 @@ mod tests { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false, + run_id: None, + work_dir: None, }; let result = engine.run(&g, &config).await; assert!(result.is_err()); @@ -2599,6 +2669,8 @@ mod tests { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: true, + run_id: None, + work_dir: None, }; let result = engine.run(&g, &config).await; assert!(result.is_err()); @@ -2622,6 +2694,8 @@ mod tests { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: true, + run_id: None, + work_dir: None, }; let result = engine.run(&g, &config).await; assert!(result.is_err()); @@ -2706,6 +2780,8 @@ mod tests { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false, + run_id: None, + work_dir: None, }; // The engine returns Err because the Fail outcome has no outgoing fail edge, diff --git a/crates/arc-attractor/src/git.rs b/crates/arc-attractor/src/git.rs new file mode 100644 index 000000000..0dd2e5c9f --- /dev/null +++ b/crates/arc-attractor/src/git.rs @@ -0,0 +1,271 @@ +use std::path::Path; +use std::process::Command; + +use crate::error::{AttractorError, Result}; + +fn git_error(msg: impl Into) -> AttractorError { + AttractorError::Engine(msg.into()) +} + +/// Assert the working directory is a clean git repo (no uncommitted changes). +pub fn ensure_clean(repo: &Path) -> Result<()> { + let output = Command::new("git") + .args(["status", "--porcelain"]) + .current_dir(repo) + .output() + .map_err(|e| git_error(format!("git status failed: {e}")))?; + + if !output.status.success() { + return Err(git_error("not a git repository")); + } + + let stdout = String::from_utf8_lossy(&output.stdout); + if !stdout.trim().is_empty() { + return Err(git_error("working directory has uncommitted changes")); + } + + Ok(()) +} + +/// Return the SHA of HEAD. +pub fn head_sha(repo: &Path) -> Result { + let output = Command::new("git") + .args(["rev-parse", "HEAD"]) + .current_dir(repo) + .output() + .map_err(|e| git_error(format!("git rev-parse failed: {e}")))?; + + if !output.status.success() { + return Err(git_error("git rev-parse HEAD failed")); + } + + Ok(String::from_utf8_lossy(&output.stdout).trim().to_string()) +} + +/// 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) + .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(()) +} + +/// 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") + .args(["worktree", "add"]) + .arg(path) + .arg(branch) + .current_dir(repo) + .output() + .map_err(|e| git_error(format!("git worktree add failed: {e}")))?; + + if !output.status.success() { + let stderr = String::from_utf8_lossy(&output.stderr); + return Err(git_error(format!("git worktree add failed: {stderr}"))); + } + + Ok(()) +} + +/// Remove a git worktree. +pub fn remove_worktree(repo: &Path, path: &Path) -> Result<()> { + let output = Command::new("git") + .args(["worktree", "remove", "--force"]) + .arg(path) + .current_dir(repo) + .output() + .map_err(|e| git_error(format!("git worktree remove failed: {e}")))?; + + if !output.status.success() { + let stderr = String::from_utf8_lossy(&output.stderr); + return Err(git_error(format!("git worktree remove failed: {stderr}"))); + } + + Ok(()) +} + +/// Stage all changes and commit in `work_dir` with a structured message. +/// Returns the new commit SHA. +pub fn checkpoint_commit( + work_dir: &Path, + run_id: &str, + node_id: &str, + status: &str, +) -> Result { + // Stage everything + let output = Command::new("git") + .args(["add", "-A"]) + .current_dir(work_dir) + .output() + .map_err(|e| git_error(format!("git add failed: {e}")))?; + + if !output.status.success() { + let stderr = String::from_utf8_lossy(&output.stderr); + return Err(git_error(format!("git add failed: {stderr}"))); + } + + // Commit with arc identity (works even if user.name/email not configured) + let message = format!("arc({run_id}): {node_id} ({status})"); + let output = Command::new("git") + .args([ + "-c", "user.name=arc", + "-c", "user.email=arc@local", + "commit", + "--allow-empty", + "-m", &message, + ]) + .current_dir(work_dir) + .output() + .map_err(|e| git_error(format!("git commit failed: {e}")))?; + + if !output.status.success() { + let stderr = String::from_utf8_lossy(&output.stderr); + return Err(git_error(format!("git commit failed: {stderr}"))); + } + + head_sha(work_dir) +} + +#[cfg(test)] +mod tests { + use super::*; + use std::fs; + + /// Create a temporary git repo with an initial commit. + fn init_repo(dir: &Path) { + Command::new("git") + .args(["init"]) + .current_dir(dir) + .output() + .unwrap(); + Command::new("git") + .args(["-c", "user.name=test", "-c", "user.email=test@test", "commit", "--allow-empty", "-m", "init"]) + .current_dir(dir) + .output() + .unwrap(); + } + + #[test] + fn ensure_clean_on_clean_repo() { + let dir = tempfile::tempdir().unwrap(); + init_repo(dir.path()); + assert!(ensure_clean(dir.path()).is_ok()); + } + + #[test] + fn ensure_clean_fails_with_dirty_file() { + let dir = tempfile::tempdir().unwrap(); + init_repo(dir.path()); + fs::write(dir.path().join("dirty.txt"), "hello").unwrap(); + let err = ensure_clean(dir.path()).unwrap_err(); + assert!(err.to_string().contains("uncommitted changes")); + } + + #[test] + fn ensure_clean_fails_on_non_repo() { + let dir = tempfile::tempdir().unwrap(); + let err = ensure_clean(dir.path()).unwrap_err(); + assert!(err.to_string().contains("not a git repository")); + } + + #[test] + fn head_sha_returns_40_char_hex() { + let dir = tempfile::tempdir().unwrap(); + init_repo(dir.path()); + let sha = head_sha(dir.path()).unwrap(); + assert_eq!(sha.len(), 40); + assert!(sha.chars().all(|c| c.is_ascii_hexdigit())); + } + + #[test] + fn create_branch_and_list() { + let dir = tempfile::tempdir().unwrap(); + init_repo(dir.path()); + create_branch(dir.path(), "test-branch").unwrap(); + + let output = Command::new("git") + .args(["branch", "--list", "test-branch"]) + .current_dir(dir.path()) + .output() + .unwrap(); + let stdout = String::from_utf8_lossy(&output.stdout); + assert!(stdout.contains("test-branch")); + } + + #[test] + fn add_and_remove_worktree() { + let dir = tempfile::tempdir().unwrap(); + init_repo(dir.path()); + create_branch(dir.path(), "wt-branch").unwrap(); + + let wt_path = dir.path().join("my-worktree"); + add_worktree(dir.path(), &wt_path, "wt-branch").unwrap(); + assert!(wt_path.join(".git").exists()); + + remove_worktree(dir.path(), &wt_path).unwrap(); + assert!(!wt_path.exists()); + } + + #[test] + fn checkpoint_commit_creates_commit() { + let dir = tempfile::tempdir().unwrap(); + init_repo(dir.path()); + create_branch(dir.path(), "run-branch").unwrap(); + + let wt_path = dir.path().join("worktree"); + add_worktree(dir.path(), &wt_path, "run-branch").unwrap(); + + // Write a file in the worktree + fs::write(wt_path.join("output.txt"), "result").unwrap(); + + let sha = checkpoint_commit(&wt_path, "run1", "nodeA", "success").unwrap(); + assert_eq!(sha.len(), 40); + assert!(sha.chars().all(|c| c.is_ascii_hexdigit())); + + // Verify commit message + let output = Command::new("git") + .args(["log", "--oneline", "-1"]) + .current_dir(&wt_path) + .output() + .unwrap(); + let log = String::from_utf8_lossy(&output.stdout); + assert!(log.contains("arc(run1): nodeA (success)")); + + remove_worktree(dir.path(), &wt_path).unwrap(); + } + + #[test] + fn checkpoint_commit_with_no_user_config() { + let dir = tempfile::tempdir().unwrap(); + // Init repo without setting global user.name/email — the -c flags on commit + // provide identity, so this should still succeed. + Command::new("git") + .args(["init"]) + .current_dir(dir.path()) + .output() + .unwrap(); + Command::new("git") + .args(["-c", "user.name=test", "-c", "user.email=test@test", "commit", "--allow-empty", "-m", "init"]) + .current_dir(dir.path()) + .output() + .unwrap(); + create_branch(dir.path(), "fallback-branch").unwrap(); + + let wt_path = dir.path().join("worktree"); + add_worktree(dir.path(), &wt_path, "fallback-branch").unwrap(); + + let sha = checkpoint_commit(&wt_path, "run2", "nodeB", "completed").unwrap(); + assert_eq!(sha.len(), 40); + + remove_worktree(dir.path(), &wt_path).unwrap(); + } +} diff --git a/crates/arc-attractor/src/handler/codergen.rs b/crates/arc-attractor/src/handler/codergen.rs index 852c87bfc..83aa7cdea 100644 --- a/crates/arc-attractor/src/handler/codergen.rs +++ b/crates/arc-attractor/src/handler/codergen.rs @@ -176,13 +176,15 @@ fn resolve_hook(node: &Node, graph: &Graph, key: &str) -> Option { } /// Execute a tool hook shell command. Returns true if the command succeeded (exit 0). -fn run_hook(command: &str, node_id: &str) -> bool { - match std::process::Command::new("sh") - .arg("-c") +fn run_hook(command: &str, node_id: &str, work_dir: Option<&Path>) -> bool { + let mut cmd = std::process::Command::new("sh"); + cmd.arg("-c") .arg(command) - .env("ATTRACTOR_NODE_ID", node_id) - .output() - { + .env("ATTRACTOR_NODE_ID", node_id); + if let Some(wd) = work_dir { + cmd.current_dir(wd); + } + match cmd.output() { Ok(output) => output.status.success(), Err(_) => false, } @@ -217,9 +219,14 @@ impl Handler for CodergenHandler { tokio::fs::create_dir_all(&stage_dir).await?; tokio::fs::write(stage_dir.join("prompt.md"), &prompt).await?; + // Resolve work_dir from context for hooks + let work_dir_str = context.get("internal.work_dir") + .and_then(|v| v.as_str().map(String::from)); + let work_dir = work_dir_str.as_deref().map(Path::new); + // 3. Execute pre-hook (spec 9.7) if let Some(pre_hook) = resolve_hook(node, graph, "tool_hooks.pre") { - if !run_hook(&pre_hook, &node.id) { + if !run_hook(&pre_hook, &node.id, work_dir) { let mut outcome = Outcome::skipped(); outcome.notes = Some("pre-hook returned non-zero, tool call skipped".to_string()); return Ok(outcome); @@ -259,7 +266,7 @@ impl Handler for CodergenHandler { // 5. Execute post-hook (spec 9.7) if let Some(post_hook) = resolve_hook(node, graph, "tool_hooks.post") { - if !run_hook(&post_hook, &node.id) { + if !run_hook(&post_hook, &node.id, work_dir) { context.append_log(format!( "post-hook failed for node {}, continuing", node.id diff --git a/crates/arc-attractor/src/lib.rs b/crates/arc-attractor/src/lib.rs index 47cfb2519..73587b4ad 100644 --- a/crates/arc-attractor/src/lib.rs +++ b/crates/arc-attractor/src/lib.rs @@ -7,6 +7,7 @@ pub mod context; pub mod engine; pub mod error; pub mod event; +pub mod git; pub mod graph; pub mod handler; pub mod interviewer; diff --git a/crates/arc-attractor/src/server.rs b/crates/arc-attractor/src/server.rs index e6eb9fc36..f0b570aeb 100644 --- a/crates/arc-attractor/src/server.rs +++ b/crates/arc-attractor/src/server.rs @@ -219,7 +219,7 @@ async fn start_pipeline( tokio::spawn(async move { let logs_root = std::env::temp_dir().join(format!("arc-{}", uuid::Uuid::new_v4())); std::fs::create_dir_all(&logs_root).expect("failed to create logs directory"); - let config = RunConfig { logs_root, cancel_token: Some(cancel_token), dry_run: state_clone.dry_run }; + let config = RunConfig { logs_root, cancel_token: Some(cancel_token), dry_run: state_clone.dry_run, run_id: None, work_dir: None }; let result = tokio::select! { result = engine.run(&graph, &config) => result, diff --git a/crates/arc-attractor/tests/daytona_integration.rs b/crates/arc-attractor/tests/daytona_integration.rs index 3d976cbce..c6a7adafb 100644 --- a/crates/arc-attractor/tests/daytona_integration.rs +++ b/crates/arc-attractor/tests/daytona_integration.rs @@ -270,6 +270,8 @@ async fn daytona_pipeline_artifact_offload_and_sync() { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false, + run_id: None, + work_dir: None, }; let outcome = engine.run(&graph, &config).await.expect("pipeline should succeed"); diff --git a/crates/arc-attractor/tests/integration.rs b/crates/arc-attractor/tests/integration.rs index 0d2a0741e..8bfc1e6d4 100644 --- a/crates/arc-attractor/tests/integration.rs +++ b/crates/arc-attractor/tests/integration.rs @@ -186,6 +186,8 @@ async fn end_to_end_linear_pipeline() { let config = RunConfig { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false, + run_id: None, + work_dir: None, }; let outcome = engine.run(&graph, &config).await.expect("run should succeed"); @@ -304,6 +306,8 @@ async fn end_to_end_branching_pipeline() { let config = RunConfig { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false, + run_id: None, + work_dir: None, }; let outcome = engine.run(&graph, &config).await.expect("run should succeed"); @@ -409,6 +413,8 @@ async fn end_to_end_human_gate_pipeline() { let config = RunConfig { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false, + run_id: None, + work_dir: None, }; let outcome = engine.run(&graph, &config).await.expect("run should succeed"); @@ -508,6 +514,8 @@ async fn goal_gate_routes_to_retry_target_on_failure() { let config = RunConfig { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false, + run_id: None, + work_dir: None, }; let result = engine.run(&graph, &config).await; @@ -616,6 +624,8 @@ async fn goal_gate_routes_to_retry_target_when_present() { let config = RunConfig { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false, + run_id: None, + work_dir: None, }; let outcome = engine @@ -924,6 +934,8 @@ async fn retry_on_failure_then_succeed() { let config = RunConfig { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false, + run_id: None, + work_dir: None, }; let outcome = engine @@ -986,6 +998,8 @@ async fn pipeline_with_many_nodes() { let config = RunConfig { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false, + run_id: None, + work_dir: None, }; let outcome = engine @@ -1300,6 +1314,8 @@ async fn smoke_test_with_mock_codergen_backend() { let config = RunConfig { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false, + run_id: None, + work_dir: None, }; let outcome = engine @@ -1394,6 +1410,8 @@ async fn end_to_end_parallel_fan_out_fan_in() { let config = RunConfig { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false, + run_id: None, + work_dir: None, }; let outcome = engine @@ -1497,6 +1515,8 @@ async fn resume_from_checkpoint_completes_pipeline() { let config = RunConfig { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false, + run_id: None, + work_dir: None, }; let outcome = engine @@ -1585,6 +1605,8 @@ async fn resume_from_checkpoint_preserves_goal_gate_outcomes() { let config = RunConfig { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false, + run_id: None, + work_dir: None, }; // This should succeed because goal gate for gated_work is satisfied @@ -1615,6 +1637,8 @@ async fn graph_goal_in_context() { let config = RunConfig { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false, + run_id: None, + work_dir: None, }; engine.run(&graph, &config).await.expect("run"); @@ -1641,6 +1665,8 @@ async fn event_streaming_lifecycle() { let config = RunConfig { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false, + run_id: None, + work_dir: None, }; engine.run(&graph, &config).await.expect("run"); @@ -1707,6 +1733,8 @@ async fn context_flow_between_stages() { let config = RunConfig { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false, + run_id: None, + work_dir: None, }; engine.run(&graph, &config).await.expect("run"); @@ -1746,6 +1774,8 @@ async fn tool_handler_e2e() { let config = RunConfig { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false, + run_id: None, + work_dir: None, }; let outcome = engine.run(&graph, &config).await.expect("run"); assert_eq!(outcome.status, StageStatus::Success); @@ -1803,6 +1833,8 @@ async fn auto_approve_interviewer_e2e() { let config = RunConfig { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false, + run_id: None, + work_dir: None, }; let outcome = engine.run(&graph, &config).await.expect("run"); assert_eq!(outcome.status, StageStatus::Success); @@ -1826,6 +1858,8 @@ async fn codergen_without_backend_simulated() { let config = RunConfig { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false, + run_id: None, + work_dir: None, }; engine.run(&graph, &config).await.expect("run"); @@ -1923,6 +1957,8 @@ async fn branching_loop_back_on_failure() { let config = RunConfig { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false, + run_id: None, + work_dir: None, }; let outcome = engine.run(&graph, &config).await.expect("run"); assert_eq!(outcome.status, StageStatus::Success); @@ -2003,6 +2039,8 @@ async fn human_gate_loops_back() { let config = RunConfig { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false, + run_id: None, + work_dir: None, }; let outcome = engine.run(&graph, &config).await.expect("run"); assert_eq!(outcome.status, StageStatus::Success); @@ -2051,6 +2089,8 @@ async fn scenario_ship_a_feature() { let config = RunConfig { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false, + run_id: None, + work_dir: None, }; let outcome = engine.run(&graph, &config).await.expect("run"); assert_eq!(outcome.status, StageStatus::Success); @@ -2127,6 +2167,8 @@ async fn scenario_parallel_expert_review() { let config = RunConfig { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false, + run_id: None, + work_dir: None, }; let outcome = engine.run(&graph, &config).await.expect("run"); assert_eq!(outcome.status, StageStatus::Success); @@ -2197,6 +2239,8 @@ async fn scenario_node_retries_on_retry_status() { let config = RunConfig { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false, + run_id: None, + work_dir: None, }; let outcome = engine.run(&graph, &config).await.expect("run"); assert_eq!(outcome.status, StageStatus::Success); @@ -2251,6 +2295,8 @@ async fn scenario_loop_restart_resets_context() { let config = RunConfig { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false, + run_id: None, + work_dir: None, }; let outcome = engine.run(&graph, &config).await.expect("run"); assert_eq!(outcome.status, StageStatus::Success); @@ -2309,6 +2355,8 @@ async fn scenario_bug_triage_router() { let config = RunConfig { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false, + run_id: None, + work_dir: None, }; let outcome = engine.run(&graph, &config).await.expect("run"); assert_eq!(outcome.status, StageStatus::Success); @@ -2355,6 +2403,8 @@ async fn scenario_crash_recovery() { let config = RunConfig { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false, + run_id: None, + work_dir: None, }; let outcome = engine .run_from_checkpoint(&graph, &config, &checkpoint) @@ -2433,6 +2483,8 @@ async fn manager_loop_stop_condition_satisfied_e2e() { let config = RunConfig { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false, + run_id: None, + work_dir: None, }; let outcome = engine.run(&graph, &config).await.expect("run"); @@ -2479,6 +2531,8 @@ async fn manager_loop_max_cycles_exceeded_e2e() { let config = RunConfig { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false, + run_id: None, + work_dir: None, }; let outcome = engine.run(&graph, &config).await.expect("run"); @@ -2605,6 +2659,8 @@ async fn conditional_branching_success_fail_paths() { let config = RunConfig { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false, + run_id: None, + work_dir: None, }; let outcome = engine.run(&graph, &config).await.expect("run"); assert_eq!(outcome.status, StageStatus::Success); @@ -2649,6 +2705,8 @@ async fn edge_selection_condition_match_wins_over_weight() { let config = RunConfig { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false, + run_id: None, + work_dir: None, }; engine.run(&graph, &config).await.expect("run"); @@ -2688,6 +2746,8 @@ async fn edge_selection_weight_breaks_ties() { let config = RunConfig { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false, + run_id: None, + work_dir: None, }; engine.run(&graph, &config).await.expect("run"); @@ -2719,6 +2779,8 @@ async fn edge_selection_lexical_tiebreak() { let config = RunConfig { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false, + run_id: None, + work_dir: None, }; engine.run(&graph, &config).await.expect("run"); @@ -2767,6 +2829,8 @@ async fn context_updates_visible_across_nodes() { let config = RunConfig { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false, + run_id: None, + work_dir: None, }; engine.run(&graph, &config).await.expect("run"); @@ -2797,6 +2861,8 @@ async fn stylesheet_applies_model_override() { let config = RunConfig { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false, + run_id: None, + work_dir: None, }; let outcome = engine.run(&graph, &config).await.expect("run"); assert_eq!(outcome.status, StageStatus::Success); @@ -2843,6 +2909,8 @@ async fn custom_handler_registration_and_execution() { let config = RunConfig { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false, + run_id: None, + work_dir: None, }; engine.run(&graph, &config).await.expect("run"); @@ -2900,6 +2968,8 @@ async fn integration_smoke_plan_implement_review_done() { let config = RunConfig { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false, + run_id: None, + work_dir: None, }; let outcome = engine.run(&graph, &config).await.expect("run"); assert_eq!(outcome.status, StageStatus::Success); @@ -3385,6 +3455,8 @@ async fn sub_pipeline_e2e_through_engine() { let config = RunConfig { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false, + run_id: None, + work_dir: None, }; let outcome = engine @@ -3527,6 +3599,8 @@ async fn manager_loop_with_child_observer_e2e() { let config = RunConfig { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false, + run_id: None, + work_dir: None, }; let outcome = engine @@ -3646,6 +3720,8 @@ async fn graph_merge_e2e_through_engine() { let config = RunConfig { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false, + run_id: None, + work_dir: None, }; let outcome = engine @@ -3786,6 +3862,8 @@ async fn fidelity_default_is_compact() { let config = RunConfig { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false, + run_id: None, + work_dir: None, }; engine.run(&graph, &config).await.expect("run"); @@ -3825,6 +3903,8 @@ async fn fidelity_graph_default_applied() { let config = RunConfig { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false, + run_id: None, + work_dir: None, }; engine.run(&graph, &config).await.expect("run"); @@ -3863,6 +3943,8 @@ async fn fidelity_node_overrides_graph_default() { let config = RunConfig { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false, + run_id: None, + work_dir: None, }; engine.run(&graph, &config).await.expect("run"); @@ -3907,6 +3989,8 @@ async fn fidelity_edge_overrides_node_and_graph() { let config = RunConfig { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false, + run_id: None, + work_dir: None, }; engine.run(&graph, &config).await.expect("run"); @@ -3941,6 +4025,8 @@ async fn fidelity_full_produces_empty_preamble() { let config = RunConfig { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false, + run_id: None, + work_dir: None, }; engine.run(&graph, &config).await.expect("run"); @@ -3982,6 +4068,8 @@ async fn fidelity_truncate_preamble_minimal() { let config = RunConfig { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false, + run_id: None, + work_dir: None, }; engine.run(&graph, &config).await.expect("run"); @@ -4039,6 +4127,8 @@ async fn fidelity_summary_low_mode() { let config = RunConfig { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false, + run_id: None, + work_dir: None, }; engine.run(&graph, &config).await.expect("run"); @@ -4091,6 +4181,8 @@ async fn fidelity_summary_medium_mode() { let config = RunConfig { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false, + run_id: None, + work_dir: None, }; engine.run(&graph, &config).await.expect("run"); @@ -4143,6 +4235,8 @@ async fn fidelity_summary_high_mode() { let config = RunConfig { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false, + run_id: None, + work_dir: None, }; engine.run(&graph, &config).await.expect("run"); @@ -4188,6 +4282,8 @@ async fn fidelity_full_sets_thread_id_in_context() { let config = RunConfig { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false, + run_id: None, + work_dir: None, }; engine.run(&graph, &config).await.expect("run"); @@ -4244,6 +4340,8 @@ async fn fidelity_full_nodes_share_thread_id() { let config = RunConfig { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false, + run_id: None, + work_dir: None, }; engine.run(&graph, &config).await.expect("run"); @@ -4307,6 +4405,8 @@ async fn fidelity_resume_degrades_full_to_summary_high() { let config = RunConfig { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false, + run_id: None, + work_dir: None, }; engine .run_from_checkpoint(&graph, &config, &checkpoint) @@ -4386,6 +4486,8 @@ async fn fidelity_resume_degrade_only_affects_first_hop() { let config = RunConfig { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false, + run_id: None, + work_dir: None, }; engine .run_from_checkpoint(&graph, &config, &checkpoint) @@ -4452,6 +4554,8 @@ async fn fidelity_resume_no_degrade_when_not_full() { let config = RunConfig { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false, + run_id: None, + work_dir: None, }; engine .run_from_checkpoint(&graph, &config, &checkpoint) @@ -4484,6 +4588,8 @@ async fn fidelity_stored_in_checkpoint_context() { let config = RunConfig { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false, + run_id: None, + work_dir: None, }; engine.run(&graph, &config).await.expect("run"); @@ -4555,6 +4661,8 @@ async fn fidelity_precedence_multi_node_pipeline() { let config = RunConfig { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false, + run_id: None, + work_dir: None, }; engine.run(&graph, &config).await.expect("run"); @@ -4608,6 +4716,8 @@ async fn fidelity_compact_preamble_includes_completed_stages_and_context() { let config = RunConfig { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false, + run_id: None, + work_dir: None, }; engine.run(&graph, &config).await.expect("run"); @@ -4653,7 +4763,7 @@ async fn fidelity_summary_low_excludes_context_values_in_pipeline() { registry_low.register("exit", Box::new(ExitHandler)); registry_low.register("fidelity_capture", Box::new(FidelityCapturingHandler { captures: captures_low.clone() })); let engine_low = PipelineEngine::new(registry_low, Arc::new(EventEmitter::new()), local_env()); - let config_low = RunConfig { logs_root: dir_low.path().to_path_buf(), cancel_token: None, dry_run: false }; + let config_low = RunConfig { logs_root: dir_low.path().to_path_buf(), cancel_token: None, dry_run: false, run_id: None, work_dir: None }; engine_low.run(&graph_low, &config_low).await.expect("run low"); { @@ -4687,7 +4797,7 @@ async fn fidelity_summary_low_excludes_context_values_in_pipeline() { registry_med.register("exit", Box::new(ExitHandler)); registry_med.register("fidelity_capture", Box::new(FidelityCapturingHandler { captures: captures_med.clone() })); let engine_med = PipelineEngine::new(registry_med, Arc::new(EventEmitter::new()), local_env()); - let config_med = RunConfig { logs_root: dir_med.path().to_path_buf(), cancel_token: None, dry_run: false }; + let config_med = RunConfig { logs_root: dir_med.path().to_path_buf(), cancel_token: None, dry_run: false, run_id: None, work_dir: None }; engine_med.run(&graph_med, &config_med).await.expect("run med"); let preambles_med = captures_med.preambles.lock().unwrap(); @@ -4740,6 +4850,8 @@ async fn fidelity_thread_id_fallback_to_previous_node_in_pipeline() { let config = RunConfig { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false, + run_id: None, + work_dir: None, }; engine.run(&graph, &config).await.expect("run"); @@ -4779,6 +4891,8 @@ async fn fidelity_thread_id_from_node_class_in_pipeline() { let config = RunConfig { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false, + run_id: None, + work_dir: None, }; engine.run(&graph, &config).await.expect("run"); @@ -4821,6 +4935,8 @@ async fn fidelity_edge_thread_id_override_in_pipeline() { let config = RunConfig { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false, + run_id: None, + work_dir: None, }; engine.run(&graph, &config).await.expect("run"); @@ -4864,6 +4980,8 @@ async fn fidelity_full_without_explicit_thread_id_uses_previous_node() { let config = RunConfig { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false, + run_id: None, + work_dir: None, }; engine.run(&graph, &config).await.expect("run"); @@ -4914,6 +5032,8 @@ async fn fidelity_from_parsed_dot_pipeline() { let config = RunConfig { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false, + run_id: None, + work_dir: None, }; engine.run(&graph, &config).await.expect("run"); @@ -4952,6 +5072,8 @@ async fn fidelity_checkpoint_roundtrip_preserves_fidelity() { let config = RunConfig { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false, + run_id: None, + work_dir: None, }; engine.run(&graph, &config).await.expect("run"); @@ -5007,6 +5129,8 @@ async fn fidelity_node_thread_id_overrides_edge_thread_id_in_pipeline() { let config = RunConfig { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false, + run_id: None, + work_dir: None, }; engine.run(&graph, &config).await.expect("run"); @@ -5076,6 +5200,8 @@ async fn fidelity_resume_preserves_context_values_across_checkpoint() { let config = RunConfig { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false, + run_id: None, + work_dir: None, }; engine .run_from_checkpoint(&graph, &config, &checkpoint) @@ -5261,6 +5387,8 @@ mod real_llm { let config = RunConfig { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false, + run_id: None, + work_dir: None, }; let outcome = tokio::time::timeout( @@ -5365,6 +5493,8 @@ mod real_llm { let config = RunConfig { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false, + run_id: None, + work_dir: None, }; let outcome = tokio::time::timeout( @@ -5499,6 +5629,8 @@ mod real_llm { let config = RunConfig { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false, + run_id: None, + work_dir: None, }; let outcome = tokio::time::timeout( @@ -5599,6 +5731,8 @@ mod real_llm { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false, + run_id: None, + work_dir: None, }; let outcome = tokio::time::timeout( @@ -5679,6 +5813,8 @@ async fn human_gate_freeform_only_routes_text() { let config = RunConfig { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false, + run_id: None, + work_dir: None, }; let outcome = engine.run(&graph, &config).await.expect("run should succeed"); @@ -5794,6 +5930,8 @@ async fn human_gate_freeform_with_fixed_choice_match() { let config = RunConfig { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false, + run_id: None, + work_dir: None, }; let outcome = engine.run(&graph, &config).await.expect("run should succeed"); @@ -5894,6 +6032,8 @@ async fn human_gate_freeform_fallback_on_unmatched_text() { let config = RunConfig { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false, + run_id: None, + work_dir: None, }; let outcome = engine.run(&graph, &config).await.expect("run should succeed"); @@ -6007,6 +6147,8 @@ async fn human_gate_freeform_sets_allow_freeform_on_question() { let config = RunConfig { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false, + run_id: None, + work_dir: None, }; let outcome = engine.run(&graph, &config).await.expect("run should succeed"); @@ -6096,6 +6238,8 @@ async fn human_gate_without_freeform_sets_allow_freeform_false() { let config = RunConfig { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false, + run_id: None, + work_dir: None, }; let outcome = engine.run(&graph, &config).await.expect("run should succeed"); @@ -6334,6 +6478,8 @@ async fn tool_hooks_pre_success_allows_pipeline_to_proceed() { let config = RunConfig { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false, + run_id: None, + work_dir: None, }; let outcome = engine.run(&graph, &config).await.expect("run should succeed"); @@ -6369,6 +6515,8 @@ async fn tool_hooks_pre_failure_skips_tool_call() { let config = RunConfig { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false, + run_id: None, + work_dir: None, }; engine.run(&graph, &config).await.expect("run should complete"); @@ -6407,6 +6555,8 @@ async fn tool_hooks_post_success_does_not_affect_outcome() { let config = RunConfig { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false, + run_id: None, + work_dir: None, }; let outcome = engine.run(&graph, &config).await.expect("run should succeed"); @@ -6437,6 +6587,8 @@ async fn tool_hooks_post_failure_does_not_block_pipeline() { let config = RunConfig { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false, + run_id: None, + work_dir: None, }; let outcome = engine.run(&graph, &config).await.expect("run should succeed"); @@ -6469,6 +6621,8 @@ async fn tool_hooks_graph_level_applies_to_all_nodes() { let config = RunConfig { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false, + run_id: None, + work_dir: None, }; let outcome = engine.run(&graph, &config).await.expect("run should succeed"); @@ -6504,6 +6658,8 @@ async fn tool_hooks_node_level_overrides_graph_level() { let config = RunConfig { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false, + run_id: None, + work_dir: None, }; let _outcome = engine.run(&graph, &config).await.expect("run should complete"); @@ -6549,6 +6705,8 @@ async fn tool_hooks_pre_receives_node_id_env_var() { let config = RunConfig { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false, + run_id: None, + work_dir: None, }; engine.run(&graph, &config).await.expect("run should succeed"); @@ -6641,6 +6799,8 @@ async fn attractor_e2e_with_real_llm() { let config = RunConfig { logs_root: logs_dir.path().to_path_buf(), cancel_token: None, dry_run: false, + run_id: None, + work_dir: None, }; let outcome = engine.run(&graph, &config).await.expect("run should succeed"); @@ -6744,6 +6904,8 @@ async fn run_fidelity_prompt_pipeline(fidelity: &str) -> String { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false, + run_id: None, + work_dir: None, }; engine.run(&graph, &config).await.expect("pipeline should succeed"); @@ -6883,6 +7045,8 @@ async fn large_context_values_are_offloaded_to_artifact_store() { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false, + run_id: None, + work_dir: None, }; let outcome = engine @@ -7054,6 +7218,8 @@ async fn artifact_pointers_rewritten_for_remote_execution_env() { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false, + run_id: None, + work_dir: None, }; let outcome = engine @@ -7164,6 +7330,8 @@ async fn node_dir_uses_visit_count_on_revisit() { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false, + run_id: None, + work_dir: None, }; let outcome = engine.run(&graph, &config).await.expect("pipeline should succeed"); @@ -7726,6 +7894,8 @@ async fn full_pipeline_with_cli_backend_node() { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false, + run_id: None, + work_dir: None, }; let outcome = engine.run(&graph, &config).await.expect("pipeline should succeed"); @@ -7810,6 +7980,8 @@ async fn stylesheet_backend_property_routes_to_cli() { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false, + run_id: None, + work_dir: None, }; let outcome = engine.run(&graph, &config).await.expect("pipeline should succeed");