From 61394ba2f1dcfba9700f044b156563b1769b1fd9 Mon Sep 17 00:00:00 2001 From: Bryan Helmkamp Date: Fri, 21 Aug 2026 19:29:14 -0400 Subject: [PATCH] Simplify run-event persistence failure plumbing - Make the failure watch channel the single record of the latched failure; drop the worker task's mirrored local state. - Replace the hand-rolled wait loop with watch::Receiver::wait_for. - Extract race_persistence/flush_or_stop helpers so the select!/flush scaffolding in RunSession::run exists once instead of three times. - Return RunEventPersistenceError from append_event_to_sink and add a From impl on Error, replacing four hand-written per-event message strings with the event name derived from the event itself. - Dedupe the RunCreated test seed literal in initialize.rs and drop the dead BlockingHandler::simulate override. Co-Authored-By: Claude Fable 5 Claude-Session: https://claude.ai/code/session_01Ryyhtbc1eNtCLw8GjrFQXZ --- lib/components/fabro-workflow/src/error.rs | 7 + .../fabro-workflow/src/event/sink.rs | 38 ++--- .../fabro-workflow/src/operations/resume.rs | 3 +- .../fabro-workflow/src/operations/start.rs | 140 ++++++++---------- .../fabro-workflow/src/pipeline/initialize.rs | 93 ++++++------ 5 files changed, 139 insertions(+), 142 deletions(-) diff --git a/lib/components/fabro-workflow/src/error.rs b/lib/components/fabro-workflow/src/error.rs index 290aacf91..7bedbda70 100644 --- a/lib/components/fabro-workflow/src/error.rs +++ b/lib/components/fabro-workflow/src/error.rs @@ -14,6 +14,7 @@ use fabro_validate::Diagnostic; use regex::Regex; use thiserror::Error as ThisError; +use crate::event::RunEventPersistenceError; use crate::outcome::{FailureDetail, Outcome, StageOutcome}; /// Classify an LLM error into a `FailureCategory` based on its structure. @@ -721,6 +722,12 @@ impl From for Error { } } +impl From for Error { + fn from(err: RunEventPersistenceError) -> Self { + Self::engine_with_source("run event persistence failed", err) + } +} + impl From for Error { fn from(err: fabro_checkpoint::MetadataError) -> Self { match err { diff --git a/lib/components/fabro-workflow/src/event/sink.rs b/lib/components/fabro-workflow/src/event/sink.rs index e27a584b0..a82f2e20f 100644 --- a/lib/components/fabro-workflow/src/event/sink.rs +++ b/lib/components/fabro-workflow/src/event/sink.rs @@ -43,9 +43,15 @@ pub async fn append_event_to_sink( sink: &RunEventSink, run_id: &RunId, event: &Event, -) -> Result<()> { +) -> Result<(), RunEventPersistenceError> { let stored = to_run_event(run_id, event); - sink.write_run_event(&stored).await + sink.write_run_event(&stored) + .await + .map_err(|err| RunEventPersistenceError::Write { + run_id: *run_id, + event: stored.body.event_name().to_string(), + source: SharedError::new(err), + }) } #[derive(Clone)] @@ -179,11 +185,12 @@ impl RunEventLogger { let (failure_tx, failure_rx) = watch::channel(None); tokio::spawn(async move { - let mut persistence_failure = None; + // The watch channel is the single record of the latched failure: + // the worker is its only writer, so borrowing it here cannot race. while let Some(command) = rx.recv().await { match command { RunEventCommand::Event(event) => { - if persistence_failure.is_some() { + if failure_tx.borrow().is_some() { continue; } if let Err(err) = sink.write_run_event(&event).await { @@ -194,17 +201,15 @@ impl RunEventLogger { error = %rendered_error, "Failed to persist run event; stopping workflow", ); - let failure = RunEventPersistenceError::Write { + failure_tx.send_replace(Some(RunEventPersistenceError::Write { run_id: event.run_id, event: event.body.event_name().to_string(), source: SharedError::new(err), - }; - persistence_failure = Some(failure.clone()); - failure_tx.send_replace(Some(failure)); + })); } } RunEventCommand::Flush(tx) => { - let result = persistence_failure.clone().map_or(Ok(()), Err); + let result = failure_tx.borrow().clone().map_or(Ok(()), Err); let _ = tx.send(result); } } @@ -229,13 +234,12 @@ impl RunEventLogger { pub async fn wait_for_failure(&self) -> RunEventPersistenceError { let mut failure_rx = self.failure_rx.clone(); - loop { - if let Some(failure) = failure_rx.borrow_and_update().clone() { - return failure; - } - if failure_rx.changed().await.is_err() { - return RunEventPersistenceError::TaskStopped; - } + let failure = failure_rx.wait_for(Option::is_some).await; + match failure { + Ok(failure) => failure + .clone() + .expect("wait_for only returns values matching the predicate"), + Err(_) => RunEventPersistenceError::TaskStopped, } } @@ -245,7 +249,7 @@ impl RunEventLogger { return Err(RunEventPersistenceError::TaskStopped); } rx.await - .map_err(|_| RunEventPersistenceError::TaskStopped)? + .unwrap_or(Err(RunEventPersistenceError::TaskStopped)) } } diff --git a/lib/components/fabro-workflow/src/operations/resume.rs b/lib/components/fabro-workflow/src/operations/resume.rs index 1911c4a66..466ae5e0c 100644 --- a/lib/components/fabro-workflow/src/operations/resume.rs +++ b/lib/components/fabro-workflow/src/operations/resume.rs @@ -43,8 +43,7 @@ pub async fn resume(run_dir: &Path, services: StartServices) -> Result Result Result Error { cancel_token.cancel(); - Error::engine_with_source("run event persistence failed", error) + error.into() +} + +/// Race a pipeline step against the first latched run-event persistence +/// failure. When the failure wins, the step future is dropped mid-flight and +/// the run token is cancelled. +async fn race_persistence( + logger: &RunEventLogger, + cancel_token: &CancellationToken, + step: impl Future, +) -> Result { + tokio::select! { + result = step => Ok(result), + failure = logger.wait_for_failure() => { + Err(stop_for_run_event_persistence_failure(cancel_token, failure)) + } + } +} + +async fn flush_or_stop( + logger: &RunEventLogger, + cancel_token: &CancellationToken, +) -> Result<(), Error> { + logger + .flush() + .await + .map_err(|failure| stop_for_run_event_persistence_failure(cancel_token, failure)) } impl RunSession { @@ -798,25 +821,16 @@ impl RunSession { seed_context: self.seed_context, fabro_run_tools: self.fabro_run_tools, }; - let mut initializing = Box::pin(pipeline::initialize(persisted, init_options)); - let initialized = tokio::select! { - result = &mut initializing => result, - failure = store_progress_logger.wait_for_failure() => { - return Err(stop_for_run_event_persistence_failure( - &run_cancel_token, - failure, - )); - } - }; - let mut initialized = match initialized { + let mut initialized = match race_persistence( + &store_progress_logger, + &run_cancel_token, + Box::pin(pipeline::initialize(persisted, init_options)), + ) + .await? + { Ok(initialized) => initialized, Err(err) => { - if let Err(failure) = store_progress_logger.flush().await { - return Err(stop_for_run_event_persistence_failure( - &run_cancel_token, - failure, - )); - } + flush_or_stop(&store_progress_logger, &run_cancel_token).await?; return Err(err); } }; @@ -843,23 +857,15 @@ impl RunSession { steering_hub_for_drain.drain_pending_at_run_end(); }); - store_progress_logger.flush().await.map_err(|failure| { - stop_for_run_event_persistence_failure(&run_cancel_token, failure) - })?; + flush_or_stop(&store_progress_logger, &run_cancel_token).await?; - let mut executing = Box::pin(pipeline::execute(initialized)); - let executed = tokio::select! { - executed = &mut executing => executed, - failure = store_progress_logger.wait_for_failure() => { - return Err(stop_for_run_event_persistence_failure( - &run_cancel_token, - failure, - )); - } - }; - store_progress_logger.flush().await.map_err(|failure| { - stop_for_run_event_persistence_failure(&run_cancel_token, failure) - })?; + let executed = race_persistence( + &store_progress_logger, + &run_cancel_token, + Box::pin(pipeline::execute(initialized)), + ) + .await?; + flush_or_stop(&store_progress_logger, &run_cancel_token).await?; let final_context = Some(executed.final_context.clone()); let finalize_opts = FinalizeOptions { @@ -879,30 +885,21 @@ impl RunSession { model: self.pr_model, }; - let mut concluding = Box::pin(async { - let concluded = Box::pin(pipeline::conclude(executed, &finalize_opts)).await?; - let published = Box::pin(pipeline::publish(concluded, &publish_opts)).await; - Box::pin(pipeline::finalize(published, &finalize_opts)).await - }); - let concluding = tokio::select! { - result = &mut concluding => result, - failure = store_progress_logger.wait_for_failure() => { - return Err(stop_for_run_event_persistence_failure( - &run_cancel_token, - failure, - )); - } - }; + let concluding = race_persistence( + &store_progress_logger, + &run_cancel_token, + Box::pin(async { + let concluded = Box::pin(pipeline::conclude(executed, &finalize_opts)).await?; + let published = Box::pin(pipeline::publish(concluded, &publish_opts)).await; + Box::pin(pipeline::finalize(published, &finalize_opts)).await + }), + ) + .await?; let finalized = match concluding { Ok(finalized) => finalized, Err(err) => { self.steering_hub.drain_pending_at_run_end(); - if let Err(failure) = store_progress_logger.flush().await { - return Err(stop_for_run_event_persistence_failure( - &run_cancel_token, - failure, - )); - } + flush_or_stop(&store_progress_logger, &run_cancel_token).await?; return Err(err); } }; @@ -911,9 +908,7 @@ impl RunSession { // scopeguard above re-runs as a no-op (drain is idempotent on an // already-empty buffer) on the way out of scope. self.steering_hub.drain_pending_at_run_end(); - store_progress_logger.flush().await.map_err(|failure| { - stop_for_run_event_persistence_failure(&run_cancel_token, failure) - })?; + flush_or_stop(&store_progress_logger, &run_cancel_token).await?; scopeguard::ScopeGuard::into_inner(cleanup_guard); @@ -1081,7 +1076,7 @@ impl Drop for DetachedRunCompletionGuard { }) .await { - let rendered_error = collect_chain(err.as_ref()).join(": "); + let rendered_error = collect_chain(&err).join(": "); tracing::warn!( error = %rendered_error, "Failed to append detached completion notice", @@ -1110,7 +1105,7 @@ async fn persist_detached_failure( exec_output_tail: None, }; if let Err(err) = append_event_to_sink(event_sink, &run_id, &event).await { - let rendered_error = collect_chain(err.as_ref()).join(": "); + let rendered_error = collect_chain(&err).join(": "); tracing::warn!( error = %rendered_error, "Failed to append detached failure notice", @@ -1229,17 +1224,6 @@ mod tests { ) -> Result { std::future::pending().await } - - async fn simulate( - &self, - _node: &fabro_graphviz::graph::Node, - _context: &Context, - _graph: &fabro_graphviz::graph::Graph, - _run_dir: &Path, - _services: &EngineServices, - ) -> Result { - std::future::pending().await - } } fn memory_store() -> Arc { diff --git a/lib/components/fabro-workflow/src/pipeline/initialize.rs b/lib/components/fabro-workflow/src/pipeline/initialize.rs index 7706ae411..d7f8a5930 100644 --- a/lib/components/fabro-workflow/src/pipeline/initialize.rs +++ b/lib/components/fabro-workflow/src/pipeline/initialize.rs @@ -668,7 +668,7 @@ mod tests { use fabro_graphviz::graph::{AttrValue, Edge, Graph, Node}; use fabro_interview::AutoApproveInterviewer; use fabro_sandbox::SandboxSpec; - use fabro_store::Database; + use fabro_store::{Database, RunDatabase}; use fabro_types::settings::run::RunModelControls; use fabro_types::{ EventBody, ForkSourceRef, RunEvent, RunId, WorkflowSettings, fixtures, test_support, @@ -713,6 +713,37 @@ mod tests { )) } + async fn seed_run_created( + run_store: &RunDatabase, + settings: serde_json::Value, + graph: serde_json::Value, + source_directory: Option, + fork_source_ref: Option, + ) { + crate::event::append_event(run_store, &test_run_id(), &Event::RunCreated { + run_id: test_run_id(), + title: None, + settings, + graph, + workflow_source: None, + labels: BTreeMap::new(), + source_directory, + workflow_slug: Some("test".to_string()), + workflow_version_id: None, + automation: None, + provenance: test_support::test_run_provenance(), + manifest_blob: None, + spec_blob: None, + git: None, + fork_source_ref, + retried_from: None, + parent_id: None, + web_url: None, + }) + .await + .unwrap(); + } + fn simple_graph() -> (Graph, String) { let source = r"digraph test { start [shape=Mdiamond]; @@ -1039,28 +1070,14 @@ mod tests { let mut run_options = test_settings(&run_dir); run_options.settings = settings; run_options.fork_source_ref = fork_source_ref; - crate::event::append_event(&run_store, &test_run_id(), &Event::RunCreated { - run_id: test_run_id(), - title: None, - settings: serde_json::to_value(&run_options.settings).unwrap(), - graph: serde_json::to_value(&graph).unwrap(), - workflow_source: None, - labels: BTreeMap::new(), - source_directory: Some(workspace.display().to_string()), - workflow_slug: Some("test".to_string()), - workflow_version_id: None, - automation: None, - provenance: test_support::test_run_provenance(), - manifest_blob: None, - spec_blob: None, - git: None, - fork_source_ref: run_options.fork_source_ref.clone(), - retried_from: None, - parent_id: None, - web_url: None, - }) - .await - .unwrap(); + seed_run_created( + &run_store, + serde_json::to_value(&run_options.settings).unwrap(), + serde_json::to_value(&graph).unwrap(), + Some(workspace.display().to_string()), + run_options.fork_source_ref.clone(), + ) + .await; initialize(persisted, InitOptions { resume: Some(ResumeState::for_test( @@ -1330,28 +1347,14 @@ mod tests { let emitter = Arc::new(crate::event::Emitter::new(test_run_id())); let store = memory_store(); let run_store = store.create_run(&test_run_id()).await.unwrap(); - crate::event::append_event(&run_store, &test_run_id(), &Event::RunCreated { - run_id: test_run_id(), - title: None, - settings: serde_json::to_value(WorkflowSettings::default()).unwrap(), - graph: serde_json::to_value(graph).unwrap(), - workflow_source: None, - labels: BTreeMap::new(), - source_directory: None, - workflow_slug: Some("test".to_string()), - workflow_version_id: None, - automation: None, - provenance: test_support::test_run_provenance(), - manifest_blob: None, - spec_blob: None, - git: None, - fork_source_ref: None, - retried_from: None, - parent_id: None, - web_url: None, - }) - .await - .unwrap(); + seed_run_created( + &run_store, + serde_json::to_value(WorkflowSettings::default()).unwrap(), + serde_json::to_value(graph).unwrap(), + None, + None, + ) + .await; let store_logger = StoreProgressLogger::new(run_store.clone()); let seen = Arc::new(std::sync::Mutex::new(Vec::new())); emitter.on_event({