From eb1933954a96023af1b119b1b9b31dfbf585da39 Mon Sep 17 00:00:00 2001 From: Bryan Helmkamp Date: Wed, 1 Apr 2026 22:27:47 -0400 Subject: [PATCH] Export checkpoint metadata from the run store --- lib/crates/fabro-workflow/src/git.rs | 155 ++++++++++++++++++ .../fabro-workflow/src/lifecycle/git.rs | 52 +++--- .../fabro-workflow/src/pipeline/finalize.rs | 8 +- 3 files changed, 187 insertions(+), 28 deletions(-) diff --git a/lib/crates/fabro-workflow/src/git.rs b/lib/crates/fabro-workflow/src/git.rs index 53d593e58..0236a1446 100644 --- a/lib/crates/fabro-workflow/src/git.rs +++ b/lib/crates/fabro-workflow/src/git.rs @@ -3,6 +3,7 @@ use std::process::Command; use fabro_checkpoint::git::Store; use fabro_config::FabroSettings; +use fabro_store::{NodeVisitRef, RunStore}; use crate::error::{FabroError, Result}; use tokio::task::{JoinError, spawn_blocking}; @@ -352,9 +353,97 @@ pub fn scan_node_files(run_dir: &Path) -> Vec<(String, Vec)> { result } +pub async fn scan_node_files_from_store(run_store: &dyn RunStore) -> Vec<(String, Vec)> { + let mut result = Vec::new(); + let Ok(node_ids) = run_store.list_node_ids().await else { + return result; + }; + + for node_id in node_ids { + let Ok(visits) = run_store.list_node_visits(&node_id).await else { + continue; + }; + for visit in visits { + let Ok(node) = run_store + .get_node(&NodeVisitRef { + node_id: &node_id, + visit, + }) + .await + else { + continue; + }; + + if let Some(prompt) = node.prompt { + result.push(( + node_file_path(&node_id, visit, "prompt.md"), + prompt.into_bytes(), + )); + } + if let Some(response) = node.response { + result.push(( + node_file_path(&node_id, visit, "response.md"), + response.into_bytes(), + )); + } + if let Some(status) = node.status { + if let Ok(bytes) = serde_json::to_vec_pretty(&status) { + result.push((node_file_path(&node_id, visit, "status.json"), bytes)); + } + } + if let Some(provider_used) = node.provider_used { + if let Ok(bytes) = serde_json::to_vec_pretty(&provider_used) { + result.push((node_file_path(&node_id, visit, "provider_used.json"), bytes)); + } + } + if let Some(diff) = node.diff { + result.push(( + node_file_path(&node_id, visit, "diff.patch"), + diff.into_bytes(), + )); + } + if let Some(script_invocation) = node.script_invocation { + if let Ok(bytes) = serde_json::to_vec_pretty(&script_invocation) { + result.push(( + node_file_path(&node_id, visit, "script_invocation.json"), + bytes, + )); + } + } + if let Some(script_timing) = node.script_timing { + if let Ok(bytes) = serde_json::to_vec_pretty(&script_timing) { + result.push((node_file_path(&node_id, visit, "script_timing.json"), bytes)); + } + } + if let Some(parallel_results) = node.parallel_results { + if let Ok(bytes) = serde_json::to_vec_pretty(¶llel_results) { + result.push(( + node_file_path(&node_id, visit, "parallel_results.json"), + bytes, + )); + } + } + } + } + + result +} + +fn node_file_path(node_id: &str, visit: u32, filename: &str) -> String { + if visit <= 1 { + format!("nodes/{node_id}/{filename}") + } else { + format!("nodes/{node_id}-visit_{visit}/{filename}") + } +} + #[cfg(test)] mod tests { use super::*; + use chrono::Utc; + use fabro_graphviz::graph::Graph; + use fabro_store::{InMemoryStore, Store}; + use fabro_types::{NodeStatusRecord, RunRecord, StageStatus, fixtures}; use std::fs; /// Create a temporary git repo with an initial commit. @@ -492,6 +581,72 @@ mod tests { assert!(files.is_empty()); } + #[tokio::test] + async fn scan_node_files_from_store_reconstructs_allowlisted_entries() { + let store = InMemoryStore::default(); + let created_at = Utc::now(); + let run = store + .create_run(&fixtures::RUN_1, created_at, None) + .await + .unwrap(); + run.put_run(&RunRecord { + run_id: fixtures::RUN_1, + created_at, + settings: fabro_config::FabroSettings::default(), + graph: Graph::new("test"), + workflow_slug: None, + working_directory: std::path::PathBuf::from("."), + host_repo_path: None, + base_branch: None, + labels: std::collections::HashMap::new(), + }) + .await + .unwrap(); + let node = NodeVisitRef { + node_id: "work", + visit: 2, + }; + run.put_node_prompt(&node, "hello").await.unwrap(); + run.put_node_response(&node, "world").await.unwrap(); + run.put_node_status( + &node, + &NodeStatusRecord { + status: StageStatus::Success, + notes: None, + failure_reason: None, + timestamp: Utc::now(), + }, + ) + .await + .unwrap(); + run.put_node_provider_used(&node, &serde_json::json!({"provider":"openai"})) + .await + .unwrap(); + run.put_node_diff(&node, "diff --git a/story.txt b/story.txt") + .await + .unwrap(); + run.put_node_script_invocation(&node, &serde_json::json!({"command":"echo hi"})) + .await + .unwrap(); + run.put_node_script_timing(&node, &serde_json::json!({"exit_code":0})) + .await + .unwrap(); + run.put_node_parallel_results(&node, &serde_json::json!([{"id":"a"}])) + .await + .unwrap(); + + let files = scan_node_files_from_store(run.as_ref()).await; + let paths: Vec<&str> = files.iter().map(|(path, _)| path.as_str()).collect(); + assert!(paths.contains(&"nodes/work-visit_2/prompt.md")); + assert!(paths.contains(&"nodes/work-visit_2/response.md")); + assert!(paths.contains(&"nodes/work-visit_2/status.json")); + assert!(paths.contains(&"nodes/work-visit_2/provider_used.json")); + assert!(paths.contains(&"nodes/work-visit_2/diff.patch")); + assert!(paths.contains(&"nodes/work-visit_2/script_invocation.json")); + assert!(paths.contains(&"nodes/work-visit_2/script_timing.json")); + assert!(paths.contains(&"nodes/work-visit_2/parallel_results.json")); + } + #[test] fn sanitize_ref_component_lowercases() { assert_eq!(sanitize_ref_component("Hello"), "hello"); diff --git a/lib/crates/fabro-workflow/src/lifecycle/git.rs b/lib/crates/fabro-workflow/src/lifecycle/git.rs index ca31f267f..3c762d167 100644 --- a/lib/crates/fabro-workflow/src/lifecycle/git.rs +++ b/lib/crates/fabro-workflow/src/lifecycle/git.rs @@ -14,7 +14,7 @@ use fabro_core::state::RunState; use crate::artifact::ArtifactStore; use crate::event::{EventEmitter, RunNoticeLevel, WorkflowRunEvent}; use crate::git::MetadataStore; -use crate::git::scan_node_files; +use crate::git::scan_node_files_from_store; use crate::graph::WorkflowGraph; use crate::graph::WorkflowNode; use crate::outcome::{Outcome, StageStatus, StageUsage}; @@ -137,16 +137,17 @@ impl RunLifecycle for GitLifecycle { let git_author = self.run_options.git_author(); let store = MetadataStore::new(repo_path, &git_author); // Build checkpoint JSON for shadow branch - self.run_store + if let Some(cp_json) = self + .run_store .get_checkpoint() .await .ok() .flatten() .and_then(|checkpoint| serde_json::to_vec_pretty(&checkpoint).ok()) - .or_else(|| std::fs::read(self.run_dir.join("checkpoint.json")).ok()) - .and_then(|cp_json| { + { + let mut extra_entries: Vec<(String, Vec)> = { let artifact_store = self.artifact_store.lock().unwrap(); - let mut extra_entries: Vec<(String, Vec)> = artifact_store + artifact_store .list() .iter() .filter_map(|info| { @@ -156,26 +157,29 @@ impl RunLifecycle for GitLifecycle { .map(|data| (format!("artifacts/{}.json", info.id), data)) }) }) - .collect(); - extra_entries.extend(scan_node_files(&self.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(&self.run_id.to_string(), &cp_json, &extra_refs) { - Ok(sha) => Some(sha), - Err(e) => { - self.emitter.emit(&WorkflowRunEvent::RunNotice { - level: RunNoticeLevel::Warn, - code: "checkpoint_metadata_write_failed".to_string(), - message: format!( - "[node: {node_id}] metadata checkpoint write failed: {e}" - ), - }); - None - } + .collect() + }; + extra_entries.extend(scan_node_files_from_store(self.run_store.as_ref()).await); + let extra_refs: Vec<(&str, &[u8])> = extra_entries + .iter() + .map(|(k, v)| (k.as_str(), v.as_slice())) + .collect(); + match store.write_checkpoint(&self.run_id.to_string(), &cp_json, &extra_refs) { + Ok(sha) => Some(sha), + Err(e) => { + self.emitter.emit(&WorkflowRunEvent::RunNotice { + level: RunNoticeLevel::Warn, + code: "checkpoint_metadata_write_failed".to_string(), + message: format!( + "[node: {node_id}] metadata checkpoint write failed: {e}" + ), + }); + None } - }) + } + } else { + None + } } else { None }; diff --git a/lib/crates/fabro-workflow/src/pipeline/finalize.rs b/lib/crates/fabro-workflow/src/pipeline/finalize.rs index f9731637a..ee96fd642 100644 --- a/lib/crates/fabro-workflow/src/pipeline/finalize.rs +++ b/lib/crates/fabro-workflow/src/pipeline/finalize.rs @@ -3,7 +3,7 @@ use std::sync::Arc; use crate::error::FabroError; use crate::event::{EventEmitter, RunNoticeLevel, WorkflowRunEvent}; -use crate::git::{MetadataStore, scan_node_files}; +use crate::git::{MetadataStore, scan_node_files_from_store}; use crate::outcome::{Outcome, OutcomeExt, StageStatus}; use crate::records::{Checkpoint, CheckpointExt, Conclusion, ConclusionExt, StageSummary}; use crate::run_options::RunOptions; @@ -264,7 +264,7 @@ pub fn persist_terminal_outcome( /// Best-effort: errors are logged as warnings. pub async fn write_finalize_commit( run_options: &RunOptions, - run_dir: &Path, + _run_dir: &Path, run_store: &dyn RunStore, ) { let (Some(meta_branch), Some(repo_path)) = ( @@ -279,10 +279,10 @@ pub async fn write_finalize_commit( let git_author = run_options.git_author(); let store = MetadataStore::new(repo_path, &git_author); - let mut entries = scan_node_files(run_dir); + let mut entries = scan_node_files_from_store(run_store).await; let retro_bytes = match run_store.get_retro().await { Ok(Some(retro)) => serde_json::to_vec_pretty(&retro).ok(), - _ => std::fs::read(run_dir.join("retro.json")).ok(), + _ => None, }; if let Some(bytes) = retro_bytes { entries.push(("retro.json".to_string(), bytes));