diff --git a/Cargo.lock b/Cargo.lock index 260efaec2..ee3bc2c61 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -3078,6 +3078,7 @@ dependencies = [ "serde_json", "serde_yaml", "sha2 0.10.9", + "slatedb", "sqlx", "strum 0.28.0", "sysinfo", diff --git a/lib/apps/fabro-server/Cargo.toml b/lib/apps/fabro-server/Cargo.toml index 8aea2bade..dd27ebb1a 100644 --- a/lib/apps/fabro-server/Cargo.toml +++ b/lib/apps/fabro-server/Cargo.toml @@ -111,6 +111,7 @@ chrono = { workspace = true } [dev-dependencies] fabro-auth = { path = "../../foundation/fabro-auth", features = ["test-support"] } +slatedb.workspace = true tokio = { workspace = true, features = ["test-util", "macros"] } tower = "0.5" http-body-util = "0.1" diff --git a/lib/apps/fabro-server/src/server.rs b/lib/apps/fabro-server/src/server.rs index bc5aeeeb3..2ec72f85b 100644 --- a/lib/apps/fabro-server/src/server.rs +++ b/lib/apps/fabro-server/src/server.rs @@ -67,7 +67,7 @@ use fabro_llm::types::{ use fabro_mcp_store::McpServerStore; use fabro_model::catalog::LlmCatalogSettings; use fabro_model::{BilledTokenCounts, Catalog, ModelRef, ModelTestMode, ProviderId}; -use fabro_redact::redact_jsonl_line; +use fabro_redact::{DisplaySafeUrl, redact_jsonl_line, redact_string}; use fabro_sandbox::daytona::{self, DaytonaSandbox}; use fabro_sandbox::details::sandbox_details; use fabro_sandbox::reconnect::reconnect_for_run; @@ -98,8 +98,8 @@ use fabro_types::settings::server::{ use fabro_types::{ AgentBackend, AskFabro, AskFabroUnavailableReason, EventBody, InterviewQuestionRecord, PairId, PairMessageId, PairTarget, PendingReason, Principal, PullRequestLink, QuestionType, RunBlobId, - RunControlAction, RunEvent, RunId, RunRunnableSource, SandboxProviderKind, ServerSettings, - SessionCapability, + RunControlAction, RunEvent, RunId, RunProjection, RunRunnableSource, SandboxProviderKind, + ServerSettings, SessionCapability, }; use fabro_util::error::{ SharedError, collect_causes, render_compact_with_causes, render_with_causes, @@ -128,8 +128,7 @@ use tokio::process::Command; use tokio::runtime::Builder as TokioRuntimeBuilder; use tokio::sync::broadcast::error::RecvError; use tokio::sync::{ - Mutex as AsyncMutex, Notify, OwnedMutexGuard, RwLock as AsyncRwLock, Semaphore, broadcast, - mpsc, oneshot, + Mutex as AsyncMutex, Notify, OwnedMutexGuard, RwLock as AsyncRwLock, Semaphore, broadcast, mpsc, }; use tokio::task::spawn_blocking; use tokio::time::{sleep, timeout}; @@ -276,10 +275,9 @@ struct ManagedRun { /// Stage IDs of currently running agent sessions that have no live /// steering capability, keyed to the session id that owns the marker. active_non_steerable_stages: HashMap, - event_tx: Option>, checkpoint: Option, - cancel_tx: Option>, cancel_token: Option, + admission_retry: AdmissionRetryState, worker_ref: Option, /// Exact worker currently covered by a cancellation escalation task. /// Prevents repeated cancel requests from arming duplicate watchdogs. @@ -288,6 +286,78 @@ struct ManagedRun { execution_mode: RunExecutionMode, } +const MAX_ADMISSION_FAILURES: u8 = 5; +const MAX_ADMISSION_RETRY_DELAY: Duration = Duration::from_secs(30); + +#[derive(Clone, Copy, Debug, Default, Eq, PartialEq)] +enum AdmissionRetryState { + #[default] + Ready, + Backoff { + failures: u8, + retry_at: Instant, + }, + GivenUp { + failures: u8, + }, +} + +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +enum AdmissionRetryOutcome { + Retry { failures: u8, delay: Duration }, + GivenUp { failures: u8 }, +} + +impl AdmissionRetryState { + fn is_eligible(self, now: Instant) -> bool { + match self { + Self::Ready => true, + Self::Backoff { retry_at, .. } => now >= retry_at, + Self::GivenUp { .. } => false, + } + } + + fn reset(&mut self) { + *self = Self::Ready; + } + + fn record_transient_failure(&mut self, now: Instant) -> AdmissionRetryOutcome { + let failures = self.failures().saturating_add(1); + if failures >= MAX_ADMISSION_FAILURES { + *self = Self::GivenUp { failures }; + return AdmissionRetryOutcome::GivenUp { failures }; + } + + let exponent = u32::from(failures.saturating_sub(1)); + let delay = Duration::from_secs(1_u64.checked_shl(exponent).unwrap_or(u64::MAX)) + .min(MAX_ADMISSION_RETRY_DELAY); + *self = Self::Backoff { + failures, + retry_at: now + delay, + }; + AdmissionRetryOutcome::Retry { failures, delay } + } + + fn record_permanent_failure(&mut self) -> AdmissionRetryOutcome { + let failures = self.failures().saturating_add(1); + *self = Self::GivenUp { failures }; + AdmissionRetryOutcome::GivenUp { failures } + } + + fn failures(self) -> u8 { + match self { + Self::Ready => 0, + Self::Backoff { failures, .. } | Self::GivenUp { failures } => failures, + } + } +} + +struct AdmittedRun { + run_store: fabro_store::RunDatabase, + run_dir: PathBuf, + cancel_token: CancellationToken, +} + impl ManagedRun { /// True if cancellation should still escalate to `worker_ref`; clears a /// stale escalation marker as a side effect. @@ -2633,9 +2703,6 @@ async fn delete_run_internal( if let Some(answer_transport) = managed_run.answer_transport.clone() { let _ = answer_transport.cancel_run().await; } - if let Some(cancel_tx) = managed_run.cancel_tx.take() { - let _ = cancel_tx.send(()); - } } // Terminal runs can still carry a stale worker ref briefly after their // completion events land, so avoid paying the full cancellation grace. @@ -2887,6 +2954,35 @@ pub(in crate::server) fn counts_toward_scheduler_capacity(status: RunStatus) -> ) } +fn select_runs_for_admission( + runs: &HashMap, + max_concurrent_runs: usize, + now: Instant, +) -> Vec { + let active = runs + .values() + .filter(|run| counts_toward_scheduler_capacity(run.status)) + .count(); + let available = max_concurrent_runs.saturating_sub(active); + if available == 0 { + return Vec::new(); + } + + let mut runnable: Vec<_> = runs + .iter() + .filter(|(_, run)| { + run.status == RunStatus::Runnable && run.admission_retry.is_eligible(now) + }) + .map(|(id, run)| (*id, run.created_at)) + .collect(); + runnable.sort_by_key(|(_, created_at)| *created_at); + runnable + .into_iter() + .take(available) + .map(|(id, _)| id) + .collect() +} + #[allow( clippy::result_large_err, reason = "Run ID parsing returns HTTP 400 responses directly." @@ -2982,8 +3078,6 @@ fn clear_live_run_state(run: &mut ManagedRun) { run.active_api_targets.clear(); run.active_steerable_stages.clear(); run.active_non_steerable_stages.clear(); - run.event_tx = None; - run.cancel_tx = None; run.cancel_token = None; run.worker_ref = None; run.cancel_escalation_worker = None; @@ -3229,6 +3323,13 @@ async fn alive_refs(state: &AppState, refs: &[WorkerRef]) -> Vec { async fn persist_cancelled_run_status(state: &AppState, run_id: RunId) -> anyhow::Result<()> { let run_store = state.stores.runs.open_run(&run_id).await?; + persist_cancelled_run_status_to_store(&run_store, run_id).await +} + +async fn persist_cancelled_run_status_to_store( + run_store: &fabro_store::RunDatabase, + run_id: RunId, +) -> anyhow::Result<()> { let run_state = run_store.state().await?; if run_state.status.is_terminal() { return Ok(()); @@ -3243,19 +3344,45 @@ async fn persist_cancelled_run_status(state: &AppState, run_id: RunId) -> anyhow None, None, ); - workflow_event::append_event(&run_store, &run_id, &failure_event).await + workflow_event::append_event(run_store, &run_id, &failure_event).await } -async fn finish_cancelled_run_before_execution(state: &Arc, run_id: RunId) { - if let Err(err) = persist_cancelled_run_status(state.as_ref(), run_id).await { - error!(run_id = %run_id, error = %err, "Failed to persist cancelled run status"); +async fn finish_cancelled_run_before_execution( + state: &Arc, + run_store: &fabro_store::RunDatabase, + run_id: RunId, +) { + let persist_result = persist_cancelled_run_status_to_store(run_store, run_id).await; + let durable_status = match run_store.state().await { + Ok(projection) => Some(projection.status), + Err(err) => { + error!( + run_id = %run_id, + error = %safe_error_chain(&anyhow::Error::new(err)), + "Failed to load durable status after pre-launch cancellation" + ); + None + } + }; + if !durable_status.is_some_and(RunStatus::is_terminal) { + if let Err(err) = &persist_result { + error!( + run_id = %run_id, + error = %safe_error_chain(err), + "Failed to persist cancelled run status" + ); + } } let mut runs = state.runs.lock().expect("runs lock poisoned"); if let Some(managed_run) = runs.get_mut(&run_id) { - managed_run.status = RunStatus::Failed { - reason: FailureReason::Cancelled, - }; + if let Some(status) = durable_status { + managed_run.status = status; + } else if persist_result.is_ok() { + managed_run.status = RunStatus::Failed { + reason: FailureReason::Cancelled, + }; + } clear_live_run_state(managed_run); } drop(runs); @@ -3267,6 +3394,7 @@ async fn finish_cancelled_run_before_execution(state: &Arc, run_id: Ru /// disabled by server policy. Returns `true` when the run was rejected. async fn reject_run_if_sandbox_provider_disabled( state: &Arc, + run_store: &fabro_store::RunDatabase, server_settings: &ServerSettings, run_id: RunId, settings: &RunNamespace, @@ -3276,36 +3404,48 @@ async fn reject_run_if_sandbox_provider_disabled( return false; }; tracing::warn!(run_id = %run_id, error = %error, "Sandbox provider disabled by server policy"); - fail_run_before_execution(state, run_id, FailureReason::LaunchFailed, error).await; + fail_run_before_execution( + state, + run_store, + run_id, + FailureReason::LaunchFailed, + WorkflowError::engine(error), + ) + .await; true } async fn fail_run_before_execution( state: &Arc, + run_store: &fabro_store::RunDatabase, run_id: RunId, reason: FailureReason, - message: String, + error: WorkflowError, ) { - match state.stores.runs.open_run(&run_id).await { - Ok(run_store) => { - let failure_event = workflow_event::Event::workflow_run_failed_from_error( - &WorkflowError::engine(message.clone()), - fabro_types::RunTiming::default(), - reason, - None, - None, - None, - None, - ); - if let Err(err) = - workflow_event::append_event(&run_store, &run_id, &failure_event).await - { - error!(run_id = %run_id, error = %err, "Failed to persist run failure status"); - } - } - Err(err) => { - error!(run_id = %run_id, error = %err, "Failed to open run store while persisting run failure"); + let message = error.display_with_causes(); + let failure_event = workflow_event::Event::workflow_run_failed_from_error( + &error, + fabro_types::RunTiming::default(), + reason, + None, + None, + None, + None, + ); + if let Err(err) = workflow_event::append_event(run_store, &run_id, &failure_event).await { + error!( + run_id = %run_id, + error = %safe_error_chain(&err), + "Failed to persist run failure status" + ); + let mut runs = state.runs.lock().expect("runs lock poisoned"); + if let Some(managed_run) = runs.get_mut(&run_id) { + clear_live_run_state(managed_run); } + drop(runs); + cleanup_worker_control_bus_for_run(state.as_ref(), run_id); + state.scheduler_notify.notify_one(); + return; } fail_managed_run(state, run_id, reason, message); @@ -3349,10 +3489,9 @@ fn managed_run( active_api_targets: HashMap::new(), active_steerable_stages: HashMap::new(), active_non_steerable_stages: HashMap::new(), - event_tx: None, checkpoint: None, - cancel_tx: None, cancel_token: None, + admission_retry: AdmissionRetryState::Ready, worker_ref: None, cancel_escalation_worker: None, run_dir: Some(run_dir), @@ -3581,19 +3720,14 @@ 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 failure_event = workflow_event::Event::workflow_run_failed_from_error( - &WorkflowError::engine_with_anyhow("Failed to spawn worker", err), - fabro_types::RunTiming::default(), + fail_run_before_execution( + state, + run_store, + run_id, FailureReason::LaunchFailed, - 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); - state.scheduler_notify.notify_one(); + WorkflowError::engine_with_anyhow("Failed to spawn worker", err), + ) + .await; } async fn append_worker_exit_failure( @@ -3898,6 +4032,265 @@ fn answer_from_request( } } +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +enum AdmissionErrorClass { + Transient(&'static str), + Permanent(&'static str), +} + +fn classify_admission_error(error: &anyhow::Error) -> AdmissionErrorClass { + let Some(store_error) = error.downcast_ref::() else { + return AdmissionErrorClass::Permanent("unrecognized"); + }; + match store_error { + fabro_store::Error::Slate(_) => AdmissionErrorClass::Transient("slate"), + fabro_store::Error::ObjectStore(_) => AdmissionErrorClass::Transient("object_store"), + fabro_store::Error::Sqlite(_) => AdmissionErrorClass::Transient("sqlite"), + fabro_store::Error::Io(_) => AdmissionErrorClass::Transient("io"), + fabro_store::Error::InvalidEvent(_) => AdmissionErrorClass::Permanent("invalid_event"), + fabro_store::Error::EventRejected { .. } => { + AdmissionErrorClass::Permanent("event_rejected") + } + fabro_store::Error::RunNotFound(_) => AdmissionErrorClass::Permanent("run_not_found"), + fabro_store::Error::SessionNotFound(_) => { + AdmissionErrorClass::Permanent("session_not_found") + } + fabro_store::Error::SessionAlreadyExists(_) => { + AdmissionErrorClass::Permanent("session_already_exists") + } + fabro_store::Error::ReadOnly => AdmissionErrorClass::Permanent("read_only"), + fabro_store::Error::EventSequenceExhausted { .. } => { + AdmissionErrorClass::Permanent("event_sequence_exhausted") + } + fabro_store::Error::InvalidKeySegment { .. } => { + AdmissionErrorClass::Permanent("invalid_key_segment") + } + fabro_store::Error::KeyParse(_) => AdmissionErrorClass::Permanent("key_parse"), + fabro_store::Error::RunSummaryMismatch { .. } => { + AdmissionErrorClass::Permanent("run_summary_mismatch") + } + fabro_store::Error::InvalidTransition(_) => { + AdmissionErrorClass::Permanent("invalid_transition") + } + fabro_store::Error::Serde(_) => AdmissionErrorClass::Permanent("serialization"), + fabro_store::Error::Other(_) => AdmissionErrorClass::Permanent("internal"), + } +} + +fn safe_error_chain(error: &anyhow::Error) -> String { + static URL_PATTERN: LazyLock = LazyLock::new(|| { + regex::Regex::new(r#"https?://[^\s\[\](){}<>\"']+"#) + .expect("embedded URL redaction regex should compile") + }); + + let rendered = render_with_causes(&error.to_string(), &collect_causes(error.as_ref())); + let url_safe = URL_PATTERN.replace_all(&rendered, |captures: ®ex::Captures<'_>| { + DisplaySafeUrl::parse(&captures[0]) + .map_or_else(|_| "".to_string(), |url| url.redacted_string()) + }); + redact_string(&url_safe) +} + +fn projection_failure_message(projection: &RunProjection) -> Option { + projection.conclusion.as_ref().and_then(|conclusion| { + conclusion.failure.as_ref().map(|failure| { + render_compact_with_causes(&failure.detail.message, &failure.detail.causes) + }) + }) +} + +fn record_admission_failure(state: &Arc, run_id: RunId, error: &anyhow::Error) { + let class = classify_admission_error(error); + let rendered_error = safe_error_chain(error); + let outcome = { + let mut runs = state.runs.lock().expect("runs lock poisoned"); + let Some(managed_run) = runs.get_mut(&run_id) else { + return; + }; + if managed_run.status == RunStatus::Starting { + managed_run.status = RunStatus::Runnable; + } + clear_live_run_state(managed_run); + match class { + AdmissionErrorClass::Transient(_) => managed_run + .admission_retry + .record_transient_failure(Instant::now()), + AdmissionErrorClass::Permanent(_) => { + managed_run.admission_retry.record_permanent_failure() + } + } + }; + + let error_kind = match class { + AdmissionErrorClass::Transient(kind) | AdmissionErrorClass::Permanent(kind) => kind, + }; + match outcome { + AdmissionRetryOutcome::Retry { failures, delay } => { + warn!( + run_id = %run_id, + failure_count = failures, + max_failures = MAX_ADMISSION_FAILURES, + retry_delay_ms = u64::try_from(delay.as_millis()).unwrap_or(u64::MAX), + error_kind, + error = %rendered_error, + "Run admission claim failed; retrying" + ); + } + AdmissionRetryOutcome::GivenUp { failures } => { + error!( + run_id = %run_id, + failure_count = failures, + max_failures = MAX_ADMISSION_FAILURES, + error_kind, + gave_up = true, + error = %rendered_error, + "Run admission claim failed; giving up" + ); + } + } + state.scheduler_notify.notify_one(); +} + +async fn admit_run(state: &Arc, run_id: RunId) -> Option { + let (run_dir, cancel_token) = { + let mut runs = state.runs.lock().expect("runs lock poisoned"); + if state.is_shutting_down() { + return None; + } + let managed_run = runs.get_mut(&run_id)?; + if managed_run.status != RunStatus::Runnable + || !managed_run.admission_retry.is_eligible(Instant::now()) + { + return None; + } + + let cancel_token = CancellationToken::new(); + managed_run.status = RunStatus::Starting; + managed_run.cancel_token = Some(cancel_token.clone()); + (managed_run.run_dir.clone(), cancel_token) + }; + + let run_store = match state.stores.runs.open_run(&run_id).await { + Ok(run_store) => run_store, + Err(error) => { + record_admission_failure(state, run_id, &anyhow::Error::new(error)); + return None; + } + }; + let run_events = run_store.subscribe(); + let mut observed_status = None; + let mut observed_error = None; + let claimed = workflow_event::append_event_if( + &run_store, + &run_id, + &workflow_event::Event::RunStarting, + |projection| { + if projection.status == RunStatus::Runnable { + true + } else { + observed_status = Some(projection.status); + observed_error = projection_failure_message(projection); + false + } + }, + ) + .await; + + match claimed { + Ok(true) => { + let mut runs = state.runs.lock().expect("runs lock poisoned"); + if let Some(managed_run) = runs.get_mut(&run_id) { + managed_run.admission_retry.reset(); + } + drop(runs); + tokio::spawn(forward_run_events_to_global( + Arc::clone(state), + run_id, + run_events, + )); + } + Ok(false) => { + let mut runs = state.runs.lock().expect("runs lock poisoned"); + if let Some(managed_run) = runs.get_mut(&run_id) { + if let Some(status) = observed_status { + managed_run.status = status; + } + managed_run.error = observed_error; + managed_run.admission_retry.reset(); + clear_live_run_state(managed_run); + } + drop(runs); + cleanup_worker_control_bus_for_run(state.as_ref(), run_id); + state.scheduler_notify.notify_one(); + return None; + } + Err(error) => { + record_admission_failure(state, run_id, &error); + return None; + } + } + + let Some(run_dir) = run_dir else { + fail_run_before_execution( + state, + &run_store, + run_id, + FailureReason::WorkflowError, + WorkflowError::engine("Run directory is unavailable before execution"), + ) + .await; + return None; + }; + + let projection = match run_store.state().await { + Ok(projection) => projection, + Err(error) => { + fail_run_before_execution( + state, + &run_store, + run_id, + FailureReason::WorkflowError, + WorkflowError::engine_with_source("Failed to load admitted run state", error), + ) + .await; + return None; + } + }; + if projection.pending_control == Some(RunControlAction::Cancel) || cancel_token.is_cancelled() { + finish_cancelled_run_before_execution(state, &run_store, run_id).await; + return None; + } + + Some(AdmittedRun { + run_store, + run_dir, + cancel_token, + }) +} + +fn execution_mode_for_run(state: &AppState, run_id: RunId) -> Option { + state + .runs + .lock() + .expect("runs lock poisoned") + .get(&run_id) + .map(|managed_run| managed_run.execution_mode) +} + +fn launch_is_allowed(state: &AppState, run_id: RunId, cancel_token: &CancellationToken) -> bool { + if cancel_token.is_cancelled() { + return false; + } + let runs = state.runs.lock().expect("runs lock poisoned"); + runs.get(&run_id).is_some_and(|managed_run| { + managed_run.status == RunStatus::Starting + && managed_run + .cancel_token + .as_ref() + .is_some_and(|token| !token.is_cancelled()) + }) +} + /// Execute a single run: transitions runnable → starting → running → /// completed/failed/cancelled. async fn execute_run(state: Arc, run_id: RunId) { @@ -3905,53 +4298,44 @@ async fn execute_run(state: Arc, run_id: RunId) { return; } + let Some(admitted) = Box::pin(admit_run(&state, run_id)).await else { + return; + }; + if state.registry_factory_override.is_some() { - Box::pin(execute_run_in_process(state, run_id)).await; + Box::pin(execute_run_in_process(state, run_id, admitted)).await; return; } - Box::pin(execute_run_subprocess(state, run_id)).await; + Box::pin(execute_run_subprocess(state, run_id, admitted)).await; } -async fn execute_run_in_process(state: Arc, run_id: RunId) { - // Transition to Starting and set up cancel infrastructure - let (cancel_rx, run_dir, event_tx, cancel_token, execution_mode) = { - let mut runs = state.runs.lock().expect("runs lock poisoned"); - let managed_run = match runs.get_mut(&run_id) { - Some(r) if r.status == RunStatus::Runnable => r, - _ => return, - }; - let Some(run_dir) = managed_run.run_dir.clone() else { - return; - }; - - let (cancel_tx, cancel_rx) = oneshot::channel::<()>(); - let cancel_token = CancellationToken::new(); - let (event_tx, _) = broadcast::channel(256); - - managed_run.status = RunStatus::Starting; - managed_run.cancel_tx = Some(cancel_tx); - managed_run.cancel_token = Some(cancel_token.clone()); - managed_run.event_tx = Some(event_tx); - - ( - cancel_rx, - run_dir, - managed_run.event_tx.clone(), - cancel_token, - managed_run.execution_mode, +async fn execute_run_in_process(state: Arc, run_id: RunId, admitted: AdmittedRun) { + let AdmittedRun { + run_store, + run_dir, + cancel_token, + } = admitted; + if !launch_is_allowed(state.as_ref(), run_id, &cancel_token) { + finish_cancelled_run_before_execution(&state, &run_store, run_id).await; + return; + } + let Some(execution_mode) = execution_mode_for_run(state.as_ref(), run_id) else { + fail_run_before_execution( + &state, + &run_store, + run_id, + FailureReason::WorkflowError, + WorkflowError::engine("Managed run disappeared before in-process execution"), ) + .await; + return; }; // Create interviewer and event plumbing (this is the "provisioning" phase) let interviewer = Arc::new(ControlInterviewer::new()); let interview_runtime: Arc = interviewer.clone(); let emitter = Emitter::new(run_id); - if let Some(tx_clone) = event_tx { - emitter.on_event(move |event| { - let _ = tx_clone.send(event.clone()); - }); - } let registry_override = state .registry_factory_override .as_ref() @@ -3959,12 +4343,12 @@ async fn execute_run_in_process(state: Arc, run_id: RunId) { let emitter = Arc::new(emitter); let steering_hub = Arc::new(fabro_workflow::SteeringHub::new(Arc::clone(&emitter))); - // Transition to Running, populate interviewer + // Publish the in-process control transport while preserving Starting + // until the workflow emits its durable Running event. let cancelled_during_setup = { let mut runs = state.runs.lock().expect("runs lock poisoned"); if let Some(managed_run) = runs.get_mut(&run_id) { - if managed_run.status == RunStatus::Starting { - managed_run.status = RunStatus::Running; + if managed_run.status == RunStatus::Starting && !cancel_token.is_cancelled() { managed_run.answer_transport = Some(RunAnswerTransport::InProcess { interviewer: Arc::clone(&interviewer), steering_hub: Arc::clone(&steering_hub), @@ -3977,46 +4361,23 @@ async fn execute_run_in_process(state: Arc, run_id: RunId) { true } } else { - false + true } }; if cancelled_during_setup { - if let Err(err) = persist_cancelled_run_status(state.as_ref(), run_id).await { - error!(run_id = %run_id, error = %err, "Failed to persist cancelled run status"); - } + finish_cancelled_run_before_execution(&state, &run_store, run_id).await; return; } - - let run_store = match state.stores.runs.open_run(&run_id).await { - Ok(run_store) => run_store, - Err(e) => { - tracing::error!(run_id = %run_id, error = %e, "Failed to open run store"); - let mut runs = state.runs.lock().expect("runs lock poisoned"); - if let Some(managed_run) = runs.get_mut(&run_id) { - managed_run.status = RunStatus::Failed { - reason: FailureReason::WorkflowError, - }; - managed_run.error = Some(format!("Failed to open run store: {e}")); - clear_live_run_state(managed_run); - } - state.scheduler_notify.notify_one(); - return; - } - }; - tokio::spawn(forward_run_events_to_global( - Arc::clone(&state), - run_id, - run_store.subscribe(), - )); let persisted = match Persisted::load_from_store(&run_store.clone().into(), &run_dir).await { Ok(persisted) => persisted, Err(e) => { tracing::error!(run_id = %run_id, error = %e, "Failed to load persisted run"); fail_run_before_execution( &state, + &run_store, run_id, FailureReason::WorkflowError, - format!("Failed to load persisted run: {e}"), + WorkflowError::engine_with_source("Failed to load persisted run", e), ) .await; return; @@ -4025,11 +4386,12 @@ async fn execute_run_in_process(state: Arc, run_id: RunId) { let server_settings = state.server_settings(); let github_settings = &server_settings.server.integrations.github; if cancel_token.is_cancelled() { - finish_cancelled_run_before_execution(&state, run_id).await; + finish_cancelled_run_before_execution(&state, &run_store, run_id).await; return; } if reject_run_if_sandbox_provider_disabled( &state, + &run_store, &server_settings, run_id, &persisted.run_spec().settings.run, @@ -4070,15 +4432,16 @@ async fn execute_run_in_process(state: Arc, run_id: RunId) { Ok(github_app) => github_app, Err(e) => { if cancel_token.is_cancelled() { - finish_cancelled_run_before_execution(&state, run_id).await; + finish_cancelled_run_before_execution(&state, &run_store, run_id).await; return; } tracing::error!(run_id = %run_id, error = %e, "Invalid GitHub credentials"); fail_run_before_execution( &state, + &run_store, run_id, FailureReason::WorkflowError, - format!("Invalid GitHub credentials: {e}"), + WorkflowError::engine_with_source("Invalid GitHub credentials", e), ) .await; return; @@ -4101,9 +4464,10 @@ async fn execute_run_in_process(state: Arc, run_id: RunId) { ); fail_run_before_execution( &state, + &run_store, run_id, FailureReason::WorkflowError, - format!("Failed to resolve GitHub permissions: {err}"), + WorkflowError::engine_with_source("Failed to resolve GitHub permissions", err), ) .await; return; @@ -4115,9 +4479,10 @@ async fn execute_run_in_process(state: Arc, run_id: RunId) { tracing::error!(run_id = %run_id, error = ?err, "Loading run secrets failed"); fail_run_before_execution( &state, + &run_store, run_id, FailureReason::WorkflowError, - "Loading run secrets failed".to_string(), + WorkflowError::engine_with_source("Loading run secrets failed", err), ) .await; return; @@ -4142,6 +4507,11 @@ async fn execute_run_in_process(state: Arc, run_id: RunId) { fabro_run_tools: None, }; + if !launch_is_allowed(state.as_ref(), run_id, &cancel_token) { + finish_cancelled_run_before_execution(&state, &run_store, run_id).await; + return; + } + let execution = async { match execution_mode { RunExecutionMode::Start => operations::start(&run_dir, services).await, @@ -4151,14 +4521,11 @@ async fn execute_run_in_process(state: Arc, run_id: RunId) { let result = tokio::select! { result = execution => ExecutionResult::Completed(Box::new(result)), - _ = cancel_rx => { - cancel_token.cancel(); - ExecutionResult::CancelledBySignal - } + () = cancel_token.cancelled() => ExecutionResult::CancelledBySignal, }; if matches!(&result, ExecutionResult::CancelledBySignal) { - if let Err(err) = persist_cancelled_run_status(state.as_ref(), run_id).await { + if let Err(err) = persist_cancelled_run_status_to_store(&run_store, run_id).await { error!(run_id = %run_id, error = %err, "Failed to persist cancelled run status"); } } @@ -4236,60 +4603,51 @@ async fn execute_run_in_process(state: Arc, run_id: RunId) { state.scheduler_notify.notify_one(); } -async fn execute_run_subprocess(state: Arc, run_id: RunId) { - let (run_dir, execution_mode) = { - let mut runs = state.runs.lock().expect("runs lock poisoned"); - if state.is_shutting_down() { - return; - } - let managed_run = match runs.get_mut(&run_id) { - Some(run) if run.status == RunStatus::Runnable => run, - _ => return, - }; - let Some(run_dir) = managed_run.run_dir.clone() else { - return; - }; - managed_run.status = RunStatus::Starting; - (run_dir, managed_run.execution_mode) +async fn execute_run_subprocess(state: Arc, run_id: RunId, admitted: AdmittedRun) { + let AdmittedRun { + run_store, + run_dir, + cancel_token, + } = admitted; + if !launch_is_allowed(state.as_ref(), run_id, &cancel_token) { + finish_cancelled_run_before_execution(&state, &run_store, run_id).await; + return; + } + let Some(execution_mode) = execution_mode_for_run(state.as_ref(), run_id) else { + fail_run_before_execution( + &state, + &run_store, + run_id, + FailureReason::WorkflowError, + WorkflowError::engine("Managed run disappeared before subprocess execution"), + ) + .await; + return; }; - let run_store = match state.stores.runs.open_run(&run_id).await { - Ok(run_store) => run_store, - Err(err) => { - tracing::error!(run_id = %run_id, error = %err, "Failed to open run store"); - fail_managed_run( - &state, - run_id, - FailureReason::WorkflowError, - format!("Failed to open run store: {err}"), - ); - state.scheduler_notify.notify_one(); - return; - } - }; - tokio::spawn(forward_run_events_to_global( - Arc::clone(&state), - run_id, - run_store.subscribe(), - )); - let run_state = match run_store.state().await { Ok(run_state) => run_state, Err(err) => { tracing::error!(run_id = %run_id, error = %err, "Failed to load run state"); - fail_managed_run( + fail_run_before_execution( &state, + &run_store, run_id, FailureReason::WorkflowError, - format!("Failed to load run state: {err}"), - ); - state.scheduler_notify.notify_one(); + WorkflowError::engine_with_source("Failed to load run state", err), + ) + .await; return; } }; let agent_fabro_tools_enabled = run_state.spec.settings.run.agent.fabro_tools; + if cancel_token.is_cancelled() { + finish_cancelled_run_before_execution(&state, &run_store, run_id).await; + return; + } if reject_run_if_sandbox_provider_disabled( &state, + &run_store, &state.server_settings(), run_id, &run_state.spec.settings.run, @@ -4304,15 +4662,19 @@ async fn execute_run_subprocess(state: Arc, run_id: RunId) { Err(err) => { fail_run_before_execution( &state, + &run_store, run_id, FailureReason::WorkflowError, - "Loading worker secrets failed".to_string(), + WorkflowError::engine_with_source("Loading worker secrets failed", err), ) .await; - tracing::error!(run_id = %run_id, error = ?err, "Loading worker secrets failed"); return; } }; + if cancel_token.is_cancelled() { + finish_cancelled_run_before_execution(&state, &run_store, run_id).await; + return; + } let state_for_build = Arc::clone(&state); let run_dir_for_build = run_dir.clone(); let start_result = spawn_blocking(move || { @@ -4329,10 +4691,28 @@ async fn execute_run_subprocess(state: Arc, run_id: RunId) { .context("worker_launch_spec task failed") .and_then(|inner| inner); - let launch_result = match start_result { - Ok(spec) => state.worker_runtime.start(spec).await, - Err(err) => Err(err), + let spec = match start_result { + Ok(spec) => spec, + Err(err) => { + fail_run_before_execution( + &state, + &run_store, + run_id, + FailureReason::LaunchFailed, + WorkflowError::engine_with_anyhow( + "Failed to build worker launch specification", + err, + ), + ) + .await; + return; + } }; + if !launch_is_allowed(state.as_ref(), run_id, &cancel_token) { + finish_cancelled_run_before_execution(&state, &run_store, run_id).await; + return; + } + let launch_result = state.worker_runtime.start(spec).await; let started_worker = match launch_result { Ok(worker) => worker, Err(err) => { @@ -4353,6 +4733,14 @@ async fn execute_run_subprocess(state: Arc, run_id: RunId) { }); } } + if cancel_token.is_cancelled() { + state.worker_runtime.request_stop(&worker_ref).await; + handler::lifecycle::schedule_worker_cancel_escalation( + Arc::clone(&state), + run_id, + worker_ref.clone(), + ); + } let stderr_task = tokio::spawn(drain_worker_stderr(run_id, started_worker.stderr)); @@ -4469,26 +4857,7 @@ pub fn spawn_scheduler(state: Arc) { } let runs_to_start = { let runs = state.runs.lock().expect("runs lock poisoned"); - let active = runs - .values() - .filter(|r| counts_toward_scheduler_capacity(r.status)) - .count(); - let available = state.max_concurrent_runs.saturating_sub(active); - if available == 0 { - Vec::new() - } else { - let mut runnable: Vec<_> = runs - .iter() - .filter(|(_, r)| r.status == RunStatus::Runnable) - .map(|(id, r)| (*id, r.created_at)) - .collect(); - runnable.sort_by_key(|(_, created_at)| *created_at); - runnable - .into_iter() - .take(available) - .map(|(id, _)| id) - .collect::>() - } + select_runs_for_admission(&runs, state.max_concurrent_runs, Instant::now()) }; for id in runs_to_start { if state.is_shutting_down() { diff --git a/lib/apps/fabro-server/src/server/handler/lifecycle.rs b/lib/apps/fabro-server/src/server/handler/lifecycle.rs index fe4020aba..f17ee3ec7 100644 --- a/lib/apps/fabro-server/src/server/handler/lifecycle.rs +++ b/lib/apps/fabro-server/src/server/handler/lifecycle.rs @@ -350,7 +350,11 @@ async fn deny_run( run_response(state.as_ref(), id, StatusCode::OK).await } -fn schedule_worker_cancel_escalation(state: Arc, run_id: RunId, worker_ref: WorkerRef) { +pub(in crate::server) fn schedule_worker_cancel_escalation( + state: Arc, + run_id: RunId, + worker_ref: WorkerRef, +) { let requested_at = Instant::now(); let armed = { let mut runs = state.runs.lock().expect("runs lock poisoned"); @@ -487,16 +491,15 @@ async fn cancel_run( reason: FailureReason::Cancelled, }; } - let cancel_tx = if should_cancel_pending_interview { + let cancel_token = if should_cancel_pending_interview { None } else { - managed_run.cancel_tx.take() + managed_run.cancel_token.clone() }; Some(( persist_cancelled_status, answer_transport, - managed_run.cancel_token.clone(), - cancel_tx, + cancel_token, managed_run.worker_ref.clone(), )) } @@ -509,7 +512,7 @@ async fn cancel_run( None => None, } }; - let Some((persist_cancelled_status, answer_transport, cancel_token, cancel_tx, worker_ref)) = + let Some((persist_cancelled_status, answer_transport, cancel_token, worker_ref)) = cancel_target else { return unmanaged_cancel_response(state.as_ref(), id, actor, pending_control).await; @@ -527,20 +530,8 @@ async fn cancel_run( if let Some(token) = &cancel_token { token.cancel(); } - let sent_in_process_cancel = if let Some(cancel_tx) = cancel_tx { - let _ = cancel_tx.send(()); - true - } else { - false - }; let delivered_control = if let Some(answer_transport) = answer_transport { - if sent_in_process_cancel - && matches!(answer_transport, RunAnswerTransport::InProcess { .. }) - { - true - } else { - answer_transport.cancel_run().await.is_ok() - } + answer_transport.cancel_run().await.is_ok() } else { false }; diff --git a/lib/apps/fabro-server/src/server/tests.rs b/lib/apps/fabro-server/src/server/tests.rs index 24c355641..c3d60aeb3 100644 --- a/lib/apps/fabro-server/src/server/tests.rs +++ b/lib/apps/fabro-server/src/server/tests.rs @@ -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; @@ -2945,6 +2945,7 @@ fn worker_token_claims(cmd: &Command, state: &AppState) -> crate::worker_token:: #[derive(Default)] struct RecordingWorkerRuntime { + starts: AtomicUsize, requested: StdMutex>, forced: StdMutex>, alive: AtomicBool, @@ -2952,6 +2953,10 @@ struct RecordingWorkerRuntime { } impl RecordingWorkerRuntime { + fn start_count(&self) -> usize { + self.starts.load(Ordering::Relaxed) + } + fn requested_refs(&self) -> Vec { self.requested .lock() @@ -2985,7 +2990,11 @@ impl RecordingWorkerRuntime { #[async_trait::async_trait] impl WorkerRuntime for RecordingWorkerRuntime { async fn start(&self, _spec: WorkerLaunchSpec) -> anyhow::Result { - anyhow::bail!("recording runtime does not start workers") + self.starts.fetch_add(1, Ordering::Relaxed); + Err( + anyhow::Error::new(std::io::Error::other("recording runtime source failure")) + .context("recording runtime does not start workers"), + ) } async fn request_stop(&self, worker_ref: &WorkerRef) { @@ -5271,8 +5280,6 @@ async fn delete_terminal_managed_run_does_not_send_cancel_signal() { RunExecutionMode::Start, ); run.cancel_token = Some(cancel_token.clone()); - let (cancel_tx, _cancel_rx) = oneshot::channel(); - run.cancel_tx = Some(cancel_tx); state .runs .lock() @@ -15541,7 +15548,6 @@ async fn cancel_durably_blocked_in_process_run_cancels_pending_interview_without let ask = tokio::spawn(async move { ask_interviewer.ask(question).await }); tokio::task::yield_now().await; - let (cancel_tx, mut cancel_rx) = oneshot::channel(); let cancel_token = CancellationToken::new(); let temp_dir = tempfile::tempdir().unwrap(); let mut run = managed_run( @@ -15557,8 +15563,7 @@ async fn cancel_durably_blocked_in_process_run_cancels_pending_interview_without fabro_workflow::event::Emitter::new(run_id), ))), }); - run.cancel_token = Some(cancel_token); - run.cancel_tx = Some(cancel_tx); + run.cancel_token = Some(cancel_token.clone()); state .runs .lock() @@ -15579,10 +15584,7 @@ async fn cancel_durably_blocked_in_process_run_cancels_pending_interview_without .expect("interview task should not panic"); assert_eq!(submission.answer.value, AnswerValue::Cancelled); assert!( - matches!( - cancel_rx.try_recv(), - Err(tokio::sync::oneshot::error::TryRecvError::Empty) - ), + !cancel_token.is_cancelled(), "blocked in-process cancellation should let the workflow unwind instead of aborting it" ); } @@ -16343,6 +16345,634 @@ fn scheduler_capacity_counts_only_runs_occupying_slots() { assert!(!counts_toward_scheduler_capacity(RunStatus::Dead)); } +async fn stage_runnable_managed_run( + state: &Arc, + run_id: RunId, + run_dir: &Path, +) -> fabro_store::RunDatabase { + create_durable_run_with_events(state, run_id, &[workflow_event::Event::RunRunnable { + source: RunRunnableSource::StartRequested, + actor: None, + }]) + .await; + state.runs.lock().expect("runs lock poisoned").insert( + run_id, + managed_run( + MINIMAL_DOT.to_string(), + RunStatus::Runnable, + Utc::now(), + run_dir.to_path_buf(), + RunExecutionMode::Start, + ), + ); + state.stores.runs.open_run(&run_id).await.unwrap() +} + +fn starting_event_count(events: &[EventEnvelope]) -> usize { + events + .iter() + .filter(|event| matches!(event.event.body, EventBody::RunStarting(_))) + .count() +} + +#[tokio::test] +async fn concurrent_durable_admission_has_one_winner() { + let state = test_app_state(); + let direct_run_id = fixtures::RUN_1; + create_durable_run_with_events(&state, direct_run_id, &[ + workflow_event::Event::RunRunnable { + source: RunRunnableSource::StartRequested, + actor: None, + }, + ]) + .await; + let run_store = state.stores.runs.open_run(&direct_run_id).await.unwrap(); + let claim = || { + workflow_event::append_event_if( + &run_store, + &direct_run_id, + &workflow_event::Event::RunStarting, + |projection| projection.status == RunStatus::Runnable, + ) + }; + let (first, second) = tokio::join!(claim(), claim()); + let mut outcomes = [first.unwrap(), second.unwrap()]; + outcomes.sort_unstable(); + assert_eq!(outcomes, [false, true]); + assert_eq!(run_store.state().await.unwrap().status, RunStatus::Starting); + assert_eq!( + starting_event_count(&run_store.list_events().await.unwrap()), + 1 + ); + + let temp = tempfile::tempdir().unwrap(); + let managed_run_id = fixtures::RUN_2; + stage_runnable_managed_run(&state, managed_run_id, temp.path()).await; + let (first, second) = tokio::join!( + Box::pin(admit_run(&state, managed_run_id)), + Box::pin(admit_run(&state, managed_run_id)), + ); + assert_eq!( + usize::from(first.is_some()) + usize::from(second.is_some()), + 1 + ); + let run_store = state.stores.runs.open_run(&managed_run_id).await.unwrap(); + assert_eq!(run_store.state().await.unwrap().status, RunStatus::Starting); + assert_eq!( + starting_event_count(&run_store.list_events().await.unwrap()), + 1 + ); +} + +#[tokio::test] +async fn admission_predicate_false_reconciles_durable_cancellation() { + let state = test_app_state(); + let run_id = fixtures::RUN_1; + create_durable_run_with_events(&state, run_id, &[ + workflow_event::Event::WorkflowRunFailed { + failure: fabro_types::RunFailure { + reason: FailureReason::Cancelled, + detail: FailureDetail::new("cancelled", FailureCategory::Canceled), + }, + timing: fabro_types::RunTiming::default(), + final_git_commit_sha: None, + final_patch: None, + diff_summary: None, + billing: None, + }, + ]) + .await; + let run_store = state.stores.runs.open_run(&run_id).await.unwrap(); + let event_count = run_store.list_events().await.unwrap().len(); + let temp = tempfile::tempdir().unwrap(); + let mut run = managed_run( + MINIMAL_DOT.to_string(), + RunStatus::Runnable, + Utc::now(), + temp.path().to_path_buf(), + RunExecutionMode::Start, + ); + run.cancel_token = Some(CancellationToken::new()); + run.worker_ref = Some(test_worker_ref(123)); + state + .runs + .lock() + .expect("runs lock poisoned") + .insert(run_id, run); + + assert!(Box::pin(admit_run(&state, run_id)).await.is_none()); + + let runs = state.runs.lock().expect("runs lock poisoned"); + let run = runs.get(&run_id).unwrap(); + assert_eq!(run.status, RunStatus::Failed { + reason: FailureReason::Cancelled, + }); + assert!(run.cancel_token.is_none()); + assert!(run.worker_ref.is_none()); + assert_eq!(run.admission_retry, AdmissionRetryState::Ready); + drop(runs); + assert_eq!(run_store.list_events().await.unwrap().len(), event_count); +} + +#[test] +fn admission_error_classifier_is_structural() { + let transient = [ + ( + fabro_store::Error::Slate(slatedb::Error::unavailable("test outage".to_string())), + "slate", + ), + ( + fabro_store::Error::ObjectStore(object_store::Error::Generic { + store: "test", + source: Box::new(std::io::Error::other("test outage")), + }), + "object_store", + ), + ( + fabro_store::Error::Sqlite(sqlx::Error::RowNotFound), + "sqlite", + ), + ( + fabro_store::Error::Io(std::io::Error::other("test outage")), + "io", + ), + ]; + for (error, kind) in transient { + assert_eq!( + classify_admission_error(&anyhow::Error::new(error)), + AdmissionErrorClass::Transient(kind), + ); + } + + let serde_error = serde_json::from_str::("{").unwrap_err(); + let permanent = [ + ( + fabro_store::Error::EventRejected { + source: Box::new(fabro_store::Error::ReadOnly), + }, + "event_rejected", + ), + ( + fabro_store::Error::EventSequenceExhausted { max_seq: 10 }, + "event_sequence_exhausted", + ), + ( + fabro_store::Error::RunNotFound("missing".to_string()), + "run_not_found", + ), + (fabro_store::Error::ReadOnly, "read_only"), + (fabro_store::Error::Serde(serde_error), "serialization"), + ]; + for (error, kind) in permanent { + assert_eq!( + classify_admission_error(&anyhow::Error::new(error)), + AdmissionErrorClass::Permanent(kind), + ); + } + assert_eq!( + classify_admission_error(&anyhow::anyhow!("unrecognized")), + AdmissionErrorClass::Permanent("unrecognized"), + ); +} + +#[test] +fn admission_retry_backoff_is_bounded_and_scheduler_selection_is_fair() { + let mut retry = AdmissionRetryState::Ready; + let mut now = Instant::now(); + for (failure_count, expected_delay) in [(1, 1), (2, 2), (3, 4), (4, 8)] { + let outcome = retry.record_transient_failure(now); + assert_eq!(outcome, AdmissionRetryOutcome::Retry { + failures: failure_count, + delay: Duration::from_secs(expected_delay), + }); + let retry_deadline = now + Duration::from_secs(expected_delay); + assert!( + !retry.is_eligible( + retry_deadline + .checked_sub(Duration::from_millis(1)) + .unwrap() + ) + ); + now = retry_deadline; + assert!(retry.is_eligible(now)); + } + assert_eq!( + retry.record_transient_failure(now), + AdmissionRetryOutcome::GivenUp { failures: 5 } + ); + assert!(!retry.is_eligible(now + Duration::from_mins(1))); + + let temp = tempfile::tempdir().unwrap(); + let mut older = managed_run( + MINIMAL_DOT.to_string(), + RunStatus::Runnable, + Utc::now() - ChronoDuration::seconds(1), + temp.path().join("older"), + RunExecutionMode::Start, + ); + older.admission_retry = AdmissionRetryState::Backoff { + failures: 1, + retry_at: now + Duration::from_secs(1), + }; + let healthy = managed_run( + MINIMAL_DOT.to_string(), + RunStatus::Runnable, + Utc::now(), + temp.path().join("healthy"), + RunExecutionMode::Start, + ); + let runs = HashMap::from([(fixtures::RUN_1, older), (fixtures::RUN_2, healthy)]); + assert_eq!(select_runs_for_admission(&runs, 1, now), vec![ + fixtures::RUN_2 + ]); +} + +#[tokio::test] +async fn transient_admission_failure_restores_runnable_and_arms_retry() { + let state = test_app_state(); + let temp = tempfile::tempdir().unwrap(); + let run_store = stage_runnable_managed_run(&state, fixtures::RUN_1, temp.path()).await; + let original_events = run_store.list_events().await.unwrap().len(); + { + let mut runs = state.runs.lock().expect("runs lock poisoned"); + let run = runs.get_mut(&fixtures::RUN_1).unwrap(); + run.status = RunStatus::Starting; + run.cancel_token = Some(CancellationToken::new()); + } + let error = anyhow::Error::new(fabro_store::Error::Io(std::io::Error::other( + "temporary outage", + ))); + + record_admission_failure(&state, fixtures::RUN_1, &error); + + let runs = state.runs.lock().expect("runs lock poisoned"); + let run = runs.get(&fixtures::RUN_1).unwrap(); + let AdmissionRetryState::Backoff { failures, retry_at } = run.admission_retry else { + panic!("transient failure should arm admission backoff") + }; + assert_eq!(failures, 1); + assert_eq!(run.status, RunStatus::Runnable); + assert!(run.cancel_token.is_none()); + assert!(select_runs_for_admission(&runs, 1, Instant::now()).is_empty()); + assert_eq!(select_runs_for_admission(&runs, 1, retry_at), vec![ + fixtures::RUN_1 + ]); + drop(runs); + assert_eq!(run_store.state().await.unwrap().status, RunStatus::Runnable); + assert_eq!( + run_store.list_events().await.unwrap().len(), + original_events + ); +} + +#[tokio::test] +async fn permanent_admission_failure_gives_up_without_mutating_durable_run() { + let state = test_app_state(); + let temp = tempfile::tempdir().unwrap(); + let run_store = stage_runnable_managed_run(&state, fixtures::RUN_1, temp.path()).await; + let original_events = run_store.list_events().await.unwrap().len(); + { + let mut runs = state.runs.lock().expect("runs lock poisoned"); + let run = runs.get_mut(&fixtures::RUN_1).unwrap(); + run.status = RunStatus::Starting; + run.cancel_token = Some(CancellationToken::new()); + } + + record_admission_failure( + &state, + fixtures::RUN_1, + &anyhow::Error::new(fabro_store::Error::EventSequenceExhausted { max_seq: 10 }), + ); + + let runs = state.runs.lock().expect("runs lock poisoned"); + let run = runs.get(&fixtures::RUN_1).unwrap(); + assert_eq!(run.status, RunStatus::Runnable); + assert_eq!(run.admission_retry, AdmissionRetryState::GivenUp { + failures: 1, + }); + assert!(run.cancel_token.is_none()); + assert!(select_runs_for_admission(&runs, 1, Instant::now()).is_empty()); + drop(runs); + assert_eq!(run_store.state().await.unwrap().status, RunStatus::Runnable); + assert_eq!( + run_store.list_events().await.unwrap().len(), + original_events + ); +} + +#[tokio::test] +async fn successful_admission_clears_retry_bookkeeping() { + let state = test_app_state(); + let temp = tempfile::tempdir().unwrap(); + let run_store = stage_runnable_managed_run(&state, fixtures::RUN_1, temp.path()).await; + { + let mut runs = state.runs.lock().expect("runs lock poisoned"); + runs.get_mut(&fixtures::RUN_1).unwrap().admission_retry = AdmissionRetryState::Backoff { + failures: 2, + retry_at: Instant::now().checked_sub(Duration::from_secs(1)).unwrap(), + }; + } + + assert!(Box::pin(admit_run(&state, fixtures::RUN_1)).await.is_some()); + let runs = state.runs.lock().expect("runs lock poisoned"); + assert_eq!( + runs.get(&fixtures::RUN_1).unwrap().admission_retry, + AdmissionRetryState::Ready + ); + drop(runs); + assert_eq!(run_store.state().await.unwrap().status, RunStatus::Starting); + assert_eq!( + starting_event_count(&run_store.list_events().await.unwrap()), + 1 + ); +} + +async fn cancel_admitted_run(app: &Router, run_id: RunId) { + let response = app + .clone() + .oneshot( + Request::builder() + .method("POST") + .uri(api(&format!("/runs/{run_id}/cancel"))) + .body(Body::empty()) + .unwrap(), + ) + .await + .unwrap(); + assert_status!(response, StatusCode::ACCEPTED).await; +} + +fn assert_cancelled_admission_history(events: &[EventEnvelope]) { + let starting = events + .iter() + .position(|event| matches!(event.event.body, EventBody::RunStarting(_))) + .expect("admission should append run.starting"); + let requested = events + .iter() + .position(|event| matches!(event.event.body, EventBody::RunCancelRequested(_))) + .expect("cancellation should append run.cancel.requested"); + let failed = events + .iter() + .position(|event| { + matches!( + event.event.body, + EventBody::RunFailed(ref props) + if props.failure.reason == FailureReason::Cancelled + ) + }) + .expect("cancellation should append a cancelled run.failed event"); + assert!(starting < requested && requested < failed); + assert!(!events.iter().any(|event| { + matches!( + event.event.body, + EventBody::RunFailed(ref props) + if props.failure.reason == FailureReason::LaunchFailed + ) + })); +} + +#[tokio::test] +async fn cancellation_after_admission_does_not_start_subprocess_worker() { + let runtime = Arc::new(RecordingWorkerRuntime::default()); + let state = TestAppStateBuilder::new() + .vault_entries([(EnvVars::OPENAI_API_KEY, "test-openai-api-key")]) + .worker_runtime(runtime.clone()) + .build(); + let app = crate::test_support::build_test_router(Arc::clone(&state)); + let temp = tempfile::tempdir().unwrap(); + let run_store = stage_runnable_managed_run(&state, fixtures::RUN_1, temp.path()).await; + let admitted = Box::pin(admit_run(&state, fixtures::RUN_1)) + .await + .expect("run should win durable admission"); + let cancel_token = admitted.cancel_token.clone(); + + cancel_admitted_run(&app, fixtures::RUN_1).await; + Box::pin(execute_run_subprocess( + Arc::clone(&state), + fixtures::RUN_1, + admitted, + )) + .await; + + assert!(cancel_token.is_cancelled()); + assert_eq!(runtime.start_count(), 0); + let runs = state.runs.lock().expect("runs lock poisoned"); + let run = runs.get(&fixtures::RUN_1).unwrap(); + assert_eq!(run.status, RunStatus::Failed { + reason: FailureReason::Cancelled, + }); + assert!(run.cancel_token.is_none()); + assert!(run.worker_ref.is_none()); + drop(runs); + assert_cancelled_admission_history(&run_store.list_events().await.unwrap()); +} + +#[tokio::test] +async fn cancellation_after_admission_does_not_start_in_process_workflow() { + let registry_calls = Arc::new(AtomicUsize::new(0)); + let calls = Arc::clone(®istry_calls); + let state = TestAppStateBuilder::new() + .vault_entries([(EnvVars::OPENAI_API_KEY, "test-openai-api-key")]) + .registry_factory(move |interviewer| { + calls.fetch_add(1, Ordering::Relaxed); + fabro_workflow::handler::default_registry(interviewer, || None) + }) + .build(); + let app = crate::test_support::build_test_router(Arc::clone(&state)); + let temp = tempfile::tempdir().unwrap(); + let run_store = stage_runnable_managed_run(&state, fixtures::RUN_1, temp.path()).await; + let admitted = Box::pin(admit_run(&state, fixtures::RUN_1)) + .await + .expect("run should win durable admission"); + + cancel_admitted_run(&app, fixtures::RUN_1).await; + Box::pin(execute_run_in_process( + Arc::clone(&state), + fixtures::RUN_1, + admitted, + )) + .await; + + assert_eq!(registry_calls.load(Ordering::Relaxed), 0); + assert_eq!(run_store.state().await.unwrap().status, RunStatus::Failed { + reason: FailureReason::Cancelled, + }); + assert_cancelled_admission_history(&run_store.list_events().await.unwrap()); +} + +fn write_worker_test_server_record(state: &AppState) { + let runtime_directory = Storage::new(state.server_storage_dir()).runtime_directory(); + ServerDaemon::new( + std::process::id(), + Bind::Tcp("127.0.0.1:32276".parse::().unwrap()), + runtime_directory.log_path(), + ) + .write(&runtime_directory) + .unwrap(); +} + +#[tokio::test] +async fn subprocess_launch_failure_follows_starting_and_preserves_source_chain() { + let runtime = Arc::new(RecordingWorkerRuntime::default()); + let state = TestAppStateBuilder::new() + .vault_entries([(EnvVars::OPENAI_API_KEY, "test-openai-api-key")]) + .worker_runtime(runtime.clone()) + .build(); + write_worker_test_server_record(state.as_ref()); + let app = crate::test_support::build_test_router(Arc::clone(&state)); + let run_id = create_and_start_run(&app, MINIMAL_DOT) + .await + .parse::() + .unwrap(); + + execute_run(Arc::clone(&state), run_id).await; + + assert_eq!(runtime.start_count(), 1); + let run_store = state.stores.runs.open_run(&run_id).await.unwrap(); + let events = run_store.list_events().await.unwrap(); + let starting = events + .iter() + .position(|event| matches!(event.event.body, EventBody::RunStarting(_))) + .unwrap(); + let (failed, failure) = events + .iter() + .enumerate() + .find_map(|(index, event)| match &event.event.body { + EventBody::RunFailed(props) => Some((index, &props.failure)), + _ => None, + }) + .expect("worker launch should persist run.failed"); + assert!(starting < failed); + assert_eq!(failure.reason, FailureReason::LaunchFailed); + assert_eq!(failure.detail.message, "Failed to spawn worker"); + assert!( + failure + .detail + .causes + .iter() + .any(|cause| cause.contains("recording runtime source failure")) + ); + assert_eq!(run_store.state().await.unwrap().status, RunStatus::Failed { + reason: FailureReason::LaunchFailed, + }); + let summaries = state + .stores + .runs + .list_runs(&fabro_store::ListRunsQuery::default(), Utc::now()) + .await + .unwrap(); + assert!(summaries.iter().any(|summary| summary.id == run_id)); + assert!( + state + .stores + .runs + .list_unreadable_runs() + .await + .unwrap() + .is_empty() + ); +} + +#[tokio::test] +async fn in_process_policy_failure_follows_durable_starting() { + let disabled_source = r#" +_version = 1 + +[server.auth] +methods = ["dev-token"] + +[server.sandbox.providers.docker] +enabled = false +"#; + let state = test_app_state_with_registry_factory(|interviewer| { + fabro_workflow::handler::default_registry(interviewer, || None) + }); + let app = crate::test_support::build_test_router(Arc::clone(&state)); + let run_id = create_and_start_run(&app, MINIMAL_DOT) + .await + .parse::() + .unwrap(); + state + .replace_runtime_settings(resolved_runtime_settings_from_toml(disabled_source)) + .unwrap(); + + execute_run(Arc::clone(&state), run_id).await; + + let run_store = state.stores.runs.open_run(&run_id).await.unwrap(); + let events = run_store.list_events().await.unwrap(); + let starting = events + .iter() + .position(|event| matches!(event.event.body, EventBody::RunStarting(_))) + .unwrap(); + let failed = events + .iter() + .position(|event| { + matches!( + event.event.body, + EventBody::RunFailed(ref props) + if props.failure.reason == FailureReason::LaunchFailed + ) + }) + .expect("disabled provider should fail the admitted run"); + assert!(starting < failed); + assert_eq!(run_store.state().await.unwrap().status, RunStatus::Failed { + reason: FailureReason::LaunchFailed, + }); +} + +#[tokio::test] +async fn failure_append_error_keeps_local_status_aligned_with_durable_starting() { + let state = test_app_state(); + let temp = tempfile::tempdir().unwrap(); + let run_store = stage_runnable_managed_run(&state, fixtures::RUN_1, temp.path()).await; + let admitted = Box::pin(admit_run(&state, fixtures::RUN_1)) + .await + .expect("run should win durable admission"); + let read_only = state + .stores + .runs + .open_run_reader(&fixtures::RUN_1) + .await + .unwrap(); + + fail_run_before_execution( + &state, + &read_only, + fixtures::RUN_1, + FailureReason::LaunchFailed, + WorkflowError::engine_with_source( + "Launch setup failed", + std::io::Error::other("source survives"), + ), + ) + .await; + + drop(admitted); + let runs = state.runs.lock().expect("runs lock poisoned"); + let run = runs.get(&fixtures::RUN_1).unwrap(); + assert_eq!(run.status, RunStatus::Starting); + assert!(run.cancel_token.is_none()); + assert!(run.worker_ref.is_none()); + drop(runs); + assert_eq!(run_store.state().await.unwrap().status, RunStatus::Starting); + assert!( + !run_store + .list_events() + .await + .unwrap() + .iter() + .any(|event| { matches!(event.event.body, EventBody::RunFailed(_)) }) + ); + + let boundary_error = anyhow::Error::new(std::io::Error::other( + "https://user:password@example.com/private", + )) + .context("failure append context"); + let rendered = safe_error_chain(&boundary_error); + assert!(rendered.contains("failure append context")); + assert!(!rendered.contains("password")); +} + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] async fn concurrency_limit_respected() { let state = test_app_state_with_options(default_test_server_settings(), RunLayer::default(), 1); diff --git a/lib/components/fabro-workflow/src/operations/start.rs b/lib/components/fabro-workflow/src/operations/start.rs index 12b7455a9..943c8b8fd 100644 --- a/lib/components/fabro-workflow/src/operations/start.rs +++ b/lib/components/fabro-workflow/src/operations/start.rs @@ -183,33 +183,47 @@ pub(super) async fn execute_persisted_run( let run_id = services.run_id; let run_store = services.run_store.clone(); let event_sink = services.event_sink.clone(); - if let Err(err) = run_store.state().await { - let error = Error::engine(err.to_string()); - let _ = persist_detached_failure( - run_id, - &run_store, - &event_sink, - run_dir, - "bootstrap", - FailureReason::BootstrapFailed, - &error, - ) - .await; - return Err(error); - } - if let Err(err) = append_event_to_sink(&event_sink, &run_id, &Event::RunStarting).await { - let error = Error::engine(err.to_string()); - let _ = persist_detached_failure( - run_id, - &run_store, - &event_sink, - run_dir, - "bootstrap", - FailureReason::BootstrapFailed, - &error, - ) - .await; - return Err(error); + let bootstrap_state = match run_store.state().await { + Ok(state) => state, + Err(err) => { + let error = Error::engine(err.to_string()); + let _ = persist_detached_failure( + run_id, + &run_store, + &event_sink, + run_dir, + "bootstrap", + FailureReason::BootstrapFailed, + &error, + ) + .await; + return Err(error); + } + }; + match bootstrap_state.status { + RunStatus::Runnable => { + if let Err(err) = append_event_to_sink(&event_sink, &run_id, &Event::RunStarting).await + { + let error = Error::engine(err.to_string()); + let _ = persist_detached_failure( + run_id, + &run_store, + &event_sink, + run_dir, + "bootstrap", + FailureReason::BootstrapFailed, + &error, + ) + .await; + return Err(error); + } + } + RunStatus::Starting => {} + status => { + return Err(Error::Precondition(format!( + "cannot execute persisted run: status is {status}, expected runnable or starting" + ))); + } } let mut bootstrap_guard = DetachedRunBootstrapGuard::arm( @@ -2083,6 +2097,50 @@ reasoning = false assert_eq!(started.finalized.conclusion.status, StageOutcome::Succeeded); let run_store = store.open_run(&fixtures::RUN_1).await.unwrap(); assert!(run_store.state().await.unwrap().conclusion.is_some()); + let starting_count = run_store + .list_events() + .await + .unwrap() + .iter() + .filter(|event| matches!(event.event.body, EventBody::RunStarting(_))) + .count(); + assert_eq!(starting_count, 1); + } + + #[tokio::test] + async fn start_from_already_starting_emits_no_duplicate_starting_event() { + let temp = tempfile::tempdir().unwrap(); + let (storage_root, run_dir) = storage_root_and_run_dir(&temp); + let emitter = Arc::new(Emitter::new(fixtures::RUN_1)); + let registry = Arc::new(test_registry()); + let (_persisted, store) = persisted_workflow(MINIMAL_DOT, &storage_root).await; + let run_store = store.open_run(&fixtures::RUN_1).await.unwrap(); + + crate::event::append_event(&run_store, &fixtures::RUN_1, &Event::RunRunnable { + source: RunRunnableSource::StartRequested, + actor: None, + }) + .await + .unwrap(); + crate::event::append_event(&run_store, &fixtures::RUN_1, &Event::RunStarting) + .await + .unwrap(); + + start( + &run_dir, + test_start_services(&store, &run_dir, emitter, registry).await, + ) + .await + .unwrap(); + + let starting_count = run_store + .list_events() + .await + .unwrap() + .iter() + .filter(|event| matches!(event.event.body, EventBody::RunStarting(_))) + .count(); + assert_eq!(starting_count, 1); } #[tokio::test]