From aefb6e01c55323e35690f7085d28f6492f75d9d8 Mon Sep 17 00:00:00 2001 From: Bryan Helmkamp Date: Sun, 29 Mar 2026 19:23:44 -0400 Subject: [PATCH] Handle projection and retro artifact write failures --- lib/crates/fabro-retro/src/retro_agent.rs | 16 +-- lib/crates/fabro-store/src/disk_projecting.rs | 121 ++++++++++++++---- lib/crates/fabro-store/src/lib.rs | 2 +- .../fabro-workflows/src/operations/start.rs | 37 +++++- 4 files changed, 135 insertions(+), 41 deletions(-) diff --git a/lib/crates/fabro-retro/src/retro_agent.rs b/lib/crates/fabro-retro/src/retro_agent.rs index 1d1b62f71..67d3e1fbf 100644 --- a/lib/crates/fabro-retro/src/retro_agent.rs +++ b/lib/crates/fabro-retro/src/retro_agent.rs @@ -276,10 +276,10 @@ async fn write_retro_prompt( prompt: &str, ) -> anyhow::Result<()> { if let Some(store) = run_store { - store - .put_retro_prompt(prompt) - .await - .map_err(|e| anyhow::anyhow!("Failed to save retro prompt to store: {e}"))?; + if let Err(err) = store.put_retro_prompt(prompt).await { + tracing::warn!(error = %err, "Failed to save retro prompt to store"); + std::fs::write(retro_dir.join("prompt.md"), prompt)?; + } } else { std::fs::write(retro_dir.join("prompt.md"), prompt)?; } @@ -292,10 +292,10 @@ async fn write_retro_response( response: &str, ) -> anyhow::Result<()> { if let Some(store) = run_store { - store - .put_retro_response(response) - .await - .map_err(|e| anyhow::anyhow!("Failed to save retro response to store: {e}"))?; + if let Err(err) = store.put_retro_response(response).await { + tracing::warn!(error = %err, "Failed to save retro response to store"); + std::fs::write(retro_dir.join("response.md"), response)?; + } } else { std::fs::write(retro_dir.join("response.md"), response)?; } diff --git a/lib/crates/fabro-store/src/disk_projecting.rs b/lib/crates/fabro-store/src/disk_projecting.rs index 932d7dc79..8f2c5212e 100644 --- a/lib/crates/fabro-store/src/disk_projecting.rs +++ b/lib/crates/fabro-store/src/disk_projecting.rs @@ -17,18 +17,39 @@ use fabro_types::{ StartRecord, }; +#[derive(Debug, Clone)] +pub struct ProjectionError { + pub path: PathBuf, + pub critical: bool, + pub error: String, +} + pub struct DiskProjectingRunStore { inner: Arc, run_dir: PathBuf, + on_projection_error: Option>, } impl DiskProjectingRunStore { #[must_use] pub fn new(inner: Arc, run_dir: PathBuf) -> Self { - Self { inner, run_dir } + Self { + inner, + run_dir, + on_projection_error: None, + } } - fn warn_projection(path: &Path, err: &std::io::Error, critical: bool) { + #[must_use] + pub fn on_projection_error( + mut self, + callback: Arc, + ) -> Self { + self.on_projection_error = Some(callback); + self + } + + fn report_projection_error(&self, path: &Path, err: &std::io::Error, critical: bool) { if critical { warn!( path = %path.display(), @@ -38,35 +59,43 @@ impl DiskProjectingRunStore { } else { warn!(path = %path.display(), error = %err, "Disk projection failed"); } - } - fn write_json_critical(path: &Path, value: &T) { - if let Err(err) = write_json(path, value) { - Self::warn_projection(path, &err, true); + if let Some(ref callback) = self.on_projection_error { + callback(ProjectionError { + path: path.to_path_buf(), + critical, + error: err.to_string(), + }); } } - fn write_json_best_effort(path: &Path, value: &T) { + fn write_json_critical(&self, path: &Path, value: &T) { if let Err(err) = write_json(path, value) { - Self::warn_projection(path, &err, false); + self.report_projection_error(path, &err, true); } } - fn write_text_best_effort(path: &Path, value: &str) { + fn write_json_best_effort(&self, path: &Path, value: &T) { + if let Err(err) = write_json(path, value) { + self.report_projection_error(path, &err, false); + } + } + + fn write_text_best_effort(&self, path: &Path, value: &str) { if let Err(err) = write_text(path, value) { - Self::warn_projection(path, &err, false); + self.report_projection_error(path, &err, false); } } fn append_jsonl_critical(&self, payload: &EventPayload) { let progress_path = self.run_dir.join("progress.jsonl"); if let Err(err) = append_jsonl(&progress_path, payload) { - Self::warn_projection(&progress_path, &err, true); + self.report_projection_error(&progress_path, &err, true); } let live_path = self.run_dir.join("live.json"); if let Err(err) = write_live_json(&live_path, payload) { - Self::warn_projection(&live_path, &err, true); + self.report_projection_error(&live_path, &err, true); } } } @@ -121,7 +150,7 @@ fn write_live_json(path: &Path, payload: &EventPayload) -> std::io::Result<()> { impl RunStore for DiskProjectingRunStore { async fn put_run(&self, record: &RunRecord) -> Result<()> { self.inner.put_run(record).await?; - Self::write_json_best_effort(&self.run_dir.join("run.json"), record); + self.write_json_best_effort(&self.run_dir.join("run.json"), record); Ok(()) } @@ -131,7 +160,7 @@ impl RunStore for DiskProjectingRunStore { async fn put_start(&self, record: &StartRecord) -> Result<()> { self.inner.put_start(record).await?; - Self::write_json_best_effort(&self.run_dir.join("start.json"), record); + self.write_json_best_effort(&self.run_dir.join("start.json"), record); Ok(()) } @@ -140,7 +169,7 @@ impl RunStore for DiskProjectingRunStore { } async fn put_status(&self, record: &RunStatusRecord) -> Result<()> { - Self::write_json_critical(&self.run_dir.join("status.json"), record); + self.write_json_critical(&self.run_dir.join("status.json"), record); self.inner.put_status(record).await } @@ -150,7 +179,7 @@ impl RunStore for DiskProjectingRunStore { async fn put_checkpoint(&self, record: &Checkpoint) -> Result<()> { self.inner.put_checkpoint(record).await?; - Self::write_json_best_effort(&self.run_dir.join("checkpoint.json"), record); + self.write_json_best_effort(&self.run_dir.join("checkpoint.json"), record); Ok(()) } @@ -167,7 +196,7 @@ impl RunStore for DiskProjectingRunStore { } async fn put_conclusion(&self, record: &Conclusion) -> Result<()> { - Self::write_json_critical(&self.run_dir.join("conclusion.json"), record); + self.write_json_critical(&self.run_dir.join("conclusion.json"), record); self.inner.put_conclusion(record).await } @@ -177,7 +206,7 @@ impl RunStore for DiskProjectingRunStore { async fn put_retro(&self, retro: &Retro) -> Result<()> { self.inner.put_retro(retro).await?; - Self::write_json_best_effort(&self.run_dir.join("retro.json"), retro); + self.write_json_best_effort(&self.run_dir.join("retro.json"), retro); Ok(()) } @@ -187,7 +216,7 @@ impl RunStore for DiskProjectingRunStore { async fn put_graph(&self, dot_source: &str) -> Result<()> { self.inner.put_graph(dot_source).await?; - Self::write_text_best_effort(&self.run_dir.join("workflow.fabro"), dot_source); + self.write_text_best_effort(&self.run_dir.join("workflow.fabro"), dot_source); Ok(()) } @@ -197,7 +226,7 @@ impl RunStore for DiskProjectingRunStore { async fn put_sandbox(&self, record: &SandboxRecord) -> Result<()> { self.inner.put_sandbox(record).await?; - Self::write_json_best_effort(&self.run_dir.join("sandbox.json"), record); + self.write_json_best_effort(&self.run_dir.join("sandbox.json"), record); Ok(()) } @@ -207,7 +236,7 @@ impl RunStore for DiskProjectingRunStore { async fn put_node_prompt(&self, node: &NodeVisitRef<'_>, prompt: &str) -> Result<()> { self.inner.put_node_prompt(node, prompt).await?; - Self::write_text_best_effort( + self.write_text_best_effort( &disk_node_dir(&self.run_dir, node.node_id, node.visit).join("prompt.md"), prompt, ); @@ -216,7 +245,7 @@ impl RunStore for DiskProjectingRunStore { async fn put_node_response(&self, node: &NodeVisitRef<'_>, response: &str) -> Result<()> { self.inner.put_node_response(node, response).await?; - Self::write_text_best_effort( + self.write_text_best_effort( &disk_node_dir(&self.run_dir, node.node_id, node.visit).join("response.md"), response, ); @@ -229,7 +258,7 @@ impl RunStore for DiskProjectingRunStore { status: &NodeStatusRecord, ) -> Result<()> { self.inner.put_node_status(node, status).await?; - Self::write_json_best_effort( + self.write_json_best_effort( &disk_node_dir(&self.run_dir, node.node_id, node.visit).join("status.json"), status, ); @@ -238,7 +267,7 @@ impl RunStore for DiskProjectingRunStore { async fn put_node_stdout(&self, node: &NodeVisitRef<'_>, log: &str) -> Result<()> { self.inner.put_node_stdout(node, log).await?; - Self::write_text_best_effort( + self.write_text_best_effort( &disk_node_dir(&self.run_dir, node.node_id, node.visit).join("stdout.log"), log, ); @@ -247,7 +276,7 @@ impl RunStore for DiskProjectingRunStore { async fn put_node_stderr(&self, node: &NodeVisitRef<'_>, log: &str) -> Result<()> { self.inner.put_node_stderr(node, log).await?; - Self::write_text_best_effort( + self.write_text_best_effort( &disk_node_dir(&self.run_dir, node.node_id, node.visit).join("stderr.log"), log, ); @@ -284,7 +313,7 @@ impl RunStore for DiskProjectingRunStore { async fn put_retro_prompt(&self, text: &str) -> Result<()> { self.inner.put_retro_prompt(text).await?; - Self::write_text_best_effort(&self.run_dir.join("retro").join("prompt.md"), text); + self.write_text_best_effort(&self.run_dir.join("retro").join("prompt.md"), text); Ok(()) } @@ -294,7 +323,7 @@ impl RunStore for DiskProjectingRunStore { async fn put_retro_response(&self, text: &str) -> Result<()> { self.inner.put_retro_response(text).await?; - Self::write_text_best_effort(&self.run_dir.join("retro").join("response.md"), text); + self.write_text_best_effort(&self.run_dir.join("retro").join("response.md"), text); Ok(()) } @@ -761,6 +790,44 @@ mod tests { assert_eq!(stored_node.prompt.as_deref(), Some("Plan the fix")); } + #[tokio::test] + async fn projection_error_callback_runs_on_disk_failure() { + let temp = TempDir::new().unwrap(); + let created_at = dt("2026-03-27T12:00:00Z"); + let inner = InMemoryStore::default() + .create_run( + "run-1", + created_at, + Some(temp.path().to_string_lossy().as_ref()), + ) + .await + .unwrap(); + let seen = Arc::new(std::sync::Mutex::new(Vec::::new())); + let seen_clone = Arc::clone(&seen); + let store = DiskProjectingRunStore::new(inner, temp.path().to_path_buf()) + .on_projection_error(Arc::new(move |error| { + seen_clone.lock().unwrap().push(error); + })); + + let mut permissions = fs::metadata(temp.path()).unwrap().permissions(); + permissions.set_readonly(true); + fs::set_permissions(temp.path(), permissions).unwrap(); + + store + .put_status(&sample_status( + RunStatus::Running, + Some(StatusReason::SandboxInitializing), + )) + .await + .unwrap(); + + let seen = seen.lock().unwrap(); + assert_eq!(seen.len(), 1); + assert!(seen[0].critical); + assert_eq!(seen[0].path, temp.path().join("status.json")); + assert!(!seen[0].error.is_empty()); + } + #[test] fn disk_node_dir_matches_legacy_layout() { let run_dir = Path::new("/tmp/fabro-run"); diff --git a/lib/crates/fabro-store/src/lib.rs b/lib/crates/fabro-store/src/lib.rs index 3064f3918..af13adc4b 100644 --- a/lib/crates/fabro-store/src/lib.rs +++ b/lib/crates/fabro-store/src/lib.rs @@ -14,7 +14,7 @@ mod runtime; mod slate; mod types; -pub use disk_projecting::DiskProjectingRunStore; +pub use disk_projecting::{DiskProjectingRunStore, ProjectionError}; pub use error::{Result, StoreError}; pub use memory::InMemoryStore; pub use runtime::RuntimeState; diff --git a/lib/crates/fabro-workflows/src/operations/start.rs b/lib/crates/fabro-workflows/src/operations/start.rs index 0cb804645..28f53cdd6 100644 --- a/lib/crates/fabro-workflows/src/operations/start.rs +++ b/lib/crates/fabro-workflows/src/operations/start.rs @@ -11,7 +11,7 @@ use fabro_config::{project as project_config, run as run_config, sandbox as sand use fabro_interview::{AutoApproveInterviewer, Interviewer}; use fabro_model::{Catalog, FallbackTarget, Provider}; use fabro_sandbox::{SandboxProvider, SandboxSpec, detect_clone_params}; -use fabro_store::{DiskProjectingRunStore, RunStore}; +use fabro_store::{DiskProjectingRunStore, ProjectionError, RunStore}; use serde::Serialize; use crate::context::Context; @@ -118,10 +118,37 @@ pub(super) async fn execute_persisted_run( mut services: StartServices, ) -> Result { let inner_store = Arc::clone(&services.run_store); - services.run_store = Arc::new(DiskProjectingRunStore::new( - inner_store, - run_dir.to_path_buf(), - )); + let projection_run_dir = run_dir.to_path_buf(); + services.run_store = Arc::new( + DiskProjectingRunStore::new(inner_store, run_dir.to_path_buf()).on_projection_error( + Arc::new(move |projection_error: ProjectionError| { + let Some(run_id) = load_run_id(&projection_run_dir) else { + return; + }; + + // Write directly to progress.jsonl/live.json so projection failures do not + // recurse back through the decorated store.append_event() path. + let _ = append_progress_event( + &projection_run_dir, + &run_id, + &WorkflowRunEvent::RunNotice { + level: RunNoticeLevel::Warn, + code: "disk_projection_failed".to_string(), + message: format!( + "{}disk projection failed for {}: {}", + if projection_error.critical { + "critical " + } else { + "" + }, + projection_error.path.display(), + projection_error.error + ), + }, + ); + }), + ), + ); let run_store = Arc::clone(&services.run_store); if let Err(err) = run_store .put_status(&run_status::RunStatusRecord::new(