From 97e2ff788203623b91b6e5753cfece77dd847378 Mon Sep 17 00:00:00 2001 From: Bryan Helmkamp Date: Thu, 26 Mar 2026 10:30:30 -0400 Subject: [PATCH] Add restore operation for resumed runs --- lib/crates/fabro-cli/src/commands/resume.rs | 109 ++++++++++-- .../fabro-workflows/src/operations/create.rs | 18 +- .../fabro-workflows/src/operations/mod.rs | 6 +- .../fabro-workflows/src/operations/restore.rs | 163 ++++++++++++++++++ 4 files changed, 261 insertions(+), 35 deletions(-) create mode 100644 lib/crates/fabro-workflows/src/operations/restore.rs diff --git a/lib/crates/fabro-cli/src/commands/resume.rs b/lib/crates/fabro-cli/src/commands/resume.rs index 8a3ac03af..252ef9a95 100644 --- a/lib/crates/fabro-cli/src/commands/resume.rs +++ b/lib/crates/fabro-cli/src/commands/resume.rs @@ -14,7 +14,7 @@ use fabro_util::terminal::Styles; use fabro_workflows::event::{EventEmitter, RunNoticeLevel}; use fabro_workflows::handler::llm::{AgentApiBackend, AgentCliBackend, BackendRouter}; use fabro_workflows::operations::{ - create_from_graph, start, RunCreateOptions, StartFinalizeOptions, StartOptions, + restore, start, RestoreOptions, RunCreateOptions, StartFinalizeOptions, StartOptions, StartPullRequestConfig, StartRetroOptions, }; use fabro_workflows::outcome::StageStatus; @@ -154,6 +154,29 @@ fn preferred_resume_repo_path( .unwrap_or_else(|| original_cwd.to_path_buf()) } +fn restore_persisted_for_resume( + rec: &RunRecord, + config: FabroConfig, + run_dir: PathBuf, + run_id: &str, + labels: HashMap, + base_branch: Option, + resume_repo_path: &std::path::Path, +) -> Result { + let mut run_record = rec.clone(); + run_record.run_id = run_id.to_string(); + run_record.config = config; + run_record.labels = labels; + run_record.base_branch = base_branch; + run_record.working_directory = resume_repo_path.to_path_buf(); + run_record.host_repo_path = Some(resume_repo_path.to_string_lossy().to_string()); + + restore(RestoreOptions { + run_dir, + run_record, + }) +} + /// Resume an interrupted workflow run. /// /// # Errors @@ -681,20 +704,14 @@ async fn prepare_from_branch( sandbox_provider, }, ); - let persisted = create_from_graph( - rec.graph.clone(), - RunCreateOptions { - config, - run_dir: Some(run_dir.clone()), - run_id: Some(run_id.clone()), - workflow_slug: rec.workflow_slug.clone(), - labels: parse_labels(&args.label), - base_branch: detected_base_branch.clone(), - working_directory: Some(resume_repo_path.clone()), - host_repo_path: Some(resume_repo_path.to_string_lossy().to_string()), - goal_override: None, - base_dir: None, - }, + let persisted = restore_persisted_for_resume( + rec, + config, + run_dir.clone(), + &run_id, + parse_labels(&args.label), + detected_base_branch.clone(), + &resume_repo_path, )?; let graph_source = persisted.source().to_string(); let run_cfg = Some(persisted.run_record().config.clone()); @@ -1450,7 +1467,7 @@ fn run_dir_name_matches_mode(name: &str, run_id_suffix: &str, dry_run: bool) -> #[cfg(test)] mod tests { use super::*; - use chrono::Utc; + use chrono::{TimeZone, Utc}; use fabro_workflows::run_status::{RunStatus, RunStatusRecord, StatusReason}; fn sample_run_record() -> RunRecord { @@ -1491,6 +1508,66 @@ mod tests { assert_eq!(selected, cwd.path()); } + #[test] + fn restore_persisted_for_resume_materializes_metadata_record_without_existing_run_dir() { + let dir = tempfile::tempdir().unwrap(); + let repo_dir = dir.path().join("repo"); + let run_dir = dir.path().join("runs").join("run-1"); + std::fs::create_dir_all(&repo_dir).unwrap(); + + let mut record = sample_run_record(); + let created_at = Utc.with_ymd_and_hms(2025, 1, 1, 0, 0, 0).single().unwrap(); + record.created_at = created_at; + record.config = fabro_config::config::FabroConfig { + llm: Some(fabro_config::run::LlmConfig { + model: Some("sonnet".to_string()), + provider: None, + fallbacks: None, + }), + dry_run: Some(true), + ..Default::default() + }; + record.labels = HashMap::from([("old".to_string(), "value".to_string())]); + record.host_repo_path = Some("/tmp/old-repo".to_string()); + + assert!(!run_dir.exists()); + + let persisted = restore_persisted_for_resume( + &record, + record.config.clone(), + run_dir.clone(), + &record.run_id, + HashMap::from([("env".to_string(), "test".to_string())]), + Some("develop".to_string()), + &repo_dir, + ) + .unwrap(); + let loaded = Persisted::load(&run_dir).unwrap(); + + assert!(run_dir.exists()); + assert_eq!(persisted.run_record().created_at, created_at); + assert_eq!(loaded.run_record().created_at, created_at); + assert_eq!(loaded.run_record().working_directory, repo_dir); + assert_eq!( + loaded.run_record().host_repo_path.as_deref(), + Some(repo_dir.to_string_lossy().as_ref()) + ); + assert_eq!(loaded.run_record().base_branch.as_deref(), Some("develop")); + assert_eq!( + loaded.run_record().labels.get("env").map(String::as_str), + Some("test") + ); + assert_eq!( + loaded + .run_record() + .config + .llm + .as_ref() + .and_then(|llm| llm.model.as_deref()), + Some("claude-sonnet-4-6") + ); + } + #[test] fn find_existing_run_dir_in_respects_dry_run_mode() { let runs = tempfile::tempdir().unwrap(); diff --git a/lib/crates/fabro-workflows/src/operations/create.rs b/lib/crates/fabro-workflows/src/operations/create.rs index 5b8ff959e..e228681b2 100644 --- a/lib/crates/fabro-workflows/src/operations/create.rs +++ b/lib/crates/fabro-workflows/src/operations/create.rs @@ -92,22 +92,6 @@ pub fn create_from_file( create(&source, options) } -/// Build a persisted workflow from an already-materialized graph. -/// -/// This is used by detached/resume CLI paths that load a graph from `RunRecord` -/// instead of re-parsing DOT source. -#[doc(hidden)] -pub fn create_from_graph( - mut graph: Graph, - options: RunCreateOptions, -) -> Result { - if let Some(goal_override) = options.goal_override.as_deref() { - apply_goal_override(&mut graph, Some(goal_override)); - } - let validated = Validated::new(graph, String::new(), vec![]); - persist_validated(validated, options) -} - fn preprocess_and_validate( dot_source: &str, base_dir: Option, @@ -193,7 +177,7 @@ fn persist_validated( ) } -fn finalize_config(config: &mut FabroConfig, graph: &Graph) { +pub(crate) fn finalize_config(config: &mut FabroConfig, graph: &Graph) { let llm_config = config.llm.as_ref(); let configured_model = llm_config.and_then(|l| l.model.as_deref()); let configured_provider = llm_config.and_then(|l| l.provider.as_deref()); diff --git a/lib/crates/fabro-workflows/src/operations/mod.rs b/lib/crates/fabro-workflows/src/operations/mod.rs index 1a513086f..114030cf7 100644 --- a/lib/crates/fabro-workflows/src/operations/mod.rs +++ b/lib/crates/fabro-workflows/src/operations/mod.rs @@ -1,13 +1,15 @@ mod create; mod fork; +mod restore; mod rewind; mod start; pub use create::{ - create, create_from_file, create_from_graph, default_run_dir, validate, validate_from_file, - RunCreateOptions, ValidateOptions, + create, create_from_file, default_run_dir, validate, validate_from_file, RunCreateOptions, + ValidateOptions, }; pub use fork::fork; +pub use restore::{restore, RestoreOptions}; pub use rewind::{ build_timeline, find_run_id_by_prefix, load_parallel_map, parse_target, resolve_target, rewind, TimelineEntry, diff --git a/lib/crates/fabro-workflows/src/operations/restore.rs b/lib/crates/fabro-workflows/src/operations/restore.rs new file mode 100644 index 000000000..3f415ce85 --- /dev/null +++ b/lib/crates/fabro-workflows/src/operations/restore.rs @@ -0,0 +1,163 @@ +use std::path::PathBuf; + +use crate::error::FabroError; +use crate::pipeline::types::PersistOptions; +use crate::pipeline::{self, Persisted, Validated}; +use crate::records::RunRecord; + +use super::create::finalize_config; + +pub struct RestoreOptions { + pub run_dir: PathBuf, + pub run_record: RunRecord, +} + +/// Materialize an existing run record to local disk. +/// +/// Unlike `create()`, this skips parsing, transforms, and validation because +/// the caller already has the resolved graph from the original run. +pub fn restore(options: RestoreOptions) -> Result { + let mut run_record = options.run_record; + finalize_config(&mut run_record.config, &run_record.graph); + let graph = run_record.graph.clone(); + let validated = Validated::new(graph, String::new(), vec![]); + + pipeline::persist( + validated, + PersistOptions { + run_dir: options.run_dir, + run_record, + }, + ) +} + +#[cfg(test)] +mod tests { + use super::*; + use std::collections::HashMap; + + use chrono::{TimeZone, Utc}; + use fabro_config::config::FabroConfig; + use fabro_graphviz::graph::{AttrValue, Graph}; + + fn sample_graph() -> Graph { + let mut graph = Graph::new("restore-test"); + graph.attrs.insert( + "goal".to_string(), + AttrValue::String("Ship feature".to_string()), + ); + graph + } + + fn sample_record() -> RunRecord { + RunRecord { + run_id: "restore-run-123".to_string(), + created_at: Utc.with_ymd_and_hms(2025, 1, 2, 3, 4, 5).single().unwrap(), + config: FabroConfig { + llm: Some(fabro_config::run::LlmConfig { + model: Some("sonnet".to_string()), + provider: None, + fallbacks: None, + }), + pull_request: Some(fabro_config::run::PullRequestConfig { + enabled: false, + ..Default::default() + }), + dry_run: Some(true), + ..Default::default() + }, + graph: sample_graph(), + workflow_slug: Some("restore-slug".to_string()), + working_directory: PathBuf::from("/tmp/original-project"), + host_repo_path: Some("/tmp/original-project".to_string()), + base_branch: Some("main".to_string()), + labels: HashMap::from([("env".to_string(), "test".to_string())]), + } + } + + #[test] + fn restore_roundtrips_and_normalizes_config() { + let temp = tempfile::tempdir().unwrap(); + let run_dir = temp.path().join("run"); + + let persisted = restore(RestoreOptions { + run_dir: run_dir.clone(), + run_record: sample_record(), + }) + .unwrap(); + let loaded = Persisted::load(&run_dir).unwrap(); + + assert_eq!(persisted.run_record().run_id, "restore-run-123"); + assert_eq!( + persisted + .run_record() + .config + .llm + .as_ref() + .and_then(|llm| llm.model.as_deref()), + Some("claude-sonnet-4-6") + ); + assert_eq!( + persisted + .run_record() + .config + .llm + .as_ref() + .and_then(|llm| llm.provider.as_deref()), + Some("anthropic") + ); + assert_eq!( + persisted.run_record().config.goal.as_deref(), + Some("Ship feature") + ); + assert!(persisted.run_record().config.pull_request.is_none()); + assert_eq!( + serde_json::to_value(loaded.run_record()).unwrap(), + serde_json::to_value(persisted.run_record()).unwrap() + ); + } + + #[test] + fn restore_preserves_run_record_fields() { + let temp = tempfile::tempdir().unwrap(); + let run_dir = temp.path().join("run"); + let record = sample_record(); + + restore(RestoreOptions { + run_dir: run_dir.clone(), + run_record: record.clone(), + }) + .unwrap(); + let loaded = Persisted::load(&run_dir).unwrap(); + + assert_eq!(loaded.run_record().run_id, record.run_id); + assert_eq!(loaded.run_record().workflow_slug, record.workflow_slug); + assert_eq!(loaded.run_record().labels, record.labels); + assert_eq!( + loaded.run_record().working_directory, + record.working_directory + ); + assert_eq!(loaded.run_record().host_repo_path, record.host_repo_path); + assert_eq!(loaded.run_record().base_branch, record.base_branch); + } + + #[test] + fn restore_preserves_created_at_and_run_lookup_uses_it_without_start_record() { + let temp = tempfile::tempdir().unwrap(); + let runs_base = temp.path().join("runs"); + let run_dir = runs_base.join("restore-run-123"); + let record = sample_record(); + + restore(RestoreOptions { + run_dir, + run_record: record.clone(), + }) + .unwrap(); + + let runs = crate::run_lookup::scan_runs(&runs_base).unwrap(); + assert_eq!(runs.len(), 1); + assert_eq!(runs[0].run_id, record.run_id); + assert_eq!(runs[0].start_time, record.created_at.to_rfc3339()); + assert_eq!(runs[0].start_time_dt, Some(record.created_at)); + } +}