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 <noreply@anthropic.com>
This commit is contained in:
Bryan Helmkamp 2026-02-28 10:51:37 -05:00
parent 172a4a89a5
commit e97c16f03f
9 changed files with 620 additions and 22 deletions

View file

@ -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<String>,
/// 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<String>,
}
impl Checkpoint {
@ -45,6 +48,7 @@ impl Checkpoint {
logs: context.logs_snapshot(),
node_outcomes,
next_node_id,
git_commit_sha: None,
}
}

View file

@ -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::<ExecutionEnvKind>())
.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]

View file

@ -467,6 +467,10 @@ pub struct RunConfig {
pub logs_root: PathBuf,
pub cancel_token: Option<Arc<AtomicBool>>,
pub dry_run: bool,
/// Pre-assigned run ID. Generated if `None`.
pub run_id: Option<String>,
/// Git worktree path for checkpoint commits.
pub work_dir: Option<PathBuf>,
}
/// The pipeline execution engine.
@ -695,12 +699,12 @@ impl PipelineEngine {
mut node_visits: HashMap<String, usize>,
) -> Result<Outcome> {
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,

View file

@ -0,0 +1,271 @@
use std::path::Path;
use std::process::Command;
use crate::error::{AttractorError, Result};
fn git_error(msg: impl Into<String>) -> 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<String> {
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<String> {
// 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();
}
}

View file

@ -176,13 +176,15 @@ fn resolve_hook(node: &Node, graph: &Graph, key: &str) -> Option<String> {
}
/// 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

View file

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

View file

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

View file

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

View file

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