Simplify durable admission internals

- Reuse fabro_redact::redacted_url_for_log in safe_error_chain instead of
  re-rolling the DisplaySafeUrl parse/redact fallback
- Reuse fabro_util::backoff::BackoffPolicy for admission retry delays
- Reuse projection_failure_message in worker-exit reconciliation
- Flatten AdmissionErrorClass into a { transient, kind } struct so the
  class is matched once instead of twice
- Return the resulting RunStatus from persist_cancelled_run_status_to_store,
  removing a redundant store read and the Result/Option reconciliation in
  finish_cancelled_run_before_execution
- Capture spec, pending_control, and execution_mode at the admission claim
  instead of re-reading the run projection up to two more times per
  admission; drop the unreachable "managed run disappeared" prologues
- Dedup test fixtures (daemon record helper, RunStarting event counting)
- Cross-reference the server claim from execute_persisted_run's twin
  bootstrap transition

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
This commit is contained in:
Scott Werner 2026-08-04 14:54:03 -04:00
parent 579acf17dc
commit 95d3216470
3 changed files with 141 additions and 178 deletions

View file

@ -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<WorkerRef> {
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<RunStatus> {
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::<fabro_store::Error>() 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: &regex::Captures<'_>| {
DisplaySafeUrl::parse(&captures[0])
.map_or_else(|_| "<invalid url>".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<AppState>, 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<AppState>, run_id: RunId, error: &anyhow
}
async fn admit_run(state: &Arc<AppState>, run_id: RunId) -> Option<AdmittedRun> {
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<AppState>, run_id: RunId) -> Option<AdmittedRun>
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<AppState>, run_id: RunId) -> Option<AdmittedRun>
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<AppState>, run_id: RunId) -> Option<AdmittedRun>
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<AppState>, run_id: RunId) -> Option<AdmittedRun>
run_store,
run_dir,
cancel_token,
spec,
execution_mode,
})
}
fn execution_mode_for_run(state: &AppState, run_id: RunId) -> Option<RunExecutionMode> {
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<AppState>, 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<AppState>, 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<AppState>, 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);

View file

@ -2859,14 +2859,7 @@ allowed_usernames = ["octocat"]
.collect::<Vec<_>>()
.join(", ")
);
let runtime_directory = Storage::new(storage_dir).runtime_directory();
ServerDaemon::new(
std::process::id(),
Bind::Tcp("127.0.0.1:32276".parse::<std::net::SocketAddr>().unwrap()),
runtime_directory.log_path(),
)
.write(&runtime_directory)
.unwrap();
write_worker_test_server_record(storage_dir);
let mut server_secret_env: HashMap<String, String> = 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::<std::net::SocketAddr>().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

View file

@ -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]