Persist pre-start worker failures

This commit is contained in:
Scott Werner 2026-09-02 17:29:41 -04:00
parent efcf8a0d93
commit 54666632c8
4 changed files with 430 additions and 13 deletions

View file

@ -3585,18 +3585,45 @@ async fn fail_worker_launch(
err: anyhow::Error,
) {
tracing::error!(run_id = %run_id, error = %err, "Failed to spawn worker");
let message = format!("Failed to spawn worker: {err}");
let cancellation_pending = match run_store.state().await {
Ok(run_state) => run_state.pending_control == Some(RunControlAction::Cancel),
Err(state_err) => {
tracing::warn!(
run_id = %run_id,
error = %render_compact_with_causes(
&state_err.to_string(),
&collect_causes(&state_err),
),
"Failed to load run state while recording worker launch failure"
);
false
}
};
let (error, reason, message) = if cancellation_pending {
(
WorkflowError::Cancelled,
FailureReason::Cancelled,
"Run cancelled before worker launch completed".to_string(),
)
} else {
let message = format!("Failed to spawn worker: {err}");
(
WorkflowError::engine_with_anyhow("Failed to spawn worker", err),
FailureReason::LaunchFailed,
message,
)
};
let failure_event = workflow_event::Event::workflow_run_failed_from_error(
&WorkflowError::engine_with_anyhow("Failed to spawn worker", err),
&error,
fabro_types::RunTiming::default(),
FailureReason::LaunchFailed,
reason,
None,
None,
None,
None,
);
let _ = workflow_event::append_event(run_store, &run_id, &failure_event).await;
fail_managed_run(state, run_id, FailureReason::LaunchFailed, message);
fail_managed_run(state, run_id, reason, message);
state.scheduler_notify.notify_one();
}

View file

@ -4,7 +4,7 @@ use std::os::unix::fs::PermissionsExt;
use std::path::{Path, PathBuf};
#[cfg(unix)]
use std::process::Stdio;
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
use std::sync::{Arc as StdArc, Mutex as StdMutex};
use async_zip::base::read::mem::ZipFileReader;
@ -58,7 +58,7 @@ use crate::worker_control::{
LocalWorkerControlBus, WorkerControlBus, WorkerControlCursor, WorkerControlReceiver,
};
use crate::worker_runtime::{
LocalWorkerRuntime, StartedWorker, WorkerLaunchSpec, WorkerRef, WorkerRuntime,
LocalWorkerRuntime, StartedWorker, WorkerExit, WorkerLaunchSpec, WorkerRef, WorkerRuntime,
};
const MINIMAL_DOT: &str = r#"digraph Test {
@ -2982,6 +2982,96 @@ impl RecordingWorkerRuntime {
}
}
#[derive(Clone, Copy)]
enum PreStartWorkerOutcome {
LaunchFailure,
EarlyExit,
}
struct PreStartWorkerRuntime {
outcome: PreStartWorkerOutcome,
starts: AtomicUsize,
hold_launch_failure: bool,
start_entered: Notify,
release_start: Notify,
}
impl PreStartWorkerRuntime {
fn new(outcome: PreStartWorkerOutcome) -> Self {
Self {
outcome,
starts: AtomicUsize::new(0),
hold_launch_failure: false,
start_entered: Notify::new(),
release_start: Notify::new(),
}
}
fn held_launch_failure() -> Self {
Self {
hold_launch_failure: true,
..Self::new(PreStartWorkerOutcome::LaunchFailure)
}
}
fn start_count(&self) -> usize {
self.starts.load(Ordering::Relaxed)
}
async fn wait_for_start(&self) {
tokio::time::timeout(std::time::Duration::from_secs(1), async {
loop {
let notified = self.start_entered.notified();
if self.start_count() > 0 {
return;
}
notified.await;
}
})
.await
.expect("test worker runtime should receive one start request");
}
fn release_held_start(&self) {
self.release_start.notify_one();
}
}
#[async_trait::async_trait]
impl WorkerRuntime for PreStartWorkerRuntime {
async fn start(&self, _spec: WorkerLaunchSpec) -> anyhow::Result<StartedWorker> {
self.starts.fetch_add(1, Ordering::Relaxed);
self.start_entered.notify_waiters();
if self.hold_launch_failure {
self.release_start.notified().await;
}
match self.outcome {
PreStartWorkerOutcome::LaunchFailure => {
anyhow::bail!("test worker launch failed")
}
PreStartWorkerOutcome::EarlyExit => Ok(StartedWorker {
worker_ref: test_worker_ref(u32::MAX),
stderr: Box::pin(tokio::io::empty()),
wait: Box::pin(async {
Ok(WorkerExit {
success: false,
detail: "test worker exited before starting".to_string(),
})
}),
}),
}
}
async fn request_stop(&self, _worker_ref: &WorkerRef) {}
async fn force_stop(&self, _worker_ref: &WorkerRef) {}
async fn is_alive(&self, _worker_ref: &WorkerRef) -> bool {
false
}
}
#[async_trait::async_trait]
impl WorkerRuntime for RecordingWorkerRuntime {
async fn start(&self, _spec: WorkerLaunchSpec) -> anyhow::Result<StartedWorker> {
@ -5700,6 +5790,232 @@ async fn create_and_start_run(app: &Router, dot_source: &str) -> String {
run_id
}
fn subprocess_pre_start_failure_state(runtime: StdArc<PreStartWorkerRuntime>) -> Arc<AppState> {
let state = TestAppStateBuilder::new()
.vault_entries([(EnvVars::OPENAI_API_KEY, "test-openai-api-key")])
.worker_runtime(runtime)
.build();
let runtime_directory = Storage::new(state.server_storage_dir()).runtime_directory();
ServerDaemon::new(
std::process::id(),
Bind::Tcp("127.0.0.1:32276".parse().expect("test bind should parse")),
runtime_directory.log_path(),
)
.write(&runtime_directory)
.expect("test server record should be written");
state
}
async fn assert_subprocess_pre_start_failure(
outcome: PreStartWorkerOutcome,
expected_reason: FailureReason,
) {
let runtime = StdArc::new(PreStartWorkerRuntime::new(outcome));
let state = subprocess_pre_start_failure_state(StdArc::clone(&runtime));
let app = crate::test_support::build_test_router(Arc::clone(&state));
let run_id = create_and_start_run(&app, MINIMAL_DOT)
.await
.parse::<RunId>()
.expect("created run id should parse");
let run_store = state
.stores
.runs
.open_run_reader(&run_id)
.await
.expect("created run should remain readable");
assert_eq!(
run_store
.state()
.await
.expect("runnable run state should load")
.status,
RunStatus::Runnable
);
execute_run(Arc::clone(&state), run_id).await;
assert_eq!(runtime.start_count(), 1);
let events = run_store
.list_events()
.await
.expect("failed run events should remain readable");
let lifecycle_events = events
.iter()
.map(|envelope| envelope.event.event_name())
.filter(|name| matches!(*name, "run.runnable" | "run.starting" | "run.failed"))
.collect::<Vec<_>>();
assert_eq!(lifecycle_events, vec!["run.runnable", "run.failed"]);
let failure_reasons = events
.iter()
.filter_map(|envelope| match &envelope.event.body {
EventBody::RunFailed(props) => Some(props.failure.reason),
_ => None,
})
.collect::<Vec<_>>();
assert_eq!(failure_reasons, vec![expected_reason]);
let expected_status = RunStatus::Failed {
reason: expected_reason,
};
assert_eq!(
run_store
.state()
.await
.expect("failed run state should load")
.status,
expected_status
);
assert_eq!(
state
.runs
.lock()
.expect("runs lock poisoned")
.get(&run_id)
.expect("managed run should remain present")
.status,
expected_status
);
let response = app
.clone()
.oneshot(
Request::builder()
.method("GET")
.uri(api(&format!("/runs/{run_id}")))
.body(Body::empty())
.expect("run request should build"),
)
.await
.expect("run request should complete");
let body = response_json!(response, StatusCode::OK).await;
assert_eq!(run_json_status(&body)["kind"], "failed");
assert_eq!(
run_json_status(&body)["reason"],
expected_reason.to_string()
);
let response = app
.oneshot(
Request::builder()
.method("GET")
.uri(api("/runs"))
.body(Body::empty())
.expect("run list request should build"),
)
.await
.expect("run list request should complete");
let body = response_json!(response, StatusCode::OK).await;
let run_id_string = run_id.to_string();
let listed = body["data"]
.as_array()
.expect("run list data should be an array")
.iter()
.find(|run| run_json_id(run) == Some(run_id_string.as_str()))
.expect("failed run should remain listed");
assert_eq!(run_json_status(listed)["kind"], "failed");
assert_eq!(
run_json_status(listed)["reason"],
expected_reason.to_string()
);
}
#[tokio::test]
async fn subprocess_pre_start_failure_persists_launch_failure_from_runnable() {
assert_subprocess_pre_start_failure(
PreStartWorkerOutcome::LaunchFailure,
FailureReason::LaunchFailed,
)
.await;
}
#[tokio::test]
async fn subprocess_pre_start_failure_persists_early_worker_exit_from_runnable() {
assert_subprocess_pre_start_failure(
PreStartWorkerOutcome::EarlyExit,
FailureReason::Terminated,
)
.await;
}
#[tokio::test]
async fn subprocess_pre_start_failure_preserves_pending_cancellation() {
let runtime = StdArc::new(PreStartWorkerRuntime::held_launch_failure());
let state = subprocess_pre_start_failure_state(StdArc::clone(&runtime));
let app = crate::test_support::build_test_router(Arc::clone(&state));
let run_id = create_and_start_run(&app, MINIMAL_DOT)
.await
.parse::<RunId>()
.expect("created run id should parse");
let execution = tokio::spawn(execute_run(Arc::clone(&state), run_id));
runtime.wait_for_start().await;
let response = app
.oneshot(
Request::builder()
.method("POST")
.uri(api(&format!("/runs/{run_id}/cancel")))
.body(Body::empty())
.expect("cancel request should build"),
)
.await
.expect("cancel request should complete");
assert_status!(response, StatusCode::ACCEPTED).await;
runtime.release_held_start();
execution.await.expect("run execution task should complete");
assert_eq!(runtime.start_count(), 1);
let run_store = state
.stores
.runs
.open_run_reader(&run_id)
.await
.expect("cancelled run should remain readable");
let events = run_store
.list_events()
.await
.expect("cancelled run events should remain readable");
let failure_reasons = events
.iter()
.filter_map(|envelope| match &envelope.event.body {
EventBody::RunFailed(props) => Some(props.failure.reason),
_ => None,
})
.collect::<Vec<_>>();
assert_eq!(failure_reasons, vec![FailureReason::Cancelled]);
assert!(!events.iter().any(|envelope| {
matches!(
&envelope.event.body,
EventBody::RunFailed(props)
if props.failure.reason == FailureReason::LaunchFailed
)
}));
let expected_status = RunStatus::Failed {
reason: FailureReason::Cancelled,
};
assert_eq!(
run_store
.state()
.await
.expect("cancelled run state should load")
.status,
expected_status
);
assert_eq!(
state
.runs
.lock()
.expect("runs lock poisoned")
.get(&run_id)
.expect("managed run should remain present")
.status,
expected_status
);
}
async fn create_durable_run_with_events(
state: &Arc<AppState>,
run_id: RunId,

View file

@ -634,6 +634,29 @@ mod tests {
)
}
fn sandbox_init_failure_payload(label: &str) -> EventPayload {
event_payload(
label,
"2026-03-27T12:00:04Z",
"run.failed",
&serde_json::json!({
"failure": {
"reason": "sandbox_init_failed",
"detail": {
"message": "sandbox initialization failed",
"category": "deterministic"
}
},
"timing": {
"wall_time_ms": 1,
"inference_time_ms": 0,
"tool_time_ms": 0,
"active_time_ms": 0
},
}),
)
}
async fn append_completed(run: &RunDatabase, label: &str, created_at: DateTime<Utc>) {
append_running(run, label, created_at).await;
run.append_event(&event_payload(
@ -977,7 +1000,7 @@ mod tests {
let events_before = run.list_events().await.unwrap();
let err = run
.append_event(&workflow_failure_payload("run-1"))
.append_event(&sandbox_init_failure_payload("run-1"))
.await
.unwrap_err();
@ -989,7 +1012,7 @@ mod tests {
Error::InvalidTransition(fabro_types::InvalidTransition {
from: RunStatus::Runnable,
to: RunStatus::Failed {
reason: FailureReason::WorkflowError,
reason: FailureReason::SandboxInitFailed,
},
})
));
@ -1009,7 +1032,7 @@ mod tests {
append_runnable(&run, "run-1", dt("2026-03-27T12:00:00Z")).await;
let err = run
.append_event(&workflow_failure_payload("run-1"))
.append_event(&sandbox_init_failure_payload("run-1"))
.await
.unwrap_err();
assert!(matches!(err, Error::EventRejected { .. }));

View file

@ -157,7 +157,18 @@ impl RunStatus {
Self::Submitted
)
| (Self::Pending { .. }, Self::Runnable)
| (Self::Runnable, Self::Starting)
// A worker can fail before it appends RunStarting, while the durable run is still
// Runnable. Keep this set aligned with the audited pre-Starting failure callers.
| (
Self::Runnable,
Self::Starting
| Self::Failed {
reason: FailureReason::LaunchFailed
| FailureReason::WorkflowError
| FailureReason::Terminated
| FailureReason::BootstrapFailed,
}
)
| (
Self::Submitted | Self::Pending { .. } | Self::Runnable,
Self::Failed {
@ -417,9 +428,6 @@ mod tests {
assert!(runnable.can_transition_to(RunStatus::Failed {
reason: FailureReason::Cancelled,
}));
assert!(!runnable.can_transition_to(RunStatus::Failed {
reason: FailureReason::Terminated,
}));
assert!(running.can_transition_to(blocked));
assert!(blocked.can_transition_to(running));
assert!(blocked.can_transition_to(paused));
@ -428,6 +436,49 @@ mod tests {
}));
}
#[test]
fn runnable_pre_start_failures_allow_only_audited_reasons() {
for reason in [
FailureReason::LaunchFailed,
FailureReason::WorkflowError,
FailureReason::Terminated,
FailureReason::BootstrapFailed,
] {
let failed = RunStatus::Failed { reason };
assert!(
RunStatus::Runnable.can_transition_to(failed),
"Runnable should accept {reason} before Starting"
);
assert!(
!RunStatus::Submitted.can_transition_to(failed),
"Submitted should reject {reason}"
);
assert!(
!RunStatus::Pending {
reason: PendingReason::ApprovalRequired,
}
.can_transition_to(failed),
"Pending should reject {reason}"
);
}
assert!(RunStatus::Runnable.can_transition_to(RunStatus::Failed {
reason: FailureReason::Cancelled,
}));
for reason in [
FailureReason::SandboxInitFailed,
FailureReason::PublishFailed,
FailureReason::TransientInfra,
FailureReason::BudgetExhausted,
FailureReason::ApprovalDenied,
] {
assert!(
!RunStatus::Runnable.can_transition_to(RunStatus::Failed { reason }),
"Runnable should reject unrelated failure {reason}"
);
}
}
#[test]
fn success_and_failure_reasons_parse_and_round_trip() {
let success = SuccessReason::from_str("completed").expect("completed should parse");