diff --git a/lib/crates/arc-workflows/src/cli/run.rs b/lib/crates/arc-workflows/src/cli/run.rs index 47d8c4fb6..1ae3f6a7d 100644 --- a/lib/crates/arc-workflows/src/cli/run.rs +++ b/lib/crates/arc-workflows/src/cli/run.rs @@ -1227,6 +1227,9 @@ pub async fn run_command( .await; } + // Write finalize commit with retro.json + final node files (captures last diff.patch) + write_finalize_commit(&config, &run_dir).await; + // Auto-create PR on successful completion let mut pushed_branch: Option = None; let mut pr_url: Option = None; @@ -1753,6 +1756,9 @@ async fn run_from_branch( .await; } + // Write finalize commit with retro.json + final node files (captures last diff.patch) + write_finalize_commit(&config, &run_dir).await; + let outcome = engine_result?; eprintln!("\n{}", styles.bold.apply_to("=== Run Result ==="),); @@ -2076,6 +2082,38 @@ async fn run_preflight( } } +/// Write a finalize commit to the shadow branch with retro.json and final node files. +/// +/// This captures the last diff.patch (written after the final checkpoint) and retro.json. +/// Best-effort: errors are logged as warnings. +async fn write_finalize_commit(config: &RunConfig, run_dir: &std::path::Path) { + let (Some(ref meta_branch), Some(ref repo_path)) = + (&config.meta_branch, &config.host_repo_path) + else { + return; + }; + + let store = crate::git::MetadataStore::new(repo_path, &config.git_author); + let mut entries = crate::git::scan_node_files(run_dir); + if let Ok(retro_bytes) = std::fs::read(run_dir.join("retro.json")) { + entries.push(("retro.json".to_string(), retro_bytes)); + } + let refs: Vec<(&str, &[u8])> = entries + .iter() + .map(|(k, v)| (k.as_str(), v.as_slice())) + .collect(); + if let Err(e) = store.write_files(&config.run_id, &refs, "finalize run") { + tracing::warn!(error = %e, "Failed to write finalize commit to metadata branch"); + return; + } + + // Push the finalize commit + let run_id_part = meta_branch.strip_prefix("refs/arc/").unwrap_or(meta_branch); + let refspec = format!("{meta_branch}:refs/heads/arc/meta/{run_id_part}"); + crate::engine::git_push_host(repo_path, &refspec, &config.github_app, "finalize metadata") + .await; +} + /// Generate a retro report for a completed workflow run. /// /// Derives a basic retro from the checkpoint, then optionally runs the retro agent diff --git a/lib/crates/arc-workflows/src/engine.rs b/lib/crates/arc-workflows/src/engine.rs index b6a73ab8e..77bec5054 100644 --- a/lib/crates/arc-workflows/src/engine.rs +++ b/lib/crates/arc-workflows/src/engine.rs @@ -637,7 +637,7 @@ pub async fn git_checkpoint( /// /// Authenticates via a GitHub App installation token so we don't depend /// on the host's ambient git credentials. -async fn git_push_host( +pub(crate) async fn git_push_host( repo_path: &Path, refspec: &str, github_app: &Option, @@ -1211,7 +1211,14 @@ impl WorkflowRunEngine { let store = crate::git::MetadataStore::new(repo_path, &config.git_author); let manifest_bytes = serde_json::to_vec_pretty(&manifest).unwrap_or_default(); let dot_source = std::fs::read(config.run_dir.join("graph.dot")).unwrap_or_default(); - if let Err(e) = store.init_run(&config.run_id, &manifest_bytes, &dot_source) { + let sandbox_json = std::fs::read(config.run_dir.join("sandbox.json")).ok(); + let mut extra_files: Vec<(&str, &[u8])> = Vec::new(); + if let Some(ref data) = sandbox_json { + extra_files.push(("sandbox.json", data)); + } + if let Err(e) = + store.init_run(&config.run_id, &manifest_bytes, &dot_source, &extra_files) + { tracing::warn!(run_id = %config.run_id, error = %e, "Metadata branch init failed"); } } @@ -1844,7 +1851,7 @@ impl WorkflowRunEngine { serde_json::to_vec_pretty(&checkpoint) .ok() .and_then(|cp_json| { - let artifact_entries: Vec<(String, Vec)> = artifact_store + let mut extra_entries: Vec<(String, Vec)> = artifact_store .list() .iter() .filter_map(|info| { @@ -1855,11 +1862,12 @@ impl WorkflowRunEngine { }) }) .collect(); - let artifact_refs: Vec<(&str, &[u8])> = artifact_entries + extra_entries.extend(crate::git::scan_node_files(&config.run_dir)); + let extra_refs: Vec<(&str, &[u8])> = extra_entries .iter() .map(|(k, v)| (k.as_str(), v.as_slice())) .collect(); - match store.write_checkpoint(&config.run_id, &cp_json, &artifact_refs) { + match store.write_checkpoint(&config.run_id, &cp_json, &extra_refs) { Ok(sha) => Some(sha), Err(e) => { context.append_log(format!( diff --git a/lib/crates/arc-workflows/src/git.rs b/lib/crates/arc-workflows/src/git.rs index 91b3bfd6e..ac2a6ff02 100644 --- a/lib/crates/arc-workflows/src/git.rs +++ b/lib/crates/arc-workflows/src/git.rs @@ -270,6 +270,54 @@ pub fn sanitize_ref_component(s: &str) -> String { result.trim_matches('-').to_string() } +/// Filenames allowed in per-node directories on the shadow branch. +const NODE_FILE_ALLOWLIST: &[&str] = &[ + "prompt.md", + "response.md", + "status.json", + "provider_used.json", + "diff.patch", + "script_invocation.json", + "script_timing.json", + "parallel_results.json", +]; + +/// Maximum size (bytes) for a single node file. Files larger than this are skipped. +const MAX_NODE_FILE_SIZE: u64 = 512 * 1024; + +/// Scan `{run_dir}/nodes/` for allowlisted files and return them as +/// `("nodes/{subdir}/{filename}", bytes)` entries suitable for the shadow tree. +pub fn scan_node_files(run_dir: &Path) -> Vec<(String, Vec)> { + let nodes_dir = run_dir.join("nodes"); + let entries = match std::fs::read_dir(&nodes_dir) { + Ok(e) => e, + Err(_) => return Vec::new(), + }; + + let mut result = Vec::new(); + for entry in entries.flatten() { + let path = entry.path(); + if !path.is_dir() { + continue; + } + let subdir_name = match path.file_name().and_then(|n| n.to_str()) { + Some(n) => n.to_string(), + None => continue, + }; + for filename in NODE_FILE_ALLOWLIST { + let file_path = path.join(filename); + match std::fs::metadata(&file_path) { + Ok(meta) if meta.is_file() && meta.len() <= MAX_NODE_FILE_SIZE => {} + _ => continue, + } + if let Ok(data) = std::fs::read(&file_path) { + result.push((format!("nodes/{subdir_name}/{filename}"), data)); + } + } + } + result +} + /// Git-native metadata storage for pipeline runs. /// /// Stores checkpoint data, manifests, and graph DOT on an orphan branch @@ -301,18 +349,39 @@ impl MetadataStore { 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<()> { + /// Initialize a run's metadata branch with manifest, graph DOT, and optional extra files. + pub fn init_run( + &self, + run_id: &str, + manifest_json: &[u8], + graph_dot: &[u8], + extra_files: &[(&str, &[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}")))?; + let mut entries: Vec<(&str, &[u8])> = + vec![("manifest.json", manifest_json), ("graph.dot", graph_dot)]; + entries.extend_from_slice(extra_files); + bs.write_entries(&entries, "init run") + .map_err(|e| git_error(format!("write_entries failed: {e}")))?; + Ok(()) + } + + /// Write arbitrary files to the metadata branch without overwriting checkpoint.json. + pub fn write_files( + &self, + run_id: &str, + entries: &[(&str, &[u8])], + message: &str, + ) -> Result<()> { + let (store, sig) = self.open_store()?; + let branch = Self::branch_name(run_id); + let bs = BranchStore::new(&store, &branch, &sig); + bs.write_entries(entries, message) + .map_err(|e| git_error(format!("write_entries failed: {e}")))?; Ok(()) } @@ -490,7 +559,7 @@ mod tests { let store = MetadataStore::new(dir.path(), &GitAuthor::default()); let manifest = br#"{"run_id":"RUN1","workflow_name":"test","goal":"g","start_time":"2025-01-01T00:00:00Z","node_count":2,"edge_count":1}"#; let dot = b"digraph { start -> end }"; - store.init_run("RUN1", manifest, dot).unwrap(); + store.init_run("RUN1", manifest, dot, &[]).unwrap(); let read_manifest = MetadataStore::read_manifest(dir.path(), "RUN1") .unwrap() @@ -510,7 +579,7 @@ mod tests { init_repo(dir.path()); let store = MetadataStore::new(dir.path(), &GitAuthor::default()); - store.init_run("RUN2", b"{}", b"digraph {}").unwrap(); + store.init_run("RUN2", b"{}", b"digraph {}", &[]).unwrap(); let ctx = crate::context::Context::new(); ctx.set("goal", serde_json::json!("test")); @@ -546,7 +615,7 @@ mod tests { init_repo(dir.path()); let store = MetadataStore::new(dir.path(), &GitAuthor::default()); - store.init_run("RUN3", b"{}", b"digraph {}").unwrap(); + store.init_run("RUN3", b"{}", b"digraph {}", &[]).unwrap(); let ctx = crate::context::Context::new(); let cp1 = crate::checkpoint::Checkpoint::from_context( @@ -599,7 +668,7 @@ mod tests { init_repo(dir.path()); let store = MetadataStore::new(dir.path(), &GitAuthor::default()); - store.init_run("RUN4", b"{}", b"digraph {}").unwrap(); + 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 @@ -617,6 +686,106 @@ mod tests { assert_eq!(read_back, artifact_data); } + #[test] + fn scan_node_files_picks_up_allowlisted() { + let dir = tempfile::tempdir().unwrap(); + let run_dir = dir.path(); + let node_dir = run_dir.join("nodes").join("work"); + fs::create_dir_all(&node_dir).unwrap(); + fs::write(node_dir.join("prompt.md"), "hello").unwrap(); + fs::write(node_dir.join("response.md"), "world").unwrap(); + fs::write(node_dir.join("not_allowed.txt"), "skip me").unwrap(); + + let files = scan_node_files(run_dir); + let paths: Vec<&str> = files.iter().map(|(p, _)| p.as_str()).collect(); + assert!(paths.contains(&"nodes/work/prompt.md")); + assert!(paths.contains(&"nodes/work/response.md")); + assert!(!paths.iter().any(|p| p.contains("not_allowed"))); + } + + #[test] + fn scan_node_files_skips_oversized() { + let dir = tempfile::tempdir().unwrap(); + let run_dir = dir.path(); + let node_dir = run_dir.join("nodes").join("big"); + fs::create_dir_all(&node_dir).unwrap(); + // Write a file just over the 512KB limit + let big_data = vec![0u8; 512 * 1024 + 1]; + fs::write(node_dir.join("prompt.md"), &big_data).unwrap(); + + let files = scan_node_files(run_dir); + assert!(files.is_empty()); + } + + #[test] + fn scan_node_files_handles_visit_suffixes() { + let dir = tempfile::tempdir().unwrap(); + let run_dir = dir.path(); + let node_dir = run_dir.join("nodes").join("work-visit_2"); + fs::create_dir_all(&node_dir).unwrap(); + fs::write(node_dir.join("status.json"), "{}").unwrap(); + + let files = scan_node_files(run_dir); + assert_eq!(files.len(), 1); + assert_eq!(files[0].0, "nodes/work-visit_2/status.json"); + } + + #[test] + fn scan_node_files_empty_when_no_nodes_dir() { + let dir = tempfile::tempdir().unwrap(); + let files = scan_node_files(dir.path()); + assert!(files.is_empty()); + } + + #[test] + fn metadata_store_write_files() { + let dir = tempfile::tempdir().unwrap(); + init_repo(dir.path()); + + let store = MetadataStore::new(dir.path(), &GitAuthor::default()); + store.init_run("RUN5", b"{}", b"digraph {}", &[]).unwrap(); + + store + .write_files( + "RUN5", + &[("retro.json", b"{\"status\":\"ok\"}")], + "finalize", + ) + .unwrap(); + + let data = MetadataStore::read_file(dir.path(), "RUN5", "retro.json") + .unwrap() + .unwrap(); + assert_eq!(data, b"{\"status\":\"ok\"}"); + + // Original files still present + let dot = MetadataStore::read_graph_dot(dir.path(), "RUN5") + .unwrap() + .unwrap(); + assert_eq!(dot, "digraph {}"); + } + + #[test] + fn metadata_store_init_run_with_extra_files() { + let dir = tempfile::tempdir().unwrap(); + init_repo(dir.path()); + + let store = MetadataStore::new(dir.path(), &GitAuthor::default()); + store + .init_run( + "RUN6", + b"{}", + b"digraph {}", + &[("sandbox.json", b"{\"type\":\"local\"}")], + ) + .unwrap(); + + let data = MetadataStore::read_file(dir.path(), "RUN6", "sandbox.json") + .unwrap() + .unwrap(); + assert_eq!(data, b"{\"type\":\"local\"}"); + } + #[test] fn sanitize_ref_component_lowercases() { assert_eq!(sanitize_ref_component("Hello"), "hello");