From e9d924fa63d59133dee0b35f75dbc7ea08de9271 Mon Sep 17 00:00:00 2001 From: Bryan Helmkamp Date: Mon, 30 Mar 2026 15:22:16 -0400 Subject: [PATCH] Rename fabro-git-storage to fabro-checkpoint --- AGENTS.md | 2 +- Cargo.lock | 29 +- .../Cargo.toml | 9 +- lib/crates/fabro-checkpoint/src/author.rs | 56 ++ .../src/branch.rs} | 4 +- .../src/error.rs | 14 + .../gitobj.rs => fabro-checkpoint/src/git.rs} | 4 +- lib/crates/fabro-checkpoint/src/lib.rs | 10 + lib/crates/fabro-checkpoint/src/metadata.rs | 436 ++++++++++++ .../src/trailer.rs} | 0 lib/crates/fabro-cli/Cargo.toml | 2 +- lib/crates/fabro-cli/src/commands/run/fork.rs | 2 +- .../fabro-cli/src/commands/run/rewind.rs | 2 +- .../fabro-cli/tests/it/scenario/recovery.rs | 4 +- lib/crates/fabro-git-storage/src/lib.rs | 7 - lib/crates/fabro-git-storage/src/snapshot.rs | 649 ------------------ lib/crates/fabro-workflow/Cargo.toml | 2 +- lib/crates/fabro-workflow/src/error.rs | 13 + lib/crates/fabro-workflow/src/git.rs | 402 +---------- .../fabro-workflow/src/operations/fork.rs | 4 +- .../src/operations/rebuild_meta.rs | 4 +- .../fabro-workflow/src/operations/rewind.rs | 4 +- .../src/operations/test_support.rs | 2 +- lib/crates/fabro-workflow/src/sandbox_git.rs | 3 +- 24 files changed, 576 insertions(+), 1088 deletions(-) rename lib/crates/{fabro-git-storage => fabro-checkpoint}/Cargo.toml (60%) create mode 100644 lib/crates/fabro-checkpoint/src/author.rs rename lib/crates/{fabro-git-storage/src/branchstore.rs => fabro-checkpoint/src/branch.rs} (99%) rename lib/crates/{fabro-git-storage => fabro-checkpoint}/src/error.rs (56%) rename lib/crates/{fabro-git-storage/src/gitobj.rs => fabro-checkpoint/src/git.rs} (99%) create mode 100644 lib/crates/fabro-checkpoint/src/lib.rs create mode 100644 lib/crates/fabro-checkpoint/src/metadata.rs rename lib/crates/{fabro-git-storage/src/trailerlink.rs => fabro-checkpoint/src/trailer.rs} (100%) delete mode 100644 lib/crates/fabro-git-storage/src/lib.rs delete mode 100644 lib/crates/fabro-git-storage/src/snapshot.rs diff --git a/AGENTS.md b/AGENTS.md index 0069a072e..9539a80c3 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -61,7 +61,7 @@ Fabro is an AI-powered workflow orchestration platform. Workflows are defined as - **fabro-mcp** — Model Context Protocol client/server - **fabro-slack** — Slack integration (socket mode, blocks API) - **fabro-devcontainer** — Parses `.devcontainer/devcontainer.json` for container setup -- **fabro-git-storage** — Git-based storage with branch store and snapshots +- **fabro-checkpoint** — Git-based checkpoint storage with branch store and metadata branches - **fabro-telemetry** — CLI analytics (Segment) and crash reporting (Sentry), with anonymous IDs, command sanitization, and detached subprocess delivery - **fabro-util** — Shared utilities (redaction, terminal formatting) diff --git a/Cargo.lock b/Cargo.lock index c701c8c72..4309f9aa5 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1446,6 +1446,20 @@ dependencies = [ "uuid", ] +[[package]] +name = "fabro-checkpoint" +version = "0.176.2" +dependencies = [ + "chrono", + "fabro-types", + "git2", + "serde", + "serde_json", + "tempfile", + "thiserror 2.0.18", + "tracing", +] + [[package]] name = "fabro-cli" version = "0.176.2" @@ -1466,9 +1480,9 @@ dependencies = [ "dirs", "dotenvy", "fabro-agent", + "fabro-checkpoint", "fabro-config", "fabro-devcontainer", - "fabro-git-storage", "fabro-github", "fabro-graphviz", "fabro-hooks", @@ -1581,17 +1595,6 @@ dependencies = [ "tracing", ] -[[package]] -name = "fabro-git-storage" -version = "0.176.2" -dependencies = [ - "git2", - "tempfile", - "thiserror 2.0.18", - "tracing", - "walkdir", -] - [[package]] name = "fabro-github" version = "0.176.2" @@ -1993,10 +1996,10 @@ dependencies = [ "dirs", "dotenvy", "fabro-agent", + "fabro-checkpoint", "fabro-config", "fabro-core", "fabro-devcontainer", - "fabro-git-storage", "fabro-github", "fabro-graphviz", "fabro-hooks", diff --git a/lib/crates/fabro-git-storage/Cargo.toml b/lib/crates/fabro-checkpoint/Cargo.toml similarity index 60% rename from lib/crates/fabro-git-storage/Cargo.toml rename to lib/crates/fabro-checkpoint/Cargo.toml index 4580c4d80..02a87e246 100644 --- a/lib/crates/fabro-git-storage/Cargo.toml +++ b/lib/crates/fabro-checkpoint/Cargo.toml @@ -1,10 +1,10 @@ [package] -name = "fabro-git-storage" +name = "fabro-checkpoint" edition.workspace = true version.workspace = true publish = false license.workspace = true -description = "Store structured data in git without touching the working directory" +description = "Git-backed checkpoint storage for Fabro workflows" repository = "https://github.com/brynary/arc" [lib] @@ -14,10 +14,13 @@ doctest = false workspace = true [dependencies] +fabro-types = { path = "../fabro-types" } git2.workspace = true +serde.workspace = true +serde_json.workspace = true thiserror.workspace = true tracing.workspace = true -walkdir.workspace = true [dev-dependencies] +chrono.workspace = true tempfile = "3" diff --git a/lib/crates/fabro-checkpoint/src/author.rs b/lib/crates/fabro-checkpoint/src/author.rs new file mode 100644 index 000000000..857ed5572 --- /dev/null +++ b/lib/crates/fabro-checkpoint/src/author.rs @@ -0,0 +1,56 @@ +use std::fmt::Write; + +use fabro_types::settings::server::GitAuthorSettings; + +/// Resolved git author identity for checkpoint commits. +#[derive(Debug, Clone, PartialEq)] +pub struct GitAuthor { + pub name: String, + pub email: String, +} + +impl Default for GitAuthor { + fn default() -> Self { + Self { + name: "Fabro".into(), + email: "noreply@fabro.sh".into(), + } + } +} + +impl GitAuthor { + /// Create a `GitAuthor` from optional name/email, falling back to defaults. + pub fn from_options(name: Option, email: Option) -> Self { + let defaults = Self::default(); + Self { + name: name.unwrap_or(defaults.name), + email: email.unwrap_or(defaults.email), + } + } + + /// Returns true when this identity matches the default Fabro identity. + pub fn is_default(&self) -> bool { + let defaults = Self::default(); + self.name == defaults.name && self.email == defaults.email + } + + /// Append the Fabro footer (and Co-Authored-By when the author is not the + /// default identity) to a commit message. + pub fn append_footer(&self, message: &mut String) { + message.push_str("\n\u{2692}\u{fe0f} Generated with [Fabro](https://fabro.sh)\n"); + if !self.is_default() { + let defaults = Self::default(); + let _ = write!( + message, + "\nCo-Authored-By: {} <{}>\n", + defaults.name, defaults.email + ); + } + } +} + +impl From<&GitAuthorSettings> for GitAuthor { + fn from(value: &GitAuthorSettings) -> Self { + Self::from_options(value.name.clone(), value.email.clone()) + } +} diff --git a/lib/crates/fabro-git-storage/src/branchstore.rs b/lib/crates/fabro-checkpoint/src/branch.rs similarity index 99% rename from lib/crates/fabro-git-storage/src/branchstore.rs rename to lib/crates/fabro-checkpoint/src/branch.rs index 408cb040f..bf849ea45 100644 --- a/lib/crates/fabro-git-storage/src/branchstore.rs +++ b/lib/crates/fabro-checkpoint/src/branch.rs @@ -2,7 +2,7 @@ use git2::{Oid, Signature}; use tracing::{debug, warn}; use crate::Result; -use crate::gitobj::{FileMode, Store, TreeEntries}; +use crate::git::{FileMode, Store, TreeEntries}; /// Metadata about a commit, returned by `log`. #[derive(Debug)] @@ -219,7 +219,7 @@ pub fn sharded_path(id: &str, prefix_len: usize) -> String { #[cfg(test)] mod tests { use super::*; - use crate::gitobj::FileMode; + use crate::git::FileMode; use git2::Repository; fn temp_repo() -> (tempfile::TempDir, Store) { diff --git a/lib/crates/fabro-git-storage/src/error.rs b/lib/crates/fabro-checkpoint/src/error.rs similarity index 56% rename from lib/crates/fabro-git-storage/src/error.rs rename to lib/crates/fabro-checkpoint/src/error.rs index 07f3ee442..8f186db80 100644 --- a/lib/crates/fabro-git-storage/src/error.rs +++ b/lib/crates/fabro-checkpoint/src/error.rs @@ -16,3 +16,17 @@ pub enum Error { #[error("branch {branch} not found")] BranchNotFound { branch: String }, } + +#[derive(Debug, thiserror::Error)] +pub enum MetadataError { + #[error(transparent)] + Storage(#[from] Error), + + #[error("deserialize {entity} on branch {branch}")] + Deserialize { + entity: &'static str, + branch: String, + #[source] + source: serde_json::Error, + }, +} diff --git a/lib/crates/fabro-git-storage/src/gitobj.rs b/lib/crates/fabro-checkpoint/src/git.rs similarity index 99% rename from lib/crates/fabro-git-storage/src/gitobj.rs rename to lib/crates/fabro-checkpoint/src/git.rs index 1a4e2cbe4..ce4d1f2cc 100644 --- a/lib/crates/fabro-git-storage/src/gitobj.rs +++ b/lib/crates/fabro-checkpoint/src/git.rs @@ -565,7 +565,7 @@ mod tests { fn read_blob_at_returns_content() { let (_dir, store) = temp_repo(); let sig = Signature::now("Test", "test@example.com").unwrap(); - let bs = crate::branchstore::BranchStore::new(&store, "test/data", &sig); + let bs = crate::branch::BranchStore::new(&store, "test/data", &sig); bs.ensure_branch().unwrap(); bs.write_entry("hello.txt", b"world", "add hello").unwrap(); @@ -580,7 +580,7 @@ mod tests { fn read_blob_at_returns_none_for_missing_path() { let (_dir, store) = temp_repo(); let sig = Signature::now("Test", "test@example.com").unwrap(); - let bs = crate::branchstore::BranchStore::new(&store, "test/data", &sig); + let bs = crate::branch::BranchStore::new(&store, "test/data", &sig); bs.ensure_branch().unwrap(); bs.write_entry("hello.txt", b"world", "add hello").unwrap(); diff --git a/lib/crates/fabro-checkpoint/src/lib.rs b/lib/crates/fabro-checkpoint/src/lib.rs new file mode 100644 index 000000000..7bdc82a5e --- /dev/null +++ b/lib/crates/fabro-checkpoint/src/lib.rs @@ -0,0 +1,10 @@ +pub mod author; +pub mod branch; +pub mod error; +pub mod git; +pub mod metadata; +pub mod trailer; + +pub const META_BRANCH_PREFIX: &str = "fabro/meta/"; + +pub use error::{Error, MetadataError, Result}; diff --git a/lib/crates/fabro-checkpoint/src/metadata.rs b/lib/crates/fabro-checkpoint/src/metadata.rs new file mode 100644 index 000000000..e73800870 --- /dev/null +++ b/lib/crates/fabro-checkpoint/src/metadata.rs @@ -0,0 +1,436 @@ +use std::path::{Path, PathBuf}; + +use fabro_types::{Checkpoint, RunRecord, StartRecord}; +use git2::{Repository, Signature}; + +use crate::META_BRANCH_PREFIX; +use crate::author::GitAuthor; +use crate::branch::BranchStore; +use crate::error::{Error, MetadataError}; +use crate::git::Store; + +/// Git-native metadata storage for pipeline runs. +/// +/// Stores checkpoint data, run records, and metadata on an orphan branch +/// (`fabro/meta/{run_id}`) so that runs can be resumed from git alone. +pub struct MetadataStore { + repo_path: PathBuf, + author: GitAuthor, +} + +impl MetadataStore { + pub fn new(repo_path: impl Into, author: &GitAuthor) -> Self { + Self { + repo_path: repo_path.into(), + author: author.clone(), + } + } + + /// Returns the branch name for a run: `fabro/meta/{run_id}`. + pub fn branch_name(run_id: &str) -> String { + format!("{META_BRANCH_PREFIX}{run_id}") + } + + /// Format a commit message with the standard Fabro footer appended. + fn commit_message(&self, subject: &str) -> String { + let mut msg = format!("{subject}\n"); + self.author.append_footer(&mut msg); + msg + } + + fn open_store(&self) -> Result<(Store, Signature<'static>), MetadataError> { + let repo = Repository::discover(&self.repo_path).map_err(Error::from)?; + let store = Store::new(repo); + let sig = Signature::now(&self.author.name, &self.author.email).map_err(Error::from)?; + Ok((store, sig)) + } + + /// Initialize a run's metadata branch with the given files. + /// + /// Callers pass all files (run.json, start.json, sandbox.json, etc.) + /// via the `files` slice. + pub fn init_run(&self, run_id: &str, files: &[(&str, &[u8])]) -> Result<(), MetadataError> { + let (store, sig) = self.open_store()?; + let branch = Self::branch_name(run_id); + let branch_store = BranchStore::new(&store, &branch, &sig); + branch_store.ensure_branch()?; + let message = self.commit_message("init run"); + branch_store.write_entries(files, &message)?; + 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<(), MetadataError> { + let (store, sig) = self.open_store()?; + let branch = Self::branch_name(run_id); + let branch_store = BranchStore::new(&store, &branch, &sig); + let message = self.commit_message(message); + branch_store.write_entries(entries, &message)?; + Ok(()) + } + + /// Write checkpoint data (and optional artifacts) to the metadata branch. + /// Returns the SHA of the new commit on the shadow branch. + pub fn write_checkpoint( + &self, + run_id: &str, + checkpoint_json: &[u8], + artifacts: &[(&str, &[u8])], + ) -> Result { + let (store, sig) = self.open_store()?; + let branch = Self::branch_name(run_id); + let branch_store = BranchStore::new(&store, &branch, &sig); + let mut entries: Vec<(&str, &[u8])> = vec![("checkpoint.json", checkpoint_json)]; + entries.extend_from_slice(artifacts); + let message = self.commit_message("checkpoint"); + let oid = branch_store.write_entries(&entries, &message)?; + 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>, MetadataError> { + let Ok(repo) = Repository::discover(repo_path) else { + return Ok(None); + }; + let store = Store::new(repo); + let sig = Signature::now("Fabro", "noreply@fabro.sh").map_err(Error::from)?; + let branch = Self::branch_name(run_id); + let branch_store = BranchStore::new(&store, &branch, &sig); + Ok(branch_store.read_entry(path)?) + } + + /// 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, MetadataError> { + let branch = Self::branch_name(run_id); + match Self::read_file(repo_path, run_id, "checkpoint.json")? { + Some(bytes) => serde_json::from_slice(&bytes).map(Some).map_err(|source| { + MetadataError::Deserialize { + entity: "checkpoint", + branch, + source, + } + }), + None => Ok(None), + } + } + + /// Read the run record from the metadata branch. Returns `None` if not found. + pub fn read_run_record( + repo_path: &Path, + run_id: &str, + ) -> Result, MetadataError> { + let branch = Self::branch_name(run_id); + match Self::read_file(repo_path, run_id, "run.json")? { + Some(bytes) => serde_json::from_slice(&bytes).map(Some).map_err(|source| { + MetadataError::Deserialize { + entity: "run record", + branch, + source, + } + }), + None => Ok(None), + } + } + + /// Read the start record from the metadata branch. Returns `None` if not found. + pub fn read_start_record( + repo_path: &Path, + run_id: &str, + ) -> Result, MetadataError> { + let branch = Self::branch_name(run_id); + match Self::read_file(repo_path, run_id, "start.json")? { + Some(bytes) => serde_json::from_slice(&bytes).map(Some).map_err(|source| { + MetadataError::Deserialize { + entity: "start record", + branch, + source, + } + }), + 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>, MetadataError> { + Self::read_file(repo_path, run_id, &format!("artifacts/{key}.json")) + } +} + +#[cfg(test)] +mod tests { + use std::collections::HashMap; + + use super::*; + use chrono::{TimeZone, Utc}; + use fabro_types::{FabroSettings, Graph, fixtures}; + + /// Create a temporary git repo with an initial commit. + fn init_repo(dir: &Path) { + std::process::Command::new("git") + .args(["init"]) + .current_dir(dir) + .output() + .unwrap(); + std::process::Command::new("git") + .args([ + "-c", + "user.name=test", + "-c", + "user.email=test@test", + "commit", + "--allow-empty", + "-m", + "init", + ]) + .current_dir(dir) + .output() + .unwrap(); + } + + fn test_run_record(run_id: fabro_types::RunId) -> RunRecord { + RunRecord { + run_id, + created_at: Utc.with_ymd_and_hms(2025, 1, 1, 0, 0, 0).single().unwrap(), + settings: FabroSettings::default(), + graph: Graph::new("test"), + workflow_slug: None, + working_directory: PathBuf::from("/tmp"), + host_repo_path: None, + base_branch: None, + labels: HashMap::new(), + } + } + + fn test_checkpoint( + current_node: &str, + completed_nodes: Vec, + next_node_id: Option, + ) -> Checkpoint { + Checkpoint { + timestamp: Utc.with_ymd_and_hms(2025, 1, 1, 0, 0, 0).single().unwrap(), + current_node: current_node.to_string(), + completed_nodes, + node_retries: HashMap::new(), + context_values: HashMap::new(), + node_outcomes: HashMap::new(), + next_node_id, + git_commit_sha: None, + loop_failure_signatures: HashMap::new(), + restart_failure_signatures: HashMap::new(), + node_visits: HashMap::new(), + } + } + + fn branch_entry(repo_dir: &Path, run_id: &str, path: &str) -> Vec { + let repo = Repository::discover(repo_dir).unwrap(); + let store = Store::new(repo); + let sig = Signature::now("Test", "test@example.com").unwrap(); + let branch = MetadataStore::branch_name(run_id); + let branch_store = BranchStore::new(&store, &branch, &sig); + branch_store.read_entry(path).unwrap().unwrap() + } + + #[test] + fn metadata_store_init_run_and_read() { + let dir = tempfile::tempdir().unwrap(); + init_repo(dir.path()); + + let store = MetadataStore::new(dir.path(), &GitAuthor::default()); + let run_id = fixtures::RUN_1.to_string(); + let run_record = serde_json::to_vec_pretty(&test_run_record(fixtures::RUN_1)).unwrap(); + store + .init_run(&run_id, &[("run.json", &run_record)]) + .unwrap(); + + let read_record = MetadataStore::read_run_record(dir.path(), &run_id) + .unwrap() + .unwrap(); + assert_eq!(read_record.run_id, fixtures::RUN_1); + assert_eq!(read_record.graph.name, "test"); + } + + #[test] + fn metadata_store_write_and_read_checkpoint() { + let dir = tempfile::tempdir().unwrap(); + init_repo(dir.path()); + + let run_id = fixtures::RUN_2.to_string(); + let store = MetadataStore::new(dir.path(), &GitAuthor::default()); + store.init_run(&run_id, &[]).unwrap(); + + let mut checkpoint = test_checkpoint( + "node_a", + vec!["start".to_string()], + Some("node_b".to_string()), + ); + checkpoint + .context_values + .insert("goal".to_string(), serde_json::json!("test")); + let checkpoint_json = serde_json::to_vec_pretty(&checkpoint).unwrap(); + store + .write_checkpoint(&run_id, &checkpoint_json, &[]) + .unwrap(); + + let loaded = MetadataStore::read_checkpoint(dir.path(), &run_id) + .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 run_id = fixtures::RUN_3.to_string(); + let store = MetadataStore::new(dir.path(), &GitAuthor::default()); + store.init_run(&run_id, &[]).unwrap(); + + let checkpoint_one = + serde_json::to_vec_pretty(&test_checkpoint("node_a", vec!["start".to_string()], None)) + .unwrap(); + store + .write_checkpoint(&run_id, &checkpoint_one, &[]) + .unwrap(); + + let checkpoint_two = serde_json::to_vec_pretty(&test_checkpoint( + "node_b", + vec!["start".to_string(), "node_a".to_string()], + Some("node_c".to_string()), + )) + .unwrap(); + store + .write_checkpoint(&run_id, &checkpoint_two, &[]) + .unwrap(); + + let loaded = MetadataStore::read_checkpoint(dir.path(), &run_id) + .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 run_id = fixtures::RUN_4.to_string(); + let store = MetadataStore::new(dir.path(), &GitAuthor::default()); + store.init_run(&run_id, &[]).unwrap(); + + let artifact_data = br#"{"large_output":"some data"}"#; + let checkpoint_json = + serde_json::to_vec_pretty(&test_checkpoint("node_a", Vec::new(), None)).unwrap(); + store + .write_checkpoint( + &run_id, + &checkpoint_json, + &[("artifacts/response.plan.json", artifact_data.as_slice())], + ) + .unwrap(); + + let read_back = MetadataStore::read_artifact(dir.path(), &run_id, "response.plan") + .unwrap() + .unwrap(); + assert_eq!(read_back, artifact_data); + } + + #[test] + fn metadata_store_write_files() { + let dir = tempfile::tempdir().unwrap(); + init_repo(dir.path()); + + let run_id = fixtures::RUN_5.to_string(); + let store = MetadataStore::new(dir.path(), &GitAuthor::default()); + let run_record = serde_json::to_vec_pretty(&test_run_record(fixtures::RUN_5)).unwrap(); + store + .init_run(&run_id, &[("run.json", &run_record)]) + .unwrap(); + + store + .write_files( + &run_id, + &[("retro.json", b"{\"status\":\"ok\"}")], + "finalize run", + ) + .unwrap(); + + let data = branch_entry(dir.path(), &run_id, "retro.json"); + assert_eq!(data, b"{\"status\":\"ok\"}"); + + let record = MetadataStore::read_run_record(dir.path(), &run_id) + .unwrap() + .unwrap(); + assert_eq!(record.run_id, fixtures::RUN_5); + } + + #[test] + fn metadata_store_init_run_with_extra_files() { + let dir = tempfile::tempdir().unwrap(); + init_repo(dir.path()); + + let run_id = fixtures::RUN_6.to_string(); + let store = MetadataStore::new(dir.path(), &GitAuthor::default()); + store + .init_run(&run_id, &[("sandbox.json", b"{\"type\":\"local\"}")]) + .unwrap(); + + let data = branch_entry(dir.path(), &run_id, "sandbox.json"); + assert_eq!(data, b"{\"type\":\"local\"}"); + } + + #[test] + fn metadata_store_read_start_record_roundtrip() { + let dir = tempfile::tempdir().unwrap(); + init_repo(dir.path()); + + let run_id = fixtures::RUN_6.to_string(); + let store = MetadataStore::new(dir.path(), &GitAuthor::default()); + let start_record = StartRecord { + run_id: fixtures::RUN_6, + start_time: Utc.with_ymd_and_hms(2025, 1, 1, 0, 0, 0).single().unwrap(), + run_branch: Some("fabro/run/test".to_string()), + base_sha: None, + }; + let bytes = serde_json::to_vec_pretty(&start_record).unwrap(); + store.init_run(&run_id, &[("start.json", &bytes)]).unwrap(); + + let loaded = MetadataStore::read_start_record(dir.path(), &run_id) + .unwrap() + .unwrap(); + assert_eq!(loaded.run_id, fixtures::RUN_6); + assert_eq!(loaded.run_branch.as_deref(), Some("fabro/run/test")); + } +} diff --git a/lib/crates/fabro-git-storage/src/trailerlink.rs b/lib/crates/fabro-checkpoint/src/trailer.rs similarity index 100% rename from lib/crates/fabro-git-storage/src/trailerlink.rs rename to lib/crates/fabro-checkpoint/src/trailer.rs diff --git a/lib/crates/fabro-cli/Cargo.toml b/lib/crates/fabro-cli/Cargo.toml index b7f2c1a7f..dc375eede 100644 --- a/lib/crates/fabro-cli/Cargo.toml +++ b/lib/crates/fabro-cli/Cargo.toml @@ -32,7 +32,7 @@ fabro-mcp = { path = "../fabro-mcp" } fabro-proctitle = { path = "../fabro-proctitle" } fabro-retro = { path = "../fabro-retro" } fabro-sandbox = { path = "../fabro-sandbox", features = ["daytona"] } -fabro-git-storage = { path = "../fabro-git-storage" } +fabro-checkpoint = { path = "../fabro-checkpoint" } fabro-graphviz = { path = "../fabro-graphviz" } fabro-validate = { path = "../fabro-validate" } fabro-workflow = { path = "../fabro-workflow" } diff --git a/lib/crates/fabro-cli/src/commands/run/fork.rs b/lib/crates/fabro-cli/src/commands/run/fork.rs index 733013509..6de549381 100644 --- a/lib/crates/fabro-cli/src/commands/run/fork.rs +++ b/lib/crates/fabro-cli/src/commands/run/fork.rs @@ -1,7 +1,7 @@ use anyhow::Context; use anyhow::Result; +use fabro_checkpoint::git::Store; use fabro_config::FabroSettingsExt; -use fabro_git_storage::gitobj::Store; use fabro_util::terminal::Styles; use fabro_workflow::operations::{ ForkRunInput, RewindTarget, build_timeline_or_rebuild, find_run_id_by_prefix_or_store, fork, diff --git a/lib/crates/fabro-cli/src/commands/run/rewind.rs b/lib/crates/fabro-cli/src/commands/run/rewind.rs index 3b8c2b4b3..02506749a 100644 --- a/lib/crates/fabro-cli/src/commands/run/rewind.rs +++ b/lib/crates/fabro-cli/src/commands/run/rewind.rs @@ -2,8 +2,8 @@ use anyhow::Context; use anyhow::Result; use cli_table::format::{Border, Separator}; use cli_table::{Cell, CellStruct, Color, Style, Table}; +use fabro_checkpoint::git::Store; use fabro_config::FabroSettingsExt; -use fabro_git_storage::gitobj::Store; use fabro_util::terminal::Styles; use fabro_workflow::operations::{ RewindInput, RewindTarget, RunTimeline, build_timeline_or_rebuild, diff --git a/lib/crates/fabro-cli/tests/it/scenario/recovery.rs b/lib/crates/fabro-cli/tests/it/scenario/recovery.rs index 8489b765e..d70131406 100644 --- a/lib/crates/fabro-cli/tests/it/scenario/recovery.rs +++ b/lib/crates/fabro-cli/tests/it/scenario/recovery.rs @@ -1,8 +1,8 @@ use std::collections::BTreeSet; use std::path::Path; -use fabro_git_storage::branchstore::BranchStore; -use fabro_git_storage::gitobj::Store as GitStore; +use fabro_checkpoint::branch::BranchStore; +use fabro_checkpoint::git::Store as GitStore; use fabro_test::{fabro_snapshot, test_context}; use fabro_types::Checkpoint; use git2::{Repository, Signature}; diff --git a/lib/crates/fabro-git-storage/src/lib.rs b/lib/crates/fabro-git-storage/src/lib.rs deleted file mode 100644 index f2fc84016..000000000 --- a/lib/crates/fabro-git-storage/src/lib.rs +++ /dev/null @@ -1,7 +0,0 @@ -pub mod branchstore; -pub mod error; -pub mod gitobj; -pub mod snapshot; -pub mod trailerlink; - -pub use error::{Error, Result}; diff --git a/lib/crates/fabro-git-storage/src/snapshot.rs b/lib/crates/fabro-git-storage/src/snapshot.rs deleted file mode 100644 index 6bfe762b0..000000000 --- a/lib/crates/fabro-git-storage/src/snapshot.rs +++ /dev/null @@ -1,649 +0,0 @@ -use std::collections::BTreeMap; -use std::path::PathBuf; - -use git2::{Oid, Signature}; -use tracing::{debug, warn}; - -use crate::Result; -use crate::gitobj::{FileMode, Store, TreeEntries}; - -/// Options for writing a snapshot. -pub struct WriteOptions<'a> { - pub branch: String, - pub base_tree: Oid, - pub changes: FileChanges, - pub metadata: BTreeMap>, - pub metadata_from_disk: Option, - pub author: Signature<'a>, - pub message: String, - pub deduplicate: bool, -} - -/// File changes to apply from the working directory. -pub struct FileChanges { - pub modified: Vec, - pub new: Vec, - pub deleted: Vec, - pub repo_root: PathBuf, -} - -/// A directory on disk to walk and embed into the tree. -pub struct DiskDir { - pub disk_path: PathBuf, - pub tree_prefix: String, -} - -/// Result of a snapshot write. -pub struct WriteResult { - pub commit_oid: Oid, - pub tree_oid: Oid, - pub skipped: bool, -} - -/// Metadata about a snapshot commit. -#[derive(Debug)] -pub struct SnapshotInfo { - pub commit_oid: Oid, - pub tree_oid: Oid, - pub message: String, - pub time: git2::Time, -} - -/// Captures full repo-state on named branches. -pub struct SnapshotStore<'a> { - objects: &'a Store, -} - -impl<'a> SnapshotStore<'a> { - pub fn new(objects: &'a Store) -> Self { - Self { objects } - } - - /// Write a snapshot to a branch. - pub fn write(&self, opts: &WriteOptions<'_>) -> Result { - debug!(branch = %opts.branch, "Writing snapshot"); - // 1. Resolve existing branch tip or use base_tree - let (base_tree_oid, parent_oid) = match self.objects.resolve_ref(&opts.branch)? { - Some(commit_oid) => { - let commit = self.objects.repo().find_commit(commit_oid)?; - (commit.tree_id(), Some(commit_oid)) - } - None => (opts.base_tree, None), - }; - - // 2. Flatten base tree - let mut entries = self.objects.read_tree(base_tree_oid)?; - - // 3. Apply FileChanges - for path in &opts.changes.deleted { - entries.remove(path); - } - for path in opts.changes.modified.iter().chain(opts.changes.new.iter()) { - let full_path = opts.changes.repo_root.join(path); - match self.objects.write_blob_from_file(&full_path) { - Ok((oid, mode)) => { - entries.set(path.clone(), oid, mode); - } - Err(crate::Error::ReadFile { .. }) => { - // File disappeared since detection — treat as deleted - warn!(path = %path, "File disappeared since detection, treating as deleted"); - entries.remove(path); - } - Err(e) => return Err(e), - } - } - - // 4. Apply in-memory metadata - for (path, content) in &opts.metadata { - let oid = self.objects.write_blob(content)?; - entries.set(path.clone(), oid, FileMode::Blob); - } - - // 5. Walk metadata_from_disk - if let Some(disk_dir) = &opts.metadata_from_disk { - self.walk_disk_dir(&mut entries, disk_dir)?; - } - - // 6. Write tree - let new_tree_oid = self.objects.write_tree(&entries)?; - - // 7. Dedup check - if opts.deduplicate { - if let Some(parent) = parent_oid { - let parent_commit = self.objects.repo().find_commit(parent)?; - if parent_commit.tree_id() == new_tree_oid { - debug!(branch = %opts.branch, "Snapshot skipped (tree unchanged)"); - return Ok(WriteResult { - commit_oid: parent, - tree_oid: new_tree_oid, - skipped: true, - }); - } - } - } - - // 8. Create commit - let parents: Vec = parent_oid.into_iter().collect(); - let commit_oid = - self.objects - .write_commit(new_tree_oid, &parents, &opts.message, &opts.author)?; - - // 9. Update ref - self.objects.update_ref(&opts.branch, commit_oid)?; - debug!(branch = %opts.branch, commit = %commit_oid, "Snapshot written"); - - Ok(WriteResult { - commit_oid, - tree_oid: new_tree_oid, - skipped: false, - }) - } - - /// Tip commit of a snapshot branch. `None` if branch doesn't exist. - pub fn latest(&self, branch: &str) -> Result> { - let Some(commit_oid) = self.objects.resolve_ref(branch)? else { - return Ok(None); - }; - let commit = self.objects.repo().find_commit(commit_oid)?; - let tree_oid = commit.tree_id(); - let message = commit.message().unwrap_or("").to_string(); - let time = commit.author().when(); - Ok(Some(SnapshotInfo { - commit_oid, - tree_oid, - message, - time, - })) - } - - /// Read a single file from a snapshot commit's tree. - pub fn read_file(&self, commit_oid: Oid, path: &str) -> Result>> { - let commit = self.objects.repo().find_commit(commit_oid)?; - let tree = commit.tree()?; - match tree.get_path(std::path::Path::new(path)) { - Ok(entry) => { - let blob = self.objects.repo().find_blob(entry.id())?; - Ok(Some(blob.content().to_vec())) - } - Err(e) if e.code() == git2::ErrorCode::NotFound => Ok(None), - Err(e) => Err(e.into()), - } - } - - /// Walk commits on a snapshot branch, newest first. - pub fn list_commits(&self, branch: &str, limit: usize) -> Result> { - let Some(commit_oid) = self.objects.resolve_ref(branch)? else { - return Ok(vec![]); - }; - let mut revwalk = self.objects.repo().revwalk()?; - revwalk.set_sorting(git2::Sort::TIME | git2::Sort::TOPOLOGICAL)?; - revwalk.push(commit_oid)?; - - let mut results = Vec::new(); - for oid_result in revwalk.take(limit) { - let oid = oid_result?; - let commit = self.objects.repo().find_commit(oid)?; - results.push(SnapshotInfo { - commit_oid: oid, - tree_oid: commit.tree_id(), - message: commit.message().unwrap_or("").to_string(), - time: commit.author().when(), - }); - } - Ok(results) - } - - /// Check if a snapshot branch exists. - pub fn exists(&self, branch: &str) -> Result { - Ok(self.objects.resolve_ref(branch)?.is_some()) - } - - /// Delete a snapshot branch. - pub fn delete(&self, branch: &str) -> Result<()> { - debug!(branch = %branch, "Deleting snapshot branch"); - self.objects.delete_ref(branch) - } - - /// Rename a snapshot branch. - pub fn rename(&self, old: &str, new: &str) -> Result<()> { - debug!(old = %old, new = %new, "Renaming snapshot branch"); - let oid = self - .objects - .resolve_ref(old)? - .ok_or_else(|| crate::Error::BranchNotFound { - branch: old.to_string(), - })?; - self.objects.update_ref(new, oid)?; - self.objects.delete_ref(old)?; - Ok(()) - } - - /// List snapshot branches matching a prefix. - pub fn list(&self, prefix: &str) -> Result> { - let full_prefix = format!("refs/heads/{prefix}"); - let mut branches = Vec::new(); - for reference in self - .objects - .repo() - .references_glob(&format!("{full_prefix}*"))? - { - let reference = reference?; - if let Some(name) = reference.name() { - if let Some(branch) = name.strip_prefix("refs/heads/") { - branches.push(branch.to_string()); - } - } - } - branches.sort(); - Ok(branches) - } - - /// Walk a directory on disk and add files to tree entries. - fn walk_disk_dir(&self, entries: &mut TreeEntries, disk_dir: &DiskDir) -> Result<()> { - let walker = walkdir::WalkDir::new(&disk_dir.disk_path) - .follow_links(false) - .into_iter() - .filter_map(std::result::Result::ok); - - for entry in walker { - // Skip symlinks - if entry.path_is_symlink() { - continue; - } - // Skip directories - if entry.file_type().is_dir() { - continue; - } - - let relative = entry - .path() - .strip_prefix(&disk_dir.disk_path) - .unwrap_or(entry.path()); - let tree_path = if disk_dir.tree_prefix.is_empty() { - relative.to_string_lossy().to_string() - } else { - format!( - "{}/{}", - disk_dir.tree_prefix.trim_end_matches('/'), - relative.to_string_lossy() - ) - }; - - let (oid, mode) = self.objects.write_blob_from_file(entry.path())?; - entries.set(tree_path, oid, mode); - } - Ok(()) - } -} - -#[cfg(test)] -mod tests { - use super::*; - use git2::Repository; - - fn temp_repo() -> (tempfile::TempDir, Store) { - let dir = tempfile::TempDir::new().unwrap(); - let repo = Repository::init(dir.path()).unwrap(); - (dir, Store::new(repo)) - } - - fn test_sig() -> Signature<'static> { - Signature::now("Test", "test@example.com").unwrap() - } - - fn empty_changes() -> FileChanges { - FileChanges { - modified: vec![], - new: vec![], - deleted: vec![], - repo_root: PathBuf::from("/tmp"), - } - } - - // -- write creates branch + commit from base tree -- - - #[test] - fn write_creates_branch_from_base_tree() { - let (_dir, store) = temp_repo(); - let sig = test_sig(); - let snap = SnapshotStore::new(&store); - - // Create a base tree with one file - let blob_oid = store.write_blob(b"base content").unwrap(); - let mut base_entries = TreeEntries::new(); - base_entries.set("existing.txt", blob_oid, FileMode::Blob); - let base_tree = store.write_tree(&base_entries).unwrap(); - - let result = snap - .write(&WriteOptions { - branch: "snap/test".to_string(), - base_tree, - changes: empty_changes(), - metadata: BTreeMap::new(), - metadata_from_disk: None, - author: sig, - message: "snapshot 1".to_string(), - deduplicate: false, - }) - .unwrap(); - - assert!(!result.skipped); - assert!(store.resolve_ref("snap/test").unwrap().is_some()); - - // Verify the file is in the snapshot - let content = snap.read_file(result.commit_oid, "existing.txt").unwrap(); - assert_eq!(content.unwrap(), b"base content"); - } - - // -- write applies file changes -- - - #[test] - fn write_applies_file_changes() { - let (dir, store) = temp_repo(); - let sig = test_sig(); - let snap = SnapshotStore::new(&store); - - // Create files on disk in the repo root - let repo_root = dir.path().to_path_buf(); - std::fs::write(repo_root.join("new_file.txt"), b"new content").unwrap(); - std::fs::write(repo_root.join("modified.txt"), b"modified content").unwrap(); - - // Create base tree with a file to delete and one to modify - let old_blob = store.write_blob(b"old content").unwrap(); - let delete_blob = store.write_blob(b"delete me").unwrap(); - let mut base_entries = TreeEntries::new(); - base_entries.set("modified.txt", old_blob, FileMode::Blob); - base_entries.set("to_delete.txt", delete_blob, FileMode::Blob); - let base_tree = store.write_tree(&base_entries).unwrap(); - - let result = snap - .write(&WriteOptions { - branch: "snap/changes".to_string(), - base_tree, - changes: FileChanges { - modified: vec!["modified.txt".to_string()], - new: vec!["new_file.txt".to_string()], - deleted: vec!["to_delete.txt".to_string()], - repo_root, - }, - metadata: BTreeMap::new(), - metadata_from_disk: None, - author: sig, - message: "apply changes".to_string(), - deduplicate: false, - }) - .unwrap(); - - assert_eq!( - snap.read_file(result.commit_oid, "modified.txt") - .unwrap() - .unwrap(), - b"modified content" - ); - assert_eq!( - snap.read_file(result.commit_oid, "new_file.txt") - .unwrap() - .unwrap(), - b"new content" - ); - assert!( - snap.read_file(result.commit_oid, "to_delete.txt") - .unwrap() - .is_none() - ); - } - - // -- write embeds in-memory metadata -- - - #[test] - fn write_embeds_metadata() { - let (_dir, store) = temp_repo(); - let sig = test_sig(); - let snap = SnapshotStore::new(&store); - let base_tree = store.write_empty_tree().unwrap(); - - let mut metadata = BTreeMap::new(); - metadata.insert( - ".meta/transcript.jsonl".to_string(), - b"line1\nline2".to_vec(), - ); - - let result = snap - .write(&WriteOptions { - branch: "snap/meta".to_string(), - base_tree, - changes: empty_changes(), - metadata, - metadata_from_disk: None, - author: sig, - message: "with metadata".to_string(), - deduplicate: false, - }) - .unwrap(); - - let content = snap - .read_file(result.commit_oid, ".meta/transcript.jsonl") - .unwrap() - .unwrap(); - assert_eq!(content, b"line1\nline2"); - } - - // -- write dedup skips when tree unchanged -- - - #[test] - fn write_dedup_skips_unchanged() { - let (_dir, store) = temp_repo(); - let snap = SnapshotStore::new(&store); - let base_tree = store.write_empty_tree().unwrap(); - let sig = test_sig(); - - // First write - let result1 = snap - .write(&WriteOptions { - branch: "snap/dedup".to_string(), - base_tree, - changes: empty_changes(), - metadata: BTreeMap::new(), - metadata_from_disk: None, - author: sig.clone(), - message: "first".to_string(), - deduplicate: true, - }) - .unwrap(); - assert!(!result1.skipped); - - // Second write with same content — should be skipped - let sig2 = test_sig(); - let result2 = snap - .write(&WriteOptions { - branch: "snap/dedup".to_string(), - base_tree, - changes: empty_changes(), - metadata: BTreeMap::new(), - metadata_from_disk: None, - author: sig2, - message: "second".to_string(), - deduplicate: true, - }) - .unwrap(); - assert!(result2.skipped); - assert_eq!(result2.commit_oid, result1.commit_oid); - } - - // -- latest / read_file / list_commits -- - - #[test] - fn latest_and_list_commits() { - let (_dir, store) = temp_repo(); - let snap = SnapshotStore::new(&store); - let base_tree = store.write_empty_tree().unwrap(); - - // Write two snapshots - let sig1 = test_sig(); - snap.write(&WriteOptions { - branch: "snap/history".to_string(), - base_tree, - changes: empty_changes(), - metadata: BTreeMap::from([("a.txt".to_string(), b"a".to_vec())]), - metadata_from_disk: None, - author: sig1, - message: "first".to_string(), - deduplicate: false, - }) - .unwrap(); - - let sig2 = test_sig(); - snap.write(&WriteOptions { - branch: "snap/history".to_string(), - base_tree, - changes: empty_changes(), - metadata: BTreeMap::from([("b.txt".to_string(), b"b".to_vec())]), - metadata_from_disk: None, - author: sig2, - message: "second".to_string(), - deduplicate: false, - }) - .unwrap(); - - let latest = snap.latest("snap/history").unwrap().unwrap(); - assert_eq!(latest.message, "second"); - - let commits = snap.list_commits("snap/history", 10).unwrap(); - assert_eq!(commits.len(), 2); - assert_eq!(commits[0].message, "second"); - assert_eq!(commits[1].message, "first"); - } - - #[test] - fn latest_nonexistent() { - let (_dir, store) = temp_repo(); - let snap = SnapshotStore::new(&store); - assert!(snap.latest("nonexistent").unwrap().is_none()); - } - - // -- exists / delete / rename / list -- - - #[test] - fn exists_and_delete() { - let (_dir, store) = temp_repo(); - let snap = SnapshotStore::new(&store); - let base_tree = store.write_empty_tree().unwrap(); - let sig = test_sig(); - - snap.write(&WriteOptions { - branch: "snap/del".to_string(), - base_tree, - changes: empty_changes(), - metadata: BTreeMap::new(), - metadata_from_disk: None, - author: sig, - message: "create".to_string(), - deduplicate: false, - }) - .unwrap(); - - assert!(snap.exists("snap/del").unwrap()); - snap.delete("snap/del").unwrap(); - assert!(!snap.exists("snap/del").unwrap()); - } - - #[test] - fn rename_branch() { - let (_dir, store) = temp_repo(); - let snap = SnapshotStore::new(&store); - let base_tree = store.write_empty_tree().unwrap(); - let sig = test_sig(); - - snap.write(&WriteOptions { - branch: "snap/old".to_string(), - base_tree, - changes: empty_changes(), - metadata: BTreeMap::from([("file.txt".to_string(), b"data".to_vec())]), - metadata_from_disk: None, - author: sig, - message: "create".to_string(), - deduplicate: false, - }) - .unwrap(); - - snap.rename("snap/old", "snap/new").unwrap(); - assert!(!snap.exists("snap/old").unwrap()); - assert!(snap.exists("snap/new").unwrap()); - - // Verify data is preserved - let info = snap.latest("snap/new").unwrap().unwrap(); - let content = snap.read_file(info.commit_oid, "file.txt").unwrap(); - assert_eq!(content.unwrap(), b"data"); - } - - #[test] - fn list_branches() { - let (_dir, store) = temp_repo(); - let snap = SnapshotStore::new(&store); - let base_tree = store.write_empty_tree().unwrap(); - - // Create several branches - for name in &["snap/a", "snap/b", "other/c"] { - let sig = test_sig(); - snap.write(&WriteOptions { - branch: name.to_string(), - base_tree, - changes: empty_changes(), - metadata: BTreeMap::new(), - metadata_from_disk: None, - author: sig, - message: "create".to_string(), - deduplicate: false, - }) - .unwrap(); - } - - let snap_branches = snap.list("snap/").unwrap(); - assert_eq!(snap_branches, vec!["snap/a", "snap/b"]); - } - - // -- metadata_from_disk -- - - #[test] - fn write_metadata_from_disk() { - let (_dir, store) = temp_repo(); - let snap = SnapshotStore::new(&store); - let base_tree = store.write_empty_tree().unwrap(); - - // Create a temp directory with files - let meta_dir = tempfile::TempDir::new().unwrap(); - std::fs::write(meta_dir.path().join("info.json"), b"{}").unwrap(); - std::fs::create_dir(meta_dir.path().join("sub")).unwrap(); - std::fs::write(meta_dir.path().join("sub/data.txt"), b"nested").unwrap(); - - let sig = test_sig(); - let result = snap - .write(&WriteOptions { - branch: "snap/disk".to_string(), - base_tree, - changes: empty_changes(), - metadata: BTreeMap::new(), - metadata_from_disk: Some(DiskDir { - disk_path: meta_dir.path().to_path_buf(), - tree_prefix: ".meta".to_string(), - }), - author: sig, - message: "from disk".to_string(), - deduplicate: false, - }) - .unwrap(); - - assert_eq!( - snap.read_file(result.commit_oid, ".meta/info.json") - .unwrap() - .unwrap(), - b"{}" - ); - assert_eq!( - snap.read_file(result.commit_oid, ".meta/sub/data.txt") - .unwrap() - .unwrap(), - b"nested" - ); - } -} diff --git a/lib/crates/fabro-workflow/Cargo.toml b/lib/crates/fabro-workflow/Cargo.toml index f4df489a1..f70831eb9 100644 --- a/lib/crates/fabro-workflow/Cargo.toml +++ b/lib/crates/fabro-workflow/Cargo.toml @@ -30,7 +30,7 @@ fabro-mcp = { path = "../fabro-mcp" } fabro-github = { path = "../fabro-github" } fabro-interview = { path = "../fabro-interview" } fabro-util = { path = "../fabro-util" } -fabro-git-storage = { path = "../fabro-git-storage" } +fabro-checkpoint = { path = "../fabro-checkpoint" } fabro-llm = { path = "../fabro-llm" } fabro-model = { path = "../fabro-model" } fabro-retro = { path = "../fabro-retro" } diff --git a/lib/crates/fabro-workflow/src/error.rs b/lib/crates/fabro-workflow/src/error.rs index 5a8545d77..d9deeccd5 100644 --- a/lib/crates/fabro-workflow/src/error.rs +++ b/lib/crates/fabro-workflow/src/error.rs @@ -343,6 +343,19 @@ impl From for FabroError { } } +impl From for FabroError { + fn from(err: fabro_checkpoint::MetadataError) -> Self { + let message = err.to_string(); + match err { + fabro_checkpoint::MetadataError::Deserialize { + entity: "checkpoint", + .. + } => Self::Checkpoint(message), + _ => Self::engine(message), + } + } +} + pub type Result = std::result::Result; #[cfg(test)] diff --git a/lib/crates/fabro-workflow/src/git.rs b/lib/crates/fabro-workflow/src/git.rs index 6d0134746..53d593e58 100644 --- a/lib/crates/fabro-workflow/src/git.rs +++ b/lib/crates/fabro-workflow/src/git.rs @@ -1,77 +1,20 @@ -use std::fmt::Write; use std::path::Path; use std::process::Command; +use fabro_checkpoint::git::Store; use fabro_config::FabroSettings; -use fabro_config::server::GitAuthorSettings; -use fabro_git_storage::branchstore::BranchStore; -use fabro_git_storage::gitobj::Store; -use git2::{Repository, Signature}; use crate::error::{FabroError, Result}; -use crate::records::{Checkpoint, RunRecord, StartRecord}; use tokio::task::{JoinError, spawn_blocking}; use tokio::time::timeout; +pub use fabro_checkpoint::META_BRANCH_PREFIX; +pub use fabro_checkpoint::author::GitAuthor; +pub use fabro_checkpoint::metadata::MetadataStore; + /// Branch prefix for workflow run branches (e.g. `fabro/run/{run_id}`). pub const RUN_BRANCH_PREFIX: &str = "fabro/run/"; -/// Branch prefix for metadata branches (e.g. `fabro/meta/{run_id}`). -pub const META_BRANCH_PREFIX: &str = "fabro/meta/"; - -/// Resolved git author identity for checkpoint commits. -#[derive(Debug, Clone, PartialEq)] -pub struct GitAuthor { - pub name: String, - pub email: String, -} - -impl Default for GitAuthor { - fn default() -> Self { - Self { - name: "Fabro".into(), - email: "noreply@fabro.sh".into(), - } - } -} - -impl GitAuthor { - /// Create a `GitAuthor` from optional name/email, falling back to defaults. - pub fn from_options(name: Option, email: Option) -> Self { - let defaults = Self::default(); - Self { - name: name.unwrap_or(defaults.name), - email: email.unwrap_or(defaults.email), - } - } - - /// Returns true when this identity matches the default Fabro identity. - pub fn is_default(&self) -> bool { - let defaults = Self::default(); - self.name == defaults.name && self.email == defaults.email - } - - /// Append the Fabro footer (and Co-Authored-By when the author is not the - /// default identity) to a commit message. - pub fn append_footer(&self, message: &mut String) { - message.push_str("\n\u{2692}\u{fe0f} Generated with [Fabro](https://fabro.sh)\n"); - if !self.is_default() { - let defaults = Self::default(); - let _ = write!( - message, - "\nCo-Authored-By: {} <{}>\n", - defaults.name, defaults.email - ); - } - } -} - -impl From<&GitAuthorSettings> for GitAuthor { - fn from(value: &GitAuthorSettings) -> Self { - Self::from_options(value.name.clone(), value.email.clone()) - } -} - pub fn git_author_from_settings(settings: &FabroSettings) -> GitAuthor { settings .git_author() @@ -409,160 +352,11 @@ pub fn scan_node_files(run_dir: &Path) -> Vec<(String, Vec)> { result } -/// Git-native metadata storage for pipeline runs. -/// -/// Stores checkpoint data, run records, and metadata on an orphan branch -/// (`fabro/meta/{run_id}`) so that runs can be resumed from git alone. -pub struct MetadataStore { - repo_path: std::path::PathBuf, - author: GitAuthor, -} - -impl MetadataStore { - pub fn new(repo_path: impl Into, author: &GitAuthor) -> Self { - Self { - repo_path: repo_path.into(), - author: author.clone(), - } - } - - /// Returns the branch name for a run: `fabro/meta/{run_id}`. - pub fn branch_name(run_id: &str) -> String { - format!("{META_BRANCH_PREFIX}{run_id}") - } - - /// Format a commit message with the standard Fabro footer appended. - fn commit_message(&self, subject: &str) -> String { - let mut msg = format!("{subject}\n"); - self.author.append_footer(&mut msg); - msg - } - - 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(&self.author.name, &self.author.email) - .map_err(|e| git_error(format!("failed to create signature: {e}")))?; - Ok((store, sig)) - } - - /// Initialize a run's metadata branch with the given files. - /// - /// Callers pass all files (run.json, start.json, sandbox.json, etc.) - /// via the `files` slice. - pub fn init_run(&self, run_id: &str, 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}")))?; - let msg = self.commit_message("init run"); - bs.write_entries(files, &msg) - .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); - let msg = self.commit_message(message); - bs.write_entries(entries, &msg) - .map_err(|e| git_error(format!("write_entries failed: {e}")))?; - Ok(()) - } - - /// Write checkpoint data (and optional artifacts) to the metadata branch. - /// Returns the SHA of the new commit on the shadow branch. - pub fn write_checkpoint( - &self, - run_id: &str, - checkpoint_json: &[u8], - artifacts: &[(&str, &[u8])], - ) -> Result { - let (store, sig) = self.open_store()?; - let branch = Self::branch_name(run_id); - let bs = BranchStore::new(&store, &branch, &sig); - let mut entries: Vec<(&str, &[u8])> = vec![("checkpoint.json", checkpoint_json)]; - entries.extend_from_slice(artifacts); - let msg = self.commit_message("checkpoint"); - let oid = bs - .write_entries(&entries, &msg) - .map_err(|e| git_error(format!("write_entries failed: {e}")))?; - Ok(oid.to_string()) - } - - /// Read a single file from the metadata branch. Returns `None` if branch or path doesn't exist. - fn read_file(repo_path: &Path, run_id: &str, path: &str) -> Result>> { - let Ok(repo) = Repository::discover(repo_path) else { - return Ok(None); - }; - let store = Store::new(repo); - let sig = Signature::now("Fabro", "noreply@fabro.sh") - .map_err(|e| git_error(format!("failed to create signature: {e}")))?; - let branch = Self::branch_name(run_id); - let bs = BranchStore::new(&store, &branch, &sig); - bs.read_entry(path) - .map_err(|e| git_error(format!("read_entry failed: {e}"))) - } - - /// Read a checkpoint from the metadata branch. Returns `None` if branch or file doesn't exist. - pub fn read_checkpoint(repo_path: &Path, run_id: &str) -> Result> { - match Self::read_file(repo_path, run_id, "checkpoint.json")? { - Some(bytes) => { - let cp: Checkpoint = serde_json::from_slice(&bytes) - .map_err(|e| FabroError::Checkpoint(format!("deserialize failed: {e}")))?; - Ok(Some(cp)) - } - None => Ok(None), - } - } - - /// Read the run record from the metadata branch. Returns `None` if not found. - pub fn read_run_record(repo_path: &Path, run_id: &str) -> Result> { - match Self::read_file(repo_path, run_id, "run.json")? { - Some(bytes) => { - let record: RunRecord = serde_json::from_slice(&bytes) - .map_err(|e| git_error(format!("run record deserialize failed: {e}")))?; - Ok(Some(record)) - } - None => Ok(None), - } - } - - /// Read the start record from the metadata branch. Returns `None` if not found. - pub fn read_start_record(repo_path: &Path, run_id: &str) -> Result> { - match Self::read_file(repo_path, run_id, "start.json")? { - Some(bytes) => { - let record: StartRecord = serde_json::from_slice(&bytes) - .map_err(|e| git_error(format!("start record deserialize failed: {e}")))?; - Ok(Some(record)) - } - None => Ok(None), - } - } - - /// Read an artifact from the metadata branch. Returns `None` if not found. - pub fn read_artifact(repo_path: &Path, run_id: &str, key: &str) -> Result>> { - Self::read_file(repo_path, run_id, &format!("artifacts/{key}.json")) - } -} - #[cfg(test)] mod tests { use super::*; - use fabro_types::fixtures; use std::fs; - use crate::records::{CheckpointExt, RunRecordExt}; - /// Create a temporary git repo with an initial commit. fn init_repo(dir: &Path) { Command::new("git") @@ -647,142 +441,6 @@ mod tests { assert!(!wt_path.exists()); } - // --- 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(), &GitAuthor::default()); - let run_id = fixtures::RUN_1.to_string(); - let run_record = format!( - r#"{{"run_id":"{run_id}","created_at":"2025-01-01T00:00:00Z","settings":{{}},"graph":{{"name":"test","nodes":{{}},"edges":[],"attrs":{{}}}},"working_directory":"/tmp"}}"# - ); - store - .init_run(&run_id, &[("run.json", run_record.as_bytes())]) - .unwrap(); - - let read_record = MetadataStore::read_run_record(dir.path(), &run_id) - .unwrap() - .unwrap(); - assert_eq!(read_record.run_id, fixtures::RUN_1); - assert_eq!(read_record.workflow_name(), "test"); - } - - #[test] - fn metadata_store_write_and_read_checkpoint() { - let dir = tempfile::tempdir().unwrap(); - init_repo(dir.path()); - - let store = MetadataStore::new(dir.path(), &GitAuthor::default()); - store.init_run("RUN2", &[]).unwrap(); - - let ctx = crate::context::Context::new(); - ctx.set("goal", serde_json::json!("test")); - let cp = crate::records::Checkpoint::from_context( - &ctx, - "node_a", - vec!["start".to_string()], - std::collections::HashMap::new(), - std::collections::HashMap::new(), - Some("node_b".to_string()), - std::collections::HashMap::new(), - std::collections::HashMap::new(), - std::collections::HashMap::new(), - ); - 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(), &GitAuthor::default()); - store.init_run("RUN3", &[]).unwrap(); - - let ctx = crate::context::Context::new(); - let cp1 = crate::records::Checkpoint::from_context( - &ctx, - "node_a", - vec!["start".to_string()], - std::collections::HashMap::new(), - std::collections::HashMap::new(), - None, - std::collections::HashMap::new(), - std::collections::HashMap::new(), - std::collections::HashMap::new(), - ); - let cp1_json = serde_json::to_vec_pretty(&cp1).unwrap(); - store.write_checkpoint("RUN3", &cp1_json, &[]).unwrap(); - - let cp2 = crate::records::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()), - std::collections::HashMap::new(), - std::collections::HashMap::new(), - std::collections::HashMap::new(), - ); - 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(), &GitAuthor::default()); - store.init_run("RUN4", &[]).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); - } - #[test] fn scan_node_files_picks_up_allowlisted() { let dir = tempfile::tempdir().unwrap(); @@ -834,56 +492,6 @@ mod tests { 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()); - let run_id = fixtures::RUN_5.to_string(); - let run_record = format!( - r#"{{"run_id":"{run_id}","created_at":"2025-01-01T00:00:00Z","settings":{{}},"graph":{{"name":"test","nodes":{{}},"edges":[],"attrs":{{}}}},"working_directory":"/tmp"}}"# - ); - store - .init_run(&run_id, &[("run.json", run_record.as_bytes())]) - .unwrap(); - - store - .write_files( - &run_id, - &[("retro.json", b"{\"status\":\"ok\"}")], - "finalize", - ) - .unwrap(); - - let data = MetadataStore::read_file(dir.path(), &run_id, "retro.json") - .unwrap() - .unwrap(); - assert_eq!(data, b"{\"status\":\"ok\"}"); - - // Original files still present - let record = MetadataStore::read_run_record(dir.path(), &run_id) - .unwrap() - .unwrap(); - assert_eq!(record.run_id, fixtures::RUN_5); - } - - #[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", &[("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"); diff --git a/lib/crates/fabro-workflow/src/operations/fork.rs b/lib/crates/fabro-workflow/src/operations/fork.rs index c4c325caa..0e47f5c30 100644 --- a/lib/crates/fabro-workflow/src/operations/fork.rs +++ b/lib/crates/fabro-workflow/src/operations/fork.rs @@ -1,6 +1,6 @@ use anyhow::{Context, Result}; -use fabro_git_storage::branchstore::BranchStore; -use fabro_git_storage::gitobj::Store; +use fabro_checkpoint::branch::BranchStore; +use fabro_checkpoint::git::Store; use fabro_types::RunId; use git2::{Oid, Signature}; diff --git a/lib/crates/fabro-workflow/src/operations/rebuild_meta.rs b/lib/crates/fabro-workflow/src/operations/rebuild_meta.rs index 4e658f2b9..f45f3d839 100644 --- a/lib/crates/fabro-workflow/src/operations/rebuild_meta.rs +++ b/lib/crates/fabro-workflow/src/operations/rebuild_meta.rs @@ -3,8 +3,8 @@ use std::fmt::Write; use std::path::PathBuf; use anyhow::{Context, Result, bail}; -use fabro_git_storage::branchstore::BranchStore; -use fabro_git_storage::gitobj::Store as GitStore; +use fabro_checkpoint::branch::BranchStore; +use fabro_checkpoint::git::Store as GitStore; use fabro_store::{ ListRunsQuery, NodeVisitRef, RunStore as DurableRunStore, Store as DurableStore, }; diff --git a/lib/crates/fabro-workflow/src/operations/rewind.rs b/lib/crates/fabro-workflow/src/operations/rewind.rs index b12978cc5..5f8beb5fa 100644 --- a/lib/crates/fabro-workflow/src/operations/rewind.rs +++ b/lib/crates/fabro-workflow/src/operations/rewind.rs @@ -3,8 +3,8 @@ use std::fmt::Write; use std::str::FromStr; use anyhow::{Context, Result, bail}; -use fabro_git_storage::branchstore::{BranchStore, CommitInfo}; -use fabro_git_storage::gitobj::Store; +use fabro_checkpoint::branch::{BranchStore, CommitInfo}; +use fabro_checkpoint::git::Store; use fabro_types::RunId; use git2::{Oid, Repository, Signature}; diff --git a/lib/crates/fabro-workflow/src/operations/test_support.rs b/lib/crates/fabro-workflow/src/operations/test_support.rs index f686bb8b9..79319fa7b 100644 --- a/lib/crates/fabro-workflow/src/operations/test_support.rs +++ b/lib/crates/fabro-workflow/src/operations/test_support.rs @@ -1,6 +1,6 @@ use std::collections::HashMap; -use fabro_git_storage::gitobj::Store; +use fabro_checkpoint::git::Store; use git2::{Repository, Signature}; pub(super) fn temp_repo() -> (tempfile::TempDir, Store) { diff --git a/lib/crates/fabro-workflow/src/sandbox_git.rs b/lib/crates/fabro-workflow/src/sandbox_git.rs index 2c39b205e..05762a4ef 100644 --- a/lib/crates/fabro-workflow/src/sandbox_git.rs +++ b/lib/crates/fabro-workflow/src/sandbox_git.rs @@ -1,7 +1,8 @@ use std::path::Path; use fabro_agent::Sandbox; -use fabro_git_storage::trailerlink::{self, Trailer}; +use fabro_checkpoint::trailer as trailerlink; +use fabro_checkpoint::trailer::Trailer; use fabro_types::RunId; use crate::asset_snapshot;