From 90dc528c997de35f5a23718984c7e1c58b0a83b0 Mon Sep 17 00:00:00 2001 From: Bryan Helmkamp Date: Sat, 28 Feb 2026 21:45:48 -0500 Subject: [PATCH] Add shadow commits for Docker and Daytona (feature parity) - Change GitCheckpointMode::Remote to Remote(PathBuf) so both variants carry a repo path for MetadataStore shadow commits - Unify init_run and shadow write logic to work with either Host or Remote mode, eliminating Host-only gates - Add trailers (Arc-Run, Arc-Completed, Arc-Checkpoint) to remote checkpoint commits via write_file + git commit -F to avoid shell escaping issues with multi-line messages - Wire up meta_branch for Daytona in run.rs (was only set for worktree) - Fix sandbox name collisions by adding random hex suffix - Fix pre-existing build_router() test compilation errors from JWT auth - Add E2E tests for Host shadow branch and Daytona shadow branch Co-Authored-By: Claude Opus 4.6 --- Cargo.lock | 1 + crates/arc-attractor/Cargo.toml | 1 + crates/arc-attractor/src/cli/mod.rs | 9 +- crates/arc-attractor/src/cli/run.rs | 188 ++++++++++- crates/arc-attractor/src/daytona_env.rs | 6 +- crates/arc-attractor/src/engine.rs | 164 ++++++++-- crates/arc-attractor/src/git.rs | 297 +++++++++++++++++- crates/arc-attractor/src/server.rs | 2 +- .../tests/daytona_integration.rs | 126 +++++++- crates/arc-attractor/tests/integration.rs | 238 +++++++++++++- 10 files changed, 977 insertions(+), 55 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index d9d8c07a7..36ca1176b 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -134,6 +134,7 @@ version = "0.1.0" dependencies = [ "anyhow", "arc-agent", + "arc-git-storage", "arc-llm", "arc-util", "assert_cmd", diff --git a/crates/arc-attractor/Cargo.toml b/crates/arc-attractor/Cargo.toml index e2e2aa4fb..e540636b2 100644 --- a/crates/arc-attractor/Cargo.toml +++ b/crates/arc-attractor/Cargo.toml @@ -22,6 +22,7 @@ anyhow.workspace = true dotenvy.workspace = true arc-agent = { path = "../arc-agent" } arc-util = { path = "../arc-util" } +arc-git-storage = { path = "../arc-git-storage" } arc-llm = { path = "../arc-llm" } thiserror.workspace = true serde.workspace = true diff --git a/crates/arc-attractor/src/cli/mod.rs b/crates/arc-attractor/src/cli/mod.rs index d0d813dbf..e6ffca372 100644 --- a/crates/arc-attractor/src/cli/mod.rs +++ b/crates/arc-attractor/src/cli/mod.rs @@ -74,8 +74,9 @@ pub enum Command { #[derive(Args)] pub struct RunArgs { - /// Path to a .dot pipeline file or .toml task config - pub pipeline: PathBuf, + /// Path to a .dot pipeline file or .toml task config (not required with --run-branch) + #[arg(required_unless_present = "run_branch")] + pub pipeline: Option, /// Log/artifact directory #[arg(long)] @@ -93,6 +94,10 @@ pub struct RunArgs { #[arg(long)] pub resume: Option, + /// Resume from a git run branch (reads checkpoint and graph from metadata branch) + #[arg(long, conflicts_with = "resume")] + pub run_branch: Option, + /// Override default LLM model #[arg(long)] pub model: Option, diff --git a/crates/arc-attractor/src/cli/run.rs b/crates/arc-attractor/src/cli/run.rs index 993d98f2a..b4ea64203 100644 --- a/crates/arc-attractor/src/cli/run.rs +++ b/crates/arc-attractor/src/cli/run.rs @@ -43,13 +43,21 @@ struct CostAccumulator { /// /// Returns an error if the pipeline cannot be read, parsed, validated, or executed. pub async fn run_command(args: RunArgs, styles: &'static Styles) -> anyhow::Result<()> { + // Handle --run-branch resume: read everything from git metadata + if let Some(branch) = args.run_branch.clone() { + return run_from_branch(args, &branch, styles).await; + } + + let pipeline_path = args.pipeline.as_ref() + .ok_or_else(|| anyhow::anyhow!("--pipeline is required unless --run-branch is provided"))?; + // 0. Load task config if TOML, resolve DOT path, run setup - let (dot_path, task_cfg) = if args.pipeline.extension().is_some_and(|ext| ext == "toml") { - let cfg = task_config::load_task_config(&args.pipeline)?; - let dot = task_config::resolve_graph_path(&args.pipeline, &cfg.graph); + let (dot_path, task_cfg) = if pipeline_path.extension().is_some_and(|ext| ext == "toml") { + let cfg = task_config::load_task_config(pipeline_path)?; + let dot = task_config::resolve_graph_path(pipeline_path, &cfg.graph); (dot, Some(cfg)) } else { - (args.pipeline.clone(), None) + (pipeline_path.clone(), None) }; if let Some(ref cfg) = task_cfg { @@ -128,8 +136,8 @@ pub async fn run_command(args: RunArgs, styles: &'static Styles) -> anyhow::Resu tokio::fs::create_dir_all(&logs_dir).await?; tokio::fs::write(logs_dir.join("graph.dot"), &source).await?; tokio::fs::write(logs_dir.join("run.pid"), std::process::id().to_string()).await?; - if args.pipeline.extension().is_some_and(|ext| ext == "toml") { - if let Ok(toml_contents) = tokio::fs::read(&args.pipeline).await { + if pipeline_path.extension().is_some_and(|ext| ext == "toml") { + if let Ok(toml_contents) = tokio::fs::read(pipeline_path).await { tokio::fs::write(logs_dir.join("task.toml"), toml_contents).await?; } } @@ -500,6 +508,12 @@ pub async fn run_command(args: RunArgs, styles: &'static Styles) -> anyhow::Resu let run_id = worktree_run_id .or(daytona_run_id) .unwrap_or_else(|| ulid::Ulid::new().to_string()); + // Set up metadata branch for git checkpointing (host or remote) + let meta_branch = if worktree_work_dir.is_some() || daytona_base_sha.is_some() { + Some(crate::git::MetadataStore::branch_name(&run_id)) + } else { + None + }; let config = RunConfig { logs_root: logs_dir.clone(), cancel_token: None, @@ -510,11 +524,12 @@ pub async fn run_command(args: RunArgs, styles: &'static Styles) -> anyhow::Resu worktree_work_dir.map(GitCheckpointMode::Host) } ExecutionEnvKind::Daytona => { - daytona_base_sha.as_ref().map(|_| GitCheckpointMode::Remote) + daytona_base_sha.as_ref().map(|_| GitCheckpointMode::Remote(original_cwd.clone())) } }, base_sha: worktree_base_sha.or(daytona_base_sha), run_branch: worktree_branch.or(daytona_branch), + meta_branch, }; let run_start = Instant::now(); @@ -669,6 +684,165 @@ async fn setup_daytona_git( Ok((run_id, base_sha, branch_name)) } +/// Resume a pipeline run from a git run branch. +/// +/// Reads the checkpoint, manifest, and graph DOT from the metadata branch +/// (`refs/arc/{run_id}`), re-attaches a worktree to the existing run branch, +/// and resumes execution via `run_from_checkpoint()`. +async fn run_from_branch(args: RunArgs, run_branch: &str, styles: &'static Styles) -> anyhow::Result<()> { + // Extract run_id from branch name: "arc/run/{run_id}" -> "{run_id}" + let run_id = run_branch + .strip_prefix("arc/run/") + .ok_or_else(|| anyhow::anyhow!( + "invalid run branch format: expected 'arc/run/', got '{run_branch}'" + ))? + .to_string(); + + let original_cwd = std::env::current_dir()?; + + // Read checkpoint from metadata branch + let checkpoint = crate::git::MetadataStore::read_checkpoint(&original_cwd, &run_id)? + .ok_or_else(|| anyhow::anyhow!("no checkpoint found on metadata branch for run {run_id}"))?; + + // Read graph DOT from metadata branch + let source = crate::git::MetadataStore::read_graph_dot(&original_cwd, &run_id)? + .ok_or_else(|| anyhow::anyhow!("no graph.dot found on metadata branch for run {run_id}"))?; + + // If --pipeline was also provided, use it instead (allows overriding) + let source = if let Some(ref pipeline_path) = args.pipeline { + super::read_dot_file(pipeline_path)? + } else { + source + }; + + let (graph, diagnostics) = crate::pipeline::PipelineBuilder::new().prepare(&source)?; + + eprintln!( + "{bold}Resuming pipeline:{reset} {} from branch {dim}{run_branch}{reset}", + graph.name, + bold = styles.bold, dim = styles.dim, reset = styles.reset, + ); + + super::print_diagnostics(&diagnostics, styles); + if diagnostics.iter().any(|d| d.severity == crate::validation::Severity::Error) { + anyhow::bail!("Validation failed"); + } + + // Set up logs directory + let logs_dir = args.logs_dir.unwrap_or_else(|| { + let base = dirs::home_dir() + .expect("could not determine home directory") + .join(".attractor") + .join("logs"); + base.join(format!( + "arc-resume-{}", + chrono::Local::now().format("%Y%m%d-%H%M%S") + )) + }); + tokio::fs::create_dir_all(&logs_dir).await?; + tokio::fs::write(logs_dir.join("graph.dot"), &source).await?; + + // Re-attach worktree to the existing run branch + let worktree_path = logs_dir.join("worktree"); + crate::git::add_worktree(&original_cwd, &worktree_path, run_branch) + .map_err(|e| anyhow::anyhow!("failed to attach worktree to {run_branch}: {e}"))?; + std::env::set_current_dir(&worktree_path)?; + + let base_sha = crate::git::MetadataStore::read_manifest(&original_cwd, &run_id)? + .and_then(|m| m.get("base_sha").and_then(|v| v.as_str()).map(String::from)); + + // Build minimal execution environment (local only for now) + let emitter = Arc::new(EventEmitter::new()); + let execution_env: Arc = { + let mut env = arc_agent::LocalExecutionEnvironment::new(worktree_path.clone()); + let emitter_cb = Arc::clone(&emitter); + env.set_event_callback(Arc::new(move |event| { + emitter_cb.emit(&crate::event::PipelineEvent::ExecutionEnv { event }); + })); + Arc::new(env) + }; + + // Build interviewer + let interviewer: Arc = if args.auto_approve { + Arc::new(crate::interviewer::auto_approve::AutoApproveInterviewer) + } else { + Arc::new(crate::interviewer::console::ConsoleInterviewer::new(styles)) + }; + + // Build engine with a backend + let dry_run_mode = args.dry_run || arc_llm::client::Client::from_env().await + .map(|c| c.provider_names().is_empty()) + .unwrap_or(true); + + let model = args.model.unwrap_or_else(|| "claude-opus-4-6".to_string()); + let provider_enum = args.provider + .as_deref() + .map(|s| s.parse::()) + .transpose() + .map_err(|e| anyhow::anyhow!("{e}"))? + .unwrap_or(arc_llm::provider::Provider::Anthropic); + + let registry = crate::handler::default_registry(interviewer.clone(), || { + if dry_run_mode { + None + } else { + let api = AgentBackend::new(model.clone(), provider_enum, args.verbose, styles); + let cli = CliBackend::new(model.clone(), provider_enum); + Some(Box::new(BackendRouter::new(Box::new(api), cli))) + } + }); + let engine = crate::engine::PipelineEngine::with_interviewer( + registry, + Arc::clone(&emitter), + interviewer, + Arc::clone(&execution_env), + ); + + let meta_branch = Some(crate::git::MetadataStore::branch_name(&run_id)); + let config = RunConfig { + logs_root: logs_dir.clone(), + cancel_token: None, + dry_run: dry_run_mode, + run_id, + git_checkpoint: Some(GitCheckpointMode::Host(worktree_path.clone())), + base_sha, + run_branch: Some(run_branch.to_string()), + meta_branch, + }; + + let run_start = Instant::now(); + let engine_result = engine.run_from_checkpoint(&graph, &config, &checkpoint).await; + let run_duration_ms = run_start.elapsed().as_millis() as u64; + + // Clean up + let _ = std::env::set_current_dir(&original_cwd); + let _ = crate::git::remove_worktree(&original_cwd, &worktree_path); + + let outcome = engine_result?; + + eprintln!( + "\n{bold}=== Pipeline Result ==={reset}", + bold = styles.bold, reset = styles.reset, + ); + let status_str = outcome.status.to_string().to_uppercase(); + let status_color = match outcome.status { + StageStatus::Success | StageStatus::PartialSuccess => styles.green, + _ => styles.red, + }; + eprintln!("Status: {status_color}{status_str}{reset}", reset = styles.reset); + eprintln!("Duration: {}", super::format_duration_human(run_duration_ms)); + eprintln!( + "{dim}Logs: {}{reset}", + logs_dir.display(), + dim = styles.dim, reset = styles.reset, + ); + + match outcome.status { + StageStatus::Success | StageStatus::PartialSuccess => Ok(()), + _ => std::process::exit(1), + } +} + #[cfg(test)] mod tests { #[test] diff --git a/crates/arc-attractor/src/daytona_env.rs b/crates/arc-attractor/src/daytona_env.rs index c363fe56e..77aff30d8 100644 --- a/crates/arc-attractor/src/daytona_env.rs +++ b/crates/arc-attractor/src/daytona_env.rs @@ -4,6 +4,7 @@ use std::time::Instant; use arc_agent::execution_env::{format_lines_numbered, DirEntry, ExecEnvEventCallback, ExecResult, ExecutionEnvEvent, ExecutionEnvironment, GrepOptions}; use async_trait::async_trait; +use rand::Rng; use serde::Deserialize; const WORKING_DIRECTORY: &str = "/home/daytona/workspace"; @@ -90,8 +91,9 @@ impl DaytonaExecutionEnvironment { /// Build `SandboxBaseParams` from config, generating a unique sandbox name. fn base_params(&self) -> daytona_sdk::SandboxBaseParams { let name = format!( - "arc-{}", - chrono::Utc::now().format("%Y%m%d-%H%M%S-%3f") + "arc-{}-{:04x}", + chrono::Utc::now().format("%Y%m%d-%H%M%S"), + rand::thread_rng().gen_range(0..0x10000u32), ); daytona_sdk::SandboxBaseParams { name: Some(name), diff --git a/crates/arc-attractor/src/engine.rs b/crates/arc-attractor/src/engine.rs index 54b500718..f94b965fe 100644 --- a/crates/arc-attractor/src/engine.rs +++ b/crates/arc-attractor/src/engine.rs @@ -10,6 +10,8 @@ use chrono::Utc; use futures::FutureExt; use rand::Rng; +use arc_git_storage::trailerlink::{self, Trailer}; + use crate::artifact::{offload_large_values, sync_artifacts_to_env, ArtifactStore}; use crate::checkpoint::Checkpoint; use crate::condition::evaluate_condition; @@ -239,27 +241,32 @@ pub fn resolve_thread_id( // --- Run directory helpers (spec 5.6) --- -/// Write manifest.json at the start of a pipeline run. -fn write_manifest(logs_root: &Path, graph: &Graph, run_branch: Option<&str>) { +/// Write manifest.json at the start of a pipeline run. Returns the manifest value. +fn write_manifest(logs_root: &Path, graph: &Graph, config: &RunConfig) -> serde_json::Value { let pipeline_name = if graph.name.is_empty() { "unnamed" } else { &graph.name }; let mut manifest = serde_json::json!({ + "run_id": config.run_id, "pipeline_name": pipeline_name, "goal": graph.goal(), "start_time": Utc::now().to_rfc3339(), "node_count": graph.nodes.len(), "edge_count": graph.edges.len(), }); - if let Some(branch) = run_branch { - manifest["run_branch"] = serde_json::Value::String(branch.to_string()); + if let Some(ref branch) = config.run_branch { + manifest["run_branch"] = serde_json::Value::String(branch.clone()); + } + if let Some(ref base) = config.base_sha { + manifest["base_sha"] = serde_json::Value::String(base.clone()); } if let Ok(json) = serde_json::to_string_pretty(&manifest) { let _ = std::fs::create_dir_all(logs_root); let _ = std::fs::write(logs_root.join("manifest.json"), json); } + manifest } /// Return the directory for a node's logs. @@ -471,13 +478,21 @@ pub enum GitCheckpointMode { /// Run git commands on the host filesystem (local & Docker bind-mount). Host(PathBuf), /// Run git commands inside the remote execution environment via `exec_command`. - Remote, + /// The `PathBuf` is the host repo path used for `MetadataStore` (shadow commits). + Remote(PathBuf), } /// Run a git checkpoint commit on the host filesystem (local/Docker bind-mount). -async fn git_checkpoint_host(work_dir: PathBuf, run_id: String, node_id: String, status: String) -> Option { +async fn git_checkpoint_host( + work_dir: PathBuf, + run_id: String, + node_id: String, + status: String, + completed_count: usize, + shadow_sha: Option, +) -> Option { match tokio::task::spawn_blocking(move || { - crate::git::checkpoint_commit(&work_dir, &run_id, &node_id, &status) + crate::git::checkpoint_commit(&work_dir, &run_id, &node_id, &status, completed_count, shadow_sha.as_deref()) }).await { Ok(Ok(sha)) => Some(sha), Ok(Err(_)) | Err(_) => None, @@ -495,19 +510,41 @@ async fn git_diff_host(work_dir: PathBuf, base: String) -> Option { } /// Run a git checkpoint commit inside a remote execution environment. -async fn git_checkpoint_remote(exec_env: &dyn ExecutionEnvironment, run_id: &str, node_id: &str, status: &str) -> Option { +async fn git_checkpoint_remote( + exec_env: &dyn ExecutionEnvironment, + run_id: &str, + node_id: &str, + status: &str, + completed_count: usize, + shadow_sha: Option, +) -> Option { // Stage everything let add_result = exec_env.exec_command("git add -A", 30_000, None, None, None).await; if add_result.as_ref().map_or(true, |r| r.exit_code != 0) { return None; } - // Commit with arc identity - let message = format!("arc({run_id}): {node_id} ({status})"); - let commit_cmd = format!( - "git -c user.name=arc -c user.email=arc@local commit --allow-empty -m '{message}'" - ); - let commit_result = exec_env.exec_command(&commit_cmd, 30_000, None, None, None).await; + // Build commit message with trailers (same format as checkpoint_commit in git.rs) + let subject = format!("arc({run_id}): {node_id} ({status})"); + let completed_str = completed_count.to_string(); + let mut trailers = vec![ + Trailer { key: "Arc-Run", value: run_id }, + Trailer { key: "Arc-Completed", value: &completed_str }, + ]; + let shadow_sha_ref = shadow_sha.as_deref().unwrap_or(""); + if shadow_sha.is_some() { + trailers.push(Trailer { key: "Arc-Checkpoint", value: shadow_sha_ref }); + } + let message = trailerlink::format_message(&subject, "", &trailers); + + // Write message to temp file in sandbox to avoid shell escaping issues + if exec_env.write_file("/tmp/arc-commit-msg", &message).await.is_err() { + return None; + } + + // Commit with arc identity using the message file + let commit_cmd = "git -c user.name=arc -c user.email=arc@local commit --allow-empty -F /tmp/arc-commit-msg"; + let commit_result = exec_env.exec_command(commit_cmd, 30_000, None, None, None).await; if commit_result.as_ref().map_or(true, |r| r.exit_code != 0) { return None; } @@ -542,6 +579,8 @@ pub struct RunConfig { pub base_sha: Option, /// Git branch name for the run (e.g. `arc/run/{run_id}`). pub run_branch: Option, + /// Metadata branch name for git-native checkpoint storage (e.g. `refs/arc/{run_id}`). + pub meta_branch: Option, } /// The pipeline execution engine. @@ -786,7 +825,23 @@ impl PipelineEngine { self.inform(&format!("Pipeline started: {}", graph.name), "pipeline"); // Write manifest.json (spec 5.6) - write_manifest(&config.logs_root, graph, config.run_branch.as_deref()); + let manifest = write_manifest(&config.logs_root, graph, config); + + // Initialize metadata branch for git-native checkpoint storage (best-effort) + if config.meta_branch.is_some() { + let store_path = match config.git_checkpoint { + Some(GitCheckpointMode::Host(ref p)) | Some(GitCheckpointMode::Remote(ref p)) => Some(p), + None => None, + }; + if let Some(repo_path) = store_path { + let store = crate::git::MetadataStore::new(repo_path); + let manifest_bytes = serde_json::to_vec_pretty(&manifest).unwrap_or_default(); + let dot_source = std::fs::read(config.logs_root.join("graph.dot")).unwrap_or_default(); + if let Err(e) = store.init_run(&config.run_id, &manifest_bytes, &dot_source) { + eprintln!("Warning: metadata branch init failed: {e}"); + } + } + } // Compute effective max-node-visits limit: // graph attr > 0 → use it; else dry_run → 10; else 0 (disabled) @@ -1096,18 +1151,51 @@ impl PipelineEngine { }); } - // Step 6b: Git checkpoint commit + // Step 6b: Write shadow branch first, then run branch commit with trailer if let Some(ref mode) = config.git_checkpoint { + // Shadow commit (best-effort): extract repo path from either variant + let shadow_sha: Option = if config.meta_branch.is_some() { + let repo_path = match mode { + GitCheckpointMode::Host(ref p) | GitCheckpointMode::Remote(ref p) => p, + }; + let store = crate::git::MetadataStore::new(repo_path); + serde_json::to_vec_pretty(&checkpoint).ok().and_then(|cp_json| { + let artifact_entries: Vec<(String, Vec)> = artifact_store.list().iter() + .filter_map(|info| { + info.file_path.as_ref().and_then(|path| { + std::fs::read(path).ok().map(|data| { + (format!("artifacts/{}.json", info.id), data) + }) + }) + }) + .collect(); + let artifact_refs: Vec<(&str, &[u8])> = artifact_entries.iter() + .map(|(k, v)| (k.as_str(), v.as_slice())) + .collect(); + match store.write_checkpoint(&config.run_id, &cp_json, &artifact_refs) { + Ok(sha) => Some(sha), + Err(e) => { + context.append_log(format!("metadata checkpoint write failed: {e}")); + None + } + } + }) + } else { + None + }; + + // Run branch commit with Arc-Meta trailer pointing to shadow commit let rid = run_id.clone(); let nid = node.id.clone(); let status_str = outcome.status.to_string(); + let completed_count = completed_nodes.len(); let commit_result = match mode { GitCheckpointMode::Host(work_dir) => { - git_checkpoint_host(work_dir.clone(), rid, nid, status_str).await + git_checkpoint_host(work_dir.clone(), rid, nid, status_str, completed_count, shadow_sha).await } - GitCheckpointMode::Remote => { - git_checkpoint_remote(&*self.services.execution_env, &run_id, &node.id, &outcome.status.to_string()).await + GitCheckpointMode::Remote(_) => { + git_checkpoint_remote(&*self.services.execution_env, &run_id, &node.id, &outcome.status.to_string(), completed_count, shadow_sha).await } }; @@ -1134,7 +1222,7 @@ impl PipelineEngine { GitCheckpointMode::Host(work_dir) => { git_diff_host(work_dir.clone(), diff_base).await } - GitCheckpointMode::Remote => { + GitCheckpointMode::Remote(_) => { git_diff_remote(&*self.services.execution_env, &diff_base).await } }; @@ -1217,7 +1305,7 @@ impl PipelineEngine { GitCheckpointMode::Host(work_dir) => { git_diff_host(work_dir.clone(), base.clone()).await } - GitCheckpointMode::Remote => { + GitCheckpointMode::Remote(_) => { git_diff_remote(&*self.services.execution_env, base).await } }; @@ -1900,6 +1988,7 @@ mod tests { git_checkpoint: None, base_sha: None, run_branch: None, + meta_branch: None, }; let outcome = engine.run(&g, &config).await.unwrap(); assert_eq!(outcome.status, StageStatus::Success); @@ -1918,6 +2007,7 @@ mod tests { git_checkpoint: None, base_sha: None, run_branch: None, + meta_branch: None, }; engine.run(&g, &config).await.unwrap(); let checkpoint_path = dir.path().join("checkpoint.json"); @@ -1945,6 +2035,7 @@ mod tests { git_checkpoint: None, base_sha: None, run_branch: None, + meta_branch: None, }; engine.run(&g, &config).await.unwrap(); @@ -1967,6 +2058,7 @@ mod tests { git_checkpoint: None, base_sha: None, run_branch: None, + meta_branch: None, }; let result = engine.run(&g, &config).await; assert!(result.is_err()); @@ -1985,6 +2077,7 @@ mod tests { git_checkpoint: None, base_sha: None, run_branch: None, + meta_branch: None, }; engine.run(&g, &config).await.unwrap(); @@ -2016,6 +2109,7 @@ mod tests { git_checkpoint: None, base_sha: None, run_branch: None, + meta_branch: None, }; let outcome = engine.run(&g, &config).await.unwrap(); assert_eq!(outcome.status, StageStatus::Success); @@ -2073,6 +2167,7 @@ mod tests { git_checkpoint: None, base_sha: None, run_branch: None, + meta_branch: None, }; engine.run(&g, &config).await.unwrap(); @@ -2148,6 +2243,7 @@ mod tests { git_checkpoint: None, base_sha: None, run_branch: None, + meta_branch: None, }; engine.run(&g, &config).await.unwrap(); @@ -2175,6 +2271,7 @@ mod tests { git_checkpoint: None, base_sha: None, run_branch: None, + meta_branch: None, }; engine.run(&g, &config).await.unwrap(); @@ -2199,6 +2296,7 @@ mod tests { git_checkpoint: None, base_sha: None, run_branch: None, + meta_branch: None, }; engine.run(&g, &config).await.unwrap(); @@ -2343,7 +2441,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, run_id: "test-run".into(), git_checkpoint: None, base_sha: None, run_branch: None }; + let config = RunConfig { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false, run_id: "test-run".into(), git_checkpoint: None, base_sha: None, run_branch: None, meta_branch: None }; engine.run(&g, &config).await.unwrap(); let manifest_path = dir.path().join("manifest.json"); @@ -2365,7 +2463,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, run_id: "test-run".into(), git_checkpoint: None, base_sha: None, run_branch: None }; + let config = RunConfig { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false, run_id: "test-run".into(), git_checkpoint: None, base_sha: None, run_branch: None, meta_branch: None }; engine.run(&g, &config).await.unwrap(); let manifest_path = dir.path().join("manifest.json"); @@ -2401,7 +2499,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, run_id: "test-run".into(), git_checkpoint: None, base_sha: None, run_branch: None }; + let config = RunConfig { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false, run_id: "test-run".into(), git_checkpoint: None, base_sha: None, run_branch: None, meta_branch: None }; let outcome = engine.run(&g, &config).await.unwrap(); assert_eq!(outcome.status, StageStatus::Success); @@ -2437,7 +2535,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, run_id: "test-run".into(), git_checkpoint: None, base_sha: None, run_branch: None }; + let config = RunConfig { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false, run_id: "test-run".into(), git_checkpoint: None, base_sha: None, run_branch: None, meta_branch: None }; let result = engine.run(&g, &config).await; assert!(result.is_ok()); @@ -2476,7 +2574,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, run_id: "test-run".into(), git_checkpoint: None, base_sha: None, run_branch: None }; + let config = RunConfig { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false, run_id: "test-run".into(), git_checkpoint: None, base_sha: None, run_branch: None, meta_branch: None }; let result = engine.run(&g, &config).await; assert!(result.is_ok()); @@ -2510,7 +2608,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, run_id: "test-run".into(), git_checkpoint: None, base_sha: None, run_branch: None }; + let config = RunConfig { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false, run_id: "test-run".into(), git_checkpoint: None, base_sha: None, run_branch: None, meta_branch: None }; let outcome = engine.run(&g, &config).await.unwrap(); assert_eq!(outcome.status, StageStatus::Success); } @@ -2541,7 +2639,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, run_id: "test-run".into(), git_checkpoint: None, base_sha: None, run_branch: None }; + let config = RunConfig { logs_root: dir.path().to_path_buf(), cancel_token: None, dry_run: false, run_id: "test-run".into(), git_checkpoint: None, base_sha: None, run_branch: None, meta_branch: None }; let outcome = engine.run(&g, &config).await.unwrap(); assert_eq!(outcome.status, StageStatus::Success); @@ -2599,6 +2697,7 @@ mod tests { git_checkpoint: None, base_sha: None, run_branch: None, + meta_branch: None, }; engine.run(&g, &config).await.unwrap(); // Give spawned inform tasks time to complete @@ -2630,6 +2729,7 @@ mod tests { git_checkpoint: None, base_sha: None, run_branch: None, + meta_branch: None, }; engine.run(&g, &config).await.unwrap(); // Give spawned inform tasks time to complete @@ -2659,6 +2759,7 @@ mod tests { git_checkpoint: None, base_sha: None, run_branch: None, + meta_branch: None, }; let outcome = engine.run(&g, &config).await.unwrap(); assert_eq!(outcome.status, StageStatus::Success); @@ -2680,6 +2781,7 @@ mod tests { git_checkpoint: None, base_sha: None, run_branch: None, + meta_branch: None, }; let result = engine.run(&g, &config).await; assert!(result.is_err()); @@ -2700,6 +2802,7 @@ mod tests { git_checkpoint: None, base_sha: None, run_branch: None, + meta_branch: None, }; let outcome = engine.run(&g, &config).await.unwrap(); assert_eq!(outcome.status, StageStatus::Success); @@ -2732,6 +2835,7 @@ mod tests { git_checkpoint: None, base_sha: None, run_branch: None, + meta_branch: None, }; // Set cancel after a short delay (while the slow handler is running) @@ -2808,6 +2912,7 @@ mod tests { git_checkpoint: None, base_sha: None, run_branch: None, + meta_branch: None, }; let result = engine.run(&g, &config).await; assert!(result.is_err()); @@ -2831,6 +2936,7 @@ mod tests { git_checkpoint: None, base_sha: None, run_branch: None, + meta_branch: None, }; let result = engine.run(&g, &config).await; assert!(result.is_err()); @@ -2858,6 +2964,7 @@ mod tests { git_checkpoint: None, base_sha: None, run_branch: None, + meta_branch: None, }; let result = engine.run(&g, &config).await; assert!(result.is_err()); @@ -2946,6 +3053,7 @@ mod tests { git_checkpoint: None, base_sha: None, run_branch: None, + meta_branch: 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 index 703030969..95c490822 100644 --- a/crates/arc-attractor/src/git.rs +++ b/crates/arc-attractor/src/git.rs @@ -1,6 +1,12 @@ use std::path::Path; use std::process::Command; +use arc_git_storage::branchstore::BranchStore; +use arc_git_storage::gitobj::Store; +use arc_git_storage::trailerlink::{self, Trailer}; +use git2::{Repository, Signature}; + +use crate::checkpoint::Checkpoint; use crate::error::{AttractorError, Result}; fn git_error(msg: impl Into) -> AttractorError { @@ -93,13 +99,16 @@ pub fn remove_worktree(repo: &Path, path: &Path) -> Result<()> { Ok(()) } -/// Stage all changes and commit in `work_dir` with a structured message. +/// Stage all changes and commit in `work_dir` with a structured message +/// including trailers for completed node count and shadow commit pointer. /// Returns the new commit SHA. pub fn checkpoint_commit( work_dir: &Path, run_id: &str, node_id: &str, status: &str, + completed_count: usize, + shadow_sha: Option<&str>, ) -> Result { // Stage everything let output = Command::new("git") @@ -113,8 +122,19 @@ pub fn checkpoint_commit( return Err(git_error(format!("git add failed: {stderr}"))); } + // Build commit message with trailers + let subject = format!("arc({run_id}): {node_id} ({status})"); + let completed_str = completed_count.to_string(); + let mut trailers = vec![ + Trailer { key: "Arc-Run", value: run_id }, + Trailer { key: "Arc-Completed", value: &completed_str }, + ]; + if let Some(sha) = shadow_sha { + trailers.push(Trailer { key: "Arc-Checkpoint", value: sha }); + } + let message = trailerlink::format_message(&subject, "", &trailers); + // 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", @@ -152,6 +172,117 @@ pub fn diff_against(work_dir: &Path, base: &str) -> Result { Ok(String::from_utf8_lossy(&output.stdout).to_string()) } +/// Git-native metadata storage for pipeline runs. +/// +/// Stores checkpoint data, manifests, and graph DOT on an orphan branch +/// (`arc/{run_id}`) so that runs can be resumed from git alone. +pub struct MetadataStore { + repo_path: std::path::PathBuf, +} + +impl MetadataStore { + pub fn new(repo_path: impl Into) -> Self { + Self { repo_path: repo_path.into() } + } + + /// Returns the branch ref name for a run: `refs/arc/{run_id}`. + pub fn branch_name(run_id: &str) -> String { + format!("refs/arc/{run_id}") + } + + fn open_store(&self) -> Result<(Store, Signature<'static>)> { + let repo = Repository::discover(&self.repo_path) + .map_err(|e| git_error(format!("failed to open repo: {e}")))?; + let store = Store::new(repo); + let sig = Signature::now("arc", "arc@local") + .map_err(|e| git_error(format!("failed to create signature: {e}")))?; + Ok((store, sig)) + } + + /// Initialize a run's metadata branch with manifest and graph DOT. + pub fn init_run(&self, run_id: &str, manifest_json: &[u8], graph_dot: &[u8]) -> Result<()> { + let (store, sig) = self.open_store()?; + let branch = Self::branch_name(run_id); + let bs = BranchStore::new(&store, &branch, &sig); + bs.ensure_branch().map_err(|e| git_error(format!("ensure_branch failed: {e}")))?; + bs.write_entries( + &[("manifest.json", manifest_json), ("graph.dot", graph_dot)], + "init run", + ).map_err(|e| git_error(format!("write_entries failed: {e}")))?; + Ok(()) + } + + /// Write checkpoint data (and optional artifacts) to the metadata branch. + /// Returns the SHA of the new commit on the shadow branch. + pub fn write_checkpoint( + &self, + run_id: &str, + checkpoint_json: &[u8], + artifacts: &[(&str, &[u8])], + ) -> Result { + let (store, sig) = self.open_store()?; + let branch = Self::branch_name(run_id); + let bs = BranchStore::new(&store, &branch, &sig); + let mut entries: Vec<(&str, &[u8])> = vec![("checkpoint.json", checkpoint_json)]; + entries.extend_from_slice(artifacts); + let oid = bs.write_entries(&entries, "checkpoint") + .map_err(|e| git_error(format!("write_entries failed: {e}")))?; + Ok(oid.to_string()) + } + + /// Read a single file from the metadata branch. Returns `None` if branch or path doesn't exist. + fn read_file(repo_path: &Path, run_id: &str, path: &str) -> Result>> { + let repo = match Repository::discover(repo_path) { + Ok(r) => r, + Err(_) => return Ok(None), + }; + let store = Store::new(repo); + let sig = Signature::now("arc", "arc@local") + .map_err(|e| git_error(format!("failed to create signature: {e}")))?; + let branch = Self::branch_name(run_id); + let bs = BranchStore::new(&store, &branch, &sig); + bs.read_entry(path) + .map_err(|e| git_error(format!("read_entry failed: {e}"))) + } + + /// Read a checkpoint from the metadata branch. Returns `None` if branch or file doesn't exist. + pub fn read_checkpoint(repo_path: &Path, run_id: &str) -> Result> { + match Self::read_file(repo_path, run_id, "checkpoint.json")? { + Some(bytes) => { + let cp: Checkpoint = serde_json::from_slice(&bytes) + .map_err(|e| AttractorError::Checkpoint(format!("deserialize failed: {e}")))?; + Ok(Some(cp)) + } + None => Ok(None), + } + } + + /// Read the manifest JSON from the metadata branch. Returns `None` if not found. + pub fn read_manifest(repo_path: &Path, run_id: &str) -> Result> { + match Self::read_file(repo_path, run_id, "manifest.json")? { + Some(bytes) => { + let val: serde_json::Value = serde_json::from_slice(&bytes) + .map_err(|e| git_error(format!("manifest deserialize failed: {e}")))?; + Ok(Some(val)) + } + None => Ok(None), + } + } + + /// Read the graph DOT source from the metadata branch. Returns `None` if not found. + pub fn read_graph_dot(repo_path: &Path, run_id: &str) -> Result> { + match Self::read_file(repo_path, run_id, "graph.dot")? { + Some(bytes) => Ok(Some(String::from_utf8_lossy(&bytes).to_string())), + None => Ok(None), + } + } + + /// Read an artifact from the metadata branch. Returns `None` if not found. + pub fn read_artifact(repo_path: &Path, run_id: &str, key: &str) -> Result>> { + Self::read_file(repo_path, run_id, &format!("artifacts/{key}.json")) + } +} + #[cfg(test)] mod tests { use super::*; @@ -233,7 +364,7 @@ mod tests { } #[test] - fn checkpoint_commit_creates_commit() { + fn checkpoint_commit_creates_commit_with_trailers() { let dir = tempfile::tempdir().unwrap(); init_repo(dir.path()); create_branch(dir.path(), "run-branch").unwrap(); @@ -244,11 +375,13 @@ mod tests { // 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(); + // Simulate a shadow commit SHA + let shadow_sha = "abcdef1234567890abcdef1234567890abcdef12"; + let sha = checkpoint_commit(&wt_path, "run1", "nodeA", "success", 3, Some(shadow_sha)).unwrap(); assert_eq!(sha.len(), 40); assert!(sha.chars().all(|c| c.is_ascii_hexdigit())); - // Verify commit message + // Verify commit message subject line let output = Command::new("git") .args(["log", "--oneline", "-1"]) .current_dir(&wt_path) @@ -257,14 +390,49 @@ mod tests { let log = String::from_utf8_lossy(&output.stdout); assert!(log.contains("arc(run1): nodeA (success)")); + // Verify trailers by reading full message (trim trailing newlines from git log) + let output = Command::new("git") + .args(["log", "--format=%B", "-1"]) + .current_dir(&wt_path) + .output() + .unwrap(); + let full_msg = String::from_utf8_lossy(&output.stdout).trim().to_string(); + assert_eq!(trailerlink::parse(&full_msg, "Arc-Run"), Some("run1")); + assert_eq!(trailerlink::parse(&full_msg, "Arc-Completed"), Some("3")); + assert_eq!(trailerlink::parse(&full_msg, "Arc-Checkpoint"), Some(shadow_sha)); + + remove_worktree(dir.path(), &wt_path).unwrap(); + } + + #[test] + fn checkpoint_commit_without_shadow_sha() { + let dir = tempfile::tempdir().unwrap(); + init_repo(dir.path()); + create_branch(dir.path(), "run-branch2").unwrap(); + + let wt_path = dir.path().join("worktree"); + add_worktree(dir.path(), &wt_path, "run-branch2").unwrap(); + + let sha = checkpoint_commit(&wt_path, "run2", "nodeB", "completed", 1, None).unwrap(); + assert_eq!(sha.len(), 40); + + // Verify Arc-Completed trailer present but no Arc-Meta + let output = Command::new("git") + .args(["log", "--format=%B", "-1"]) + .current_dir(&wt_path) + .output() + .unwrap(); + let full_msg = String::from_utf8_lossy(&output.stdout).trim().to_string(); + assert_eq!(trailerlink::parse(&full_msg, "Arc-Run"), Some("run2")); + assert_eq!(trailerlink::parse(&full_msg, "Arc-Completed"), Some("1")); + assert_eq!(trailerlink::parse(&full_msg, "Arc-Checkpoint"), None); + 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()) @@ -280,7 +448,7 @@ mod tests { 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(); + let sha = checkpoint_commit(&wt_path, "run2", "nodeB", "completed", 0, None).unwrap(); assert_eq!(sha.len(), 40); remove_worktree(dir.path(), &wt_path).unwrap(); @@ -318,4 +486,117 @@ mod tests { let patch = diff_against(dir.path(), &base).unwrap(); assert!(patch.is_empty()); } + + // --- MetadataStore tests --- + + #[test] + fn metadata_store_init_run_and_read() { + let dir = tempfile::tempdir().unwrap(); + init_repo(dir.path()); + + let store = MetadataStore::new(dir.path()); + let manifest = br#"{"run_id":"RUN1","pipeline":"test"}"#; + let dot = b"digraph { start -> end }"; + store.init_run("RUN1", manifest, dot).unwrap(); + + let read_manifest = MetadataStore::read_manifest(dir.path(), "RUN1").unwrap().unwrap(); + assert_eq!(read_manifest["run_id"], "RUN1"); + assert_eq!(read_manifest["pipeline"], "test"); + + let read_dot = MetadataStore::read_graph_dot(dir.path(), "RUN1").unwrap().unwrap(); + assert_eq!(read_dot, "digraph { start -> end }"); + } + + #[test] + fn metadata_store_write_and_read_checkpoint() { + let dir = tempfile::tempdir().unwrap(); + init_repo(dir.path()); + + let store = MetadataStore::new(dir.path()); + store.init_run("RUN2", b"{}", b"digraph {}").unwrap(); + + let ctx = crate::context::Context::new(); + ctx.set("goal", serde_json::json!("test")); + let cp = crate::checkpoint::Checkpoint::from_context( + &ctx, + "node_a", + vec!["start".to_string()], + std::collections::HashMap::new(), + std::collections::HashMap::new(), + Some("node_b".to_string()), + ); + let cp_json = serde_json::to_vec_pretty(&cp).unwrap(); + store.write_checkpoint("RUN2", &cp_json, &[]).unwrap(); + + let loaded = MetadataStore::read_checkpoint(dir.path(), "RUN2").unwrap().unwrap(); + assert_eq!(loaded.current_node, "node_a"); + assert_eq!(loaded.completed_nodes, vec!["start"]); + assert_eq!(loaded.next_node_id.as_deref(), Some("node_b")); + assert_eq!(loaded.context_values.get("goal"), Some(&serde_json::json!("test"))); + } + + #[test] + fn metadata_store_write_checkpoint_overwrites() { + let dir = tempfile::tempdir().unwrap(); + init_repo(dir.path()); + + let store = MetadataStore::new(dir.path()); + store.init_run("RUN3", b"{}", b"digraph {}").unwrap(); + + let ctx = crate::context::Context::new(); + let cp1 = crate::checkpoint::Checkpoint::from_context( + &ctx, + "node_a", + vec!["start".to_string()], + std::collections::HashMap::new(), + std::collections::HashMap::new(), + None, + ); + let cp1_json = serde_json::to_vec_pretty(&cp1).unwrap(); + store.write_checkpoint("RUN3", &cp1_json, &[]).unwrap(); + + let cp2 = crate::checkpoint::Checkpoint::from_context( + &ctx, + "node_b", + vec!["start".to_string(), "node_a".to_string()], + std::collections::HashMap::new(), + std::collections::HashMap::new(), + Some("node_c".to_string()), + ); + let cp2_json = serde_json::to_vec_pretty(&cp2).unwrap(); + store.write_checkpoint("RUN3", &cp2_json, &[]).unwrap(); + + let loaded = MetadataStore::read_checkpoint(dir.path(), "RUN3").unwrap().unwrap(); + assert_eq!(loaded.current_node, "node_b"); + assert_eq!(loaded.completed_nodes.len(), 2); + } + + #[test] + fn metadata_store_read_checkpoint_missing_branch() { + let dir = tempfile::tempdir().unwrap(); + init_repo(dir.path()); + + let result = MetadataStore::read_checkpoint(dir.path(), "NONEXISTENT").unwrap(); + assert!(result.is_none()); + } + + #[test] + fn metadata_store_artifact_roundtrip() { + let dir = tempfile::tempdir().unwrap(); + init_repo(dir.path()); + + let store = MetadataStore::new(dir.path()); + store.init_run("RUN4", b"{}", b"digraph {}").unwrap(); + + let artifact_data = br#"{"large_output":"some data"}"#; + let cp_json = b"{}"; // minimal checkpoint for the test + store.write_checkpoint( + "RUN4", + cp_json, + &[("artifacts/response.plan.json", artifact_data.as_slice())], + ).unwrap(); + + let read_back = MetadataStore::read_artifact(dir.path(), "RUN4", "response.plan").unwrap().unwrap(); + assert_eq!(read_back, artifact_data); + } } diff --git a/crates/arc-attractor/src/server.rs b/crates/arc-attractor/src/server.rs index 6867b6d14..a9bf0b0e7 100644 --- a/crates/arc-attractor/src/server.rs +++ b/crates/arc-attractor/src/server.rs @@ -222,7 +222,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, run_id: run_id_clone.clone(), git_checkpoint: None, base_sha: None, run_branch: None }; + let config = RunConfig { logs_root, cancel_token: Some(cancel_token), dry_run: state_clone.dry_run, run_id: run_id_clone.clone(), git_checkpoint: None, base_sha: None, run_branch: None, meta_branch: 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 985148d4d..b9f94f466 100644 --- a/crates/arc-attractor/tests/daytona_integration.rs +++ b/crates/arc-attractor/tests/daytona_integration.rs @@ -274,6 +274,7 @@ async fn daytona_pipeline_artifact_offload_and_sync() { git_checkpoint: None, base_sha: None, run_branch: None, + meta_branch: None, }; let outcome = engine.run(&graph, &config).await.expect("pipeline should succeed"); @@ -432,9 +433,10 @@ async fn daytona_git_checkpoint_remote_emits_events() { cancel_token: None, dry_run: false, run_id, - git_checkpoint: Some(GitCheckpointMode::Remote), + git_checkpoint: Some(GitCheckpointMode::Remote(dir.path().to_path_buf())), base_sha: Some(base_sha), run_branch: Some(branch_name), + meta_branch: None, }; let outcome = engine.run(&graph, &config).await.expect("pipeline should succeed"); @@ -595,3 +597,125 @@ async fn daytona_cli_gemini() { ) .await; } + +// --------------------------------------------------------------------------- +// Daytona shadow commit E2E — Remote mode with MetadataStore +// --------------------------------------------------------------------------- + +use arc_attractor::git::MetadataStore; + +/// End-to-end test: pipeline with `GitCheckpointMode::Remote(host_repo_path)` + `meta_branch` +/// writes shadow branch on the host repo and includes `Arc-Checkpoint` trailer in sandbox commits. +#[tokio::test] +#[ignore] +async fn daytona_git_checkpoint_with_shadow_branch() { + let env = create_env().await; + env.initialize().await.unwrap(); + let env: Arc = Arc::new(env); + + // Install git if not available + let git_check = env.exec_command("git --version", 10_000, None, None, None).await; + if git_check.as_ref().map_or(true, |r| r.exit_code != 0) { + let install = env.exec_command( + "apt-get update -qq && apt-get install -y -qq git >/dev/null 2>&1", + 120_000, None, None, None, + ).await.expect("apt-get install git should not error"); + assert_eq!(install.exit_code, 0, "git install failed: {}", install.stderr); + } + + // Set up git in the sandbox + let (run_id, base_sha, branch_name) = setup_daytona_git(&*env).await; + + // Create a temp git repo on the host for MetadataStore + let host_repo = tempfile::tempdir().unwrap(); + std::process::Command::new("git") + .args(["init"]) + .current_dir(host_repo.path()) + .output() + .unwrap(); + std::process::Command::new("git") + .args(["-c", "user.name=test", "-c", "user.email=test@test", + "commit", "--allow-empty", "-m", "init"]) + .current_dir(host_repo.path()) + .output() + .unwrap(); + + // Pipeline: start -> work -> exit + let mut graph = Graph::new("DaytonaShadowBranch"); + graph.attrs.insert( + "goal".to_string(), + AttrValue::String("Test Daytona shadow branch".to_string()), + ); + + let mut start = Node::new("start"); + start.attrs.insert("shape".to_string(), AttrValue::String("Mdiamond".to_string())); + graph.nodes.insert("start".to_string(), start); + + let mut exit = Node::new("exit"); + exit.attrs.insert("shape".to_string(), AttrValue::String("Msquare".to_string())); + graph.nodes.insert("exit".to_string(), exit); + + let mut work = Node::new("work"); + work.attrs.insert("label".to_string(), AttrValue::String("Work".to_string())); + graph.nodes.insert("work".to_string(), work); + + graph.edges.push(Edge::new("start", "work")); + graph.edges.push(Edge::new("work", "exit")); + + let dir = tempfile::tempdir().unwrap(); + // Write graph.dot so init_run can read it + std::fs::write(dir.path().join("graph.dot"), "digraph {}").unwrap(); + + let mut registry = HandlerRegistry::new(Box::new(FileWriterHandler)); + registry.register("start", Box::new(StartHandler)); + registry.register("exit", Box::new(ExitHandler)); + + let meta_branch = MetadataStore::branch_name(&run_id); + let engine = PipelineEngine::new(registry, Arc::new(EventEmitter::new()), env.clone()); + let config = RunConfig { + logs_root: dir.path().to_path_buf(), + cancel_token: None, + dry_run: false, + run_id: run_id.clone(), + git_checkpoint: Some(GitCheckpointMode::Remote(host_repo.path().to_path_buf())), + base_sha: Some(base_sha), + run_branch: Some(branch_name), + meta_branch: Some(meta_branch), + }; + + let outcome = engine.run(&graph, &config).await.expect("pipeline should succeed"); + assert_eq!(outcome.status, StageStatus::Success); + + // Assert shadow branch on host has checkpoint data + let checkpoint = MetadataStore::read_checkpoint(host_repo.path(), &run_id) + .expect("read_checkpoint should not error") + .expect("shadow branch should contain checkpoint data"); + assert!( + !checkpoint.completed_nodes.is_empty(), + "checkpoint should have completed nodes" + ); + assert!( + checkpoint.completed_nodes.contains(&"work".to_string()), + "checkpoint should contain the 'work' node" + ); + + // Assert sandbox commit has Arc-Checkpoint trailer + let log_result = env.exec_command("git log --format=%B -1", 10_000, None, None, None).await + .expect("git log should succeed"); + assert_eq!(log_result.exit_code, 0); + let commit_msg = log_result.stdout.trim().to_string(); + assert!( + commit_msg.contains("Arc-Checkpoint:"), + "sandbox commit should have Arc-Checkpoint trailer, got:\n{commit_msg}" + ); + assert!( + commit_msg.contains("Arc-Run:"), + "sandbox commit should have Arc-Run trailer, got:\n{commit_msg}" + ); + + // Assert final.patch exists + let final_patch = dir.path().join("final.patch"); + assert!(final_patch.exists(), "final.patch should exist in logs_root"); + + env.cleanup().await.unwrap(); +} diff --git a/crates/arc-attractor/tests/integration.rs b/crates/arc-attractor/tests/integration.rs index 63ae8a174..305931901 100644 --- a/crates/arc-attractor/tests/integration.rs +++ b/crates/arc-attractor/tests/integration.rs @@ -190,6 +190,7 @@ async fn end_to_end_linear_pipeline() { git_checkpoint: None, base_sha: None, run_branch: None, + meta_branch: None, }; let outcome = engine.run(&graph, &config).await.expect("run should succeed"); @@ -312,6 +313,7 @@ async fn end_to_end_branching_pipeline() { git_checkpoint: None, base_sha: None, run_branch: None, + meta_branch: None, }; let outcome = engine.run(&graph, &config).await.expect("run should succeed"); @@ -421,6 +423,7 @@ async fn end_to_end_human_gate_pipeline() { git_checkpoint: None, base_sha: None, run_branch: None, + meta_branch: None, }; let outcome = engine.run(&graph, &config).await.expect("run should succeed"); @@ -524,6 +527,7 @@ async fn goal_gate_routes_to_retry_target_on_failure() { git_checkpoint: None, base_sha: None, run_branch: None, + meta_branch: None, }; let result = engine.run(&graph, &config).await; @@ -636,6 +640,7 @@ async fn goal_gate_routes_to_retry_target_when_present() { git_checkpoint: None, base_sha: None, run_branch: None, + meta_branch: None, }; let outcome = engine @@ -948,6 +953,7 @@ async fn retry_on_failure_then_succeed() { git_checkpoint: None, base_sha: None, run_branch: None, + meta_branch: None, }; let outcome = engine @@ -1014,6 +1020,7 @@ async fn pipeline_with_many_nodes() { git_checkpoint: None, base_sha: None, run_branch: None, + meta_branch: None, }; let outcome = engine @@ -1332,6 +1339,7 @@ async fn smoke_test_with_mock_codergen_backend() { git_checkpoint: None, base_sha: None, run_branch: None, + meta_branch: None, }; let outcome = engine @@ -1430,6 +1438,7 @@ async fn end_to_end_parallel_fan_out_fan_in() { git_checkpoint: None, base_sha: None, run_branch: None, + meta_branch: None, }; let outcome = engine @@ -1537,6 +1546,7 @@ async fn resume_from_checkpoint_completes_pipeline() { git_checkpoint: None, base_sha: None, run_branch: None, + meta_branch: None, }; let outcome = engine @@ -1629,6 +1639,7 @@ async fn resume_from_checkpoint_preserves_goal_gate_outcomes() { git_checkpoint: None, base_sha: None, run_branch: None, + meta_branch: None, }; // This should succeed because goal gate for gated_work is satisfied @@ -1663,6 +1674,7 @@ async fn graph_goal_in_context() { git_checkpoint: None, base_sha: None, run_branch: None, + meta_branch: None, }; engine.run(&graph, &config).await.expect("run"); @@ -1693,6 +1705,7 @@ async fn event_streaming_lifecycle() { git_checkpoint: None, base_sha: None, run_branch: None, + meta_branch: None, }; engine.run(&graph, &config).await.expect("run"); @@ -1763,6 +1776,7 @@ async fn context_flow_between_stages() { git_checkpoint: None, base_sha: None, run_branch: None, + meta_branch: None, }; engine.run(&graph, &config).await.expect("run"); @@ -1806,6 +1820,7 @@ async fn tool_handler_e2e() { git_checkpoint: None, base_sha: None, run_branch: None, + meta_branch: None, }; let outcome = engine.run(&graph, &config).await.expect("run"); assert_eq!(outcome.status, StageStatus::Success); @@ -1867,6 +1882,7 @@ async fn auto_approve_interviewer_e2e() { git_checkpoint: None, base_sha: None, run_branch: None, + meta_branch: None, }; let outcome = engine.run(&graph, &config).await.expect("run"); assert_eq!(outcome.status, StageStatus::Success); @@ -1894,6 +1910,7 @@ async fn codergen_without_backend_simulated() { git_checkpoint: None, base_sha: None, run_branch: None, + meta_branch: None, }; engine.run(&graph, &config).await.expect("run"); @@ -1995,6 +2012,7 @@ async fn branching_loop_back_on_failure() { git_checkpoint: None, base_sha: None, run_branch: None, + meta_branch: None, }; let outcome = engine.run(&graph, &config).await.expect("run"); assert_eq!(outcome.status, StageStatus::Success); @@ -2079,6 +2097,7 @@ async fn human_gate_loops_back() { git_checkpoint: None, base_sha: None, run_branch: None, + meta_branch: None, }; let outcome = engine.run(&graph, &config).await.expect("run"); assert_eq!(outcome.status, StageStatus::Success); @@ -2131,6 +2150,7 @@ async fn scenario_ship_a_feature() { git_checkpoint: None, base_sha: None, run_branch: None, + meta_branch: None, }; let outcome = engine.run(&graph, &config).await.expect("run"); assert_eq!(outcome.status, StageStatus::Success); @@ -2211,6 +2231,7 @@ async fn scenario_parallel_expert_review() { git_checkpoint: None, base_sha: None, run_branch: None, + meta_branch: None, }; let outcome = engine.run(&graph, &config).await.expect("run"); assert_eq!(outcome.status, StageStatus::Success); @@ -2285,6 +2306,7 @@ async fn scenario_node_retries_on_retry_status() { git_checkpoint: None, base_sha: None, run_branch: None, + meta_branch: None, }; let outcome = engine.run(&graph, &config).await.expect("run"); assert_eq!(outcome.status, StageStatus::Success); @@ -2343,6 +2365,7 @@ async fn scenario_loop_restart_resets_context() { git_checkpoint: None, base_sha: None, run_branch: None, + meta_branch: None, }; let outcome = engine.run(&graph, &config).await.expect("run"); assert_eq!(outcome.status, StageStatus::Success); @@ -2405,6 +2428,7 @@ async fn scenario_bug_triage_router() { git_checkpoint: None, base_sha: None, run_branch: None, + meta_branch: None, }; let outcome = engine.run(&graph, &config).await.expect("run"); assert_eq!(outcome.status, StageStatus::Success); @@ -2455,6 +2479,7 @@ async fn scenario_crash_recovery() { git_checkpoint: None, base_sha: None, run_branch: None, + meta_branch: None, }; let outcome = engine .run_from_checkpoint(&graph, &config, &checkpoint) @@ -2537,6 +2562,7 @@ async fn manager_loop_stop_condition_satisfied_e2e() { git_checkpoint: None, base_sha: None, run_branch: None, + meta_branch: None, }; let outcome = engine.run(&graph, &config).await.expect("run"); @@ -2587,6 +2613,7 @@ async fn manager_loop_max_cycles_exceeded_e2e() { git_checkpoint: None, base_sha: None, run_branch: None, + meta_branch: None, }; let outcome = engine.run(&graph, &config).await.expect("run"); @@ -2717,6 +2744,7 @@ async fn conditional_branching_success_fail_paths() { git_checkpoint: None, base_sha: None, run_branch: None, + meta_branch: None, }; let outcome = engine.run(&graph, &config).await.expect("run"); assert_eq!(outcome.status, StageStatus::Success); @@ -2765,6 +2793,7 @@ async fn edge_selection_condition_match_wins_over_weight() { git_checkpoint: None, base_sha: None, run_branch: None, + meta_branch: None, }; engine.run(&graph, &config).await.expect("run"); @@ -2808,6 +2837,7 @@ async fn edge_selection_weight_breaks_ties() { git_checkpoint: None, base_sha: None, run_branch: None, + meta_branch: None, }; engine.run(&graph, &config).await.expect("run"); @@ -2843,6 +2873,7 @@ async fn edge_selection_lexical_tiebreak() { git_checkpoint: None, base_sha: None, run_branch: None, + meta_branch: None, }; engine.run(&graph, &config).await.expect("run"); @@ -2895,6 +2926,7 @@ async fn context_updates_visible_across_nodes() { git_checkpoint: None, base_sha: None, run_branch: None, + meta_branch: None, }; engine.run(&graph, &config).await.expect("run"); @@ -2929,6 +2961,7 @@ async fn stylesheet_applies_model_override() { git_checkpoint: None, base_sha: None, run_branch: None, + meta_branch: None, }; let outcome = engine.run(&graph, &config).await.expect("run"); assert_eq!(outcome.status, StageStatus::Success); @@ -2979,6 +3012,7 @@ async fn custom_handler_registration_and_execution() { git_checkpoint: None, base_sha: None, run_branch: None, + meta_branch: None, }; engine.run(&graph, &config).await.expect("run"); @@ -3040,6 +3074,7 @@ async fn integration_smoke_plan_implement_review_done() { git_checkpoint: None, base_sha: None, run_branch: None, + meta_branch: None, }; let outcome = engine.run(&graph, &config).await.expect("run"); assert_eq!(outcome.status, StageStatus::Success); @@ -3117,7 +3152,7 @@ mod server_lifecycle { #[tokio::test] async fn full_http_lifecycle_approve_and_complete() { let state = create_app_state(gate_registry); - let app = build_router(Arc::clone(&state)); + let app = build_router(Arc::clone(&state), arc_attractor::jwt_auth::AuthMode::Disabled); // 1. Start pipeline let req = Request::builder() @@ -3216,7 +3251,7 @@ mod server_lifecycle { #[tokio::test] async fn full_http_lifecycle_cancel() { let state = create_app_state(gate_registry); - let app = build_router(Arc::clone(&state)); + let app = build_router(Arc::clone(&state), arc_attractor::jwt_auth::AuthMode::Disabled); // Start a pipeline that will block at the human gate let req = Request::builder() @@ -3296,7 +3331,7 @@ mod sse_events { #[tokio::test] async fn sse_stream_contains_expected_event_types() { let state = create_app_state(simple_registry); - let app = build_router(Arc::clone(&state)); + let app = build_router(Arc::clone(&state), arc_attractor::jwt_auth::AuthMode::Disabled); // Start pipeline let req = Request::builder() @@ -3429,7 +3464,7 @@ mod serve_dry_run { default_registry(interviewer, || None) }; let state = create_app_state(factory); - build_router(state) + build_router(state, arc_attractor::jwt_auth::AuthMode::Disabled) } async fn body_json(body: Body) -> serde_json::Value { @@ -3529,6 +3564,7 @@ async fn sub_pipeline_e2e_through_engine() { git_checkpoint: None, base_sha: None, run_branch: None, + meta_branch: None, }; let outcome = engine @@ -3675,6 +3711,7 @@ async fn manager_loop_with_child_observer_e2e() { git_checkpoint: None, base_sha: None, run_branch: None, + meta_branch: None, }; let outcome = engine @@ -3798,6 +3835,7 @@ async fn graph_merge_e2e_through_engine() { git_checkpoint: None, base_sha: None, run_branch: None, + meta_branch: None, }; let outcome = engine @@ -3942,6 +3980,7 @@ async fn fidelity_default_is_compact() { git_checkpoint: None, base_sha: None, run_branch: None, + meta_branch: None, }; engine.run(&graph, &config).await.expect("run"); @@ -3985,6 +4024,7 @@ async fn fidelity_graph_default_applied() { git_checkpoint: None, base_sha: None, run_branch: None, + meta_branch: None, }; engine.run(&graph, &config).await.expect("run"); @@ -4027,6 +4067,7 @@ async fn fidelity_node_overrides_graph_default() { git_checkpoint: None, base_sha: None, run_branch: None, + meta_branch: None, }; engine.run(&graph, &config).await.expect("run"); @@ -4075,6 +4116,7 @@ async fn fidelity_edge_overrides_node_and_graph() { git_checkpoint: None, base_sha: None, run_branch: None, + meta_branch: None, }; engine.run(&graph, &config).await.expect("run"); @@ -4113,6 +4155,7 @@ async fn fidelity_full_produces_empty_preamble() { git_checkpoint: None, base_sha: None, run_branch: None, + meta_branch: None, }; engine.run(&graph, &config).await.expect("run"); @@ -4158,6 +4201,7 @@ async fn fidelity_truncate_preamble_minimal() { git_checkpoint: None, base_sha: None, run_branch: None, + meta_branch: None, }; engine.run(&graph, &config).await.expect("run"); @@ -4219,6 +4263,7 @@ async fn fidelity_summary_low_mode() { git_checkpoint: None, base_sha: None, run_branch: None, + meta_branch: None, }; engine.run(&graph, &config).await.expect("run"); @@ -4275,6 +4320,7 @@ async fn fidelity_summary_medium_mode() { git_checkpoint: None, base_sha: None, run_branch: None, + meta_branch: None, }; engine.run(&graph, &config).await.expect("run"); @@ -4331,6 +4377,7 @@ async fn fidelity_summary_high_mode() { git_checkpoint: None, base_sha: None, run_branch: None, + meta_branch: None, }; engine.run(&graph, &config).await.expect("run"); @@ -4380,6 +4427,7 @@ async fn fidelity_full_sets_thread_id_in_context() { git_checkpoint: None, base_sha: None, run_branch: None, + meta_branch: None, }; engine.run(&graph, &config).await.expect("run"); @@ -4440,6 +4488,7 @@ async fn fidelity_full_nodes_share_thread_id() { git_checkpoint: None, base_sha: None, run_branch: None, + meta_branch: None, }; engine.run(&graph, &config).await.expect("run"); @@ -4507,6 +4556,7 @@ async fn fidelity_resume_degrades_full_to_summary_high() { git_checkpoint: None, base_sha: None, run_branch: None, + meta_branch: None, }; engine .run_from_checkpoint(&graph, &config, &checkpoint) @@ -4590,6 +4640,7 @@ async fn fidelity_resume_degrade_only_affects_first_hop() { git_checkpoint: None, base_sha: None, run_branch: None, + meta_branch: None, }; engine .run_from_checkpoint(&graph, &config, &checkpoint) @@ -4660,6 +4711,7 @@ async fn fidelity_resume_no_degrade_when_not_full() { git_checkpoint: None, base_sha: None, run_branch: None, + meta_branch: None, }; engine .run_from_checkpoint(&graph, &config, &checkpoint) @@ -4696,6 +4748,7 @@ async fn fidelity_stored_in_checkpoint_context() { git_checkpoint: None, base_sha: None, run_branch: None, + meta_branch: None, }; engine.run(&graph, &config).await.expect("run"); @@ -4771,6 +4824,7 @@ async fn fidelity_precedence_multi_node_pipeline() { git_checkpoint: None, base_sha: None, run_branch: None, + meta_branch: None, }; engine.run(&graph, &config).await.expect("run"); @@ -4828,6 +4882,7 @@ async fn fidelity_compact_preamble_includes_completed_stages_and_context() { git_checkpoint: None, base_sha: None, run_branch: None, + meta_branch: None, }; engine.run(&graph, &config).await.expect("run"); @@ -4873,7 +4928,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, run_id: "test-run".into(), git_checkpoint: None, base_sha: None, run_branch: None }; + let config_low = RunConfig { logs_root: dir_low.path().to_path_buf(), cancel_token: None, dry_run: false, run_id: "test-run".into(), git_checkpoint: None, base_sha: None, run_branch: None, meta_branch: None }; engine_low.run(&graph_low, &config_low).await.expect("run low"); { @@ -4907,7 +4962,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, run_id: "test-run".into(), git_checkpoint: None, base_sha: None, run_branch: None }; + let config_med = RunConfig { logs_root: dir_med.path().to_path_buf(), cancel_token: None, dry_run: false, run_id: "test-run".into(), git_checkpoint: None, base_sha: None, run_branch: None, meta_branch: None }; engine_med.run(&graph_med, &config_med).await.expect("run med"); let preambles_med = captures_med.preambles.lock().unwrap(); @@ -4964,6 +5019,7 @@ async fn fidelity_thread_id_fallback_to_previous_node_in_pipeline() { git_checkpoint: None, base_sha: None, run_branch: None, + meta_branch: None, }; engine.run(&graph, &config).await.expect("run"); @@ -5007,6 +5063,7 @@ async fn fidelity_thread_id_from_node_class_in_pipeline() { git_checkpoint: None, base_sha: None, run_branch: None, + meta_branch: None, }; engine.run(&graph, &config).await.expect("run"); @@ -5053,6 +5110,7 @@ async fn fidelity_edge_thread_id_override_in_pipeline() { git_checkpoint: None, base_sha: None, run_branch: None, + meta_branch: None, }; engine.run(&graph, &config).await.expect("run"); @@ -5100,6 +5158,7 @@ async fn fidelity_full_without_explicit_thread_id_uses_previous_node() { git_checkpoint: None, base_sha: None, run_branch: None, + meta_branch: None, }; engine.run(&graph, &config).await.expect("run"); @@ -5154,6 +5213,7 @@ async fn fidelity_from_parsed_dot_pipeline() { git_checkpoint: None, base_sha: None, run_branch: None, + meta_branch: None, }; engine.run(&graph, &config).await.expect("run"); @@ -5196,6 +5256,7 @@ async fn fidelity_checkpoint_roundtrip_preserves_fidelity() { git_checkpoint: None, base_sha: None, run_branch: None, + meta_branch: None, }; engine.run(&graph, &config).await.expect("run"); @@ -5255,6 +5316,7 @@ async fn fidelity_node_thread_id_overrides_edge_thread_id_in_pipeline() { git_checkpoint: None, base_sha: None, run_branch: None, + meta_branch: None, }; engine.run(&graph, &config).await.expect("run"); @@ -5328,6 +5390,7 @@ async fn fidelity_resume_preserves_context_values_across_checkpoint() { git_checkpoint: None, base_sha: None, run_branch: None, + meta_branch: None, }; engine .run_from_checkpoint(&graph, &config, &checkpoint) @@ -5517,6 +5580,7 @@ mod real_llm { git_checkpoint: None, base_sha: None, run_branch: None, + meta_branch: None, }; let outcome = tokio::time::timeout( @@ -5625,6 +5689,7 @@ mod real_llm { git_checkpoint: None, base_sha: None, run_branch: None, + meta_branch: None, }; let outcome = tokio::time::timeout( @@ -5763,6 +5828,7 @@ mod real_llm { git_checkpoint: None, base_sha: None, run_branch: None, + meta_branch: None, }; let outcome = tokio::time::timeout( @@ -5867,6 +5933,7 @@ mod real_llm { git_checkpoint: None, base_sha: None, run_branch: None, + meta_branch: None, }; let outcome = tokio::time::timeout( @@ -5951,6 +6018,7 @@ async fn human_gate_freeform_only_routes_text() { git_checkpoint: None, base_sha: None, run_branch: None, + meta_branch: None, }; let outcome = engine.run(&graph, &config).await.expect("run should succeed"); @@ -6070,6 +6138,7 @@ async fn human_gate_freeform_with_fixed_choice_match() { git_checkpoint: None, base_sha: None, run_branch: None, + meta_branch: None, }; let outcome = engine.run(&graph, &config).await.expect("run should succeed"); @@ -6174,6 +6243,7 @@ async fn human_gate_freeform_fallback_on_unmatched_text() { git_checkpoint: None, base_sha: None, run_branch: None, + meta_branch: None, }; let outcome = engine.run(&graph, &config).await.expect("run should succeed"); @@ -6291,6 +6361,7 @@ async fn human_gate_freeform_sets_allow_freeform_on_question() { git_checkpoint: None, base_sha: None, run_branch: None, + meta_branch: None, }; let outcome = engine.run(&graph, &config).await.expect("run should succeed"); @@ -6384,6 +6455,7 @@ async fn human_gate_without_freeform_sets_allow_freeform_false() { git_checkpoint: None, base_sha: None, run_branch: None, + meta_branch: None, }; let outcome = engine.run(&graph, &config).await.expect("run should succeed"); @@ -6626,6 +6698,7 @@ async fn tool_hooks_pre_success_allows_pipeline_to_proceed() { git_checkpoint: None, base_sha: None, run_branch: None, + meta_branch: None, }; let outcome = engine.run(&graph, &config).await.expect("run should succeed"); @@ -6665,6 +6738,7 @@ async fn tool_hooks_pre_failure_skips_tool_call() { git_checkpoint: None, base_sha: None, run_branch: None, + meta_branch: None, }; engine.run(&graph, &config).await.expect("run should complete"); @@ -6707,6 +6781,7 @@ async fn tool_hooks_post_success_does_not_affect_outcome() { git_checkpoint: None, base_sha: None, run_branch: None, + meta_branch: None, }; let outcome = engine.run(&graph, &config).await.expect("run should succeed"); @@ -6741,6 +6816,7 @@ async fn tool_hooks_post_failure_does_not_block_pipeline() { git_checkpoint: None, base_sha: None, run_branch: None, + meta_branch: None, }; let outcome = engine.run(&graph, &config).await.expect("run should succeed"); @@ -6777,6 +6853,7 @@ async fn tool_hooks_graph_level_applies_to_all_nodes() { git_checkpoint: None, base_sha: None, run_branch: None, + meta_branch: None, }; let outcome = engine.run(&graph, &config).await.expect("run should succeed"); @@ -6816,6 +6893,7 @@ async fn tool_hooks_node_level_overrides_graph_level() { git_checkpoint: None, base_sha: None, run_branch: None, + meta_branch: None, }; let _outcome = engine.run(&graph, &config).await.expect("run should complete"); @@ -6865,6 +6943,7 @@ async fn tool_hooks_pre_receives_node_id_env_var() { git_checkpoint: None, base_sha: None, run_branch: None, + meta_branch: None, }; engine.run(&graph, &config).await.expect("run should succeed"); @@ -6961,6 +7040,7 @@ async fn attractor_e2e_with_real_llm() { git_checkpoint: None, base_sha: None, run_branch: None, + meta_branch: None, }; let outcome = engine.run(&graph, &config).await.expect("run should succeed"); @@ -7068,6 +7148,7 @@ async fn run_fidelity_prompt_pipeline(fidelity: &str) -> String { git_checkpoint: None, base_sha: None, run_branch: None, + meta_branch: None, }; engine.run(&graph, &config).await.expect("pipeline should succeed"); @@ -7211,6 +7292,7 @@ async fn large_context_values_are_offloaded_to_artifact_store() { git_checkpoint: None, base_sha: None, run_branch: None, + meta_branch: None, }; let outcome = engine @@ -7386,6 +7468,7 @@ async fn artifact_pointers_rewritten_for_remote_execution_env() { git_checkpoint: None, base_sha: None, run_branch: None, + meta_branch: None, }; let outcome = engine @@ -7500,6 +7583,7 @@ async fn node_dir_uses_visit_count_on_revisit() { git_checkpoint: None, base_sha: None, run_branch: None, + meta_branch: None, }; let outcome = engine.run(&graph, &config).await.expect("pipeline should succeed"); @@ -8066,6 +8150,7 @@ async fn full_pipeline_with_cli_backend_node() { git_checkpoint: None, base_sha: None, run_branch: None, + meta_branch: None, }; let outcome = engine.run(&graph, &config).await.expect("pipeline should succeed"); @@ -8154,6 +8239,7 @@ async fn stylesheet_backend_property_routes_to_cli() { git_checkpoint: None, base_sha: None, run_branch: None, + meta_branch: None, }; let outcome = engine.run(&graph, &config).await.expect("pipeline should succeed"); @@ -8353,6 +8439,7 @@ async fn git_checkpoint_host_emits_events_and_diff_patch() { git_checkpoint: Some(GitCheckpointMode::Host(worktree_path.clone())), base_sha: Some(base_sha.clone()), run_branch: Some("arc/run/test-docker".to_string()), + meta_branch: None, }; // 5. Run pipeline @@ -8401,6 +8488,145 @@ async fn git_checkpoint_host_emits_events_and_diff_patch() { let patch_content = std::fs::read_to_string(&final_patch).unwrap(); assert!(patch_content.contains("hello.txt"), "final.patch should contain hello.txt changes"); + // Cleanup worktree + let _ = std::process::Command::new("git") + .args(["worktree", "remove", "--force"]) + .arg(&worktree_path) + .current_dir(repo.path()) + .output(); +} + +/// End-to-end test: pipeline with `GitCheckpointMode::Host` + `meta_branch` writes +/// shadow branch with checkpoint data and includes `Arc-Checkpoint` trailer in run-branch commits. +#[tokio::test] +async fn git_checkpoint_host_writes_shadow_branch() { + use arc_attractor::git::MetadataStore; + + // 1. Create a temporary git repo with an initial commit + let repo = tempfile::tempdir().unwrap(); + std::process::Command::new("git") + .args(["init"]) + .current_dir(repo.path()) + .output() + .unwrap(); + std::process::Command::new("git") + .args(["-c", "user.name=test", "-c", "user.email=test@test", + "commit", "--allow-empty", "-m", "init"]) + .current_dir(repo.path()) + .output() + .unwrap(); + + // 2. Create a branch and worktree + let run_id = "test-shadow"; + let base_sha = { + let out = std::process::Command::new("git") + .args(["rev-parse", "HEAD"]) + .current_dir(repo.path()) + .output() + .unwrap(); + String::from_utf8_lossy(&out.stdout).trim().to_string() + }; + std::process::Command::new("git") + .args(["branch", &format!("arc/run/{run_id}"), "HEAD"]) + .current_dir(repo.path()) + .output() + .unwrap(); + let worktree_path = repo.path().join("worktree"); + std::process::Command::new("git") + .args(["worktree", "add"]) + .arg(&worktree_path) + .arg(format!("arc/run/{run_id}")) + .current_dir(repo.path()) + .output() + .unwrap(); + + // Write a file in the worktree so there's something to commit + std::fs::write(worktree_path.join("shadow_test.txt"), "shadow branch test").unwrap(); + + // 3. Build a simple pipeline: start -> work -> exit + let mut graph = Graph::new("ShadowBranchTest"); + graph.attrs.insert("goal".to_string(), AttrValue::String("Test shadow branch".to_string())); + let mut start = Node::new("start"); + start.attrs.insert("shape".to_string(), AttrValue::String("Mdiamond".to_string())); + graph.nodes.insert("start".to_string(), start); + let mut exit = Node::new("exit"); + exit.attrs.insert("shape".to_string(), AttrValue::String("Msquare".to_string())); + graph.nodes.insert("exit".to_string(), exit); + let mut work = Node::new("work"); + work.attrs.insert("label".to_string(), AttrValue::String("Work".to_string())); + graph.nodes.insert("work".to_string(), work); + graph.edges.push(Edge::new("start", "work")); + graph.edges.push(Edge::new("work", "exit")); + + // 4. Set up engine with meta_branch + let logs_dir = tempfile::tempdir().unwrap(); + // Write graph.dot so init_run can read it + std::fs::write(logs_dir.path().join("graph.dot"), "digraph {}").unwrap(); + let emitter = EventEmitter::new(); + + let env: Arc = Arc::new( + arc_agent::LocalExecutionEnvironment::new(worktree_path.clone()), + ); + let mut registry = HandlerRegistry::new(Box::new(ContextSetterHandler)); + registry.register("start", Box::new(StartHandler)); + registry.register("exit", Box::new(ExitHandler)); + let engine = PipelineEngine::new(registry, Arc::new(emitter), env); + + let meta_branch = MetadataStore::branch_name(run_id); + let config = RunConfig { + logs_root: logs_dir.path().to_path_buf(), + cancel_token: None, + dry_run: false, + run_id: run_id.into(), + git_checkpoint: Some(GitCheckpointMode::Host(worktree_path.clone())), + base_sha: Some(base_sha), + run_branch: Some(format!("arc/run/{run_id}")), + meta_branch: Some(meta_branch), + }; + + // 5. Run pipeline + let outcome = engine.run(&graph, &config).await.expect("pipeline should succeed"); + assert_eq!(outcome.status, StageStatus::Success); + + // 6. Assert shadow branch has checkpoint data on the host repo + let checkpoint = MetadataStore::read_checkpoint(repo.path(), run_id) + .expect("read_checkpoint should not error") + .expect("shadow branch should contain checkpoint data"); + assert!( + !checkpoint.completed_nodes.is_empty(), + "checkpoint should have completed nodes" + ); + assert!( + checkpoint.completed_nodes.contains(&"work".to_string()), + "checkpoint should contain the 'work' node" + ); + + // 7. Assert run-branch commit has Arc-Checkpoint trailer pointing to shadow SHA + let output = std::process::Command::new("git") + .args(["log", "--format=%B", "-1"]) + .current_dir(&worktree_path) + .output() + .unwrap(); + let commit_msg = String::from_utf8_lossy(&output.stdout).trim().to_string(); + assert!( + commit_msg.contains("Arc-Checkpoint:"), + "run-branch commit should have Arc-Checkpoint trailer, got:\n{commit_msg}" + ); + assert!( + commit_msg.contains("Arc-Run:"), + "run-branch commit should have Arc-Run trailer, got:\n{commit_msg}" + ); + assert!( + commit_msg.contains("Arc-Completed:"), + "run-branch commit should have Arc-Completed trailer, got:\n{commit_msg}" + ); + + // 8. Verify round-trip: shadow checkpoint's completed_nodes matches expected + let manifest = MetadataStore::read_manifest(repo.path(), run_id) + .expect("read_manifest should not error") + .expect("shadow branch should contain manifest"); + assert_eq!(manifest["run_id"], run_id); + // Cleanup worktree let _ = std::process::Command::new("git") .args(["worktree", "remove", "--force"])