Add per-node files and run-level metadata to shadow branch

Store execution trace data (prompts, responses, status, diffs, etc.) on
the shadow branch alongside checkpoint/manifest data. This makes the
shadow branch a complete, self-contained record of each run.

- scan_node_files() scans nodes/ for allowlisted files (<512KB)
- Checkpoint commits now include per-node files
- init_run includes sandbox.json
- Finalize commit after workflow writes retro.json + final node files

Co-Authored-By: Claude Opus 4.6 (1M context) <noreply@anthropic.com>
This commit is contained in:
Bryan Helmkamp 2026-03-11 00:37:52 -04:00
parent a2288b7598
commit 0bdb8e93b1
3 changed files with 231 additions and 16 deletions

View file

@ -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<String> = None;
let mut pr_url: Option<String> = 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

View file

@ -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<arc_github::GitHubAppCredentials>,
@ -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<u8>)> = artifact_store
let mut extra_entries: Vec<(String, Vec<u8>)> = 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!(

View file

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