mirror of
https://github.com/fabro-sh/fabro.git
synced 2026-10-09 03:20:56 +00:00
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 <noreply@anthropic.com>
This commit is contained in:
parent
77b7bf7481
commit
90dc528c99
10 changed files with 977 additions and 55 deletions
1
Cargo.lock
generated
1
Cargo.lock
generated
|
|
@ -134,6 +134,7 @@ version = "0.1.0"
|
|||
dependencies = [
|
||||
"anyhow",
|
||||
"arc-agent",
|
||||
"arc-git-storage",
|
||||
"arc-llm",
|
||||
"arc-util",
|
||||
"assert_cmd",
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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<PathBuf>,
|
||||
|
||||
/// Log/artifact directory
|
||||
#[arg(long)]
|
||||
|
|
@ -93,6 +94,10 @@ pub struct RunArgs {
|
|||
#[arg(long)]
|
||||
pub resume: Option<PathBuf>,
|
||||
|
||||
/// Resume from a git run branch (reads checkpoint and graph from metadata branch)
|
||||
#[arg(long, conflicts_with = "resume")]
|
||||
pub run_branch: Option<String>,
|
||||
|
||||
/// Override default LLM model
|
||||
#[arg(long)]
|
||||
pub model: Option<String>,
|
||||
|
|
|
|||
|
|
@ -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/<run_id>', 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<dyn arc_agent::ExecutionEnvironment> = {
|
||||
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<dyn crate::interviewer::Interviewer> = 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::<arc_llm::provider::Provider>())
|
||||
.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]
|
||||
|
|
|
|||
|
|
@ -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),
|
||||
|
|
|
|||
|
|
@ -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<String> {
|
||||
async fn git_checkpoint_host(
|
||||
work_dir: PathBuf,
|
||||
run_id: String,
|
||||
node_id: String,
|
||||
status: String,
|
||||
completed_count: usize,
|
||||
shadow_sha: Option<String>,
|
||||
) -> Option<String> {
|
||||
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<String> {
|
|||
}
|
||||
|
||||
/// 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<String> {
|
||||
async fn git_checkpoint_remote(
|
||||
exec_env: &dyn ExecutionEnvironment,
|
||||
run_id: &str,
|
||||
node_id: &str,
|
||||
status: &str,
|
||||
completed_count: usize,
|
||||
shadow_sha: Option<String>,
|
||||
) -> Option<String> {
|
||||
// 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<String>,
|
||||
/// Git branch name for the run (e.g. `arc/run/{run_id}`).
|
||||
pub run_branch: Option<String>,
|
||||
/// Metadata branch name for git-native checkpoint storage (e.g. `refs/arc/{run_id}`).
|
||||
pub meta_branch: Option<String>,
|
||||
}
|
||||
|
||||
/// 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<String> = 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<u8>)> = 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,
|
||||
|
|
|
|||
|
|
@ -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<String>) -> 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<String> {
|
||||
// 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<String> {
|
|||
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<std::path::PathBuf>) -> 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<String> {
|
||||
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<Option<Vec<u8>>> {
|
||||
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<Option<Checkpoint>> {
|
||||
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<Option<serde_json::Value>> {
|
||||
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<Option<String>> {
|
||||
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<Option<Vec<u8>>> {
|
||||
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);
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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,
|
||||
|
|
|
|||
|
|
@ -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<dyn ExecutionEnvironment> = Arc::new(env);
|
||||
|
||||
// Install git if not available
|
||||
let git_check = env.exec_command("git --version", 10_000, None, None, None).await;
|
||||
if git_check.as_ref().map_or(true, |r| r.exit_code != 0) {
|
||||
let install = env.exec_command(
|
||||
"apt-get update -qq && apt-get install -y -qq git >/dev/null 2>&1",
|
||||
120_000, None, None, None,
|
||||
).await.expect("apt-get install git should not error");
|
||||
assert_eq!(install.exit_code, 0, "git install failed: {}", install.stderr);
|
||||
}
|
||||
|
||||
// Set up git in the sandbox
|
||||
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();
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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<dyn arc_agent::ExecutionEnvironment> = 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"])
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue