diff --git a/lib/crates/fabro-types/src/pull_request.rs b/lib/crates/fabro-types/src/pull_request.rs index e1139ab77..c1170abdc 100644 --- a/lib/crates/fabro-types/src/pull_request.rs +++ b/lib/crates/fabro-types/src/pull_request.rs @@ -1,5 +1,3 @@ -use std::path::Path; - use serde::{Deserialize, Serialize}; /// Record of a pull request created for a workflow run. @@ -13,11 +11,3 @@ pub struct PullRequestRecord { pub head_branch: String, pub title: String, } - -impl PullRequestRecord { - pub fn save(&self, path: &Path) -> Result<(), String> { - let json = serde_json::to_string_pretty(self) - .map_err(|e| format!("Failed to serialize pull_request.json: {e}"))?; - std::fs::write(path, json).map_err(|e| format!("Failed to write pull_request.json: {e}")) - } -} diff --git a/lib/crates/fabro-workflow/src/asset_snapshot.rs b/lib/crates/fabro-workflow/src/asset_snapshot.rs index 5ab1a46aa..ac1ee4380 100644 --- a/lib/crates/fabro-workflow/src/asset_snapshot.rs +++ b/lib/crates/fabro-workflow/src/asset_snapshot.rs @@ -275,10 +275,13 @@ fn compute_asset_info( }) } -fn write_asset_manifest(stage_dir: &Path, summary: &AssetCollectionSummary) -> Result<(), String> { +fn write_asset_manifest( + asset_capture_dir: &Path, + summary: &AssetCollectionSummary, +) -> Result<(), String> { let json = serde_json::to_string_pretty(summary) .map_err(|e| format!("failed to serialize manifest: {e}"))?; - let manifest_path = stage_dir.join("manifest.json"); + let manifest_path = asset_capture_dir.join("manifest.json"); if let Some(parent) = manifest_path.parent() { std::fs::create_dir_all(parent).map_err(|e| { format!( @@ -292,18 +295,18 @@ fn write_asset_manifest(stage_dir: &Path, summary: &AssetCollectionSummary) -> R Ok(()) } -fn cleanup_asset_stage_dir(stage_dir: &Path) -> Result<(), String> { - if !stage_dir.exists() { +fn cleanup_asset_capture_dir(asset_capture_dir: &Path) -> Result<(), String> { + if !asset_capture_dir.exists() { return Ok(()); } - std::fs::remove_dir_all(stage_dir) - .map_err(|e| format!("failed to clean up {}: {e}", stage_dir.display())) + std::fs::remove_dir_all(asset_capture_dir) + .map_err(|e| format!("failed to clean up {}: {e}", asset_capture_dir.display())) } /// Collect asset files matching the configured globs that were created during this stage. pub async fn collect_assets( sandbox: &dyn Sandbox, - stage_dir: &Path, + asset_capture_dir: &Path, globs: &[String], command_start_epoch: f64, ) -> Result { @@ -330,7 +333,7 @@ pub async fn collect_assets( let mut captured_assets: Vec = Vec::new(); for file in &to_collect { - let dest = stage_dir.join(&file.relative_path); + let dest = asset_capture_dir.join(&file.relative_path); match sandbox .download_file_to_local(&file.relative_path, &dest) .await @@ -373,8 +376,8 @@ pub async fn collect_assets( }; if files_copied > 0 { - if let Err(e) = write_asset_manifest(stage_dir, &summary) { - let cleanup_suffix = match cleanup_asset_stage_dir(stage_dir) { + if let Err(e) = write_asset_manifest(asset_capture_dir, &summary) { + let cleanup_suffix = match cleanup_asset_capture_dir(asset_capture_dir) { Ok(()) => String::new(), Err(cleanup_err) => format!("; cleanup failed: {cleanup_err}"), }; @@ -786,7 +789,7 @@ mod tests { #[cfg(unix)] #[test] - fn write_asset_manifest_failure_cleans_up_stage_dir() { + fn write_asset_manifest_failure_cleans_up_asset_capture_dir() { use std::fs::Permissions; use std::os::unix::fs::PermissionsExt; @@ -816,7 +819,7 @@ mod tests { assert!(err.contains("failed to write")); fs::set_permissions(&stage_dir, Permissions::from_mode(0o755)).unwrap(); - cleanup_asset_stage_dir(&stage_dir).unwrap(); + cleanup_asset_capture_dir(&stage_dir).unwrap(); assert!(!stage_dir.exists()); } diff --git a/lib/crates/fabro-workflow/src/handler/agent.rs b/lib/crates/fabro-workflow/src/handler/agent.rs index 2de148dca..888d6de27 100644 --- a/lib/crates/fabro-workflow/src/handler/agent.rs +++ b/lib/crates/fabro-workflow/src/handler/agent.rs @@ -13,10 +13,8 @@ use crate::event::EventEmitter; use crate::outcome::{ FailureCategory, FailureDetail, Outcome, OutcomeExt, StageStatus, StageUsage, }; -use crate::run_dir::{node_dir, visit_from_context}; use crate::vars::expand_vars; use fabro_graphviz::graph::{Graph, Node}; -use tokio::fs; use super::{EngineServices, Handler}; @@ -43,7 +41,6 @@ pub trait CodergenBackend: Send + Sync { context: &Context, thread_id: Option<&str>, emitter: &Arc, - stage_dir: &Path, sandbox: &Arc, tool_hooks: Option>, ) -> Result; @@ -54,7 +51,6 @@ pub trait CodergenBackend: Send + Sync { _node: &Node, _prompt: &str, _system_prompt: Option<&str>, - _stage_dir: &Path, ) -> Result { Err(FabroError::Validation( "one_shot mode not supported by this backend".into(), @@ -236,7 +232,7 @@ impl Handler for AgentHandler { node: &Node, context: &Context, graph: &Graph, - run_dir: &Path, + _run_dir: &Path, services: &EngineServices, ) -> Result { // 1. Build prompt (prepend fidelity preamble if present) @@ -252,10 +248,6 @@ impl Handler for AgentHandler { format!("{preamble}\n\n{expanded}") }; - let visit = visit_from_context(context); - let stage_dir = node_dir(run_dir, &node.id, visit); - fs::create_dir_all(&stage_dir).await?; - // 3. Call LLM backend (agent loop) let thread_id = context.thread_id(); let run_id = context @@ -282,7 +274,6 @@ impl Handler for AgentHandler { context, thread_id.as_deref(), &services.emitter, - &stage_dir, &services.sandbox, tool_hooks, ) @@ -563,7 +554,6 @@ mod tests { _context: &Context, _thread_id: Option<&str>, _emitter: &Arc, - _stage_dir: &Path, _sandbox: &Arc, _tool_hooks: Option>, ) -> Result { @@ -620,7 +610,6 @@ mod tests { _context: &Context, _thread_id: Option<&str>, _emitter: &Arc, - _stage_dir: &Path, _sandbox: &Arc, _tool_hooks: Option>, ) -> Result { @@ -679,7 +668,6 @@ mod tests { context: &Context, _thread_id: Option<&str>, emitter: &Arc, - _stage_dir: &Path, _sandbox: &Arc, _tool_hooks: Option>, ) -> Result { @@ -791,7 +779,6 @@ mod tests { _context: &Context, thread_id: Option<&str>, _emitter: &Arc, - _stage_dir: &Path, _sandbox: &Arc, _tool_hooks: Option>, ) -> Result { @@ -844,7 +831,6 @@ mod tests { _context: &Context, thread_id: Option<&str>, _emitter: &Arc, - _stage_dir: &Path, _sandbox: &Arc, _tool_hooks: Option>, ) -> Result { @@ -892,7 +878,6 @@ mod tests { _context: &Context, _thread_id: Option<&str>, _emitter: &Arc, - _stage_dir: &Path, _sandbox: &Arc, _tool_hooks: Option>, ) -> Result { @@ -1036,7 +1021,6 @@ Some text in between. _context: &Context, _thread_id: Option<&str>, _emitter: &Arc, - _stage_dir: &Path, _sandbox: &Arc, _tool_hooks: Option>, ) -> Result { @@ -1075,7 +1059,6 @@ Some text in between. _context: &Context, _thread_id: Option<&str>, _emitter: &Arc, - _stage_dir: &std::path::Path, _sandbox: &Arc, _tool_hooks: Option>, ) -> Result { @@ -1145,7 +1128,6 @@ Some text in between. _context: &Context, _thread_id: Option<&str>, _emitter: &Arc, - _stage_dir: &std::path::Path, _sandbox: &Arc, _tool_hooks: Option>, ) -> Result { diff --git a/lib/crates/fabro-workflow/src/handler/fan_in.rs b/lib/crates/fabro-workflow/src/handler/fan_in.rs index 43dd3ca91..52e1d3ebc 100644 --- a/lib/crates/fabro-workflow/src/handler/fan_in.rs +++ b/lib/crates/fabro-workflow/src/handler/fan_in.rs @@ -6,12 +6,11 @@ use crate::context::keys; use crate::error::FabroError; use crate::event::{EventEmitter, WorkflowRunEvent}; use crate::outcome::{Outcome, OutcomeExt}; -use crate::run_dir::{node_dir, visit_from_context}; +use crate::run_dir::visit_from_context; use crate::sandbox_git::git_merge_ff_only; use async_trait::async_trait; use fabro_agent::Sandbox; use fabro_graphviz::graph::{Graph, Node}; -use tokio::fs; use super::agent::{CodergenBackend, CodergenResult}; use super::{EngineServices, Handler}; @@ -219,7 +218,7 @@ async fn llm_evaluate( prompt: &str, results: &serde_json::Value, context: &Context, - run_dir: &Path, + _run_dir: &Path, node_id: &str, emitter: &Arc, sandbox: &Arc, @@ -232,10 +231,7 @@ async fn llm_evaluate( Respond with the ID of the best candidate." ); - let visit = visit_from_context(context); - let visit_u32 = u32::try_from(visit).unwrap_or(u32::MAX); - let stage_dir = node_dir(run_dir, node_id, visit); - fs::create_dir_all(&stage_dir).await?; + let visit_u32 = u32::try_from(visit_from_context(context)).unwrap_or(u32::MAX); emitter.emit(&WorkflowRunEvent::Prompt { stage: node_id.to_string(), @@ -257,7 +253,6 @@ async fn llm_evaluate( context, None, emitter, - &stage_dir, sandbox, None, ) @@ -464,7 +459,6 @@ mod tests { _context: &Context, _thread_id: Option<&str>, _emitter: &Arc, - _stage_dir: &std::path::Path, _sandbox: &Arc, _tool_hooks: Option>, ) -> Result { @@ -505,17 +499,6 @@ mod tests { outcome.context_updates.get(keys::PARALLEL_FAN_IN_BEST_ID), Some(&serde_json::json!("branch_b")) ); - - // Verify prompt and response files were written - let prompt_path = tmp.path().join("nodes").join("fan_in").join("prompt.md"); - assert!(prompt_path.exists()); - let prompt_content = std::fs::read_to_string(&prompt_path).unwrap(); - assert!(prompt_content.contains("Pick the best branch")); - - let response_path = tmp.path().join("nodes").join("fan_in").join("response.md"); - assert!(response_path.exists()); - let response_content = std::fs::read_to_string(&response_path).unwrap(); - assert!(response_content.contains("branch_b")); } #[tokio::test] diff --git a/lib/crates/fabro-workflow/src/handler/llm/api.rs b/lib/crates/fabro-workflow/src/handler/llm/api.rs index 56449fc56..d1e93538c 100644 --- a/lib/crates/fabro-workflow/src/handler/llm/api.rs +++ b/lib/crates/fabro-workflow/src/handler/llm/api.rs @@ -268,7 +268,6 @@ impl CodergenBackend for AgentApiBackend { node: &Node, prompt: &str, system_prompt: Option<&str>, - _stage_dir: &std::path::Path, ) -> Result { let client = Client::from_env() .await @@ -412,7 +411,6 @@ impl CodergenBackend for AgentApiBackend { context: &Context, thread_id: Option<&str>, emitter: &Arc, - _stage_dir: &std::path::Path, sandbox: &Arc, tool_hooks: Option>, ) -> Result { diff --git a/lib/crates/fabro-workflow/src/handler/llm/cli.rs b/lib/crates/fabro-workflow/src/handler/llm/cli.rs index 7c9a0e845..25b58a589 100644 --- a/lib/crates/fabro-workflow/src/handler/llm/cli.rs +++ b/lib/crates/fabro-workflow/src/handler/llm/cli.rs @@ -1,5 +1,4 @@ use std::collections::HashMap; -use std::path::Path; use std::sync::Arc; use async_trait::async_trait; @@ -465,7 +464,6 @@ impl CodergenBackend for AgentCliBackend { _context: &Context, _thread_id: Option<&str>, emitter: &Arc, - _stage_dir: &Path, sandbox: &Arc, _tool_hooks: Option>, ) -> Result { @@ -755,20 +753,19 @@ impl CodergenBackend for BackendRouter { context: &Context, thread_id: Option<&str>, emitter: &Arc, - stage_dir: &Path, sandbox: &Arc, tool_hooks: Option>, ) -> Result { if self.should_use_cli(node) { self.cli_backend .run( - node, prompt, context, thread_id, emitter, stage_dir, sandbox, tool_hooks, + node, prompt, context, thread_id, emitter, sandbox, tool_hooks, ) .await } else { self.api_backend .run( - node, prompt, context, thread_id, emitter, stage_dir, sandbox, tool_hooks, + node, prompt, context, thread_id, emitter, sandbox, tool_hooks, ) .await } @@ -779,12 +776,9 @@ impl CodergenBackend for BackendRouter { node: &Node, prompt: &str, system_prompt: Option<&str>, - stage_dir: &Path, ) -> Result { // CLI backend doesn't support one_shot, always route to API - self.api_backend - .one_shot(node, prompt, system_prompt, stage_dir) - .await + self.api_backend.one_shot(node, prompt, system_prompt).await } } @@ -792,6 +786,7 @@ impl CodergenBackend for BackendRouter { mod tests { use super::*; use fabro_graphviz::graph::AttrValue; + use std::path::Path; // -- AgentCli -- @@ -1196,7 +1191,6 @@ mod tests { _context: &Context, _thread_id: Option<&str>, _emitter: &Arc, - _stage_dir: &Path, _sandbox: &Arc, _tool_hooks: Option>, ) -> Result { diff --git a/lib/crates/fabro-workflow/src/handler/prompt.rs b/lib/crates/fabro-workflow/src/handler/prompt.rs index 7607156c5..e0662baaf 100644 --- a/lib/crates/fabro-workflow/src/handler/prompt.rs +++ b/lib/crates/fabro-workflow/src/handler/prompt.rs @@ -7,10 +7,9 @@ use crate::context::{Context, WorkflowContext}; use crate::error::FabroError; use crate::event::WorkflowRunEvent; use crate::outcome::Outcome; -use crate::run_dir::{node_dir, visit_from_context}; +use crate::run_dir::visit_from_context; use fabro_graphviz::graph::{Graph, Node}; use fabro_model::Provider; -use tokio::fs; use super::agent::{ CodergenBackend, CodergenResult, expand_variables, extract_status_fields, truncate, @@ -47,7 +46,7 @@ impl Handler for PromptHandler { node: &Node, context: &Context, graph: &Graph, - run_dir: &Path, + _run_dir: &Path, services: &EngineServices, ) -> Result { // 1. Build prompt (prepend fidelity preamble if present) @@ -87,10 +86,6 @@ impl Handler for PromptHandler { None }; - let visit = visit_from_context(context); - let stage_dir = node_dir(run_dir, &node.id, visit); - fs::create_dir_all(&stage_dir).await?; - let prompt_provider = node .provider() .map(String::from) @@ -98,7 +93,7 @@ impl Handler for PromptHandler { let prompt_model = node.model().map(String::from); services.emitter.emit(&WorkflowRunEvent::Prompt { stage: node.id.clone(), - visit: u32::try_from(visit).unwrap_or(u32::MAX), + visit: u32::try_from(visit_from_context(context)).unwrap_or(u32::MAX), text: prompt.clone(), mode: Some("prompt".to_string()), provider: prompt_provider.clone(), @@ -109,7 +104,7 @@ impl Handler for PromptHandler { let (response_text, stage_usage, backend_files_touched) = if let Some(backend) = &self.backend { let result = backend - .one_shot(node, &prompt, system_prompt.as_deref(), &stage_dir) + .one_shot(node, &prompt, system_prompt.as_deref()) .await; match result { Ok(CodergenResult::Full(outcome)) => return Ok(outcome), @@ -268,7 +263,6 @@ mod tests { _context: &Context, _thread_id: Option<&str>, _emitter: &Arc, - _stage_dir: &Path, _sandbox: &Arc, _tool_hooks: Option>, ) -> Result { @@ -280,7 +274,6 @@ mod tests { _node: &Node, _prompt: &str, _system_prompt: Option<&str>, - _stage_dir: &Path, ) -> Result { Ok(CodergenResult::Text { text: "one-shot response".to_string(), @@ -307,14 +300,12 @@ mod tests { .unwrap(); assert_eq!(outcome.status, crate::outcome::StageStatus::Success); - let response_content = std::fs::read_to_string( - tmp.path() - .join("nodes") - .join("classify") - .join("response.md"), - ) - .unwrap(); - assert_eq!(response_content, "one-shot response"); + assert_eq!( + outcome + .context_updates + .get(&crate::context::keys::response_key("classify")), + Some(&serde_json::json!("one-shot response")) + ); } #[tokio::test] @@ -332,7 +323,6 @@ mod tests { _context: &Context, _thread_id: Option<&str>, _emitter: &Arc, - _stage_dir: &Path, _sandbox: &Arc, _tool_hooks: Option>, ) -> Result { @@ -344,7 +334,6 @@ mod tests { _node: &Node, _prompt: &str, _system_prompt: Option<&str>, - _stage_dir: &Path, ) -> Result { Ok(CodergenResult::Text { text: "one-shot response".to_string(), @@ -396,7 +385,6 @@ mod tests { _context: &Context, _thread_id: Option<&str>, _emitter: &Arc, - _stage_dir: &Path, _sandbox: &Arc, _tool_hooks: Option>, ) -> Result { @@ -408,7 +396,6 @@ mod tests { _node: &Node, prompt: &str, system_prompt: Option<&str>, - _stage_dir: &Path, ) -> Result { *self.captured_prompt.lock().unwrap() = Some(prompt.to_string()); *self.captured_system_prompt.lock().unwrap() = Some(system_prompt.map(String::from)); diff --git a/lib/crates/fabro-workflow/src/lib.rs b/lib/crates/fabro-workflow/src/lib.rs index 026b00bf9..fdf39bfb9 100644 --- a/lib/crates/fabro-workflow/src/lib.rs +++ b/lib/crates/fabro-workflow/src/lib.rs @@ -19,7 +19,6 @@ use std::sync::Arc; use fabro_retro::retro::CompletedStage; use fabro_store::EventEnvelope; -use serde::de::DeserializeOwned; /// Callback invoked when a workflow node starts executing. pub type OnNodeCallback = Option>; @@ -29,16 +28,6 @@ pub(crate) fn millis_u64(d: std::time::Duration) -> u64 { u64::try_from(d.as_millis()).unwrap_or(u64::MAX) } -/// Load a value from a JSON file. -pub(crate) fn load_json( - path: &std::path::Path, - label: &str, -) -> error::Result { - let data = std::fs::read_to_string(path)?; - serde_json::from_str(&data) - .map_err(|e| error::FabroError::Checkpoint(format!("{label} deserialize failed: {e}"))) -} - /// Build `Vec` from a `Checkpoint`, mapping workflow-engine /// types into the flat struct expected by `fabro_retro::retro::derive_retro`. pub fn build_completed_stages(cp: &records::Checkpoint, run_failed: bool) -> Vec { diff --git a/lib/crates/fabro-workflow/src/lifecycle/artifact.rs b/lib/crates/fabro-workflow/src/lifecycle/artifact.rs index 5c1f18dcc..76ca11c7f 100644 --- a/lib/crates/fabro-workflow/src/lifecycle/artifact.rs +++ b/lib/crates/fabro-workflow/src/lifecycle/artifact.rs @@ -95,13 +95,13 @@ impl RunLifecycle for ArtifactLifecycle { } else { format!("{node_id}-visit_{visit}") }; - let stage_dir = self + let asset_capture_dir = self .assets_dir .join(&node_slug) .join(format!("retry_{}", ctx.attempt)); - let _ = std::fs::create_dir_all(&stage_dir); + let _ = std::fs::create_dir_all(&asset_capture_dir); - match collect_assets(&*self.sandbox, &stage_dir, &self.asset_globs, epoch).await { + match collect_assets(&*self.sandbox, &asset_capture_dir, &self.asset_globs, epoch).await { Ok(summary) if summary.files_copied > 0 => { for asset in &summary.captured_assets { self.emitter.emit(&WorkflowRunEvent::AssetCaptured { diff --git a/lib/crates/fabro-workflow/src/operations/start.rs b/lib/crates/fabro-workflow/src/operations/start.rs index e2f69d139..7c00d7b98 100644 --- a/lib/crates/fabro-workflow/src/operations/start.rs +++ b/lib/crates/fabro-workflow/src/operations/start.rs @@ -1119,7 +1119,11 @@ mod tests { HashMap::new(), HashMap::new(), ); - checkpoint.save(&run_dir.join("checkpoint.json")).unwrap(); + std::fs::write( + run_dir.join("checkpoint.json"), + serde_json::to_string_pretty(&checkpoint).unwrap(), + ) + .unwrap(); let conclusion = crate::records::Conclusion { timestamp: Utc::now(), diff --git a/lib/crates/fabro-workflow/src/pipeline/pull_request.rs b/lib/crates/fabro-workflow/src/pipeline/pull_request.rs index 5784ad737..ad7477de7 100644 --- a/lib/crates/fabro-workflow/src/pipeline/pull_request.rs +++ b/lib/crates/fabro-workflow/src/pipeline/pull_request.rs @@ -1431,32 +1431,6 @@ mod tests { assert_eq!(pr_title_from_goal("Fix bug"), "Fix bug"); } - #[test] - fn pull_request_record_save_writes_json() { - let tmp = tempfile::tempdir().unwrap(); - let path = tmp.path().join("pull_request.json"); - let record = PullRequestRecord { - html_url: "https://github.com/owner/repo/pull/42".to_string(), - number: 42, - owner: "owner".to_string(), - repo: "repo".to_string(), - base_branch: "main".to_string(), - head_branch: "fabro/run/abc".to_string(), - title: "Fix the thing".to_string(), - }; - record.save(&path).unwrap(); - - let content: serde_json::Value = - serde_json::from_str(&std::fs::read_to_string(&path).unwrap()).unwrap(); - assert_eq!(content["html_url"], "https://github.com/owner/repo/pull/42"); - assert_eq!(content["number"], 42); - assert_eq!(content["owner"], "owner"); - assert_eq!(content["repo"], "repo"); - assert_eq!(content["base_branch"], "main"); - assert_eq!(content["head_branch"], "fabro/run/abc"); - assert_eq!(content["title"], "Fix the thing"); - } - #[tokio::test] async fn empty_diff_returns_none() { let tmp = tempfile::tempdir().unwrap(); diff --git a/lib/crates/fabro-workflow/src/records/checkpoint.rs b/lib/crates/fabro-workflow/src/records/checkpoint.rs index 7ccb5bf45..c463062a3 100644 --- a/lib/crates/fabro-workflow/src/records/checkpoint.rs +++ b/lib/crates/fabro-workflow/src/records/checkpoint.rs @@ -1,8 +1,6 @@ use std::collections::HashMap; -use std::path::Path; use crate::context::Context; -use crate::error::{FabroError, Result as CrateResult}; use crate::outcome::Outcome; pub use fabro_types::checkpoint::Checkpoint; use fabro_types::failure_signature::FailureSignature; @@ -19,11 +17,6 @@ pub trait CheckpointExt { restart_failure_signatures: HashMap, node_visits: HashMap, ) -> Self; - - fn save(&self, path: &Path) -> CrateResult<()>; - fn load(path: &Path) -> CrateResult - where - Self: Sized; } impl CheckpointExt for Checkpoint { @@ -52,15 +45,4 @@ impl CheckpointExt for Checkpoint { node_visits, } } - - fn save(&self, path: &Path) -> CrateResult<()> { - let json = serde_json::to_string_pretty(self) - .map_err(|e| FabroError::Checkpoint(format!("checkpoint serialize failed: {e}")))?; - std::fs::write(path, json)?; - Ok(()) - } - - fn load(path: &Path) -> CrateResult { - crate::load_json(path, "checkpoint") - } } diff --git a/lib/crates/fabro-workflow/src/run_dir.rs b/lib/crates/fabro-workflow/src/run_dir.rs index a1028e1c6..9acd5e698 100644 --- a/lib/crates/fabro-workflow/src/run_dir.rs +++ b/lib/crates/fabro-workflow/src/run_dir.rs @@ -1,21 +1,5 @@ -use std::path::{Path, PathBuf}; - use crate::context::Context; -/// Return the directory for a node's logs. -/// -/// First visit (`visit <= 1`): `{run_dir}/nodes/{node_id}` -/// Subsequent visits: `{run_dir}/nodes/{node_id}-visit_{visit}` -pub(crate) fn node_dir(run_dir: &Path, node_id: &str, visit: usize) -> PathBuf { - if visit <= 1 { - run_dir.join("nodes").join(node_id) - } else { - run_dir - .join("nodes") - .join(format!("{node_id}-visit_{visit}")) - } -} - /// Read the workflow visit ordinal from context. /// /// The raw context value is `0` when unset; workflow execution code treats @@ -45,28 +29,4 @@ mod tests { ); assert_eq!(visit_from_context(&ctx), 3); } - - #[test] - fn node_dir_first_visit() { - let root = Path::new("/tmp/logs"); - assert_eq!(node_dir(root, "work", 1), root.join("nodes").join("work")); - } - - #[test] - fn node_dir_second_visit() { - let root = Path::new("/tmp/logs"); - assert_eq!( - node_dir(root, "work", 2), - root.join("nodes").join("work-visit_2") - ); - } - - #[test] - fn node_dir_fifth_visit() { - let root = Path::new("/tmp/logs"); - assert_eq!( - node_dir(root, "work", 5), - root.join("nodes").join("work-visit_5") - ); - } } diff --git a/lib/crates/fabro-workflow/tests/it/daytona_integration.rs b/lib/crates/fabro-workflow/tests/it/daytona_integration.rs index 336bfb14b..1495cca8b 100644 --- a/lib/crates/fabro-workflow/tests/it/daytona_integration.rs +++ b/lib/crates/fabro-workflow/tests/it/daytona_integration.rs @@ -31,7 +31,7 @@ use fabro_workflow::handler::exit::ExitHandler; use fabro_workflow::handler::start::StartHandler; use fabro_workflow::handler::{Handler, HandlerRegistry}; use fabro_workflow::outcome::{Outcome, OutcomeExt, StageStatus}; -use fabro_workflow::records::{Checkpoint, CheckpointExt}; +use fabro_workflow::records::Checkpoint; use fabro_workflow::run_options::{GitCheckpointOptions, RunOptions}; use fabro_workflow::test_support::WorkflowRunner; use ulid::Ulid; @@ -42,6 +42,11 @@ fn test_run_id(label: &str) -> RunId { RunId::from(Ulid(u128::from(hasher.finish()))) } +fn load_checkpoint(path: &Path) -> Result> { + let data = std::fs::read_to_string(path)?; + Ok(serde_json::from_str(&data)?) +} + async fn create_env() -> DaytonaSandbox { let creds = load_github_app_credentials(); create_env_with_github_app(Some(creds)).await @@ -410,7 +415,7 @@ async fn daytona_pipeline_artifact_offload_and_sync() { // Checkpoint should have a pointer rewritten for Daytona let checkpoint = - Checkpoint::load(&dir.path().join("checkpoint.json")).expect("checkpoint should load"); + load_checkpoint(&dir.path().join("checkpoint.json")).expect("checkpoint should load"); let pointer_value = checkpoint .context_values .get("response.big_output") @@ -638,7 +643,7 @@ async fn daytona_git_checkpoint_remote_emits_events() { // Verify checkpoint.json has git_commit_sha let checkpoint = - Checkpoint::load(&dir.path().join("checkpoint.json")).expect("checkpoint should load"); + load_checkpoint(&dir.path().join("checkpoint.json")).expect("checkpoint should load"); assert!( checkpoint.git_commit_sha.is_some(), "checkpoint should have git_commit_sha" @@ -789,7 +794,7 @@ async fn daytona_parallel_git_branching_e2e() { // Verify parallel.results has head_sha for each branch let checkpoint = - Checkpoint::load(&run_tmp.path().join("checkpoint.json")).expect("checkpoint should load"); + load_checkpoint(&run_tmp.path().join("checkpoint.json")).expect("checkpoint should load"); let parallel_results = checkpoint .context_values .get("parallel.results") @@ -958,7 +963,6 @@ async fn run_daytona_cli_test(provider: Provider, model: &str, install_command: &context, None, &emitter, - dir.path(), &env, None, ) diff --git a/lib/crates/fabro-workflow/tests/it/integration.rs b/lib/crates/fabro-workflow/tests/it/integration.rs index fad936176..832da0cc0 100644 --- a/lib/crates/fabro-workflow/tests/it/integration.rs +++ b/lib/crates/fabro-workflow/tests/it/integration.rs @@ -62,6 +62,15 @@ fn test_run_id(label: &str) -> RunId { RunId::from(Ulid(u128::from(hasher.finish()))) } +fn load_checkpoint(path: &Path) -> Result> { + let data = std::fs::read_to_string(path)?; + Ok(serde_json::from_str(&data)?) +} + +fn save_checkpoint(path: &Path, checkpoint: &Checkpoint) { + std::fs::write(path, serde_json::to_string_pretty(checkpoint).unwrap()).unwrap(); +} + // --------------------------------------------------------------------------- // 1. Parse and validate all 3 spec examples (Section 2.13) // --------------------------------------------------------------------------- @@ -236,7 +245,7 @@ async fn end_to_end_linear_pipeline() { let checkpoint_path = dir.path().join("checkpoint.json"); assert!(checkpoint_path.exists(), "checkpoint.json should exist"); - let checkpoint = Checkpoint::load(&checkpoint_path).expect("checkpoint should load"); + let checkpoint = load_checkpoint(&checkpoint_path).expect("checkpoint should load"); assert!(checkpoint.completed_nodes.contains(&"start".to_string())); assert!( checkpoint @@ -372,7 +381,7 @@ async fn end_to_end_branching_pipeline() { .expect("run should succeed"); assert_eq!(outcome.status, StageStatus::Success); - let checkpoint = Checkpoint::load(&dir.path().join("checkpoint.json")).unwrap(); + let checkpoint = load_checkpoint(&dir.path().join("checkpoint.json")).unwrap(); assert!( checkpoint .completed_nodes @@ -491,7 +500,7 @@ async fn end_to_end_human_gate_pipeline() { .expect("run should succeed"); assert_eq!(outcome.status, StageStatus::Success); - let checkpoint = Checkpoint::load(&dir.path().join("checkpoint.json")).unwrap(); + let checkpoint = load_checkpoint(&dir.path().join("checkpoint.json")).unwrap(); assert!( checkpoint.completed_nodes.contains(&"reject".to_string()), "should have traversed reject path" @@ -597,7 +606,7 @@ async fn human_gate_aborted_input_fails_closed_without_fail_route() { "unexpected outcome: {outcome:?}" ); - let checkpoint = Checkpoint::load(&dir.path().join("checkpoint.json")).unwrap(); + let checkpoint = load_checkpoint(&dir.path().join("checkpoint.json")).unwrap(); assert!( checkpoint.node_outcomes.contains_key("gate"), "gate outcome should be checkpointed before termination" @@ -696,7 +705,7 @@ async fn human_gate_aborted_input_routes_via_outcome_fail_condition() { .expect("aborted human gate should follow explicit fail route"); assert_eq!(outcome.status, StageStatus::Success); - let checkpoint = Checkpoint::load(&dir.path().join("checkpoint.json")).unwrap(); + let checkpoint = load_checkpoint(&dir.path().join("checkpoint.json")).unwrap(); assert!( checkpoint .completed_nodes @@ -928,7 +937,7 @@ async fn goal_gate_routes_to_retry_target_when_present() { .expect("run should eventually succeed after retry"); assert_eq!(outcome.status, StageStatus::Success); - let checkpoint = Checkpoint::load(&dir.path().join("checkpoint.json")).unwrap(); + let checkpoint = load_checkpoint(&dir.path().join("checkpoint.json")).unwrap(); // gated_work should appear in completed nodes (at least twice -- first fail, then succeed) let gated_work_count = checkpoint .completed_nodes @@ -1313,7 +1322,7 @@ async fn pipeline_with_many_nodes() { .expect("large pipeline should succeed"); assert_eq!(outcome.status, StageStatus::Success); - let checkpoint = Checkpoint::load(&dir.path().join("checkpoint.json")).unwrap(); + let checkpoint = load_checkpoint(&dir.path().join("checkpoint.json")).unwrap(); // All 10 step nodes should be in completed_nodes for name in &node_names { assert!( @@ -1349,9 +1358,9 @@ fn checkpoint_save_and_resume_roundtrip() { std::collections::HashMap::new(), ); - checkpoint.save(&path).expect("save should succeed"); + save_checkpoint(&path, &checkpoint); - let loaded = Checkpoint::load(&path).expect("load should succeed"); + let loaded = load_checkpoint(&path).expect("load should succeed"); assert_eq!(loaded.current_node, "step_2"); assert_eq!(loaded.completed_nodes.len(), 2); assert!(loaded.completed_nodes.contains(&"start".to_string())); @@ -1382,7 +1391,6 @@ impl CodergenBackend for MockCodergenBackend { _context: &Context, _thread_id: Option<&str>, _emitter: &Arc, - _stage_dir: &std::path::Path, _sandbox: &Arc, _tool_hooks: Option>, ) -> Result { @@ -1634,7 +1642,7 @@ async fn smoke_test_with_mock_codergen_backend() { .expect("smoke test should succeed"); assert_eq!(outcome.status, StageStatus::Success); - let checkpoint = Checkpoint::load(&dir.path().join("checkpoint.json")).unwrap(); + let checkpoint = load_checkpoint(&dir.path().join("checkpoint.json")).unwrap(); assert!( checkpoint.completed_nodes.contains(&"plan".to_string()), "plan should have executed" @@ -1734,7 +1742,7 @@ async fn end_to_end_parallel_fan_out_fan_in() { .expect("parallel pipeline should succeed"); assert_eq!(outcome.status, StageStatus::Success); - let checkpoint = Checkpoint::load(&dir.path().join("checkpoint.json")).unwrap(); + let checkpoint = load_checkpoint(&dir.path().join("checkpoint.json")).unwrap(); // The parallel node (fan_out) and fan_in_node should be in completed_nodes. // Branch nodes run inside the parallel handler, so they are not recorded @@ -1847,7 +1855,7 @@ async fn resume_from_checkpoint_completes_pipeline() { assert_eq!(outcome.status, StageStatus::Success); // Verify checkpoint written after resume contains step_b - let final_cp = Checkpoint::load(&dir.path().join("checkpoint.json")).unwrap(); + let final_cp = load_checkpoint(&dir.path().join("checkpoint.json")).unwrap(); assert!( final_cp.completed_nodes.contains(&"step_b".to_string()), "step_b should have been executed after resume" @@ -1982,7 +1990,7 @@ async fn graph_goal_in_context() { }; engine.run(&graph, &run_options).await.expect("run"); - let cp = Checkpoint::load(&dir.path().join("checkpoint.json")).unwrap(); + let cp = load_checkpoint(&dir.path().join("checkpoint.json")).unwrap(); assert_eq!( cp.context_values.get("graph.goal"), Some(&serde_json::json!("Ship the widget")) @@ -2092,7 +2100,7 @@ async fn context_flow_between_stages() { }; engine.run(&graph, &run_options).await.expect("run"); - let cp = Checkpoint::load(&dir.path().join("checkpoint.json")).unwrap(); + let cp = load_checkpoint(&dir.path().join("checkpoint.json")).unwrap(); assert_eq!( cp.context_values.get("last_stage"), Some(&serde_json::json!("step_b")) @@ -2145,7 +2153,7 @@ async fn tool_handler_e2e() { let outcome = engine.run(&graph, &run_options).await.expect("run"); assert_eq!(outcome.status, StageStatus::Success); - let cp = Checkpoint::load(&dir.path().join("checkpoint.json")).unwrap(); + let cp = load_checkpoint(&dir.path().join("checkpoint.json")).unwrap(); let command_output = cp .context_values .get("command.output") @@ -2216,7 +2224,7 @@ async fn auto_approve_interviewer_e2e() { let outcome = engine.run(&graph, &run_options).await.expect("run"); assert_eq!(outcome.status, StageStatus::Success); - let cp = Checkpoint::load(&dir.path().join("checkpoint.json")).unwrap(); + let cp = load_checkpoint(&dir.path().join("checkpoint.json")).unwrap(); assert!(cp.completed_nodes.contains(&"approve".to_string())); assert!(!cp.completed_nodes.contains(&"reject".to_string())); } @@ -2251,7 +2259,7 @@ async fn codergen_without_backend_simulated() { }; engine.run(&graph, &run_options).await.expect("run"); - let cp = Checkpoint::load(&dir.path().join("checkpoint.json")).unwrap(); + let cp = load_checkpoint(&dir.path().join("checkpoint.json")).unwrap(); let last_response = cp .context_values .get("last_response") @@ -2353,7 +2361,7 @@ async fn branching_loop_back_on_failure() { let outcome = engine.run(&graph, &run_options).await.expect("run"); assert_eq!(outcome.status, StageStatus::Success); - let cp = Checkpoint::load(&dir.path().join("checkpoint.json")).unwrap(); + let cp = load_checkpoint(&dir.path().join("checkpoint.json")).unwrap(); let implement_count = cp .completed_nodes .iter() @@ -2435,7 +2443,7 @@ async fn human_gate_loops_back() { let outcome = engine.run(&graph, &run_options).await.expect("run"); assert_eq!(outcome.status, StageStatus::Success); - let cp = Checkpoint::load(&dir.path().join("checkpoint.json")).unwrap(); + let cp = load_checkpoint(&dir.path().join("checkpoint.json")).unwrap(); let gate_count = cp.completed_nodes.iter().filter(|n| *n == "gate").count(); assert!( gate_count >= 2, @@ -2492,7 +2500,7 @@ async fn scenario_ship_a_feature() { let outcome = engine.run(&graph, &run_options).await.expect("run"); assert_eq!(outcome.status, StageStatus::Success); - let cp = Checkpoint::load(&dir.path().join("checkpoint.json")).unwrap(); + let cp = load_checkpoint(&dir.path().join("checkpoint.json")).unwrap(); let command_output = cp .context_values .get("command.output") @@ -2573,7 +2581,7 @@ async fn scenario_parallel_expert_review() { let outcome = engine.run(&graph, &run_options).await.expect("run"); assert_eq!(outcome.status, StageStatus::Success); - let cp = Checkpoint::load(&dir.path().join("checkpoint.json")).unwrap(); + let cp = load_checkpoint(&dir.path().join("checkpoint.json")).unwrap(); let results = cp .context_values .get("parallel.results") @@ -2656,7 +2664,7 @@ async fn scenario_node_retries_on_retry_status() { let outcome = engine.run(&graph, &run_options).await.expect("run"); assert_eq!(outcome.status, StageStatus::Success); - let cp = Checkpoint::load(&dir.path().join("checkpoint.json")).unwrap(); + let cp = load_checkpoint(&dir.path().join("checkpoint.json")).unwrap(); let retry_count = cp .node_retries .get("flaky") @@ -2784,7 +2792,7 @@ async fn scenario_bug_triage_router() { let outcome = engine.run(&graph, &run_options).await.expect("run"); assert_eq!(outcome.status, StageStatus::Success); - let cp = Checkpoint::load(&dir.path().join("checkpoint.json")).unwrap(); + let cp = load_checkpoint(&dir.path().join("checkpoint.json")).unwrap(); assert!( cp.completed_nodes.contains(&"critical".to_string()), "critical should be selected (highest weight)" @@ -2845,7 +2853,7 @@ async fn scenario_crash_recovery() { .expect("run"); assert_eq!(outcome.status, StageStatus::Success); - let cp = Checkpoint::load(&dir.path().join("checkpoint.json")).unwrap(); + let cp = load_checkpoint(&dir.path().join("checkpoint.json")).unwrap(); assert!(cp.completed_nodes.contains(&"b".to_string())); assert!(cp.completed_nodes.contains(&"c".to_string())); assert!(cp.completed_nodes.contains(&"a".to_string())); @@ -2949,7 +2957,7 @@ async fn manager_loop_stop_condition_satisfied_e2e() { }; let outcome = engine.run(&graph, &run_options).await.expect("run"); - let cp = Checkpoint::load(&dir.path().join("checkpoint.json")).unwrap(); + let cp = load_checkpoint(&dir.path().join("checkpoint.json")).unwrap(); let manager_outcome = cp.node_outcomes.get("manager").expect("manager outcome"); assert_eq!(manager_outcome.status, StageStatus::Success); assert!( @@ -3027,7 +3035,7 @@ async fn manager_loop_max_cycles_exceeded_e2e() { }; let outcome = engine.run(&graph, &run_options).await.expect("run"); - let cp = Checkpoint::load(&dir.path().join("checkpoint.json")).unwrap(); + let cp = load_checkpoint(&dir.path().join("checkpoint.json")).unwrap(); let manager_outcome = cp.node_outcomes.get("manager").expect("manager outcome"); assert_eq!(manager_outcome.status, StageStatus::Fail); assert!( @@ -3165,7 +3173,7 @@ async fn conditional_branching_success_fail_paths() { let outcome = engine.run(&graph, &run_options).await.expect("run"); assert_eq!(outcome.status, StageStatus::Success); - let cp = Checkpoint::load(&dir.path().join("checkpoint.json")).unwrap(); + let cp = load_checkpoint(&dir.path().join("checkpoint.json")).unwrap(); assert!(cp.completed_nodes.contains(&"fail_path".to_string())); assert!(!cp.completed_nodes.contains(&"success_path".to_string())); } @@ -3216,7 +3224,7 @@ async fn edge_selection_condition_match_wins_over_weight() { }; engine.run(&graph, &run_options).await.expect("run"); - let cp = Checkpoint::load(&dir.path().join("checkpoint.json")).unwrap(); + let cp = load_checkpoint(&dir.path().join("checkpoint.json")).unwrap(); assert!(cp.completed_nodes.contains(&"cond_target".to_string())); assert!(!cp.completed_nodes.contains(&"weighted_target".to_string())); } @@ -3262,7 +3270,7 @@ async fn edge_selection_weight_breaks_ties() { }; engine.run(&graph, &run_options).await.expect("run"); - let cp = Checkpoint::load(&dir.path().join("checkpoint.json")).unwrap(); + let cp = load_checkpoint(&dir.path().join("checkpoint.json")).unwrap(); assert!(cp.completed_nodes.contains(&"high".to_string())); assert!(!cp.completed_nodes.contains(&"low".to_string())); } @@ -3300,7 +3308,7 @@ async fn edge_selection_lexical_tiebreak() { }; engine.run(&graph, &run_options).await.expect("run"); - let cp = Checkpoint::load(&dir.path().join("checkpoint.json")).unwrap(); + let cp = load_checkpoint(&dir.path().join("checkpoint.json")).unwrap(); assert!(cp.completed_nodes.contains(&"alpha".to_string())); assert!(!cp.completed_nodes.contains(&"beta".to_string())); } @@ -3357,7 +3365,7 @@ async fn context_updates_visible_across_nodes() { }; engine.run(&graph, &run_options).await.expect("run"); - let cp = Checkpoint::load(&dir.path().join("checkpoint.json")).unwrap(); + let cp = load_checkpoint(&dir.path().join("checkpoint.json")).unwrap(); assert!(cp.completed_nodes.contains(&"yes".to_string())); assert!(!cp.completed_nodes.contains(&"no".to_string())); } @@ -3455,7 +3463,7 @@ async fn custom_handler_registration_and_execution() { }; engine.run(&graph, &run_options).await.expect("run"); - let cp = Checkpoint::load(&dir.path().join("checkpoint.json")).unwrap(); + let cp = load_checkpoint(&dir.path().join("checkpoint.json")).unwrap(); assert_eq!( cp.context_values.get("custom.ran"), Some(&serde_json::json!("true")) @@ -3527,7 +3535,7 @@ async fn integration_smoke_plan_implement_review_done() { assert_eq!(outcome.status, StageStatus::Success); // Verify all nodes completed - let cp = Checkpoint::load(&dir.path().join("checkpoint.json")).unwrap(); + let cp = load_checkpoint(&dir.path().join("checkpoint.json")).unwrap(); assert!(cp.completed_nodes.contains(&"plan".to_string())); assert!(cp.completed_nodes.contains(&"implement".to_string())); assert!(cp.completed_nodes.contains(&"review".to_string())); @@ -3630,7 +3638,7 @@ async fn manager_loop_runs_child_engine_e2e() { .expect("manager loop E2E should succeed"); assert_eq!(outcome.status, StageStatus::Success); - let checkpoint = Checkpoint::load(&dir.path().join("checkpoint.json")).unwrap(); + let checkpoint = load_checkpoint(&dir.path().join("checkpoint.json")).unwrap(); assert!( checkpoint .completed_nodes @@ -3761,7 +3769,7 @@ async fn manager_loop_context_flows_e2e() { assert_eq!(outcome.status, StageStatus::Success); // Check that child's context updates were propagated through the manager - let checkpoint = Checkpoint::load(&dir.path().join("checkpoint.json")).unwrap(); + let checkpoint = load_checkpoint(&dir.path().join("checkpoint.json")).unwrap(); let sup_outcome = checkpoint.node_outcomes.get("supervisor").unwrap(); assert_eq!( sup_outcome.context_updates.get("review.result"), @@ -3936,7 +3944,7 @@ async fn import_e2e_through_engine() { .expect("import E2E should succeed"); assert_eq!(outcome.status, StageStatus::Success); - let checkpoint = Checkpoint::load(&dir.path().join("checkpoint.json")).unwrap(); + let checkpoint = load_checkpoint(&dir.path().join("checkpoint.json")).unwrap(); assert!( checkpoint .completed_nodes @@ -4989,7 +4997,7 @@ async fn fidelity_stored_in_checkpoint_context() { }; engine.run(&graph, &run_options).await.expect("run"); - let cp = Checkpoint::load(&dir.path().join("checkpoint.json")).unwrap(); + let cp = load_checkpoint(&dir.path().join("checkpoint.json")).unwrap(); assert_eq!( cp.context_values.get("internal.fidelity"), Some(&serde_json::json!("summary:low")), @@ -5635,15 +5643,15 @@ async fn fidelity_checkpoint_roundtrip_preserves_fidelity() { // Load, save, load again to verify roundtrip let checkpoint_path = dir.path().join("checkpoint.json"); - let cp1 = Checkpoint::load(&checkpoint_path).expect("first load"); + let cp1 = load_checkpoint(&checkpoint_path).expect("first load"); assert_eq!( cp1.context_values.get("internal.fidelity"), Some(&serde_json::json!("summary:high")), ); let roundtrip_path = dir.path().join("checkpoint_roundtrip.json"); - cp1.save(&roundtrip_path).expect("save"); - let cp2 = Checkpoint::load(&roundtrip_path).expect("second load"); + save_checkpoint(&roundtrip_path, &cp1); + let cp2 = load_checkpoint(&roundtrip_path).expect("second load"); assert_eq!( cp2.context_values.get("internal.fidelity"), Some(&serde_json::json!("summary:high")), @@ -5799,7 +5807,7 @@ async fn fidelity_resume_preserves_context_values_across_checkpoint() { ); // Verify the final checkpoint still has the fidelity - let final_cp = Checkpoint::load(&dir.path().join("checkpoint.json")).unwrap(); + let final_cp = load_checkpoint(&dir.path().join("checkpoint.json")).unwrap(); assert_eq!( final_cp.context_values.get("internal.fidelity"), Some(&serde_json::json!("summary:low")), @@ -5841,7 +5849,6 @@ mod real_llm { _context: &Context, _thread_id: Option<&str>, _emitter: &Arc, - _stage_dir: &std::path::Path, _sandbox: &Arc, _tool_hooks: Option>, ) -> Result { @@ -5853,7 +5860,6 @@ mod real_llm { _node: &Node, prompt: &str, _system_prompt: Option<&str>, - _stage_dir: &std::path::Path, ) -> Result { self.complete(prompt).await } @@ -5938,7 +5944,7 @@ mod real_llm { }) } - use super::{local_env, test_run_id}; + use super::{load_checkpoint, local_env, test_run_id}; use fabro_graphviz::graph::{AttrValue, Edge, Graph}; use fabro_interview::AutoApproveInterviewer; use fabro_workflow::event::EventEmitter; @@ -5947,7 +5953,6 @@ mod real_llm { use fabro_workflow::handler::human::HumanHandler; use fabro_workflow::handler::start::StartHandler; use fabro_workflow::outcome::StageStatus; - use fabro_workflow::records::{Checkpoint, CheckpointExt}; use fabro_workflow::run_options::RunOptions; use fabro_workflow::test_support::WorkflowRunner; @@ -6036,7 +6041,7 @@ mod real_llm { assert_eq!(outcome.status, StageStatus::Success); - let checkpoint = Checkpoint::load(&dir.path().join("checkpoint.json")).unwrap(); + let checkpoint = load_checkpoint(&dir.path().join("checkpoint.json")).unwrap(); assert!(checkpoint.completed_nodes.contains(&"plan".to_string())); assert!(checkpoint.completed_nodes.contains(&"review".to_string())); @@ -6143,7 +6148,7 @@ mod real_llm { assert_eq!(outcome.status, StageStatus::Success); - let checkpoint = Checkpoint::load(&dir.path().join("checkpoint.json")).unwrap(); + let checkpoint = load_checkpoint(&dir.path().join("checkpoint.json")).unwrap(); let last_stage = checkpoint .context_values .get("last_stage") @@ -6275,7 +6280,7 @@ mod real_llm { assert_eq!(outcome.status, StageStatus::Success); - let checkpoint = Checkpoint::load(&dir.path().join("checkpoint.json")).unwrap(); + let checkpoint = load_checkpoint(&dir.path().join("checkpoint.json")).unwrap(); assert!( checkpoint.completed_nodes.contains(&"write".to_string()), "write should be completed" @@ -6466,7 +6471,7 @@ async fn human_gate_freeform_only_routes_text() { .expect("run should succeed"); assert_eq!(outcome.status, StageStatus::Success); - let checkpoint = Checkpoint::load(&dir.path().join("checkpoint.json")).unwrap(); + let checkpoint = load_checkpoint(&dir.path().join("checkpoint.json")).unwrap(); assert!( checkpoint .completed_nodes @@ -6595,7 +6600,7 @@ async fn human_gate_freeform_with_fixed_choice_match() { .expect("run should succeed"); assert_eq!(outcome.status, StageStatus::Success); - let checkpoint = Checkpoint::load(&dir.path().join("checkpoint.json")).unwrap(); + let checkpoint = load_checkpoint(&dir.path().join("checkpoint.json")).unwrap(); assert!( checkpoint.completed_nodes.contains(&"approve".to_string()), "fixed choice match should route to approve" @@ -6709,7 +6714,7 @@ async fn human_gate_freeform_fallback_on_unmatched_text() { .expect("run should succeed"); assert_eq!(outcome.status, StageStatus::Success); - let checkpoint = Checkpoint::load(&dir.path().join("checkpoint.json")).unwrap(); + let checkpoint = load_checkpoint(&dir.path().join("checkpoint.json")).unwrap(); assert!( checkpoint .completed_nodes @@ -8431,8 +8436,8 @@ async fn large_context_values_are_offloaded_to_artifact_store() { assert_eq!(outcome.status, StageStatus::Success); // The checkpoint context should contain an artifact pointer, not the full value - let checkpoint = fabro_workflow::records::Checkpoint::load(&dir.path().join("checkpoint.json")) - .expect("checkpoint should load"); + let checkpoint = + load_checkpoint(&dir.path().join("checkpoint.json")).expect("checkpoint should load"); let pointer_value = checkpoint .context_values .get("response.big_output") @@ -8650,7 +8655,7 @@ async fn artifact_pointers_rewritten_for_remote_sandbox() { // The checkpoint context should contain a pointer rewritten for the remote env let checkpoint = - Checkpoint::load(&dir.path().join("checkpoint.json")).expect("checkpoint should load"); + load_checkpoint(&dir.path().join("checkpoint.json")).expect("checkpoint should load"); let pointer_value = checkpoint .context_values .get("response.big_output") @@ -9049,7 +9054,6 @@ async fn cli_backend_run_writes_prompt_and_calls_exec() { &context, None, &emitter, - dir.path(), &env, None, ) @@ -9123,7 +9127,6 @@ async fn cli_backend_run_detects_changed_files() { &context, None, &emitter, - dir.path(), &env, None, ) @@ -9152,16 +9155,7 @@ async fn cli_backend_run_with_codex_provider() { let dir = tempfile::tempdir().unwrap(); let result = backend - .run( - &node, - "Build the API", - &context, - None, - &emitter, - dir.path(), - &env, - None, - ) + .run(&node, "Build the API", &context, None, &emitter, &env, None) .await .expect("CLI backend should succeed"); @@ -9326,7 +9320,6 @@ async fn cli_backend_run_fails_on_nonzero_exit() { &context, None, &emitter, - dir.path(), &failing_env, None, ) @@ -9359,16 +9352,7 @@ async fn cli_backend_run_fails_on_unparseable_output() { let dir = tempfile::tempdir().unwrap(); let result = backend - .run( - &node, - "do something", - &context, - None, - &emitter, - dir.path(), - &env, - None, - ) + .run(&node, "do something", &context, None, &emitter, &env, None) .await; let err = match result { @@ -9402,16 +9386,7 @@ async fn cli_backend_run_uses_node_model_override() { let dir = tempfile::tempdir().unwrap(); backend - .run( - &node, - "test", - &context, - None, - &emitter, - dir.path(), - &env, - None, - ) + .run(&node, "test", &context, None, &emitter, &env, None) .await .expect("should succeed"); @@ -9453,16 +9428,7 @@ async fn cli_backend_run_uses_node_provider_override() { let dir = tempfile::tempdir().unwrap(); backend - .run( - &node, - "test", - &context, - None, - &emitter, - dir.path(), - &env, - None, - ) + .run(&node, "test", &context, None, &emitter, &env, None) .await .expect("should succeed"); @@ -9488,16 +9454,7 @@ async fn cli_backend_run_returns_text_and_usage() { let dir = tempfile::tempdir().unwrap(); let result = backend - .run( - &node, - "test", - &context, - None, - &emitter, - dir.path(), - &env, - None, - ) + .run(&node, "test", &context, None, &emitter, &env, None) .await .expect("should succeed"); @@ -9538,16 +9495,7 @@ async fn backend_router_delegates_to_cli_for_cli_node() { let dir = tempfile::tempdir().unwrap(); let result = router - .run( - &node, - "Fix the bug", - &context, - None, - &emitter, - dir.path(), - &env, - None, - ) + .run(&node, "Fix the bug", &context, None, &emitter, &env, None) .await .expect("router should succeed"); @@ -9582,16 +9530,7 @@ async fn backend_router_delegates_to_api_for_normal_node() { let dir = tempfile::tempdir().unwrap(); let result = router - .run( - &node, - "Plan the work", - &context, - None, - &emitter, - dir.path(), - &env, - None, - ) + .run(&node, "Plan the work", &context, None, &emitter, &env, None) .await .expect("router should succeed"); @@ -9629,16 +9568,7 @@ async fn backend_router_delegates_to_cli_for_backend_attr() { let dir = tempfile::tempdir().unwrap(); let result = router - .run( - &node, - "Build it", - &context, - None, - &emitter, - dir.path(), - &env, - None, - ) + .run(&node, "Build it", &context, None, &emitter, &env, None) .await .expect("router should succeed"); @@ -10115,7 +10045,7 @@ async fn git_checkpoint_host_emits_events_and_diff_patch() { // 8. Verify checkpoint.json has git_commit_sha let checkpoint = - Checkpoint::load(&run_dir.path().join("checkpoint.json")).expect("checkpoint should load"); + load_checkpoint(&run_dir.path().join("checkpoint.json")).expect("checkpoint should load"); assert!( checkpoint.git_commit_sha.is_some(), "checkpoint should have git_commit_sha" @@ -10460,7 +10390,7 @@ async fn parallel_git_branching_host_e2e() { // 6. Verify parallel.results has head_sha for each branch let checkpoint = - Checkpoint::load(&run_dir.path().join("checkpoint.json")).expect("checkpoint should load"); + load_checkpoint(&run_dir.path().join("checkpoint.json")).expect("checkpoint should load"); let parallel_results = checkpoint .context_values .get("parallel.results") @@ -11323,7 +11253,7 @@ async fn e2e_failure_signature_persisted_in_context() { assert_eq!(outcome.status, StageStatus::Success); // Verify checkpoint has failure_signature in context - let cp = Checkpoint::load(&dir.path().join("checkpoint.json")).unwrap(); + let cp = load_checkpoint(&dir.path().join("checkpoint.json")).unwrap(); let sig_value = cp .context_values .get("failure_signature") @@ -11382,7 +11312,7 @@ async fn e2e_failure_signature_hint_overrides_reason_in_context() { }; let _outcome = engine.run(&graph, &run_options).await.unwrap(); - let cp = Checkpoint::load(&dir.path().join("checkpoint.json")).unwrap(); + let cp = load_checkpoint(&dir.path().join("checkpoint.json")).unwrap(); let sig_str = cp .context_values .get("failure_signature") @@ -11439,7 +11369,7 @@ async fn e2e_signature_maps_persist_in_checkpoint() { assert_eq!(outcome.status, StageStatus::Success); // Load checkpoint and verify signature maps - let cp = Checkpoint::load(&dir.path().join("checkpoint.json")).unwrap(); + let cp = load_checkpoint(&dir.path().join("checkpoint.json")).unwrap(); // The pipeline had 3 deterministic failures at "work" before succeeding. // loop_failure_signatures should have recorded them. assert!( @@ -11520,9 +11450,9 @@ fn e2e_checkpoint_signatures_roundtrip() { restart_sigs, std::collections::HashMap::new(), ); - cp.save(&path).unwrap(); + save_checkpoint(&path, &cp); - let loaded = Checkpoint::load(&path).unwrap(); + let loaded = load_checkpoint(&path).unwrap(); assert_eq!(loaded.loop_failure_signatures.len(), 1); assert_eq!(loaded.restart_failure_signatures.len(), 1); assert_eq!(loaded.loop_failure_signatures.get(&sig1), Some(&2)); @@ -11633,7 +11563,7 @@ async fn e2e_circuit_breaker_does_not_fire_below_limit() { ); // Verify signatures were tracked but didn't trigger abort - let cp = Checkpoint::load(&dir.path().join("checkpoint.json")).unwrap(); + let cp = load_checkpoint(&dir.path().join("checkpoint.json")).unwrap(); let total_failures: usize = cp.loop_failure_signatures.values().sum(); assert_eq!( total_failures, 4,