refactor(run): remove legacy checkpoint and stage dir I/O

This commit is contained in:
Bryan Helmkamp 2026-04-03 15:23:45 -07:00
parent 336a0eece6
commit 4d7c9c55a7
No known key found for this signature in database
15 changed files with 126 additions and 346 deletions

View file

@ -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}"))
}
}

View file

@ -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<AssetCollectionSummary, String> {
@ -330,7 +333,7 @@ pub async fn collect_assets(
let mut captured_assets: Vec<CapturedAssetInfo> = 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());
}

View file

@ -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<EventEmitter>,
stage_dir: &Path,
sandbox: &Arc<dyn Sandbox>,
tool_hooks: Option<Arc<dyn fabro_agent::ToolHookCallback>>,
) -> Result<CodergenResult, FabroError>;
@ -54,7 +51,6 @@ pub trait CodergenBackend: Send + Sync {
_node: &Node,
_prompt: &str,
_system_prompt: Option<&str>,
_stage_dir: &Path,
) -> Result<CodergenResult, FabroError> {
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<Outcome, FabroError> {
// 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<EventEmitter>,
_stage_dir: &Path,
_sandbox: &Arc<dyn fabro_agent::Sandbox>,
_tool_hooks: Option<Arc<dyn fabro_agent::ToolHookCallback>>,
) -> Result<CodergenResult, FabroError> {
@ -620,7 +610,6 @@ mod tests {
_context: &Context,
_thread_id: Option<&str>,
_emitter: &Arc<EventEmitter>,
_stage_dir: &Path,
_sandbox: &Arc<dyn fabro_agent::Sandbox>,
_tool_hooks: Option<Arc<dyn fabro_agent::ToolHookCallback>>,
) -> Result<CodergenResult, FabroError> {
@ -679,7 +668,6 @@ mod tests {
context: &Context,
_thread_id: Option<&str>,
emitter: &Arc<EventEmitter>,
_stage_dir: &Path,
_sandbox: &Arc<dyn fabro_agent::Sandbox>,
_tool_hooks: Option<Arc<dyn fabro_agent::ToolHookCallback>>,
) -> Result<CodergenResult, FabroError> {
@ -791,7 +779,6 @@ mod tests {
_context: &Context,
thread_id: Option<&str>,
_emitter: &Arc<EventEmitter>,
_stage_dir: &Path,
_sandbox: &Arc<dyn Sandbox>,
_tool_hooks: Option<Arc<dyn fabro_agent::ToolHookCallback>>,
) -> Result<CodergenResult, FabroError> {
@ -844,7 +831,6 @@ mod tests {
_context: &Context,
thread_id: Option<&str>,
_emitter: &Arc<EventEmitter>,
_stage_dir: &Path,
_sandbox: &Arc<dyn Sandbox>,
_tool_hooks: Option<Arc<dyn fabro_agent::ToolHookCallback>>,
) -> Result<CodergenResult, FabroError> {
@ -892,7 +878,6 @@ mod tests {
_context: &Context,
_thread_id: Option<&str>,
_emitter: &Arc<EventEmitter>,
_stage_dir: &Path,
_sandbox: &Arc<dyn Sandbox>,
_tool_hooks: Option<Arc<dyn fabro_agent::ToolHookCallback>>,
) -> Result<CodergenResult, FabroError> {
@ -1036,7 +1021,6 @@ Some text in between.
_context: &Context,
_thread_id: Option<&str>,
_emitter: &Arc<EventEmitter>,
_stage_dir: &Path,
_sandbox: &Arc<dyn Sandbox>,
_tool_hooks: Option<Arc<dyn fabro_agent::ToolHookCallback>>,
) -> Result<CodergenResult, FabroError> {
@ -1075,7 +1059,6 @@ Some text in between.
_context: &Context,
_thread_id: Option<&str>,
_emitter: &Arc<EventEmitter>,
_stage_dir: &std::path::Path,
_sandbox: &Arc<dyn Sandbox>,
_tool_hooks: Option<Arc<dyn fabro_agent::ToolHookCallback>>,
) -> Result<CodergenResult, FabroError> {
@ -1145,7 +1128,6 @@ Some text in between.
_context: &Context,
_thread_id: Option<&str>,
_emitter: &Arc<EventEmitter>,
_stage_dir: &std::path::Path,
_sandbox: &Arc<dyn Sandbox>,
_tool_hooks: Option<Arc<dyn fabro_agent::ToolHookCallback>>,
) -> Result<CodergenResult, FabroError> {

View file

@ -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<EventEmitter>,
sandbox: &Arc<dyn Sandbox>,
@ -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<EventEmitter>,
_stage_dir: &std::path::Path,
_sandbox: &Arc<dyn Sandbox>,
_tool_hooks: Option<Arc<dyn fabro_agent::ToolHookCallback>>,
) -> Result<CodergenResult, FabroError> {
@ -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]

View file

@ -268,7 +268,6 @@ impl CodergenBackend for AgentApiBackend {
node: &Node,
prompt: &str,
system_prompt: Option<&str>,
_stage_dir: &std::path::Path,
) -> Result<CodergenResult, FabroError> {
let client = Client::from_env()
.await
@ -412,7 +411,6 @@ impl CodergenBackend for AgentApiBackend {
context: &Context,
thread_id: Option<&str>,
emitter: &Arc<EventEmitter>,
_stage_dir: &std::path::Path,
sandbox: &Arc<dyn Sandbox>,
tool_hooks: Option<Arc<dyn fabro_agent::ToolHookCallback>>,
) -> Result<CodergenResult, FabroError> {

View file

@ -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<EventEmitter>,
_stage_dir: &Path,
sandbox: &Arc<dyn Sandbox>,
_tool_hooks: Option<Arc<dyn fabro_agent::ToolHookCallback>>,
) -> Result<CodergenResult, FabroError> {
@ -755,20 +753,19 @@ impl CodergenBackend for BackendRouter {
context: &Context,
thread_id: Option<&str>,
emitter: &Arc<EventEmitter>,
stage_dir: &Path,
sandbox: &Arc<dyn Sandbox>,
tool_hooks: Option<Arc<dyn fabro_agent::ToolHookCallback>>,
) -> Result<CodergenResult, FabroError> {
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<CodergenResult, FabroError> {
// 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<EventEmitter>,
_stage_dir: &Path,
_sandbox: &Arc<dyn Sandbox>,
_tool_hooks: Option<Arc<dyn fabro_agent::ToolHookCallback>>,
) -> Result<CodergenResult, FabroError> {

View file

@ -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<Outcome, FabroError> {
// 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<crate::event::EventEmitter>,
_stage_dir: &Path,
_sandbox: &Arc<dyn Sandbox>,
_tool_hooks: Option<Arc<dyn fabro_agent::ToolHookCallback>>,
) -> Result<CodergenResult, FabroError> {
@ -280,7 +274,6 @@ mod tests {
_node: &Node,
_prompt: &str,
_system_prompt: Option<&str>,
_stage_dir: &Path,
) -> Result<CodergenResult, FabroError> {
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<crate::event::EventEmitter>,
_stage_dir: &Path,
_sandbox: &Arc<dyn Sandbox>,
_tool_hooks: Option<Arc<dyn fabro_agent::ToolHookCallback>>,
) -> Result<CodergenResult, FabroError> {
@ -344,7 +334,6 @@ mod tests {
_node: &Node,
_prompt: &str,
_system_prompt: Option<&str>,
_stage_dir: &Path,
) -> Result<CodergenResult, FabroError> {
Ok(CodergenResult::Text {
text: "one-shot response".to_string(),
@ -396,7 +385,6 @@ mod tests {
_context: &Context,
_thread_id: Option<&str>,
_emitter: &Arc<crate::event::EventEmitter>,
_stage_dir: &Path,
_sandbox: &Arc<dyn fabro_agent::Sandbox>,
_tool_hooks: Option<Arc<dyn fabro_agent::ToolHookCallback>>,
) -> Result<CodergenResult, FabroError> {
@ -408,7 +396,6 @@ mod tests {
_node: &Node,
prompt: &str,
system_prompt: Option<&str>,
_stage_dir: &Path,
) -> Result<CodergenResult, FabroError> {
*self.captured_prompt.lock().unwrap() = Some(prompt.to_string());
*self.captured_system_prompt.lock().unwrap() = Some(system_prompt.map(String::from));

View file

@ -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<Arc<dyn Fn(&str) + Send + Sync>>;
@ -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<T: DeserializeOwned>(
path: &std::path::Path,
label: &str,
) -> error::Result<T> {
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<CompletedStage>` 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<CompletedStage> {

View file

@ -95,13 +95,13 @@ impl RunLifecycle<WorkflowGraph> 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 {

View file

@ -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(),

View file

@ -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();

View file

@ -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<FailureSignature, usize>,
node_visits: HashMap<String, usize>,
) -> Self;
fn save(&self, path: &Path) -> CrateResult<()>;
fn load(path: &Path) -> CrateResult<Self>
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<Self> {
crate::load_json(path, "checkpoint")
}
}

View file

@ -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")
);
}
}

View file

@ -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<Checkpoint, Box<dyn std::error::Error>> {
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,
)

View file

@ -62,6 +62,15 @@ fn test_run_id(label: &str) -> RunId {
RunId::from(Ulid(u128::from(hasher.finish())))
}
fn load_checkpoint(path: &Path) -> Result<Checkpoint, Box<dyn std::error::Error>> {
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<EventEmitter>,
_stage_dir: &std::path::Path,
_sandbox: &Arc<dyn fabro_agent::Sandbox>,
_tool_hooks: Option<Arc<dyn fabro_agent::ToolHookCallback>>,
) -> Result<CodergenResult, FabroError> {
@ -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<EventEmitter>,
_stage_dir: &std::path::Path,
_sandbox: &Arc<dyn fabro_agent::Sandbox>,
_tool_hooks: Option<Arc<dyn fabro_agent::ToolHookCallback>>,
) -> Result<CodergenResult, FabroError> {
@ -5853,7 +5860,6 @@ mod real_llm {
_node: &Node,
prompt: &str,
_system_prompt: Option<&str>,
_stage_dir: &std::path::Path,
) -> Result<CodergenResult, FabroError> {
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,