diff --git a/lib/apps/fabro-server/src/server.rs b/lib/apps/fabro-server/src/server.rs index 2ec72f85b..0835077ca 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::{DisplaySafeUrl, redact_jsonl_line, redact_string}; +use fabro_redact::{redact_jsonl_line, redact_string, redacted_url_for_log}; use fabro_sandbox::daytona::{self, DaytonaSandbox}; use fabro_sandbox::details::sandbox_details; use fabro_sandbox::reconnect::reconnect_for_run; @@ -98,9 +98,10 @@ use fabro_types::settings::server::{ use fabro_types::{ AgentBackend, AskFabro, AskFabroUnavailableReason, EventBody, InterviewQuestionRecord, PairId, PairMessageId, PairTarget, PendingReason, Principal, PullRequestLink, QuestionType, RunBlobId, - RunControlAction, RunEvent, RunId, RunProjection, RunRunnableSource, SandboxProviderKind, - ServerSettings, SessionCapability, + RunControlAction, RunEvent, RunId, RunProjection, RunRunnableSource, RunSpec, + SandboxProviderKind, ServerSettings, SessionCapability, }; +use fabro_util::backoff::BackoffPolicy; use fabro_util::error::{ SharedError, collect_causes, render_compact_with_causes, render_with_causes, }; @@ -287,7 +288,12 @@ struct ManagedRun { } const MAX_ADMISSION_FAILURES: u8 = 5; -const MAX_ADMISSION_RETRY_DELAY: Duration = Duration::from_secs(30); +const ADMISSION_RETRY_BACKOFF: BackoffPolicy = BackoffPolicy { + initial_delay: Duration::from_secs(1), + factor: 2.0, + max_delay: Duration::from_secs(30), + jitter: false, +}; #[derive(Clone, Copy, Debug, Default, Eq, PartialEq)] enum AdmissionRetryState { @@ -328,9 +334,7 @@ impl AdmissionRetryState { 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); + let delay = ADMISSION_RETRY_BACKOFF.delay_for_attempt(u32::from(failures)); *self = Self::Backoff { failures, retry_at: now + delay, @@ -353,9 +357,12 @@ impl AdmissionRetryState { } struct AdmittedRun { - run_store: fabro_store::RunDatabase, - run_dir: PathBuf, - cancel_token: CancellationToken, + run_store: fabro_store::RunDatabase, + run_dir: PathBuf, + cancel_token: CancellationToken, + /// Spec observed by the successful claim; immutable after `RunStarting`. + spec: RunSpec, + execution_mode: RunExecutionMode, } impl ManagedRun { @@ -3323,16 +3330,21 @@ 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 + persist_cancelled_run_status_to_store(&run_store, run_id) + .await + .map(|_| ()) } +/// Returns the run's durable status after the persist: the pre-existing +/// terminal status when nothing needed to be appended, or cancelled-failed +/// after a successful append. async fn persist_cancelled_run_status_to_store( run_store: &fabro_store::RunDatabase, run_id: RunId, -) -> anyhow::Result<()> { +) -> anyhow::Result { let run_state = run_store.state().await?; if run_state.status.is_terminal() { - return Ok(()); + return Ok(run_state.status); } let failure_event = workflow_event::Event::workflow_run_failed_from_error( @@ -3344,7 +3356,10 @@ async fn persist_cancelled_run_status_to_store( None, None, ); - workflow_event::append_event(run_store, &run_id, &failure_event).await + workflow_event::append_event(run_store, &run_id, &failure_event).await?; + Ok(RunStatus::Failed { + reason: FailureReason::Cancelled, + }) } async fn finish_cancelled_run_before_execution( @@ -3352,36 +3367,37 @@ async fn finish_cancelled_run_before_execution( 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 + let durable_status = match persist_cancelled_run_status_to_store(run_store, run_id).await { + Ok(status) => Some(status), + Err(persist_err) => { + 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 + } + }; + // A concurrent transition to a terminal status makes the persist + // failure moot; only report it otherwise. + if !durable_status.is_some_and(RunStatus::is_terminal) { + error!( + run_id = %run_id, + error = %safe_error_chain(&persist_err), + "Failed to persist cancelled run status" + ); + } + durable_status } }; - 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) { 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); } @@ -4033,47 +4049,63 @@ fn answer_from_request( } #[derive(Clone, Copy, Debug, Eq, PartialEq)] -enum AdmissionErrorClass { - Transient(&'static str), - Permanent(&'static str), +struct AdmissionErrorClass { + transient: bool, + kind: &'static str, +} + +impl AdmissionErrorClass { + const fn transient(kind: &'static str) -> Self { + Self { + transient: true, + kind, + } + } + + const fn permanent(kind: &'static str) -> Self { + Self { + transient: false, + kind, + } + } } fn classify_admission_error(error: &anyhow::Error) -> AdmissionErrorClass { let Some(store_error) = error.downcast_ref::() else { - return AdmissionErrorClass::Permanent("unrecognized"); + 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::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") + AdmissionErrorClass::permanent("event_rejected") } - fabro_store::Error::RunNotFound(_) => AdmissionErrorClass::Permanent("run_not_found"), + fabro_store::Error::RunNotFound(_) => AdmissionErrorClass::permanent("run_not_found"), fabro_store::Error::SessionNotFound(_) => { - AdmissionErrorClass::Permanent("session_not_found") + AdmissionErrorClass::permanent("session_not_found") } fabro_store::Error::SessionAlreadyExists(_) => { - AdmissionErrorClass::Permanent("session_already_exists") + AdmissionErrorClass::permanent("session_already_exists") } - fabro_store::Error::ReadOnly => AdmissionErrorClass::Permanent("read_only"), + fabro_store::Error::ReadOnly => AdmissionErrorClass::permanent("read_only"), fabro_store::Error::EventSequenceExhausted { .. } => { - AdmissionErrorClass::Permanent("event_sequence_exhausted") + AdmissionErrorClass::permanent("event_sequence_exhausted") } fabro_store::Error::InvalidKeySegment { .. } => { - AdmissionErrorClass::Permanent("invalid_key_segment") + AdmissionErrorClass::permanent("invalid_key_segment") } - fabro_store::Error::KeyParse(_) => AdmissionErrorClass::Permanent("key_parse"), + fabro_store::Error::KeyParse(_) => AdmissionErrorClass::permanent("key_parse"), fabro_store::Error::RunSummaryMismatch { .. } => { - AdmissionErrorClass::Permanent("run_summary_mismatch") + AdmissionErrorClass::permanent("run_summary_mismatch") } fabro_store::Error::InvalidTransition(_) => { - AdmissionErrorClass::Permanent("invalid_transition") + AdmissionErrorClass::permanent("invalid_transition") } - fabro_store::Error::Serde(_) => AdmissionErrorClass::Permanent("serialization"), - fabro_store::Error::Other(_) => AdmissionErrorClass::Permanent("internal"), + fabro_store::Error::Serde(_) => AdmissionErrorClass::permanent("serialization"), + fabro_store::Error::Other(_) => AdmissionErrorClass::permanent("internal"), } } @@ -4085,8 +4117,7 @@ fn safe_error_chain(error: &anyhow::Error) -> String { 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()) + redacted_url_for_log(&captures[0]) }); redact_string(&url_safe) } @@ -4111,19 +4142,16 @@ fn record_admission_failure(state: &Arc, run_id: RunId, error: &anyhow managed_run.status = RunStatus::Runnable; } clear_live_run_state(managed_run); - match class { - AdmissionErrorClass::Transient(_) => managed_run + if class.transient { + managed_run .admission_retry - .record_transient_failure(Instant::now()), - AdmissionErrorClass::Permanent(_) => { - managed_run.admission_retry.record_permanent_failure() - } + .record_transient_failure(Instant::now()) + } else { + managed_run.admission_retry.record_permanent_failure() } }; - let error_kind = match class { - AdmissionErrorClass::Transient(kind) | AdmissionErrorClass::Permanent(kind) => kind, - }; + let error_kind = class.kind; match outcome { AdmissionRetryOutcome::Retry { failures, delay } => { warn!( @@ -4152,7 +4180,7 @@ fn record_admission_failure(state: &Arc, run_id: RunId, error: &anyhow } async fn admit_run(state: &Arc, run_id: RunId) -> Option { - let (run_dir, cancel_token) = { + let (run_dir, cancel_token, execution_mode) = { let mut runs = state.runs.lock().expect("runs lock poisoned"); if state.is_shutting_down() { return None; @@ -4167,7 +4195,11 @@ async fn admit_run(state: &Arc, run_id: RunId) -> Option 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) + ( + managed_run.run_dir.clone(), + cancel_token, + managed_run.execution_mode, + ) }; let run_store = match state.stores.runs.open_run(&run_id).await { @@ -4180,12 +4212,14 @@ async fn admit_run(state: &Arc, run_id: RunId) -> Option let run_events = run_store.subscribe(); let mut observed_status = None; let mut observed_error = None; + let mut claimed_run = None; let claimed = workflow_event::append_event_if( &run_store, &run_id, &workflow_event::Event::RunStarting, |projection| { if projection.status == RunStatus::Runnable { + claimed_run = Some((projection.spec.clone(), projection.pending_control)); true } else { observed_status = Some(projection.status); @@ -4242,21 +4276,9 @@ async fn admit_run(state: &Arc, run_id: RunId) -> Option 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() { + let (spec, pending_control) = + claimed_run.expect("successful claim must have evaluated the admission predicate"); + if pending_control == Some(RunControlAction::Cancel) || cancel_token.is_cancelled() { finish_cancelled_run_before_execution(state, &run_store, run_id).await; return None; } @@ -4265,18 +4287,11 @@ async fn admit_run(state: &Arc, run_id: RunId) -> Option run_store, run_dir, cancel_token, + spec, + execution_mode, }) } -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; @@ -4315,22 +4330,13 @@ async fn execute_run_in_process(state: Arc, run_id: RunId, admitted: A run_store, run_dir, cancel_token, + execution_mode, + .. } = 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()); @@ -4608,49 +4614,20 @@ async fn execute_run_subprocess(state: Arc, run_id: RunId, admitted: A run_store, run_dir, cancel_token, + spec, + execution_mode, } = 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_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_run_before_execution( - &state, - &run_store, - run_id, - FailureReason::WorkflowError, - 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; - } + let agent_fabro_tools_enabled = spec.settings.run.agent.fabro_tools; if reject_run_if_sandbox_provider_disabled( &state, &run_store, &state.server_settings(), run_id, - &run_state.spec.settings.run, + &spec.settings.run, ) .await { @@ -4827,15 +4804,8 @@ async fn execute_run_subprocess(state: Arc, run_id: RunId, admitted: A reason: FailureReason::Terminated, }; } - managed_run.error = final_state - .conclusion - .as_ref() - .and_then(|conclusion| { - conclusion.failure.as_ref().map(|failure| { - render_compact_with_causes(&failure.detail.message, &failure.detail.causes) - }) - }) - .or_else(|| managed_run.error.clone()); + managed_run.error = + projection_failure_message(&final_state).or_else(|| managed_run.error.clone()); managed_run.checkpoint = final_state.current_checkpoint().cloned(); managed_run.run_dir = Some(run_dir); clear_live_run_state(managed_run); diff --git a/lib/apps/fabro-server/src/server/tests.rs b/lib/apps/fabro-server/src/server/tests.rs index c3d60aeb3..50aa2acf0 100644 --- a/lib/apps/fabro-server/src/server/tests.rs +++ b/lib/apps/fabro-server/src/server/tests.rs @@ -2859,14 +2859,7 @@ allowed_usernames = ["octocat"] .collect::>() .join(", ") ); - let runtime_directory = Storage::new(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(); + write_worker_test_server_record(storage_dir); let mut server_secret_env: HashMap = dev_token .map(|token| HashMap::from([("FABRO_DEV_TOKEN".to_string(), token)])) @@ -16500,7 +16493,7 @@ fn admission_error_classifier_is_structural() { for (error, kind) in transient { assert_eq!( classify_admission_error(&anyhow::Error::new(error)), - AdmissionErrorClass::Transient(kind), + AdmissionErrorClass::transient(kind), ); } @@ -16526,12 +16519,12 @@ fn admission_error_classifier_is_structural() { for (error, kind) in permanent { assert_eq!( classify_admission_error(&anyhow::Error::new(error)), - AdmissionErrorClass::Permanent(kind), + AdmissionErrorClass::permanent(kind), ); } assert_eq!( classify_admission_error(&anyhow::anyhow!("unrecognized")), - AdmissionErrorClass::Permanent("unrecognized"), + AdmissionErrorClass::permanent("unrecognized"), ); } @@ -16800,8 +16793,8 @@ async fn cancellation_after_admission_does_not_start_in_process_workflow() { 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(); +fn write_worker_test_server_record(storage_dir: &Path) { + let runtime_directory = Storage::new(storage_dir).runtime_directory(); ServerDaemon::new( std::process::id(), Bind::Tcp("127.0.0.1:32276".parse::().unwrap()), @@ -16818,7 +16811,7 @@ async fn subprocess_launch_failure_follows_starting_and_preserves_source_chain() .vault_entries([(EnvVars::OPENAI_API_KEY, "test-openai-api-key")]) .worker_runtime(runtime.clone()) .build(); - write_worker_test_server_record(state.as_ref()); + write_worker_test_server_record(&state.server_storage_dir()); let app = crate::test_support::build_test_router(Arc::clone(&state)); let run_id = create_and_start_run(&app, MINIMAL_DOT) .await diff --git a/lib/components/fabro-workflow/src/operations/start.rs b/lib/components/fabro-workflow/src/operations/start.rs index 943c8b8fd..ce794b5e2 100644 --- a/lib/components/fabro-workflow/src/operations/start.rs +++ b/lib/components/fabro-workflow/src/operations/start.rs @@ -200,6 +200,10 @@ pub(super) async fn execute_persisted_run( return Err(error); } }; + // Non-atomic twin of the server's durable admission claim (`admit_run`'s + // conditional `RunStarting` append): it goes through the event sink so + // fan-out sinks observe the event. Keep the accepted-status set in sync + // with the server's claim predicate. match bootstrap_state.status { RunStatus::Runnable => { if let Err(err) = append_event_to_sink(&event_sink, &run_id, &Event::RunStarting).await @@ -1828,6 +1832,16 @@ reasoning = false .unwrap(); } + async fn starting_event_count(run_store: &fabro_store::RunDatabase) -> usize { + run_store + .list_events() + .await + .unwrap() + .iter() + .filter(|event| matches!(event.event.body, EventBody::RunStarting(_))) + .count() + } + async fn wait_for_conclusion( run_store: &fabro_store::RunDatabase, ) -> crate::records::Conclusion { @@ -2097,14 +2111,7 @@ 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); + assert_eq!(starting_event_count(&run_store).await, 1); } #[tokio::test] @@ -2133,14 +2140,7 @@ reasoning = false .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); + assert_eq!(starting_event_count(&run_store).await, 1); } #[tokio::test]