From 5a8eff0d63b7bb8843cf9e2e079882b4ba8709f2 Mon Sep 17 00:00:00 2001 From: Bryan Helmkamp Date: Wed, 25 Mar 2026 14:06:19 -0400 Subject: [PATCH] Fix persisted resume boundary gaps --- lib/crates/fabro-cli/src/commands/resume.rs | 275 ++++++++++-------- .../fabro-workflows/src/operations/start.rs | 36 +++ .../src/pipeline/initialize.rs | 61 ++++ 3 files changed, 250 insertions(+), 122 deletions(-) diff --git a/lib/crates/fabro-cli/src/commands/resume.rs b/lib/crates/fabro-cli/src/commands/resume.rs index 3042d1b50..ec98ca607 100644 --- a/lib/crates/fabro-cli/src/commands/resume.rs +++ b/lib/crates/fabro-cli/src/commands/resume.rs @@ -463,18 +463,17 @@ async fn prepare_from_checkpoint( cancel_token: None, dry_run: args.dry_run, run_id: run_id.clone(), - host_repo_path: None, + host_repo_path: persisted + .run_record() + .host_repo_path + .as_deref() + .map(PathBuf::from), git: None, - labels: args - .label - .iter() - .filter_map(|s| s.split_once('=')) - .map(|(k, v)| (k.to_string(), v.to_string())) - .collect(), + labels: persisted.run_record().labels.clone(), github_app: github_app.clone(), git_author, - base_branch: None, - workflow_slug, + base_branch: persisted.run_record().base_branch.clone(), + workflow_slug: persisted.run_record().workflow_slug.clone(), }; let devcontainer_env = devcontainer_config @@ -561,50 +560,6 @@ async fn prepare_from_branch( .or_else(|| repo_info.as_ref().and_then(|(_, branch)| branch.clone())); let base_sha = start_record.as_ref().and_then(|s| s.base_sha.clone()); - let (mut validated, mut graph_source, mut run_cfg, mut sandbox_provider, mut workflow_slug) = - if let Some(ref workflow_path) = args.workflow { - let prepared = prepare_workflow_with_project_config( - &resume_as_run_args(args, workflow_path.clone()), - run_defaults.clone(), - styles, - true, - false, - )?; - ( - prepared.validated, - prepared.raw_source, - prepared.run_cfg, - prepared.sandbox_provider, - prepared.workflow_slug, - ) - } else if let Some(ref rec) = record { - // Use the fully transformed graph from the RunRecord - let validated = create_from_graph(rec.graph.clone(), String::new()); - let source = String::new(); // no DOT source needed — graph is from RunRecord - let run_cfg = Some(rec.config.clone()); - let sandbox_provider = if args.dry_run { - SandboxProvider::Local - } else { - let sp = rec - .config - .sandbox - .as_ref() - .and_then(|s| s.provider.as_deref()) - .and_then(|s| s.parse::().ok()) - .unwrap_or_default(); - args.sandbox.map(Into::into).unwrap_or(sp) - }; - ( - validated, - source, - run_cfg, - sandbox_provider, - rec.workflow_slug.clone(), - ) - } else { - bail!("no run.json found on metadata branch for run {run_id}"); - }; - // Set up logs directory — reuse existing run dir for this run_id to avoid // "ambiguous prefix" errors when the resume happens on a different day. // Skip reuse when dry-running to avoid corrupting real run data. @@ -617,28 +572,145 @@ async fn prepare_from_branch( }; tokio::fs::create_dir_all(&run_dir).await?; let run_dir = tokio::fs::canonicalize(&run_dir).await.unwrap_or(run_dir); - if args.workflow.is_none() { - if let Ok(loaded) = Persisted::load(&run_dir) { - graph_source = loaded.source().to_string(); - run_cfg = Some(loaded.run_record().config.clone()); - workflow_slug = loaded.run_record().workflow_slug.clone(); - validated = create_from_graph(loaded.graph().clone(), loaded.source().to_string()); - if !args.dry_run { - if let Some(provider) = loaded + let cli_flags = super::create::CliFlags { + dry_run: args.dry_run, + auto_approve: args.auto_approve, + no_retro: args.no_retro, + verbose: args.verbose, + preserve_sandbox: args.preserve_sandbox, + }; + + let (persisted, _run_cfg, mut sandbox_provider, graph_source) = + if let Some(ref workflow_path) = args.workflow { + let prepared = prepare_workflow_with_project_config( + &resume_as_run_args(args, workflow_path.clone()), + run_defaults.clone(), + styles, + true, + false, + )?; + let graph = prepared.validated.graph().clone(); + let (model_str, provider_str) = resolve_model_provider( + args.model.as_deref(), + args.provider.as_deref(), + prepared.run_cfg.as_ref(), + run_defaults, + &graph, + ); + let normalized = super::create::normalize_config( + prepared.run_cfg.as_ref(), + run_defaults, + &model_str, + provider_str.as_deref(), + prepared.sandbox_provider, + &graph, + cli_flags, + ); + let persisted = fabro_workflows::pipeline::persist( + prepared.validated, + PersistOptions { + run_dir: run_dir.clone(), + run_record: fabro_workflows::records::RunRecord { + run_id: run_id.clone(), + created_at: chrono::Utc::now(), + config: normalized, + graph, + workflow_slug: prepared.workflow_slug, + working_directory: resume_repo_path.clone(), + host_repo_path: Some(resume_repo_path.to_string_lossy().to_string()), + base_branch: detected_base_branch.clone(), + labels: args + .label + .iter() + .filter_map(|s| s.split_once('=')) + .map(|(k, v)| (k.to_string(), v.to_string())) + .collect(), + }, + }, + )?; + let graph_source = persisted.source().to_string(); + let run_cfg = Some(persisted.run_record().config.clone()); + let sandbox_provider = prepared.sandbox_provider; + (persisted, run_cfg, sandbox_provider, graph_source) + } else if let Ok(loaded) = Persisted::load(&run_dir) { + let sandbox_provider = if args.dry_run { + SandboxProvider::Local + } else { + let persisted_provider = loaded .run_record() .config .sandbox .as_ref() .and_then(|s| s.provider.as_deref()) .and_then(|s| s.parse::().ok()) - { - sandbox_provider = args.sandbox.map(Into::into).unwrap_or(provider); - } - } - } - } + .unwrap_or_default(); + args.sandbox.map(Into::into).unwrap_or(persisted_provider) + }; + let graph_source = loaded.source().to_string(); + let run_cfg = Some(loaded.run_record().config.clone()); + (loaded, run_cfg, sandbox_provider, graph_source) + } else if let Some(ref rec) = record { + let sandbox_provider = if args.dry_run { + SandboxProvider::Local + } else { + let sp = rec + .config + .sandbox + .as_ref() + .and_then(|s| s.provider.as_deref()) + .and_then(|s| s.parse::().ok()) + .unwrap_or_default(); + args.sandbox.map(Into::into).unwrap_or(sp) + }; + let validated = create_from_graph(rec.graph.clone(), String::new()); + let graph = validated.graph().clone(); + let run_cfg = Some(rec.config.clone()); + let (model_str, provider_str) = resolve_model_provider( + args.model.as_deref(), + args.provider.as_deref(), + run_cfg.as_ref(), + run_defaults, + &graph, + ); + let normalized = super::create::normalize_config( + run_cfg.as_ref(), + run_defaults, + &model_str, + provider_str.as_deref(), + sandbox_provider, + &graph, + cli_flags, + ); + let persisted = fabro_workflows::pipeline::persist( + validated, + PersistOptions { + run_dir: run_dir.clone(), + run_record: fabro_workflows::records::RunRecord { + run_id: run_id.clone(), + created_at: chrono::Utc::now(), + config: normalized, + graph, + workflow_slug: rec.workflow_slug.clone(), + working_directory: resume_repo_path.clone(), + host_repo_path: Some(resume_repo_path.to_string_lossy().to_string()), + base_branch: detected_base_branch.clone(), + labels: args + .label + .iter() + .filter_map(|s| s.split_once('=')) + .map(|(k, v)| (k.to_string(), v.to_string())) + .collect(), + }, + }, + )?; + let graph_source = persisted.source().to_string(); + let run_cfg = Some(persisted.run_record().config.clone()); + (persisted, run_cfg, sandbox_provider, graph_source) + } else { + bail!("no run.json found on metadata branch for run {run_id}"); + }; - let graph = validated.graph().clone(); + let graph = persisted.graph().clone(); eprintln!( "{} {} from branch {} ({})", @@ -648,51 +720,6 @@ async fn prepare_from_branch( run_id, ); - let (model_str, provider_str) = resolve_model_provider( - args.model.as_deref(), - args.provider.as_deref(), - run_cfg.as_ref(), - run_defaults, - &graph, - ); - let cli_flags = super::create::CliFlags { - dry_run: args.dry_run, - auto_approve: args.auto_approve, - no_retro: args.no_retro, - verbose: args.verbose, - preserve_sandbox: args.preserve_sandbox, - }; - let normalized = super::create::normalize_config( - run_cfg.as_ref(), - run_defaults, - &model_str, - provider_str.as_deref(), - sandbox_provider, - &graph, - cli_flags, - ); - let persisted = fabro_workflows::pipeline::persist( - validated, - PersistOptions { - run_dir: run_dir.clone(), - run_record: fabro_workflows::records::RunRecord { - run_id: run_id.clone(), - created_at: chrono::Utc::now(), - config: normalized, - graph: graph.clone(), - workflow_slug: workflow_slug.clone(), - working_directory: resume_repo_path.clone(), - host_repo_path: Some(resume_repo_path.to_string_lossy().to_string()), - base_branch: detected_base_branch.clone(), - labels: args - .label - .iter() - .filter_map(|s| s.split_once('=')) - .map(|(k, v)| (k.to_string(), v.to_string())) - .collect(), - }, - }, - )?; let run_cfg: Option = Some(persisted.run_record().config.clone()); let settings_config = persisted.run_record().config.clone(); @@ -926,22 +953,26 @@ async fn prepare_from_branch( cancel_token: None, dry_run: args.dry_run, run_id: run_id.clone(), - host_repo_path: Some(resume_repo_path.clone()), + host_repo_path: persisted + .run_record() + .host_repo_path + .as_deref() + .map(PathBuf::from) + .or_else(|| Some(resume_repo_path.clone())), git: Some(GitCheckpointSettings { base_sha, run_branch: Some(run_branch), meta_branch: Some(fabro_workflows::git::MetadataStore::branch_name(&run_id)), }), - labels: args - .label - .iter() - .filter_map(|s| s.split_once('=')) - .map(|(k, v)| (k.to_string(), v.to_string())) - .collect(), + labels: persisted.run_record().labels.clone(), github_app: github_app.clone(), git_author, - base_branch: detected_base_branch, - workflow_slug, + base_branch: persisted + .run_record() + .base_branch + .clone() + .or(detected_base_branch), + workflow_slug: persisted.run_record().workflow_slug.clone(), }; let devcontainer_env = devcontainer_config diff --git a/lib/crates/fabro-workflows/src/operations/start.rs b/lib/crates/fabro-workflows/src/operations/start.rs index 0956bb11e..59826f348 100644 --- a/lib/crates/fabro-workflows/src/operations/start.rs +++ b/lib/crates/fabro-workflows/src/operations/start.rs @@ -509,4 +509,40 @@ mod tests { assert_eq!(started.finalized.conclusion.status, StageStatus::Success); assert!(started.retro.is_none()); } + + #[tokio::test] + async fn start_runs_loaded_persisted_workflow() { + let temp = tempfile::tempdir().unwrap(); + let run_dir = temp.path().join("run"); + let wrong_run_dir = temp.path().join("wrong-run-dir"); + let emitter = Arc::new(EventEmitter::new()); + let registry = Arc::new(test_registry()); + let sandbox: Arc = + Arc::new(LocalSandbox::new(std::env::current_dir().unwrap())); + + persisted_workflow(MINIMAL_DOT, &run_dir); + let loaded = Persisted::load(&run_dir).unwrap(); + + let started = start( + loaded, + test_start_options( + &wrong_run_dir, + sandbox, + emitter, + registry, + LifecycleConfig { + setup_commands: vec![], + setup_command_timeout_ms: 1_000, + devcontainer_phases: vec![], + }, + true, + ), + ) + .await + .unwrap(); + + assert_eq!(started.finalized.conclusion.status, StageStatus::Success); + assert!(run_dir.join("conclusion.json").exists()); + assert!(!wrong_run_dir.join("conclusion.json").exists()); + } } diff --git a/lib/crates/fabro-workflows/src/pipeline/initialize.rs b/lib/crates/fabro-workflows/src/pipeline/initialize.rs index 27fef880a..79a8cf1f9 100644 --- a/lib/crates/fabro-workflows/src/pipeline/initialize.rs +++ b/lib/crates/fabro-workflows/src/pipeline/initialize.rs @@ -197,6 +197,7 @@ mod tests { use super::*; use crate::handler::default_registry; + use crate::pipeline::PersistOptions; use crate::records::RunRecord; use crate::run_settings::RunSettings; @@ -344,4 +345,64 @@ mod tests { assert!(initialized.source.is_empty()); } + + #[tokio::test] + async fn initialize_uses_loaded_persisted_run_state() { + let temp = tempfile::tempdir().unwrap(); + let run_dir = temp.path().join("run"); + let wrong_run_dir = temp.path().join("wrong-run-dir"); + let (graph, source) = simple_graph(); + + crate::pipeline::persist( + crate::pipeline::Validated::new(graph.clone(), source.clone(), vec![]), + PersistOptions { + run_dir: run_dir.clone(), + run_record: RunRecord { + run_id: "run-test".to_string(), + created_at: Utc::now(), + config: FabroConfig::default(), + graph, + workflow_slug: Some("test".to_string()), + working_directory: std::env::current_dir().unwrap(), + host_repo_path: Some(std::env::current_dir().unwrap().display().to_string()), + base_branch: Some("main".to_string()), + labels: HashMap::new(), + }, + }, + ) + .unwrap(); + + let loaded = Persisted::load(&run_dir).unwrap(); + let emitter = Arc::new(crate::event::EventEmitter::new()); + let sandbox = Arc::new(fabro_agent::LocalSandbox::new( + std::env::current_dir().unwrap(), + )); + let registry = Arc::new(default_registry(Arc::new(AutoApproveInterviewer), || None)); + + let initialized = initialize( + loaded, + InitOptions { + run_id: "run-test".to_string(), + dry_run: false, + emitter, + sandbox, + registry, + lifecycle: crate::run_settings::LifecycleConfig { + setup_commands: vec![], + setup_command_timeout_ms: 1_000, + devcontainer_phases: vec![], + }, + run_settings: test_settings(&wrong_run_dir), + hooks: fabro_hooks::HookConfig { hooks: vec![] }, + sandbox_env: HashMap::new(), + checkpoint: None, + seed_context: None, + }, + ) + .await + .unwrap(); + + assert_eq!(initialized.settings.run_dir, run_dir); + assert_eq!(initialized.source, source); + } }