diff --git a/lib/crates/fabro-store/src/slate/run_store.rs b/lib/crates/fabro-store/src/slate/run_store.rs index 4b5df4e6b..7e2ea363b 100644 --- a/lib/crates/fabro-store/src/slate/run_store.rs +++ b/lib/crates/fabro-store/src/slate/run_store.rs @@ -89,7 +89,7 @@ impl SlateRunStore { self.inner.record.clone() } - pub(crate) fn created_at(&self) -> DateTime { + pub fn created_at(&self) -> DateTime { self.inner.record.created_at } diff --git a/lib/crates/fabro-workflow/src/event.rs b/lib/crates/fabro-workflow/src/event.rs index 778698f0a..1a6333d9d 100644 --- a/lib/crates/fabro-workflow/src/event.rs +++ b/lib/crates/fabro-workflow/src/event.rs @@ -189,6 +189,8 @@ pub enum WorkflowRunEvent { delay_ms: u64, }, ParallelStarted { + node_id: String, + visit: u32, branch_count: usize, join_policy: String, }, @@ -205,6 +207,8 @@ pub enum WorkflowRunEvent { head_sha: Option, }, ParallelCompleted { + node_id: String, + visit: u32, duration_ms: u64, success_count: usize, failure_count: usize, @@ -680,6 +684,7 @@ impl WorkflowRunEvent { Self::ParallelStarted { branch_count, join_policy, + .. } => { debug!(branch_count, join_policy, "Parallel execution started"); } @@ -703,6 +708,7 @@ impl WorkflowRunEvent { success_count, failure_count, results, + .. } => { debug!( duration_ms, @@ -1302,6 +1308,8 @@ fn extract_envelope_fields(event: &WorkflowRunEvent) -> EnvelopeFields { | WorkflowRunEvent::SubgraphCompleted { .. } | WorkflowRunEvent::AssetCaptured { .. } | WorkflowRunEvent::PromptCompleted { .. } + | WorkflowRunEvent::ParallelStarted { .. } + | WorkflowRunEvent::ParallelCompleted { .. } | WorkflowRunEvent::CommandStarted { .. } | WorkflowRunEvent::CommandCompleted { .. } | WorkflowRunEvent::AgentCliStarted { .. } diff --git a/lib/crates/fabro-workflow/src/handler/agent.rs b/lib/crates/fabro-workflow/src/handler/agent.rs index 3d6924d08..66ea4b667 100644 --- a/lib/crates/fabro-workflow/src/handler/agent.rs +++ b/lib/crates/fabro-workflow/src/handler/agent.rs @@ -4,7 +4,10 @@ use std::sync::Arc; use async_trait::async_trait; use fabro_agent::Sandbox; +use fabro_model::Provider; +use fabro_store::NodeVisitRef; use fabro_types::RunId; +use tokio::fs; use crate::context::keys; use crate::context::{Context, WorkflowContext}; @@ -13,6 +16,7 @@ use crate::event::EventEmitter; use crate::outcome::{ FailureCategory, FailureDetail, Outcome, OutcomeExt, StageStatus, StageUsage, }; +use crate::run_dir::visit_from_context; use crate::transforms::variable_expansion::expand_vars; use fabro_graphviz::graph::{Graph, Node}; @@ -195,6 +199,91 @@ pub(crate) fn truncate(s: &str, max_chars: usize) -> &str { } } +pub(crate) fn stage_dir(run_dir: &Path, node_id: &str, visit: u32) -> std::path::PathBuf { + let node_dir = if visit <= 1 { + node_id.to_string() + } else { + format!("{node_id}-visit_{visit}") + }; + run_dir.join("nodes").join(node_dir) +} + +pub(crate) fn status_json_value(outcome: &Outcome) -> serde_json::Value { + let mut status = serde_json::Map::new(); + status.insert( + "status".to_string(), + serde_json::Value::String(outcome.status.to_string()), + ); + status.insert( + "outcome".to_string(), + serde_json::Value::String(outcome.status.to_string()), + ); + if let Some(label) = outcome.preferred_label.as_ref() { + status.insert( + "preferred_next_label".to_string(), + serde_json::Value::String(label.clone()), + ); + } + if !outcome.suggested_next_ids.is_empty() { + status.insert( + "suggested_next_ids".to_string(), + serde_json::json!(outcome.suggested_next_ids), + ); + } + if let Some(failure) = outcome.failure.as_ref() { + status.insert( + "failure_reason".to_string(), + serde_json::Value::String(failure.message.clone()), + ); + } + if !outcome.context_updates.is_empty() { + status.insert( + "context_updates".to_string(), + serde_json::Value::Object( + outcome + .context_updates + .iter() + .map(|(key, value)| (key.clone(), value.clone())) + .collect(), + ), + ); + } + serde_json::Value::Object(status) +} + +pub(crate) async fn write_provider_used_file( + services: &EngineServices, + stage_dir: &Path, + node_id: &str, + visit: u32, + fallback: Option, +) -> Result<(), FabroError> { + let provider_used = services + .run_store + .state() + .await + .ok() + .and_then(|state| { + let node = NodeVisitRef { node_id, visit }; + state + .node(&node) + .and_then(|node_state| node_state.provider_used.clone()) + }) + .or(fallback); + let Some(provider_used) = provider_used else { + return Ok(()); + }; + fs::write( + stage_dir.join("provider_used.json"), + serde_json::to_vec_pretty(&provider_used).map_err(|err| { + FabroError::handler(format!("failed to serialize provider_used.json: {err}")) + })?, + ) + .await + .map_err(|err| FabroError::handler(format!("failed to write provider_used.json: {err}")))?; + Ok(()) +} + /// Shared simulate implementation for LLM-backed handlers (agent & prompt). /// Produces a simulated outcome with standard context updates. pub(crate) fn simulate_llm_handler(node: &Node) -> Outcome { @@ -248,6 +337,15 @@ impl Handler for AgentHandler { format!("{preamble}\n\n{expanded}") }; + let visit = u32::try_from(visit_from_context(context)).unwrap_or(u32::MAX); + let stage_dir = stage_dir(_run_dir, &node.id, visit); + fs::create_dir_all(&stage_dir) + .await + .map_err(|err| FabroError::handler(format!("failed to create stage dir: {err}")))?; + fs::write(stage_dir.join("prompt.md"), &prompt) + .await + .map_err(|err| FabroError::handler(format!("failed to write prompt file: {err}")))?; + // 3. Call LLM backend (agent loop) let thread_id = context.thread_id(); let run_id = context @@ -350,6 +448,31 @@ impl Handler for AgentHandler { } outcome.usage = stage_usage; outcome.files_touched = backend_files_touched; + fs::write(stage_dir.join("response.md"), &response_text) + .await + .map_err(|err| FabroError::handler(format!("failed to write response file: {err}")))?; + fs::write( + stage_dir.join("status.json"), + serde_json::to_vec_pretty(&status_json_value(&outcome)) + .map_err(|err| FabroError::handler(format!("failed to serialize status: {err}")))?, + ) + .await + .map_err(|err| FabroError::handler(format!("failed to write status file: {err}")))?; + write_provider_used_file( + services, + &stage_dir, + &node.id, + visit, + Some(serde_json::json!({ + "mode": if node.backend() == Some("cli") { "cli" } else { "agent" }, + "provider": node + .provider() + .map(String::from) + .unwrap_or_else(|| Provider::default_from_env().as_str().to_string()), + "model": node.model().map(String::from).unwrap_or_default(), + })), + ) + .await?; Ok(outcome) } @@ -390,6 +513,7 @@ mod tests { .await .unwrap(); let services = EngineServices { + emitter: Arc::new(crate::event::EventEmitter::new(fixtures::RUN_1)), run_store: run_store.clone(), ..EngineServices::test_default() }; diff --git a/lib/crates/fabro-workflow/src/handler/command.rs b/lib/crates/fabro-workflow/src/handler/command.rs index 89bfd6d81..2383f59ad 100644 --- a/lib/crates/fabro-workflow/src/handler/command.rs +++ b/lib/crates/fabro-workflow/src/handler/command.rs @@ -7,6 +7,7 @@ use crate::event::WorkflowRunEvent; use crate::outcome::{Outcome, OutcomeExt}; use async_trait::async_trait; use fabro_graphviz::graph::{Graph, Node}; +use tokio::fs; use super::{EngineServices, Handler}; @@ -59,7 +60,7 @@ impl Handler for CommandHandler { node: &Node, _context: &Context, _graph: &Graph, - _run_dir: &Path, + run_dir: &Path, services: &EngineServices, ) -> Result { let script = node @@ -97,6 +98,22 @@ impl Handler for CommandHandler { } else { script.to_string() }; + let stage_dir = run_dir.join("nodes").join(&node.id); + fs::create_dir_all(&stage_dir) + .await + .map_err(|err| FabroError::handler(format!("failed to create stage dir: {err}")))?; + fs::write( + stage_dir.join("script_invocation.json"), + serde_json::to_vec_pretty(&serde_json::json!({ + "script": script, + "command": command, + "language": language, + "timeout_ms": timeout_ms(node), + })) + .map_err(|err| FabroError::handler(format!("failed to serialize invocation: {err}")))?, + ) + .await + .map_err(|err| FabroError::handler(format!("failed to write invocation file: {err}")))?; let timeout_ms = node .timeout() @@ -121,6 +138,23 @@ impl Handler for CommandHandler { duration_ms: result.duration_ms, timed_out: result.timed_out, }); + fs::write(stage_dir.join("stdout.log"), &result.stdout) + .await + .map_err(|err| FabroError::handler(format!("failed to write stdout log: {err}")))?; + fs::write(stage_dir.join("stderr.log"), &result.stderr) + .await + .map_err(|err| FabroError::handler(format!("failed to write stderr log: {err}")))?; + fs::write( + stage_dir.join("script_timing.json"), + serde_json::to_vec_pretty(&serde_json::json!({ + "duration_ms": result.duration_ms, + "exit_code": (!result.timed_out).then_some(result.exit_code), + "timed_out": result.timed_out, + })) + .map_err(|err| FabroError::handler(format!("failed to serialize timing: {err}")))?, + ) + .await + .map_err(|err| FabroError::handler(format!("failed to write timing file: {err}")))?; if result.timed_out { return Err(FabroError::handler(format!( diff --git a/lib/crates/fabro-workflow/src/handler/mod.rs b/lib/crates/fabro-workflow/src/handler/mod.rs index 2e77eef0e..7b2bafdf9 100644 --- a/lib/crates/fabro-workflow/src/handler/mod.rs +++ b/lib/crates/fabro-workflow/src/handler/mod.rs @@ -86,12 +86,22 @@ impl EngineServices { sandbox: Arc::new(fabro_agent::LocalSandbox::new( std::env::current_dir().unwrap_or_else(|_| std::path::PathBuf::from(".")), )), - run_store: futures::executor::block_on(async { - store - .create_run(&fabro_types::RunId::new(), chrono::Utc::now(), None) - .await - .expect("slate-backed test run store should initialize") - }), + // Build the test run store on a dedicated runtime so this helper + // remains safe to call from both sync tests and #[tokio::test]. + run_store: std::thread::spawn(move || { + tokio::runtime::Builder::new_current_thread() + .enable_all() + .build() + .expect("test runtime should initialize") + .block_on(async { + store + .create_run(&fabro_types::RunId::new(), chrono::Utc::now(), None) + .await + .expect("slate-backed test run store should initialize") + }) + }) + .join() + .expect("test run store thread should join"), git_state: std::sync::RwLock::new(None), hook_runner: None, env: HashMap::new(), diff --git a/lib/crates/fabro-workflow/src/handler/parallel.rs b/lib/crates/fabro-workflow/src/handler/parallel.rs index 31004d2ef..9854428bc 100644 --- a/lib/crates/fabro-workflow/src/handler/parallel.rs +++ b/lib/crates/fabro-workflow/src/handler/parallel.rs @@ -5,6 +5,7 @@ use std::time::Instant; use async_trait::async_trait; use fabro_agent::{Sandbox, WorktreeOptions, WorktreeSandbox}; use fabro_types::RunId; +use tokio::fs; use tokio::sync::Semaphore; use crate::context::keys; @@ -151,6 +152,8 @@ impl Handler for ParallelHandler { ); services.emitter.emit(&WorkflowRunEvent::ParallelStarted { + node_id: node.id.clone(), + visit: u32::try_from(visit_from_context(context)).unwrap_or(u32::MAX), branch_count: branches.len(), join_policy: join_policy.to_string(), }); @@ -473,10 +476,26 @@ impl Handler for ParallelHandler { entry }) .collect(); + let stage_dir = run_dir.join("nodes").join(&node.id); + fs::create_dir_all(&stage_dir) + .await + .map_err(|err| FabroError::handler(format!("failed to create stage dir: {err}")))?; + fs::write( + stage_dir.join("parallel_results.json"), + serde_json::to_vec_pretty(&results_json).map_err(|err| { + FabroError::handler(format!("failed to serialize parallel results: {err}")) + })?, + ) + .await + .map_err(|err| { + FabroError::handler(format!("failed to write parallel results file: {err}")) + })?; context.set(keys::PARALLEL_RESULTS, serde_json::json!(results_json)); context.set(keys::PARALLEL_BRANCH_COUNT, serde_json::json!(total)); services.emitter.emit(&WorkflowRunEvent::ParallelCompleted { + node_id: node.id.clone(), + visit: u32::try_from(visit_from_context(context)).unwrap_or(u32::MAX), duration_ms: millis_u64(parallel_start.elapsed()), success_count, failure_count: fail_count, @@ -678,9 +697,12 @@ mod tests { .await .unwrap(); let services = EngineServices { + emitter: Arc::new(crate::event::EventEmitter::new(fixtures::RUN_1)), run_store: run_store.clone(), ..EngineServices::test_default() }; + let logger = crate::event::StoreProgressLogger::new(run_store.clone()); + logger.register(services.emitter.as_ref()); let mut node = Node::new("par"); node.attrs.insert( "shape".to_string(), @@ -703,6 +725,7 @@ mod tests { .execute(&node, &context, &graph, tmp.path(), &services) .await .unwrap(); + logger.flush().await; let state = run_store.state().await.unwrap(); let node_state = state diff --git a/lib/crates/fabro-workflow/src/handler/prompt.rs b/lib/crates/fabro-workflow/src/handler/prompt.rs index d5c1b344d..54f8e1e89 100644 --- a/lib/crates/fabro-workflow/src/handler/prompt.rs +++ b/lib/crates/fabro-workflow/src/handler/prompt.rs @@ -1,6 +1,7 @@ use std::path::Path; use async_trait::async_trait; +use tokio::fs; use crate::context::keys; use crate::context::{Context, WorkflowContext}; @@ -12,7 +13,8 @@ use fabro_graphviz::graph::{Graph, Node}; use fabro_model::Provider; use super::agent::{ - CodergenBackend, CodergenResult, expand_variables, extract_status_fields, truncate, + CodergenBackend, CodergenResult, expand_variables, extract_status_fields, stage_dir, + status_json_value, truncate, write_provider_used_file, }; use super::{EngineServices, Handler}; @@ -46,7 +48,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) @@ -61,6 +63,14 @@ impl Handler for PromptHandler { } else { format!("{preamble}\n\n{expanded}") }; + let visit = u32::try_from(visit_from_context(context)).unwrap_or(u32::MAX); + let stage_dir = stage_dir(run_dir, &node.id, visit); + fs::create_dir_all(&stage_dir) + .await + .map_err(|err| FabroError::handler(format!("failed to create stage dir: {err}")))?; + fs::write(stage_dir.join("prompt.md"), &prompt) + .await + .map_err(|err| FabroError::handler(format!("failed to write prompt file: {err}")))?; // 1b. Discover project docs for system prompt when project_memory is enabled let system_prompt = if node.project_memory() { @@ -166,6 +176,28 @@ impl Handler for PromptHandler { extract_status_fields(&response_text, &mut outcome); outcome.usage = stage_usage; outcome.files_touched = backend_files_touched; + fs::write(stage_dir.join("response.md"), &response_text) + .await + .map_err(|err| FabroError::handler(format!("failed to write response file: {err}")))?; + fs::write( + stage_dir.join("status.json"), + serde_json::to_vec_pretty(&status_json_value(&outcome)) + .map_err(|err| FabroError::handler(format!("failed to serialize status: {err}")))?, + ) + .await + .map_err(|err| FabroError::handler(format!("failed to write status file: {err}")))?; + write_provider_used_file( + services, + &stage_dir, + &node.id, + visit, + Some(serde_json::json!({ + "mode": "prompt", + "provider": prompt_provider.unwrap_or_default(), + "model": prompt_model.unwrap_or_default(), + })), + ) + .await?; Ok(outcome) } @@ -205,6 +237,7 @@ mod tests { .await .unwrap(); let services = EngineServices { + emitter: Arc::new(crate::event::EventEmitter::new(fixtures::RUN_1)), run_store: run_store.clone(), ..EngineServices::test_default() }; diff --git a/lib/crates/fabro-workflow/src/lifecycle/disk.rs b/lib/crates/fabro-workflow/src/lifecycle/disk.rs index 6f6eee3d5..a027eeee8 100644 --- a/lib/crates/fabro-workflow/src/lifecycle/disk.rs +++ b/lib/crates/fabro-workflow/src/lifecycle/disk.rs @@ -1,10 +1,12 @@ use std::path::PathBuf; -use std::sync::Arc; +use std::sync::{Arc, Mutex}; use async_trait::async_trait; use fabro_store::SlateRunStore; use fabro_types::RunId; +use tokio::fs; +use fabro_core::error::CoreError; use fabro_core::error::Result as CoreResult; use fabro_core::graph::NodeSpec; use fabro_core::lifecycle::RunLifecycle; @@ -12,10 +14,11 @@ use fabro_core::outcome::NodeResult; use fabro_core::state::ExecutionState; use super::circuit_breaker::CircuitBreakerLifecycle; +use super::git::GitCheckpointResult; use crate::event::{EventEmitter, RunNoticeLevel, WorkflowRunEvent, append_workflow_event}; use crate::graph::WorkflowGraph; use crate::graph::WorkflowNode; -use crate::outcome::StageUsage; +use crate::outcome::{OutcomeExt, StageUsage}; use crate::run_options::RunOptions; use fabro_graphviz::graph::types::Graph as GvGraph; @@ -30,10 +33,38 @@ pub(crate) struct DiskLifecycle { pub graph: Arc, pub run_options: Arc, pub emitter: Arc, + pub checkpoint_git_result: Arc>>, pub circuit_breaker: Arc, pub checkpoint_enabled: bool, } +pub(super) fn build_checkpoint( + node: &WorkflowNode, + result: &WfNodeResult, + next_node_id: Option<&str>, + state: &WfRunState, + loop_failure_signatures: std::collections::HashMap, + restart_failure_signatures: std::collections::HashMap, + git_commit_sha: Option, +) -> fabro_types::Checkpoint { + let mut node_outcomes = state.node_outcomes.clone(); + node_outcomes.insert(node.id().to_string(), result.outcome.clone()); + + fabro_types::Checkpoint { + timestamp: chrono::Utc::now(), + current_node: node.id().to_string(), + completed_nodes: state.completed_nodes.clone(), + node_outcomes, + node_retries: state.node_retries.clone(), + context_values: state.context.snapshot(), + next_node_id: next_node_id.map(String::from), + git_commit_sha, + node_visits: state.node_visits.clone(), + loop_failure_signatures, + restart_failure_signatures, + } +} + #[async_trait] impl RunLifecycle for DiskLifecycle { async fn on_run_start(&self, _graph: &WorkflowGraph, _state: &WfRunState) -> CoreResult<()> { @@ -55,10 +86,45 @@ impl RunLifecycle for DiskLifecycle { async fn after_node( &self, - _node: &WorkflowNode, - _result: &mut WfNodeResult, - _state: &WfRunState, + node: &WorkflowNode, + result: &mut WfNodeResult, + state: &WfRunState, ) -> CoreResult<()> { + let visit = + u32::try_from(*state.node_visits.get(node.id()).unwrap_or(&1)).unwrap_or(u32::MAX); + let node_dir = if visit <= 1 { + self.run_dir.join("nodes").join(node.id()) + } else { + self.run_dir + .join("nodes") + .join(format!("{}-visit_{visit}", node.id())) + }; + fs::create_dir_all(&node_dir) + .await + .map_err(|err| CoreError::Other(format!("failed to create node dir: {err}")))?; + let mut status = serde_json::Map::new(); + status.insert( + "status".to_string(), + serde_json::Value::String(result.outcome.status.to_string()), + ); + status.insert( + "outcome".to_string(), + serde_json::Value::String(result.outcome.status.to_string()), + ); + if let Some(failure_reason) = result.outcome.failure_reason() { + status.insert( + "failure_reason".to_string(), + serde_json::Value::String(failure_reason.to_string()), + ); + } + fs::write( + node_dir.join("status.json"), + serde_json::to_vec_pretty(&serde_json::Value::Object(status)).map_err(|err| { + CoreError::Other(format!("failed to serialize status.json: {err}")) + })?, + ) + .await + .map_err(|err| CoreError::Other(format!("failed to write status.json: {err}")))?; Ok(()) } @@ -73,25 +139,27 @@ impl RunLifecycle for DiskLifecycle { return Ok(()); } + let git_commit_sha = self + .checkpoint_git_result + .lock() + .unwrap() + .as_ref() + .and_then(|result| result.commit_sha.clone()); let (loop_sigs, restart_sigs) = self.circuit_breaker.snapshot(); - - // Build checkpoint from state - let mut node_outcomes = state.node_outcomes.clone(); - node_outcomes.insert(node.id().to_string(), result.outcome.clone()); - - let _checkpoint = fabro_types::Checkpoint { - timestamp: chrono::Utc::now(), - current_node: node.id().to_string(), - completed_nodes: state.completed_nodes.clone(), - node_outcomes, - node_retries: state.node_retries.clone(), - context_values: state.context.snapshot(), - next_node_id: next_node_id.map(String::from), - git_commit_sha: None, - node_visits: state.node_visits.clone(), - loop_failure_signatures: loop_sigs, - restart_failure_signatures: restart_sigs, - }; + let checkpoint = build_checkpoint( + node, + result, + next_node_id, + state, + loop_sigs, + restart_sigs, + git_commit_sha, + ); + let checkpoint_bytes = serde_json::to_vec_pretty(&checkpoint) + .map_err(|err| CoreError::Other(format!("failed to serialize checkpoint: {err}")))?; + fs::write(self.run_dir.join("checkpoint.json"), checkpoint_bytes) + .await + .map_err(|err| CoreError::Other(format!("failed to write checkpoint.json: {err}")))?; Ok(()) } } diff --git a/lib/crates/fabro-workflow/src/lifecycle/git.rs b/lib/crates/fabro-workflow/src/lifecycle/git.rs index 8dae8b56d..f48719a07 100644 --- a/lib/crates/fabro-workflow/src/lifecycle/git.rs +++ b/lib/crates/fabro-workflow/src/lifecycle/git.rs @@ -1,9 +1,11 @@ +use std::collections::HashMap; use std::path::PathBuf; use std::sync::{Arc, Mutex}; use async_trait::async_trait; use fabro_store::SlateRunStore; use fabro_types::RunId; +use tokio::fs; use fabro_core::error::{CoreError, Result as CoreResult}; use fabro_core::graph::NodeSpec; @@ -16,6 +18,7 @@ use crate::event::{EventEmitter, RunNoticeLevel, WorkflowRunEvent}; use crate::git::MetadataStore; use crate::graph::WorkflowGraph; use crate::graph::WorkflowNode; +use crate::lifecycle::disk::build_checkpoint; use crate::outcome::{Outcome, StageStatus, StageUsage}; use crate::run_dump::RunDump; use crate::run_options::RunOptions; @@ -92,7 +95,7 @@ impl RunLifecycle for GitLifecycle { &self, node: &WorkflowNode, result: &WfNodeResult, - _next_node_id: Option<&str>, + next_node_id: Option<&str>, state: &WfRunState, ) -> CoreResult<()> { let node_id = node.id(); @@ -113,15 +116,16 @@ impl RunLifecycle for GitLifecycle { ) { let git_author = self.run_options.git_author(); let store = MetadataStore::new(repo_path, &git_author); - // Build checkpoint JSON for shadow branch - if let Some(cp_json) = self - .run_store - .state() - .await - .ok() - .and_then(|state| state.checkpoint) - .and_then(|checkpoint| serde_json::to_vec_pretty(&checkpoint).ok()) - { + let checkpoint = build_checkpoint( + node, + result, + next_node_id, + state, + HashMap::new(), + HashMap::new(), + None, + ); + if let Ok(cp_json) = serde_json::to_vec_pretty(&checkpoint) { let mut extra_entries: Vec<(String, Vec)> = { let artifact_store = self.artifact_store.lock().unwrap(); artifact_store @@ -296,6 +300,13 @@ impl RunLifecycle for GitLifecycle { match git_diff(&*self.sandbox, &base_sha).await { Ok(patch) if !patch.is_empty() => { *self.final_patch.lock().unwrap() = Some(patch.clone()); + if let Err(err) = fs::write(self.run_dir.join("final.patch"), patch).await { + self.emitter.emit(&WorkflowRunEvent::RunNotice { + level: RunNoticeLevel::Warn, + code: "final_patch_write_failed".to_string(), + message: format!("failed to write final.patch: {err}"), + }); + } } Ok(_) => { *self.final_patch.lock().unwrap() = None; diff --git a/lib/crates/fabro-workflow/src/lifecycle/mod.rs b/lib/crates/fabro-workflow/src/lifecycle/mod.rs index 5805fe90d..4a5155250 100644 --- a/lib/crates/fabro-workflow/src/lifecycle/mod.rs +++ b/lib/crates/fabro-workflow/src/lifecycle/mod.rs @@ -150,6 +150,7 @@ impl WorkflowLifecycle { graph: Arc::clone(&graph), run_options: Arc::clone(run_options), emitter: Arc::clone(emitter), + checkpoint_git_result: Arc::clone(&checkpoint_git_result), circuit_breaker: Arc::clone(&circuit_breaker), checkpoint_enabled: true, }; @@ -407,10 +408,10 @@ impl RunLifecycle for WorkflowLifecycle { next_node_id: Option<&str>, state: &WfRunState, ) -> CoreResult<()> { - self.disk + self.git .on_checkpoint(node, result, next_node_id, state) .await?; - self.git + self.disk .on_checkpoint(node, result, next_node_id, state) .await?; self.event diff --git a/lib/crates/fabro-workflow/src/pipeline/persist.rs b/lib/crates/fabro-workflow/src/pipeline/persist.rs index d8e6d01b8..6f04be230 100644 --- a/lib/crates/fabro-workflow/src/pipeline/persist.rs +++ b/lib/crates/fabro-workflow/src/pipeline/persist.rs @@ -33,9 +33,10 @@ pub(crate) async fn load_from_store( .state() .await .map_err(|err| FabroError::engine(err.to_string()))?; - let run_record = state + let mut run_record = state .run .ok_or_else(|| FabroError::Precondition("run record missing from store".to_string()))?; + run_record.created_at = run_store.created_at(); let graph = run_record.graph.clone(); let source = state.graph_source.unwrap_or_default();