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 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01Ryyhtbc1eNtCLw8GjrFQXZ
This commit is contained in:
Bryan Helmkamp 2026-08-21 19:29:14 -04:00
parent 129fa0ea0c
commit 61394ba2f1
No known key found for this signature in database
5 changed files with 139 additions and 142 deletions

View file

@ -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<fabro_validate::ValidationError> for Error {
}
}
impl From<RunEventPersistenceError> for Error {
fn from(err: RunEventPersistenceError) -> Self {
Self::engine_with_source("run event persistence failed", err)
}
}
impl From<fabro_checkpoint::MetadataError> for Error {
fn from(err: fabro_checkpoint::MetadataError) -> Self {
match err {

View file

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

View file

@ -43,8 +43,7 @@ pub async fn resume(run_dir: &Path, services: StartServices) -> Result<Started,
&services.run_id,
&Event::RunSubmitted { definition_blob },
)
.await
.map_err(|err| Error::engine_with_source("failed to persist run.submitted event", err))?;
.await?;
Box::pin(execute_persisted_run(run_dir, Some(resume_state), services)).await
}

View file

@ -1,4 +1,5 @@
use std::collections::{HashMap, HashSet};
use std::future::Future;
use std::path::Path;
use std::sync::{Arc, Mutex};
use std::time::{Duration, Instant};
@ -159,10 +160,7 @@ pub async fn start(run_dir: &Path, services: StartServices) -> Result<Started, E
actor: None,
},
)
.await
.map_err(|err| {
Error::engine_with_source("failed to persist run.start_requested event", err)
})?;
.await?;
append_event_to_sink(
&services.event_sink,
&services.run_id,
@ -171,8 +169,7 @@ pub async fn start(run_dir: &Path, services: StartServices) -> Result<Started, E
actor: None,
},
)
.await
.map_err(|err| Error::engine_with_source("failed to persist run.runnable event", err))?;
.await?;
}
Box::pin(execute_persisted_run(run_dir, None, services)).await
@ -202,7 +199,7 @@ pub(super) async fn execute_persisted_run(
return Err(error);
}
if let Err(err) = append_event_to_sink(&event_sink, &run_id, &Event::RunStarting).await {
let error = Error::engine_with_source("failed to persist run.starting event", err);
let error = Error::from(err);
let _ = persist_detached_failure(
run_id,
&run_store,
@ -322,7 +319,7 @@ async fn emit_workflow_run_failed(
conclusion.billing,
);
if let Err(err) = append_event_to_sink(event_sink, &run_id, &failure_event).await {
let rendered_error = collect_chain(err.as_ref()).join(": ");
let rendered_error = collect_chain(&err).join(": ");
tracing::error!(
run_id = %run_id,
event = "run.failed",
@ -362,7 +359,33 @@ fn stop_for_run_event_persistence_failure(
error: RunEventPersistenceError,
) -> 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<T>(
logger: &RunEventLogger,
cancel_token: &CancellationToken,
step: impl Future<Output = T>,
) -> Result<T, Error> {
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<Outcome, Error> {
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<Outcome, Error> {
std::future::pending().await
}
}
fn memory_store() -> Arc<Database> {

View file

@ -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<String>,
fork_source_ref: Option<ForkSourceRef>,
) {
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({