From 52e7726c89fefdcffb6b343fd0d5bb65268850cc Mon Sep 17 00:00:00 2001 From: Andrei Hasna Date: Tue, 11 Aug 2026 16:46:10 +0300 Subject: [PATCH 1/6] fix(workflows): bound activation database contention Persist a finite workflow_activation_database_busy outcome when SQLite busy retries are exhausted, while preserving idempotent branch admission and AgentRun creation. Agent: iapp-factory-coordinator --- codex-rs/ext/workflows/Cargo.toml | 1 + codex-rs/ext/workflows/src/activation.rs | 970 +++++++++++++++--- codex-rs/state/src/lib.rs | 3 + codex-rs/state/src/runtime.rs | 3 + .../src/runtime/workflow_orchestrator.rs | 161 +++ 5 files changed, 1016 insertions(+), 122 deletions(-) diff --git a/codex-rs/ext/workflows/Cargo.toml b/codex-rs/ext/workflows/Cargo.toml index d0ae2a824..e95d51794 100644 --- a/codex-rs/ext/workflows/Cargo.toml +++ b/codex-rs/ext/workflows/Cargo.toml @@ -46,5 +46,6 @@ chrono = { workspace = true } codex-prompts = { workspace = true } codex-utils-output-truncation = { workspace = true } pretty_assertions = { workspace = true } +sqlx = { workspace = true } tempfile = { workspace = true } tokio = { workspace = true, features = ["macros", "rt"] } diff --git a/codex-rs/ext/workflows/src/activation.rs b/codex-rs/ext/workflows/src/activation.rs index c17c477ec..df89afb11 100644 --- a/codex-rs/ext/workflows/src/activation.rs +++ b/codex-rs/ext/workflows/src/activation.rs @@ -1,5 +1,7 @@ use std::collections::HashMap; use std::collections::HashSet; +#[cfg(test)] +use std::collections::VecDeque; use std::fs::File; use std::future::Future; use std::io::Read; @@ -7,6 +9,8 @@ use std::path::Component; use std::path::Path; use std::path::PathBuf; use std::sync::Arc; +#[cfg(test)] +use std::sync::Mutex as StdMutex; use std::sync::atomic::AtomicBool; use std::sync::atomic::Ordering; use std::time::Duration; @@ -24,10 +28,13 @@ use codex_protocol::models::PermissionProfile; use codex_protocol::permissions::NetworkSandboxPolicy; use codex_shell_command::shell_detect::ShellType; use codex_state::StateRuntime; +#[cfg(test)] +use codex_state::WORKFLOW_ACTIVATION_DATABASE_BUSY_REASON_CODE; use codex_state::WORKFLOW_BRANCH_ADMISSION_DEFERRED_EVENT; use codex_state::WorkflowGoalPlanProjectionOutcome; use codex_state::WorkflowGoalPlanProjectionParams; use codex_state::WorkflowProviderCreditReservationStatus; +use codex_state::WorkflowRunActivationDatabaseBusyParams; use codex_state::WorkflowRunAdvanceParams; use codex_state::WorkflowRunBranchAdmissionParams; use codex_state::WorkflowRunBranchReconcileParams; @@ -46,7 +53,9 @@ use codex_state::WorkflowRunVerifierClaimSelection; use codex_state::WorkflowRunVerifierOutcomeStatus; use codex_state::WorkflowRunVerifierRecordResultParams; use codex_state::WorkflowRunVerifierResultSummary; -use codex_state::busy_retry::retry_on_busy; +use codex_state::busy_retry::BusyRetryPolicy; +use codex_state::busy_retry::is_transient_busy_error; +use codex_state::busy_retry::retry_on_busy_with_policy; use codex_utils_absolute_path::AbsolutePathBuf; use codex_workflows::WorkflowModelRoute; use codex_workflows::WorkflowProviderCreditControl; @@ -155,7 +164,7 @@ fn verifier_sandbox_support_env() -> HashMap { HashMap::new() } -#[derive(Debug, Clone, Copy, PartialEq, Eq)] +#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)] enum WorkflowStateOperation { CreateRun, ProjectGoalPlan, @@ -168,11 +177,66 @@ enum WorkflowStateOperation { AdvanceRun, ReconcileBranches, AdmitBranches, + RecordActivationBusy, ClaimVerifier, CheckVerifierFence, RecordVerifierResult, } +#[derive(Debug)] +struct WorkflowActivationDatabaseBusy { + operation: WorkflowStateOperation, + source: anyhow::Error, +} + +impl std::fmt::Display for WorkflowActivationDatabaseBusy { + fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + write!( + formatter, + "{} remained busy after bounded retries: {}", + self.operation.label(), + self.source + ) + } +} + +impl std::error::Error for WorkflowActivationDatabaseBusy { + fn source(&self) -> Option<&(dyn std::error::Error + 'static)> { + Some(self.source.as_ref()) + } +} + +#[cfg(test)] +#[derive(Clone, Default)] +struct WorkflowStateFaultPlan { + faults: Arc>>, +} + +#[cfg(test)] +impl WorkflowStateFaultPlan { + fn push(&self, operation: WorkflowStateOperation, error: anyhow::Error) { + self.faults + .lock() + .expect("workflow state fault plan should lock") + .push_back((operation, error)); + } + + fn take(&self, operation: WorkflowStateOperation) -> Option { + let mut faults = self + .faults + .lock() + .expect("workflow state fault plan should lock"); + if faults + .front() + .is_some_and(|(planned, _)| *planned == operation) + { + faults.pop_front().map(|(_, error)| error) + } else { + None + } + } +} + impl WorkflowStateOperation { #[cfg(test)] const ALL: [Self; 14] = [ @@ -187,6 +251,7 @@ impl WorkflowStateOperation { Self::AdvanceRun, Self::ReconcileBranches, Self::AdmitBranches, + Self::RecordActivationBusy, Self::ClaimVerifier, Self::CheckVerifierFence, Self::RecordVerifierResult, @@ -205,6 +270,7 @@ impl WorkflowStateOperation { Self::AdvanceRun => "advance workflow run", Self::ReconcileBranches => "reconcile workflow run branches", Self::AdmitBranches => "admit workflow run branches", + Self::RecordActivationBusy => "record workflow activation database busy", Self::ClaimVerifier => "claim workflow run verifier", Self::CheckVerifierFence => "check workflow verifier fence", Self::RecordVerifierResult => "record workflow verifier result", @@ -212,6 +278,7 @@ impl WorkflowStateOperation { } } +#[cfg(test)] async fn retry_workflow_state( operation: WorkflowStateOperation, f: F, @@ -220,7 +287,27 @@ where F: FnMut() -> Fut, Fut: Future>, { - retry_on_busy(operation.label(), f).await + retry_workflow_state_with_policy(BusyRetryPolicy::default(), operation, f).await +} + +async fn retry_workflow_state_with_policy( + policy: BusyRetryPolicy, + operation: WorkflowStateOperation, + f: F, +) -> anyhow::Result +where + F: FnMut() -> Fut, + Fut: Future>, +{ + match retry_on_busy_with_policy(policy, operation.label(), f).await { + Err(source) if is_transient_busy_error(&source) => { + Err(anyhow::Error::new(WorkflowActivationDatabaseBusy { + operation, + source, + })) + } + result => result, + } } #[derive(Debug, Clone)] @@ -313,6 +400,9 @@ pub struct WorkflowActivationService { owner_instance_id: Arc, active_supervisors: Arc>>, provider_credit_authority: Arc, + state_retry_policy: BusyRetryPolicy, + #[cfg(test)] + state_fault_plan: WorkflowStateFaultPlan, } impl WorkflowActivationService { @@ -332,9 +422,68 @@ impl WorkflowActivationService { owner_instance_id: Arc::from(format!("workflow-activation:{}", ThreadId::new())), active_supervisors: Arc::new(Mutex::new(HashSet::new())), provider_credit_authority, + state_retry_policy: BusyRetryPolicy::default(), + #[cfg(test)] + state_fault_plan: WorkflowStateFaultPlan::default(), } } + #[cfg(test)] + fn with_state_fault_plan(mut self, state_fault_plan: WorkflowStateFaultPlan) -> Self { + self.state_fault_plan = state_fault_plan; + self + } + + #[cfg(test)] + fn with_state_retry_policy(mut self, state_retry_policy: BusyRetryPolicy) -> Self { + self.state_retry_policy = state_retry_policy; + self + } + + async fn retry_state( + &self, + operation: WorkflowStateOperation, + mut f: F, + ) -> anyhow::Result + where + F: FnMut() -> Fut, + Fut: Future>, + { + retry_workflow_state_with_policy(self.state_retry_policy, operation, || { + #[cfg(test)] + let planned_fault = self.state_fault_plan.take(operation); + let future = f(); + async move { + #[cfg(test)] + if let Some(error) = planned_fault { + return Err(error); + } + future.await + } + }) + .await + } + + async fn persist_activation_database_busy( + &self, + run_id: &str, + generation: i64, + operation: WorkflowStateOperation, + ) -> anyhow::Result> { + let params = WorkflowRunActivationDatabaseBusyParams { + run_id: run_id.to_string(), + owner_id: self.owner_instance_id.to_string(), + generation, + operation: operation.label().to_string(), + }; + self.retry_state(WorkflowStateOperation::RecordActivationBusy, || { + self.state_db + .block_workflow_run_activation_database_busy(params.clone()) + }) + .await + .map(|outcome| outcome.map(|outcome| outcome.snapshot)) + } + pub async fn start_workflow_run( &self, request: WorkflowStartRequest, @@ -346,7 +495,7 @@ impl WorkflowActivationService { idempotency_key: request.idempotency_key.clone(), }; let snapshot = if create_params.idempotency_key.is_some() { - retry_workflow_state(WorkflowStateOperation::CreateRun, || { + self.retry_state(WorkflowStateOperation::CreateRun, || { self.state_db .workflows() .create_workflow_run(create_params.clone()) @@ -364,11 +513,12 @@ impl WorkflowActivationService { thread_id: request.source_thread_id, idempotency_key: request.idempotency_key.clone(), }; - let goal_plan = match retry_workflow_state(WorkflowStateOperation::ProjectGoalPlan, || { - self.state_db - .project_workflow_run_to_goal_plan(projection_params.clone()) - }) - .await + let goal_plan = match self + .retry_state(WorkflowStateOperation::ProjectGoalPlan, || { + self.state_db + .project_workflow_run_to_goal_plan(projection_params.clone()) + }) + .await { Ok(goal_plan) => goal_plan, Err(error) => { @@ -388,7 +538,7 @@ impl WorkflowActivationService { .await); } }; - if snapshot.run.status.is_terminal() { + if snapshot.run.status.is_terminal() || snapshot.run.status == WorkflowRunStatus::Blocked { return Ok(WorkflowStartOutcome { snapshot, goal_plan, @@ -400,10 +550,11 @@ impl WorkflowActivationService { owner_id: self.owner_instance_id.to_string(), lease_duration_ms: None, }; - let initial_claim = match retry_workflow_state(WorkflowStateOperation::ClaimRun, || { - self.state_db.claim_workflow_run(claim_params.clone()) - }) - .await + let initial_claim = match self + .retry_state(WorkflowStateOperation::ClaimRun, || { + self.state_db.claim_workflow_run(claim_params.clone()) + }) + .await { Ok(initial_claim) => initial_claim, Err(error) => { @@ -454,13 +605,14 @@ impl WorkflowActivationService { }; } let owns_initial_claim = initial_claim.is_some(); - let snapshot = retry_workflow_state(WorkflowStateOperation::LoadSnapshot, || { - self.state_db - .workflows() - .get_workflow_run_snapshot(snapshot.run.run_id.as_str()) - }) - .await? - .ok_or_else(|| anyhow::anyhow!("workflow run disappeared after credit reservation"))?; + let snapshot = self + .retry_state(WorkflowStateOperation::LoadSnapshot, || { + self.state_db + .workflows() + .get_workflow_run_snapshot(snapshot.run.run_id.as_str()) + }) + .await? + .ok_or_else(|| anyhow::anyhow!("workflow run disappeared after credit reservation"))?; let initial_admission_errors = if let Some(claim) = initial_claim { self.activate_with_claim( snapshot.run.run_id.clone(), @@ -514,11 +666,12 @@ impl WorkflowActivationService { ) -> anyhow::Result { tokio::time::timeout(WORKFLOW_INITIAL_ADMISSION_TIMEOUT, async { loop { - let snapshot = retry_workflow_state(WorkflowStateOperation::LoadSnapshot, || { - self.state_db.workflows().get_workflow_run_snapshot(run_id) - }) - .await? - .ok_or_else(|| anyhow::anyhow!("workflow run disappeared during admission"))?; + let snapshot = self + .retry_state(WorkflowStateOperation::LoadSnapshot, || { + self.state_db.workflows().get_workflow_run_snapshot(run_id) + }) + .await? + .ok_or_else(|| anyhow::anyhow!("workflow run disappeared during admission"))?; if initial_admission_is_observable(&snapshot) { return Ok(snapshot); } @@ -660,14 +813,15 @@ impl WorkflowActivationService { ) -> anyhow::Result<()> { let mut cursor = None; loop { - let page = retry_workflow_state(WorkflowStateOperation::ListThreadRuns, || { - self.state_db.workflows().list_thread_workflow_runs_page( - thread_id, - cursor, - codex_state::DEFAULT_THREAD_WORKFLOW_RUN_LIST_LIMIT, - ) - }) - .await?; + let page = self + .retry_state(WorkflowStateOperation::ListThreadRuns, || { + self.state_db.workflows().list_thread_workflow_runs_page( + thread_id, + cursor, + codex_state::DEFAULT_THREAD_WORKFLOW_RUN_LIST_LIMIT, + ) + }) + .await?; for snapshot in page.data { self.reconcile_failed_goal_plan_projection(&snapshot) .await?; @@ -678,7 +832,8 @@ impl WorkflowActivationService { .is_some_and(|credit| !credit.status.is_terminal()); if needs_terminal_reconciliation || (!snapshot.run.status.is_terminal() - && snapshot.run.status != WorkflowRunStatus::Paused) + && snapshot.run.status != WorkflowRunStatus::Paused + && snapshot.run.status != WorkflowRunStatus::Blocked) { self.activate(snapshot.run.run_id, config.clone()).await; } @@ -695,7 +850,7 @@ impl WorkflowActivationService { snapshot: &WorkflowRunSnapshot, ) -> anyhow::Result<()> { if snapshot.run.status == WorkflowRunStatus::Failed { - retry_workflow_state(WorkflowStateOperation::ReconcileFailedProjection, || { + self.retry_state(WorkflowStateOperation::ReconcileFailedProjection, || { self.state_db .reconcile_failed_workflow_goal_plan_projection(snapshot.run.run_id.as_str()) }) @@ -721,24 +876,26 @@ impl WorkflowActivationService { .as_ref() .is_some_and(|snapshot| snapshot.run.status == WorkflowRunStatus::CancelRequested) { - let claim = retry_workflow_state(WorkflowStateOperation::ClaimRun, || { - self.state_db.claim_workflow_run(WorkflowRunClaimParams { - run_id: run_id.clone(), - owner_id: self.owner_instance_id.to_string(), - lease_duration_ms: None, - }) - }) - .await?; - if let Some(claim) = claim { - let advanced = retry_workflow_state(WorkflowStateOperation::AdvanceRun, || { - self.state_db - .advance_workflow_run(WorkflowRunAdvanceParams { - run_id: run_id.clone(), - owner_id: self.owner_instance_id.to_string(), - generation: claim.generation, - }) + let claim = self + .retry_state(WorkflowStateOperation::ClaimRun, || { + self.state_db.claim_workflow_run(WorkflowRunClaimParams { + run_id: run_id.clone(), + owner_id: self.owner_instance_id.to_string(), + lease_duration_ms: None, + }) }) .await?; + if let Some(claim) = claim { + let advanced = self + .retry_state(WorkflowStateOperation::AdvanceRun, || { + self.state_db + .advance_workflow_run(WorkflowRunAdvanceParams { + run_id: run_id.clone(), + owner_id: self.owner_instance_id.to_string(), + generation: claim.generation, + }) + }) + .await?; if let Some(advanced) = advanced && advanced.snapshot.run.status == WorkflowRunStatus::Cancelled { @@ -747,12 +904,13 @@ impl WorkflowActivationService { } self.activate(run_id.clone(), config).await; for _ in 0..20 { - let current = retry_workflow_state(WorkflowStateOperation::LoadSnapshot, || { - self.state_db - .workflows() - .get_workflow_run_snapshot(run_id.as_str()) - }) - .await?; + let current = self + .retry_state(WorkflowStateOperation::LoadSnapshot, || { + self.state_db + .workflows() + .get_workflow_run_snapshot(run_id.as_str()) + }) + .await?; if current .as_ref() .is_none_or(|snapshot| snapshot.run.status.is_terminal()) @@ -762,7 +920,7 @@ impl WorkflowActivationService { tokio::time::sleep(WORKFLOW_SUPERVISOR_POLL_INTERVAL).await; } } - retry_workflow_state(WorkflowStateOperation::LoadSnapshot, || { + self.retry_state(WorkflowStateOperation::LoadSnapshot, || { self.state_db .workflows() .get_workflow_run_snapshot(run_id.as_str()) @@ -806,6 +964,33 @@ impl WorkflowActivationService { else { break; }; + if let Some(operation) = err + .downcast_ref::() + .map(|busy| busy.operation) + { + let persisted = service + .persist_activation_database_busy( + run_id.as_str(), + generation, + operation, + ) + .await; + if let Some(sender) = initial_admission_sender.take() { + let result = match persisted { + Ok(Some(snapshot)) => Ok(snapshot), + Ok(None) => Err(err), + Err(persist_error) => Err(persist_error), + }; + let _ = sender.send(result); + } + tracing::warn!( + workflow_run_id = %run_id, + operation = operation.label(), + "workflow activation stopped after bounded database busy retries" + ); + service.active_supervisors.lock().await.remove(&run_id); + return; + } if let Some(sender) = initial_admission_sender.take() { tracing::warn!( workflow_run_id = %run_id, @@ -849,6 +1034,17 @@ impl WorkflowActivationService { ); break; } + Err(err) + if err + .downcast_ref::() + .is_some() => + { + tracing::warn!( + workflow_run_id = %run_id, + "workflow activation supervisor stopped after bounded database busy retries: {err}" + ); + break; + } Err(err) => { tracing::warn!( workflow_run_id = %run_id, @@ -869,10 +1065,11 @@ impl WorkflowActivationService { mut config: WorkflowActivationConfig, ) -> anyhow::Result<()> { loop { - let Some(snapshot) = retry_workflow_state(WorkflowStateOperation::LoadSnapshot, || { - self.state_db.workflows().get_workflow_run_snapshot(run_id) - }) - .await? + let Some(snapshot) = self + .retry_state(WorkflowStateOperation::LoadSnapshot, || { + self.state_db.workflows().get_workflow_run_snapshot(run_id) + }) + .await? else { return Ok(()); }; @@ -882,6 +1079,9 @@ impl WorkflowActivationService { tokio::time::sleep(WORKFLOW_SUPERVISOR_POLL_INTERVAL).await; continue; } + if snapshot.run.status == WorkflowRunStatus::Blocked { + return Ok(()); + } if snapshot.run.status.is_terminal() && snapshot .provider_credit @@ -937,10 +1137,11 @@ impl WorkflowActivationService { owner_id: self.owner_instance_id.to_string(), lease_duration_ms: None, }; - let Some(claim) = retry_workflow_state(WorkflowStateOperation::ClaimRun, || { - self.state_db.claim_workflow_run(claim_params.clone()) - }) - .await? + let Some(claim) = self + .retry_state(WorkflowStateOperation::ClaimRun, || { + self.state_db.claim_workflow_run(claim_params.clone()) + }) + .await? else { tokio::time::sleep(WORKFLOW_SUPERVISOR_POLL_INTERVAL).await; continue; @@ -967,14 +1168,26 @@ impl WorkflowActivationService { .await?; let mut initial_admission_sender = None; - self.drive_generation( - run_id, - claim.generation, - &spec, - config.clone(), - &mut initial_admission_sender, - ) - .await?; + if let Err(err) = self + .drive_generation( + run_id, + claim.generation, + &spec, + config.clone(), + &mut initial_admission_sender, + ) + .await + { + if let Some(operation) = err + .downcast_ref::() + .map(|busy| busy.operation) + { + self.persist_activation_database_busy(run_id, claim.generation, operation) + .await?; + return Ok(()); + } + return Err(err); + } } } @@ -1009,12 +1222,13 @@ impl WorkflowActivationService { generation, lease_duration_ms: None, }; - if retry_workflow_state(WorkflowStateOperation::HeartbeatRun, || { - self.state_db - .heartbeat_workflow_run(heartbeat_params.clone()) - }) - .await? - .is_none() + if self + .retry_state(WorkflowStateOperation::HeartbeatRun, || { + self.state_db + .heartbeat_workflow_run(heartbeat_params.clone()) + }) + .await? + .is_none() { return Ok(()); } @@ -1029,10 +1243,11 @@ impl WorkflowActivationService { owner_id: self.owner_instance_id.to_string(), generation, }; - let Some(advanced) = retry_workflow_state(WorkflowStateOperation::AdvanceRun, || { - self.state_db.advance_workflow_run(advance_params.clone()) - }) - .await? + let Some(advanced) = self + .retry_state(WorkflowStateOperation::AdvanceRun, || { + self.state_db.advance_workflow_run(advance_params.clone()) + }) + .await? else { return Ok(()); }; @@ -1053,8 +1268,8 @@ impl WorkflowActivationService { owner_id: self.owner_instance_id.to_string(), generation, }; - let Some(reconciled) = - retry_workflow_state(WorkflowStateOperation::ReconcileBranches, || { + let Some(reconciled) = self + .retry_state(WorkflowStateOperation::ReconcileBranches, || { self.state_db .reconcile_workflow_run_branches(reconcile_params.clone()) }) @@ -1094,8 +1309,8 @@ impl WorkflowActivationService { parent_agent_run_id: None, max_active_background_agent_runs: config.max_active_background_agent_runs, }; - let Some(admitted) = - retry_workflow_state(WorkflowStateOperation::AdmitBranches, || { + let Some(admitted) = self + .retry_state(WorkflowStateOperation::AdmitBranches, || { self.state_db .admit_workflow_run_branches(admission_params.clone()) }) @@ -1124,8 +1339,8 @@ impl WorkflowActivationService { generation, selection: WorkflowRunVerifierClaimSelection::NextRunCommands, }; - if let Some(claimed_verifier) = - retry_workflow_state(WorkflowStateOperation::ClaimVerifier, || { + if let Some(claimed_verifier) = self + .retry_state(WorkflowStateOperation::ClaimVerifier, || { self.state_db .claim_workflow_run_verifier(verifier_claim_params.clone()) }) @@ -1147,8 +1362,8 @@ impl WorkflowActivationService { generation, selection: WorkflowRunVerifierClaimSelection::VerifierRunId(verifier_run_id), }; - if let Some(claimed_verifier) = - retry_workflow_state(WorkflowStateOperation::ClaimVerifier, || { + if let Some(claimed_verifier) = self + .retry_state(WorkflowStateOperation::ClaimVerifier, || { self.state_db .claim_workflow_run_verifier(verifier_claim_params.clone()) }) @@ -1404,11 +1619,12 @@ impl WorkflowActivationService { generation, }; if fence_lost.load(Ordering::Relaxed) - || !retry_workflow_state(WorkflowStateOperation::CheckVerifierFence, || { - self.state_db - .workflow_run_fence_is_current(fence_params.clone()) - }) - .await? + || !self + .retry_state(WorkflowStateOperation::CheckVerifierFence, || { + self.state_db + .workflow_run_fence_is_current(fence_params.clone()) + }) + .await? { return Ok(false); } @@ -1433,11 +1649,12 @@ impl WorkflowActivationService { output_truncated, }, }; - let recorded = retry_workflow_state(WorkflowStateOperation::RecordVerifierResult, || { - self.state_db - .record_workflow_run_verifier_result(record_params.clone()) - }) - .await?; + let recorded = self + .retry_state(WorkflowStateOperation::RecordVerifierResult, || { + self.state_db + .record_workflow_run_verifier_result(record_params.clone()) + }) + .await?; Ok(recorded.is_some()) } @@ -1489,11 +1706,12 @@ impl WorkflowActivationService { owner_id: self.owner_instance_id.to_string(), generation, }; - if !retry_workflow_state(WorkflowStateOperation::CheckVerifierFence, || { - self.state_db - .workflow_run_fence_is_current(fence_params.clone()) - }) - .await? + if !self + .retry_state(WorkflowStateOperation::CheckVerifierFence, || { + self.state_db + .workflow_run_fence_is_current(fence_params.clone()) + }) + .await? { return Ok(false); } @@ -1514,11 +1732,12 @@ impl WorkflowActivationService { output_truncated: false, }, }; - let recorded = retry_workflow_state(WorkflowStateOperation::RecordVerifierResult, || { - self.state_db - .record_workflow_run_verifier_result(record_params.clone()) - }) - .await?; + let recorded = self + .retry_state(WorkflowStateOperation::RecordVerifierResult, || { + self.state_db + .record_workflow_run_verifier_result(record_params.clone()) + }) + .await?; Ok(recorded.is_some()) } @@ -1546,11 +1765,12 @@ impl WorkflowActivationService { output_truncated: false, }, }; - let recorded = retry_workflow_state(WorkflowStateOperation::RecordVerifierResult, || { - self.state_db - .record_workflow_run_verifier_result(record_params.clone()) - }) - .await?; + let recorded = self + .retry_state(WorkflowStateOperation::RecordVerifierResult, || { + self.state_db + .record_workflow_run_verifier_result(record_params.clone()) + }) + .await?; Ok(recorded.is_some()) } @@ -1575,10 +1795,11 @@ impl WorkflowActivationService { generation, lease_duration_ms: None, }; - match retry_workflow_state(WorkflowStateOperation::HeartbeatRun, || { - state_db.heartbeat_workflow_run(heartbeat_params.clone()) - }) - .await + match self + .retry_state(WorkflowStateOperation::HeartbeatRun, || { + state_db.heartbeat_workflow_run(heartbeat_params.clone()) + }) + .await { Ok(Some(_)) => {} Ok(None) | Err(_) => { @@ -2330,6 +2551,11 @@ mod tests { use codex_state::WorkflowSpecCreateParams; use pretty_assertions::assert_eq; use serde_json::json; + use sqlx::ConnectOptions; + use sqlx::Connection; + use sqlx::sqlite::SqliteConnectOptions; + use sqlx::sqlite::SqliteConnection; + use sqlx::sqlite::SqliteJournalMode; use std::sync::Arc; use std::sync::Mutex as StdMutex; use std::sync::atomic::AtomicUsize; @@ -2377,6 +2603,154 @@ mod tests { ) -> anyhow::Result<()> { Ok(()) } + + struct ActivationTestFixture { + _temp_dir: tempfile::TempDir, + state_db: Arc, + source_thread_id: ThreadId, + workflow_record_id: String, + } + + async fn activation_test_fixture() -> ActivationTestFixture { + let temp_dir = tempfile::tempdir().expect("create activation fixture"); + let codex_home = temp_dir.path().join("codex-home"); + let source_cwd = temp_dir.path().join("project-workspace"); + std::fs::create_dir_all(&source_cwd).expect("create non-repository source directory"); + let state_db = StateRuntime::init(codex_home.clone(), "activation-test".to_string()) + .await + .expect("state db should initialize"); + let source_thread_id = ThreadId::new(); + let mut thread = ThreadMetadataBuilder::new( + source_thread_id, + codex_home.join(format!("rollout-{source_thread_id}.jsonl")), + Utc::now(), + SessionSource::Cli, + ); + thread.cwd = source_cwd; + state_db + .upsert_thread(&thread.build("activation-test")) + .await + .expect("source thread should be persisted"); + let spec = state_db + .workflows() + .save_workflow_spec_yaml(WorkflowSpecCreateParams { + source_thread_id: Some(source_thread_id), + source_yaml: ROUTE_ACTIVATION_WORKFLOW_YAML.to_string(), + }) + .await + .expect("workflow spec should save"); + ActivationTestFixture { + _temp_dir: temp_dir, + state_db, + source_thread_id, + workflow_record_id: spec.workflow_record_id, + } + } + + async fn real_sqlite_contention_errors( + busy_count: usize, + busy_snapshot_count: usize, + ) -> Vec { + let temp_dir = tempfile::tempdir().expect("create SQLite contention fixture"); + let database_path = temp_dir.path().join("contention.sqlite"); + let options = SqliteConnectOptions::new() + .filename(&database_path) + .create_if_missing(true) + .journal_mode(SqliteJournalMode::Wal) + .busy_timeout(Duration::ZERO) + .disable_statement_logging(); + let mut holder = SqliteConnection::connect_with(&options) + .await + .expect("open lock holder"); + let mut contender = SqliteConnection::connect_with(&options) + .await + .expect("open lock contender"); + sqlx::query("CREATE TABLE contention (id INTEGER PRIMARY KEY, value INTEGER NOT NULL)") + .execute(&mut holder) + .await + .expect("create contention table"); + sqlx::query("INSERT INTO contention (id, value) VALUES (1, 0)") + .execute(&mut holder) + .await + .expect("seed contention table"); + + sqlx::query("BEGIN IMMEDIATE") + .execute(&mut holder) + .await + .expect("hold write transaction"); + let mut errors = Vec::with_capacity(busy_count + busy_snapshot_count); + for _ in 0..busy_count { + let error = sqlx::query("UPDATE contention SET value = value + 1 WHERE id = 1") + .execute(&mut contender) + .await + .expect_err("contending writer should receive SQLITE_BUSY"); + assert_eq!( + Some("5"), + error + .as_database_error() + .and_then(|database_error| database_error.code()) + .as_deref() + ); + errors.push(anyhow::Error::new(error)); + } + sqlx::query("ROLLBACK") + .execute(&mut holder) + .await + .expect("release write transaction"); + + let mut reader = SqliteConnection::connect_with(&options) + .await + .expect("open snapshot reader"); + let mut writer = SqliteConnection::connect_with(&options) + .await + .expect("open snapshot writer"); + sqlx::query("BEGIN") + .execute(&mut reader) + .await + .expect("begin read transaction"); + let _: i64 = sqlx::query_scalar("SELECT value FROM contention WHERE id = 1") + .fetch_one(&mut reader) + .await + .expect("establish read snapshot"); + sqlx::query("BEGIN IMMEDIATE") + .execute(&mut writer) + .await + .expect("begin concurrent writer"); + sqlx::query("UPDATE contention SET value = value + 1 WHERE id = 1") + .execute(&mut writer) + .await + .expect("advance WAL"); + sqlx::query("COMMIT") + .execute(&mut writer) + .await + .expect("commit concurrent writer"); + for _ in 0..busy_snapshot_count { + let error = sqlx::query("UPDATE contention SET value = value + 1 WHERE id = 1") + .execute(&mut reader) + .await + .expect_err("stale reader should receive SQLITE_BUSY_SNAPSHOT"); + assert_eq!( + Some("517"), + error + .as_database_error() + .and_then(|database_error| database_error.code()) + .as_deref() + ); + errors.push(anyhow::Error::new(error)); + } + sqlx::query("ROLLBACK") + .execute(&mut reader) + .await + .expect("release stale read transaction"); + errors + } + + fn tiny_state_retry_policy() -> BusyRetryPolicy { + BusyRetryPolicy { + initial_delay: Duration::from_millis(1), + max_delay: Duration::from_millis(1), + total_delay_budget: Duration::from_millis(3), + } } #[async_trait] @@ -3021,6 +3395,357 @@ cleanup: assert_eq!(WorkflowRunStatus::Cancelled, cancelled.run.status); } + #[tokio::test] + async fn initial_start_retries_real_busy_and_busy_snapshot_without_duplicate_admission() { + let fixture = activation_test_fixture().await; + let fault_plan = WorkflowStateFaultPlan::default(); + for error in + real_sqlite_contention_errors(/*busy_count*/ 1, /*busy_snapshot_count*/ 1).await + { + fault_plan.push(WorkflowStateOperation::AdmitBranches, error); + } + let service = WorkflowActivationService::new(Arc::clone(&fixture.state_db)) + .with_state_fault_plan(fault_plan) + .with_state_retry_policy(tiny_state_retry_policy()); + let activation_config = WorkflowActivationConfig { + route_runtime: supported_route_runtime(), + max_active_background_agent_runs: Some(10), + ..Default::default() + }; + let idempotency_key = "initial-start-real-sqlite-contention"; + let outcome = service + .start_workflow_run(WorkflowStartRequest { + workflow_record_id: fixture.workflow_record_id.clone(), + source_thread_id: fixture.source_thread_id, + idempotency_key: Some(idempotency_key.to_string()), + activation_config: activation_config.clone(), + }) + .await + .expect("bounded contention should admit the initial branch"); + let admitted_step = outcome + .snapshot + .steps + .iter() + .find(|step| step.step_id == "run") + .expect("initial step should exist"); + let background_agent_run_id = admitted_step + .background_agent_run_id + .clone() + .expect("initial step should admit one background AgentRun"); + assert_eq!(WorkflowRunStepStatus::Active, admitted_step.status); + assert_eq!( + 1, + outcome + .snapshot + .events + .iter() + .filter(|event| event.event_type == "branch_admitted") + .count() + ); + assert_eq!( + 1, + fixture + .state_db + .list_background_agent_runs(/*limit*/ None) + .await + .expect("background AgentRuns should list") + .len() + ); + + let replay = service + .start_workflow_run(WorkflowStartRequest { + workflow_record_id: fixture.workflow_record_id, + source_thread_id: fixture.source_thread_id, + idempotency_key: Some(idempotency_key.to_string()), + activation_config: activation_config.clone(), + }) + .await + .expect("idempotent replay should preserve the admitted branch"); + let replayed_step = replay + .snapshot + .steps + .iter() + .find(|step| step.step_id == "run") + .expect("replayed initial step should exist"); + assert_eq!( + Some(background_agent_run_id.as_str()), + replayed_step.background_agent_run_id.as_deref() + ); + assert_eq!( + 1, + replay + .snapshot + .events + .iter() + .filter(|event| event.event_type == "branch_admitted") + .count() + ); + assert_eq!( + 1, + fixture + .state_db + .list_background_agent_runs(/*limit*/ None) + .await + .expect("background AgentRuns should list after replay") + .len() + ); + + let cancelled = service + .cancel_workflow_run( + WorkflowRunCancelParams { + run_id: outcome.snapshot.run.run_id, + reason: "initial SQLite contention regression cleanup".to_string(), + }, + activation_config, + ) + .await + .expect("workflow cancellation should succeed") + .expect("workflow run should exist"); + assert_eq!(WorkflowRunStatus::Cancelled, cancelled.run.status); + } + + #[tokio::test] + async fn thread_start_recovery_retries_real_busy_and_busy_snapshot_once() { + let fixture = activation_test_fixture().await; + let idempotency_key = "thread-start-real-sqlite-contention"; + let created = fixture + .state_db + .workflows() + .create_workflow_run(WorkflowRunCreateParams { + workflow_record_id: fixture.workflow_record_id, + source_thread_id: Some(fixture.source_thread_id), + idempotency_key: Some(idempotency_key.to_string()), + }) + .await + .expect("workflow run should create without activation"); + fixture + .state_db + .project_workflow_run_to_goal_plan(WorkflowGoalPlanProjectionParams { + workflow_run_id: created.run.run_id.clone(), + thread_id: fixture.source_thread_id, + idempotency_key: Some(idempotency_key.to_string()), + }) + .await + .expect("workflow run should project") + .expect("workflow projection should exist"); + + let fault_plan = WorkflowStateFaultPlan::default(); + for error in + real_sqlite_contention_errors(/*busy_count*/ 1, /*busy_snapshot_count*/ 1).await + { + fault_plan.push(WorkflowStateOperation::AdmitBranches, error); + } + let service = WorkflowActivationService::new(Arc::clone(&fixture.state_db)) + .with_state_fault_plan(fault_plan) + .with_state_retry_policy(tiny_state_retry_policy()); + let activation_config = WorkflowActivationConfig { + route_runtime: supported_route_runtime(), + max_active_background_agent_runs: Some(10), + ..Default::default() + }; + service + .activate_thread_runs(fixture.source_thread_id, activation_config.clone()) + .await + .expect("thread start should schedule workflow recovery"); + let recovered = tokio::time::timeout(Duration::from_secs(5), async { + loop { + let snapshot = fixture + .state_db + .workflows() + .get_workflow_run_snapshot(created.run.run_id.as_str()) + .await + .expect("workflow snapshot should load") + .expect("workflow run should exist"); + if snapshot + .steps + .iter() + .any(|step| step.step_id == "run" && step.background_agent_run_id.is_some()) + { + return snapshot; + } + tokio::time::sleep(WORKFLOW_INITIAL_ADMISSION_POLL_INTERVAL).await; + } + }) + .await + .expect("thread start recovery should admit after bounded contention"); + let recovered_step = recovered + .steps + .iter() + .find(|step| step.step_id == "run") + .expect("recovered initial step should exist"); + assert_eq!(WorkflowRunStepStatus::Active, recovered_step.status); + assert!(recovered_step.background_agent_run_id.is_some()); + assert_eq!( + 1, + recovered + .events + .iter() + .filter(|event| event.event_type == "branch_admitted") + .count() + ); + assert_eq!( + 1, + fixture + .state_db + .list_background_agent_runs(/*limit*/ None) + .await + .expect("background AgentRuns should list") + .len() + ); + + service + .activate_thread_runs(fixture.source_thread_id, activation_config.clone()) + .await + .expect("repeated thread start should not duplicate the supervisor"); + let after_replay = fixture + .state_db + .workflows() + .get_workflow_run_snapshot(created.run.run_id.as_str()) + .await + .expect("workflow snapshot should load after repeated thread start") + .expect("workflow run should exist"); + assert_eq!( + 1, + after_replay + .events + .iter() + .filter(|event| event.event_type == "branch_admitted") + .count() + ); + assert_eq!( + 1, + fixture + .state_db + .list_background_agent_runs(/*limit*/ None) + .await + .expect("background AgentRuns should list after repeated thread start") + .len() + ); + + let cancelled = service + .cancel_workflow_run( + WorkflowRunCancelParams { + run_id: created.run.run_id, + reason: "thread start SQLite contention regression cleanup".to_string(), + }, + activation_config, + ) + .await + .expect("workflow cancellation should succeed") + .expect("workflow run should exist"); + assert_eq!(WorkflowRunStatus::Cancelled, cancelled.run.status); + } + + #[tokio::test] + async fn always_busy_initial_start_stops_with_one_structured_block() { + let fixture = activation_test_fixture().await; + let fault_plan = WorkflowStateFaultPlan::default(); + for error in + real_sqlite_contention_errors(/*busy_count*/ 8, /*busy_snapshot_count*/ 0).await + { + fault_plan.push(WorkflowStateOperation::AdmitBranches, error); + } + let service = WorkflowActivationService::new(Arc::clone(&fixture.state_db)) + .with_state_fault_plan(fault_plan) + .with_state_retry_policy(tiny_state_retry_policy()); + let request = WorkflowStartRequest { + workflow_record_id: fixture.workflow_record_id, + source_thread_id: fixture.source_thread_id, + idempotency_key: Some("always-busy-initial-start".to_string()), + activation_config: WorkflowActivationConfig { + route_runtime: supported_route_runtime(), + max_active_background_agent_runs: Some(10), + ..Default::default() + }, + }; + let outcome = tokio::time::timeout( + Duration::from_secs(2), + service.start_workflow_run(request.clone()), + ) + .await + .expect("always-busy start should terminate within the finite retry bound") + .expect("always-busy start should return the durable blocked snapshot"); + let blocked_step = outcome + .snapshot + .steps + .iter() + .find(|step| step.step_id == "run") + .expect("blocked initial step should exist"); + assert_eq!(WorkflowRunStatus::Blocked, outcome.snapshot.run.status); + assert_eq!( + Some(WORKFLOW_ACTIVATION_DATABASE_BUSY_REASON_CODE), + outcome.snapshot.run.reason_code.as_deref() + ); + assert_eq!(WorkflowRunStepStatus::Blocked, blocked_step.status); + assert_eq!( + Some(WORKFLOW_ACTIVATION_DATABASE_BUSY_REASON_CODE), + blocked_step.reason_code.as_deref() + ); + assert_eq!(None, blocked_step.background_agent_run_id); + assert_eq!( + 1, + outcome + .snapshot + .events + .iter() + .filter(|event| { + event.event_type == WORKFLOW_ACTIVATION_DATABASE_BUSY_REASON_CODE + }) + .count() + ); + assert_eq!( + 0, + fixture + .state_db + .list_background_agent_runs(/*limit*/ None) + .await + .expect("background AgentRuns should list") + .len() + ); + + let replay = service + .start_workflow_run(request) + .await + .expect("blocked idempotent replay should return without reactivation"); + assert_eq!(outcome.snapshot.run.run_id, replay.snapshot.run.run_id); + assert_eq!(WorkflowRunStatus::Blocked, replay.snapshot.run.status); + assert_eq!( + 1, + replay + .snapshot + .events + .iter() + .filter(|event| { + event.event_type == WORKFLOW_ACTIVATION_DATABASE_BUSY_REASON_CODE + }) + .count() + ); + assert_eq!( + 0, + fixture + .state_db + .list_background_agent_runs(/*limit*/ None) + .await + .expect("background AgentRuns should remain empty") + .len() + ); + tokio::time::timeout(Duration::from_secs(1), async { + loop { + if !service + .active_supervisors + .lock() + .await + .contains(outcome.snapshot.run.run_id.as_str()) + { + return; + } + tokio::task::yield_now().await; + } + }) + .await + .expect("always-busy supervisor should stop instead of retrying forever"); + } + #[tokio::test] async fn initial_admission_failure_relinquishes_supervisor_before_recovery() { let temp_dir = tempfile::tempdir().expect("create state home"); @@ -3547,6 +4272,7 @@ cleanup: "advance workflow run", "reconcile workflow run branches", "admit workflow run branches", + "record workflow activation database busy", "claim workflow run verifier", "check workflow verifier fence", "record workflow verifier result", diff --git a/codex-rs/state/src/lib.rs b/codex-rs/state/src/lib.rs index 4fd19c3be..259a5e228 100644 --- a/codex-rs/state/src/lib.rs +++ b/codex-rs/state/src/lib.rs @@ -265,6 +265,7 @@ pub use runtime::UsageProfileLeaseReleaseParams; pub use runtime::UsageProfileLeaseRenewParams; pub use runtime::UsageProfileLeaseValidateParams; pub use runtime::WEBHOOK_EVENT_DEDUPE_CONFLICT_MESSAGE; +pub use runtime::WORKFLOW_ACTIVATION_DATABASE_BUSY_REASON_CODE; pub use runtime::WORKFLOW_BRANCH_ADMISSION_DEFERRED_EVENT; pub use runtime::WORKFLOW_STEP_APPROVAL_APPROVED; pub use runtime::WORKFLOW_STEP_APPROVAL_PENDING; @@ -286,6 +287,8 @@ pub use runtime::WorkflowProviderCreditReconcileParams; pub use runtime::WorkflowProviderCreditRecoveryClaimParams; pub use runtime::WorkflowProviderCreditRefreshParams; pub use runtime::WorkflowProviderCreditStageReleaseParams; +pub use runtime::WorkflowRunActivationDatabaseBusyOutcome; +pub use runtime::WorkflowRunActivationDatabaseBusyParams; pub use runtime::WorkflowRunAdvanceOutcome; pub use runtime::WorkflowRunAdvanceParams; pub use runtime::WorkflowRunBranchAdmission; diff --git a/codex-rs/state/src/runtime.rs b/codex-rs/state/src/runtime.rs index 34ca71816..503004170 100644 --- a/codex-rs/state/src/runtime.rs +++ b/codex-rs/state/src/runtime.rs @@ -286,7 +286,10 @@ pub use workflow_automation::WorkflowTimerFireCompleteParams; pub use workflow_effects::WorkflowRunVerifierEffectPublishParams; pub use workflow_goal_plan_projections::WorkflowGoalPlanProjectionOutcome; pub use workflow_goal_plan_projections::WorkflowGoalPlanProjectionParams; +pub use workflow_orchestrator::WORKFLOW_ACTIVATION_DATABASE_BUSY_REASON_CODE; pub use workflow_orchestrator::WORKFLOW_BRANCH_ADMISSION_DEFERRED_EVENT; +pub use workflow_orchestrator::WorkflowRunActivationDatabaseBusyOutcome; +pub use workflow_orchestrator::WorkflowRunActivationDatabaseBusyParams; pub use workflow_orchestrator::WorkflowRunAdvanceOutcome; pub use workflow_orchestrator::WorkflowRunAdvanceParams; pub use workflow_orchestrator::WorkflowRunBranchAdmission; diff --git a/codex-rs/state/src/runtime/workflow_orchestrator.rs b/codex-rs/state/src/runtime/workflow_orchestrator.rs index 6739206f4..3316df1f5 100644 --- a/codex-rs/state/src/runtime/workflow_orchestrator.rs +++ b/codex-rs/state/src/runtime/workflow_orchestrator.rs @@ -59,6 +59,9 @@ const WORKFLOW_BRANCH_ADMITTED_REASON: &str = "workflow branch admitted"; const WORKFLOW_BRANCH_ADMITTED_REASON_CODE: &str = "workflow_branch_admitted"; /// Durable evidence that branch admission ran but capacity was unavailable. pub const WORKFLOW_BRANCH_ADMISSION_DEFERRED_EVENT: &str = "branch_admission_deferred"; +/// Durable reason persisted when workflow activation exhausts its bounded +/// SQLite busy retry budget before a ready branch can be admitted. +pub const WORKFLOW_ACTIVATION_DATABASE_BUSY_REASON_CODE: &str = "workflow_activation_database_busy"; const WORKFLOW_BRANCH_PROVIDER_ENV_MISSING_REASON: &str = "OpenRouter workflow branch requires OPENROUTER_API_KEY before admission"; const WORKFLOW_BRANCH_PROVIDER_ENV_MISSING_REASON_CODE: &str = @@ -205,6 +208,20 @@ pub struct WorkflowRunBranchAdmissionOutcome { pub changed: bool, } +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct WorkflowRunActivationDatabaseBusyParams { + pub run_id: String, + pub owner_id: String, + pub generation: i64, + pub operation: String, +} + +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct WorkflowRunActivationDatabaseBusyOutcome { + pub snapshot: crate::WorkflowRunSnapshot, + pub changed: bool, +} + #[derive(Debug, Clone, PartialEq, Eq)] pub struct WorkflowRunBranchReconcileParams { pub run_id: String, @@ -593,6 +610,150 @@ WHERE run_id = ? })) } + pub async fn block_workflow_run_activation_database_busy( + &self, + params: WorkflowRunActivationDatabaseBusyParams, + ) -> anyhow::Result> { + validate_owner_id(¶ms.owner_id)?; + let operation = params.operation.trim(); + if operation.is_empty() { + anyhow::bail!("workflow activation database busy operation must not be empty"); + } + let status_reason = format!( + "workflow activation database remained busy after bounded retries during {operation}" + ); + let now_ms = datetime_to_epoch_millis(Utc::now()); + let mut tx = self.pool.begin().await?; + let Some(run) = claim_checked_workflow_run_in_tx( + &mut tx, + params.run_id.as_str(), + params.owner_id.as_str(), + self.workflow_owner_instance_id.as_str(), + params.generation, + now_ms, + ) + .await? + else { + tx.commit().await?; + return Ok(None); + }; + if run.status.is_terminal() + || run.status == crate::WorkflowRunStatus::CancelRequested + || run.status == crate::WorkflowRunStatus::Paused + || run.status == crate::WorkflowRunStatus::Blocked + { + let snapshot = snapshot_workflow_run_in_tx(&mut tx, params.run_id.as_str()).await?; + tx.commit().await?; + return Ok(Some(WorkflowRunActivationDatabaseBusyOutcome { + snapshot, + changed: false, + })); + } + + let blocked_steps = sqlx::query( + r#" +UPDATE workflow_run_steps +SET + status = ?, + status_reason = ?, + reason_code = ?, + updated_at_ms = ? +WHERE run_id = ? + AND status = 'ready' + AND background_agent_run_id IS NULL +RETURNING step_run_id, step_id + "#, + ) + .bind(crate::WorkflowRunStepStatus::Blocked.as_str()) + .bind(status_reason.as_str()) + .bind(WORKFLOW_ACTIVATION_DATABASE_BUSY_REASON_CODE) + .bind(now_ms) + .bind(params.run_id.as_str()) + .fetch_all(&mut *tx) + .await?; + if !blocked_steps.is_empty() { + for row in &blocked_steps { + let step_id: String = row.try_get("step_id")?; + sqlx::query( + r#" +UPDATE workflow_run_step_verifiers +SET + status = ?, + status_reason = ?, + reason_code = ?, + completed_at_ms = ?, + updated_at_ms = ? +WHERE run_id = ? + AND step_id = ? + AND status NOT IN ('passed', 'failed', 'skipped', 'blocked') + "#, + ) + .bind(crate::WorkflowRunStepVerifierStatus::Blocked.as_str()) + .bind(status_reason.as_str()) + .bind(WORKFLOW_ACTIVATION_DATABASE_BUSY_REASON_CODE) + .bind(now_ms) + .bind(now_ms) + .bind(params.run_id.as_str()) + .bind(step_id) + .execute(&mut *tx) + .await?; + } + sqlx::query( + r#" +UPDATE workflow_runs +SET + status = ?, + status_reason = ?, + reason_code = ?, + updated_at_ms = ? +WHERE run_id = ? + "#, + ) + .bind(crate::WorkflowRunStatus::Blocked.as_str()) + .bind(status_reason.as_str()) + .bind(WORKFLOW_ACTIVATION_DATABASE_BUSY_REASON_CODE) + .bind(now_ms) + .bind(params.run_id.as_str()) + .execute(&mut *tx) + .await?; + let blocked_step_run_ids = blocked_steps + .iter() + .map(|row| row.try_get::("step_run_id")) + .collect::, _>>()?; + append_workflow_run_event_in_tx( + &mut tx, + params.run_id.as_str(), + WorkflowRunEventAppend { + event_type: WORKFLOW_ACTIVATION_DATABASE_BUSY_REASON_CODE, + actor_kind: "orchestrator", + actor_id: Some(params.owner_id.clone()), + step_run_id: None, + verifier_run_id: None, + visibility: "internal", + payload: json!({ + "operation": operation, + "reasonCode": WORKFLOW_ACTIVATION_DATABASE_BUSY_REASON_CODE, + "stepRunIds": blocked_step_run_ids, + }), + now_ms, + }, + ) + .await?; + } + let changed = !blocked_steps.is_empty(); + let snapshot = snapshot_workflow_run_in_tx(&mut tx, params.run_id.as_str()).await?; + tx.commit().await?; + if changed { + self.thread_goals + .block_workflow_goal_plan_projection(params.run_id.as_str()) + .await?; + } + Ok(Some(WorkflowRunActivationDatabaseBusyOutcome { + snapshot, + changed, + })) + } + pub async fn reconcile_workflow_run_branches( &self, params: WorkflowRunBranchReconcileParams, From 3192435f081fac9d267b1bf1b19fc0fb439b6dbb Mon Sep 17 00:00:00 2001 From: Andrei Hasna Date: Tue, 11 Aug 2026 16:51:17 +0300 Subject: [PATCH 2/6] test(workflows): keep busy fixtures lockfile-neutral Inject the exact SQLite code 5 and 517 error shapes already recognized by the shared state retry layer, without adding a test-only dependency that changes Cargo.lock metadata. Agent: iapp-factory-coordinator --- codex-rs/ext/workflows/Cargo.toml | 1 - codex-rs/ext/workflows/src/activation.rs | 111 ++--------------------- 2 files changed, 10 insertions(+), 102 deletions(-) diff --git a/codex-rs/ext/workflows/Cargo.toml b/codex-rs/ext/workflows/Cargo.toml index e95d51794..d0ae2a824 100644 --- a/codex-rs/ext/workflows/Cargo.toml +++ b/codex-rs/ext/workflows/Cargo.toml @@ -46,6 +46,5 @@ chrono = { workspace = true } codex-prompts = { workspace = true } codex-utils-output-truncation = { workspace = true } pretty_assertions = { workspace = true } -sqlx = { workspace = true } tempfile = { workspace = true } tokio = { workspace = true, features = ["macros", "rt"] } diff --git a/codex-rs/ext/workflows/src/activation.rs b/codex-rs/ext/workflows/src/activation.rs index df89afb11..7eee9dd0f 100644 --- a/codex-rs/ext/workflows/src/activation.rs +++ b/codex-rs/ext/workflows/src/activation.rs @@ -2551,11 +2551,6 @@ mod tests { use codex_state::WorkflowSpecCreateParams; use pretty_assertions::assert_eq; use serde_json::json; - use sqlx::ConnectOptions; - use sqlx::Connection; - use sqlx::sqlite::SqliteConnectOptions; - use sqlx::sqlite::SqliteConnection; - use sqlx::sqlite::SqliteJournalMode; use std::sync::Arc; use std::sync::Mutex as StdMutex; use std::sync::atomic::AtomicUsize; @@ -2647,101 +2642,21 @@ mod tests { } } - async fn real_sqlite_contention_errors( + fn sqlite_contention_errors( busy_count: usize, busy_snapshot_count: usize, ) -> Vec { - let temp_dir = tempfile::tempdir().expect("create SQLite contention fixture"); - let database_path = temp_dir.path().join("contention.sqlite"); - let options = SqliteConnectOptions::new() - .filename(&database_path) - .create_if_missing(true) - .journal_mode(SqliteJournalMode::Wal) - .busy_timeout(Duration::ZERO) - .disable_statement_logging(); - let mut holder = SqliteConnection::connect_with(&options) - .await - .expect("open lock holder"); - let mut contender = SqliteConnection::connect_with(&options) - .await - .expect("open lock contender"); - sqlx::query("CREATE TABLE contention (id INTEGER PRIMARY KEY, value INTEGER NOT NULL)") - .execute(&mut holder) - .await - .expect("create contention table"); - sqlx::query("INSERT INTO contention (id, value) VALUES (1, 0)") - .execute(&mut holder) - .await - .expect("seed contention table"); - - sqlx::query("BEGIN IMMEDIATE") - .execute(&mut holder) - .await - .expect("hold write transaction"); let mut errors = Vec::with_capacity(busy_count + busy_snapshot_count); for _ in 0..busy_count { - let error = sqlx::query("UPDATE contention SET value = value + 1 WHERE id = 1") - .execute(&mut contender) - .await - .expect_err("contending writer should receive SQLITE_BUSY"); - assert_eq!( - Some("5"), - error - .as_database_error() - .and_then(|database_error| database_error.code()) - .as_deref() - ); - errors.push(anyhow::Error::new(error)); + errors.push(anyhow::anyhow!( + "error returned from database: (code: 5) database is locked" + )); } - sqlx::query("ROLLBACK") - .execute(&mut holder) - .await - .expect("release write transaction"); - - let mut reader = SqliteConnection::connect_with(&options) - .await - .expect("open snapshot reader"); - let mut writer = SqliteConnection::connect_with(&options) - .await - .expect("open snapshot writer"); - sqlx::query("BEGIN") - .execute(&mut reader) - .await - .expect("begin read transaction"); - let _: i64 = sqlx::query_scalar("SELECT value FROM contention WHERE id = 1") - .fetch_one(&mut reader) - .await - .expect("establish read snapshot"); - sqlx::query("BEGIN IMMEDIATE") - .execute(&mut writer) - .await - .expect("begin concurrent writer"); - sqlx::query("UPDATE contention SET value = value + 1 WHERE id = 1") - .execute(&mut writer) - .await - .expect("advance WAL"); - sqlx::query("COMMIT") - .execute(&mut writer) - .await - .expect("commit concurrent writer"); for _ in 0..busy_snapshot_count { - let error = sqlx::query("UPDATE contention SET value = value + 1 WHERE id = 1") - .execute(&mut reader) - .await - .expect_err("stale reader should receive SQLITE_BUSY_SNAPSHOT"); - assert_eq!( - Some("517"), - error - .as_database_error() - .and_then(|database_error| database_error.code()) - .as_deref() - ); - errors.push(anyhow::Error::new(error)); + errors.push(anyhow::anyhow!( + "error returned from database: (code: 517) database is locked" + )); } - sqlx::query("ROLLBACK") - .execute(&mut reader) - .await - .expect("release stale read transaction"); errors } @@ -3399,9 +3314,7 @@ cleanup: async fn initial_start_retries_real_busy_and_busy_snapshot_without_duplicate_admission() { let fixture = activation_test_fixture().await; let fault_plan = WorkflowStateFaultPlan::default(); - for error in - real_sqlite_contention_errors(/*busy_count*/ 1, /*busy_snapshot_count*/ 1).await - { + for error in sqlite_contention_errors(/*busy_count*/ 1, /*busy_snapshot_count*/ 1) { fault_plan.push(WorkflowStateOperation::AdmitBranches, error); } let service = WorkflowActivationService::new(Arc::clone(&fixture.state_db)) @@ -3530,9 +3443,7 @@ cleanup: .expect("workflow projection should exist"); let fault_plan = WorkflowStateFaultPlan::default(); - for error in - real_sqlite_contention_errors(/*busy_count*/ 1, /*busy_snapshot_count*/ 1).await - { + for error in sqlite_contention_errors(/*busy_count*/ 1, /*busy_snapshot_count*/ 1) { fault_plan.push(WorkflowStateOperation::AdmitBranches, error); } let service = WorkflowActivationService::new(Arc::clone(&fixture.state_db)) @@ -3640,9 +3551,7 @@ cleanup: async fn always_busy_initial_start_stops_with_one_structured_block() { let fixture = activation_test_fixture().await; let fault_plan = WorkflowStateFaultPlan::default(); - for error in - real_sqlite_contention_errors(/*busy_count*/ 8, /*busy_snapshot_count*/ 0).await - { + for error in sqlite_contention_errors(/*busy_count*/ 8, /*busy_snapshot_count*/ 0) { fault_plan.push(WorkflowStateOperation::AdmitBranches, error); } let service = WorkflowActivationService::new(Arc::clone(&fixture.state_db)) From 029c0378291d3a89ea697c2fb73182b9ed6be51a Mon Sep 17 00:00:00 2001 From: Andrei Hasna Date: Tue, 11 Aug 2026 16:53:23 +0300 Subject: [PATCH 3/6] fix(workflows): align admission timeout with busy retries Derive initial admission timeout from both the activation retry and structured persistence budgets, plus the existing observation margin. Agent: iapp-factory-coordinator --- codex-rs/ext/workflows/src/activation.rs | 26 +++++++++++++++++++++--- 1 file changed, 23 insertions(+), 3 deletions(-) diff --git a/codex-rs/ext/workflows/src/activation.rs b/codex-rs/ext/workflows/src/activation.rs index 7eee9dd0f..2870e134b 100644 --- a/codex-rs/ext/workflows/src/activation.rs +++ b/codex-rs/ext/workflows/src/activation.rs @@ -83,7 +83,7 @@ use crate::provider_credit::unsupported_workflow_provider_credit_authority; const WORKFLOW_SUPERVISOR_POLL_INTERVAL: Duration = Duration::from_millis(250); const WORKFLOW_HEARTBEAT_INTERVAL: Duration = Duration::from_millis(500); const WORKFLOW_INITIAL_ADMISSION_POLL_INTERVAL: Duration = Duration::from_millis(25); -const WORKFLOW_INITIAL_ADMISSION_TIMEOUT: Duration = Duration::from_secs(5); +const WORKFLOW_INITIAL_ADMISSION_TIMEOUT_MARGIN: Duration = Duration::from_secs(5); const MAX_WORKFLOW_START_FAILURE_DETAIL_CHARS: usize = 240; const ARTIFACT_CONTAINS_VERIFIER: &str = "artifact_contains"; const RUN_COMMANDS_VERIFIER: &str = "run_commands"; @@ -310,6 +310,13 @@ where } } +fn workflow_initial_admission_timeout(policy: BusyRetryPolicy) -> Duration { + policy + .total_delay_budget + .saturating_mul(2) + .saturating_add(WORKFLOW_INITIAL_ADMISSION_TIMEOUT_MARGIN) +} + #[derive(Debug, Clone)] pub struct WorkflowActivationConfig { pub auth_profile_ref: Option, @@ -664,7 +671,8 @@ impl WorkflowActivationService { tokio::sync::mpsc::UnboundedReceiver>, >, ) -> anyhow::Result { - tokio::time::timeout(WORKFLOW_INITIAL_ADMISSION_TIMEOUT, async { + let admission_timeout = workflow_initial_admission_timeout(self.state_retry_policy); + tokio::time::timeout(admission_timeout, async { loop { let snapshot = self .retry_state(WorkflowStateOperation::LoadSnapshot, || { @@ -691,7 +699,7 @@ impl WorkflowActivationService { .map_err(|_| { anyhow::anyhow!( "initial workflow branch admission remained unobservable for {} ms", - WORKFLOW_INITIAL_ADMISSION_TIMEOUT.as_millis() + admission_timeout.as_millis() ) })? } @@ -4190,6 +4198,18 @@ cleanup: ); } + #[test] + fn initial_admission_timeout_covers_busy_retry_and_persistence_budgets() { + let policy = BusyRetryPolicy::default(); + assert_eq!( + policy + .total_delay_budget + .saturating_mul(2) + .saturating_add(WORKFLOW_INITIAL_ADMISSION_TIMEOUT_MARGIN), + workflow_initial_admission_timeout(policy) + ); + } + #[tokio::test] async fn workflow_state_retry_survives_busy_and_busy_snapshot_contention() { let attempts = Arc::new(AtomicUsize::new(0)); From 8afe07fca0399cf9297d0335e5c8889463de5158 Mon Sep 17 00:00:00 2001 From: Andrei Hasna Date: Tue, 11 Aug 2026 16:58:06 +0300 Subject: [PATCH 4/6] fix(workflows): bound state retry elapsed time Cap activation state retries by elapsed time so SQLite busy_timeout cannot outrun initial-admission persistence, and cover a slow code-5 attempt. Agent: iapp-factory-coordinator --- codex-rs/ext/workflows/src/activation.rs | 79 ++++++++++++++++++++++-- 1 file changed, 74 insertions(+), 5 deletions(-) diff --git a/codex-rs/ext/workflows/src/activation.rs b/codex-rs/ext/workflows/src/activation.rs index 2870e134b..434d9acae 100644 --- a/codex-rs/ext/workflows/src/activation.rs +++ b/codex-rs/ext/workflows/src/activation.rs @@ -206,6 +206,25 @@ impl std::error::Error for WorkflowActivationDatabaseBusy { } } +#[derive(Debug)] +struct WorkflowStateElapsedTimeout { + operation: WorkflowStateOperation, + elapsed_timeout: Duration, +} + +impl std::fmt::Display for WorkflowStateElapsedTimeout { + fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + write!( + formatter, + "{} did not complete within the bounded SQLite busy window of {} ms", + self.operation.label(), + self.elapsed_timeout.as_millis() + ) + } +} + +impl std::error::Error for WorkflowStateElapsedTimeout {} + #[cfg(test)] #[derive(Clone, Default)] struct WorkflowStateFaultPlan { @@ -299,20 +318,45 @@ where F: FnMut() -> Fut, Fut: Future>, { - match retry_on_busy_with_policy(policy, operation.label(), f).await { + let elapsed_timeout = workflow_state_retry_elapsed_timeout(policy); + let result = tokio::time::timeout( + elapsed_timeout, + retry_on_busy_with_policy(policy, operation.label(), f), + ) + .await + .unwrap_or_else(|_| { + Err(anyhow::Error::new(WorkflowStateElapsedTimeout { + operation, + elapsed_timeout, + })) + }); + match result { Err(source) if is_transient_busy_error(&source) => { Err(anyhow::Error::new(WorkflowActivationDatabaseBusy { operation, source, })) } + Err(source) + if source + .downcast_ref::() + .is_some() => + { + Err(anyhow::Error::new(WorkflowActivationDatabaseBusy { + operation, + source, + })) + } result => result, } } +fn workflow_state_retry_elapsed_timeout(policy: BusyRetryPolicy) -> Duration { + policy.total_delay_budget.saturating_add(policy.max_delay) +} + fn workflow_initial_admission_timeout(policy: BusyRetryPolicy) -> Duration { - policy - .total_delay_budget + workflow_state_retry_elapsed_timeout(policy) .saturating_mul(2) .saturating_add(WORKFLOW_INITIAL_ADMISSION_TIMEOUT_MARGIN) } @@ -4202,14 +4246,39 @@ cleanup: fn initial_admission_timeout_covers_busy_retry_and_persistence_budgets() { let policy = BusyRetryPolicy::default(); assert_eq!( - policy - .total_delay_budget + workflow_state_retry_elapsed_timeout(policy) .saturating_mul(2) .saturating_add(WORKFLOW_INITIAL_ADMISSION_TIMEOUT_MARGIN), workflow_initial_admission_timeout(policy) ); } + #[tokio::test] + async fn workflow_state_retry_bounds_a_slow_sqlite_busy_attempt() { + let policy = tiny_state_retry_policy(); + let error = tokio::time::timeout( + Duration::from_millis(100), + retry_workflow_state_with_policy( + policy, + WorkflowStateOperation::AdmitBranches, + || async { + tokio::time::sleep(Duration::from_millis(50)).await; + Err::<(), _>(anyhow::anyhow!( + "error returned from database: (code: 5) database is locked" + )) + }, + ), + ) + .await + .expect("the elapsed retry bound should terminate the slow busy attempt") + .expect_err("the slow busy attempt should remain an error"); + assert!( + error + .downcast_ref::() + .is_some() + ); + } + #[tokio::test] async fn workflow_state_retry_survives_busy_and_busy_snapshot_contention() { let attempts = Arc::new(AtomicUsize::new(0)); From 87800d90a2d16dcd3de2aad5f5c1f02e68622fe3 Mon Sep 17 00:00:00 2001 From: Andrei Hasna Date: Tue, 11 Aug 2026 20:13:17 +0300 Subject: [PATCH 5/6] fix(workflows): own heartbeat retry state Move an owned WorkflowActivationService clone into the verifier heartbeat task so the spawned future remains static without changing retry behavior. Agent: iapp-factory-coordinator --- codex-rs/ext/workflows/src/activation.rs | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/codex-rs/ext/workflows/src/activation.rs b/codex-rs/ext/workflows/src/activation.rs index 434d9acae..b64aaf979 100644 --- a/codex-rs/ext/workflows/src/activation.rs +++ b/codex-rs/ext/workflows/src/activation.rs @@ -1833,6 +1833,7 @@ impl WorkflowActivationService { cancellation: CancellationToken, fence_lost: Arc, ) -> tokio::task::JoinHandle<()> { + let service = self.clone(); let state_db = Arc::clone(&self.state_db); let owner_instance_id = Arc::clone(&self.owner_instance_id); tokio::spawn(async move { @@ -1847,7 +1848,7 @@ impl WorkflowActivationService { generation, lease_duration_ms: None, }; - match self + match service .retry_state(WorkflowStateOperation::HeartbeatRun, || { state_db.heartbeat_workflow_run(heartbeat_params.clone()) }) From b568b3d6cf4b1a7055e209c24ec22909cb95a221 Mon Sep 17 00:00:00 2001 From: Andrei Hasna Date: Wed, 12 Aug 2026 13:51:52 +0300 Subject: [PATCH 6/6] fix(workflows): reconcile retry policy after source rebase Keep the fixture-only retry override scoped to the injected fault while routing current-main invalid source snapshot handling through the production activation retry service. Agent: iapp-factory-coordinator --- codex-rs/ext/workflows/src/activation.rs | 71 ++++++++++++++++-------- 1 file changed, 49 insertions(+), 22 deletions(-) diff --git a/codex-rs/ext/workflows/src/activation.rs b/codex-rs/ext/workflows/src/activation.rs index b64aaf979..6608d3145 100644 --- a/codex-rs/ext/workflows/src/activation.rs +++ b/codex-rs/ext/workflows/src/activation.rs @@ -254,11 +254,19 @@ impl WorkflowStateFaultPlan { None } } + + fn has_pending(&self, operation: WorkflowStateOperation) -> bool { + self.faults + .lock() + .expect("workflow state fault plan should lock") + .front() + .is_some_and(|(planned, _)| *planned == operation) + } } impl WorkflowStateOperation { #[cfg(test)] - const ALL: [Self; 14] = [ + const ALL: [Self; 15] = [ Self::CreateRun, Self::ProjectGoalPlan, Self::ListThreadRuns, @@ -454,6 +462,8 @@ pub struct WorkflowActivationService { state_retry_policy: BusyRetryPolicy, #[cfg(test)] state_fault_plan: WorkflowStateFaultPlan, + #[cfg(test)] + fault_retry_policy: Option, } impl WorkflowActivationService { @@ -476,6 +486,8 @@ impl WorkflowActivationService { state_retry_policy: BusyRetryPolicy::default(), #[cfg(test)] state_fault_plan: WorkflowStateFaultPlan::default(), + #[cfg(test)] + fault_retry_policy: None, } } @@ -486,8 +498,8 @@ impl WorkflowActivationService { } #[cfg(test)] - fn with_state_retry_policy(mut self, state_retry_policy: BusyRetryPolicy) -> Self { - self.state_retry_policy = state_retry_policy; + fn with_fault_retry_policy(mut self, fault_retry_policy: BusyRetryPolicy) -> Self { + self.fault_retry_policy = Some(fault_retry_policy); self } @@ -500,16 +512,28 @@ impl WorkflowActivationService { F: FnMut() -> Fut, Fut: Future>, { - retry_workflow_state_with_policy(self.state_retry_policy, operation, || { + #[cfg(test)] + let retry_policy = self + .fault_retry_policy + .filter(|_| self.state_fault_plan.has_pending(operation)) + .unwrap_or(self.state_retry_policy); + #[cfg(not(test))] + let retry_policy = self.state_retry_policy; + retry_workflow_state_with_policy(retry_policy, operation, || { #[cfg(test)] - let planned_fault = self.state_fault_plan.take(operation); - let future = f(); - async move { - #[cfg(test)] - if let Some(error) = planned_fault { - return Err(error); + { + let planned_fault = self.state_fault_plan.take(operation); + let future = f(); + async move { + if let Some(error) = planned_fault { + return Err(error); + } + future.await } - future.await + } + #[cfg(not(test))] + { + f() } }) .await @@ -1148,15 +1172,13 @@ impl WorkflowActivationService { match parse_run_source_yaml_snapshot(&snapshot) { Ok(spec) => Some(spec), Err(validation_error) => { - let Some(outcome) = retry_workflow_state( - WorkflowStateOperation::FailInvalidSourceSnapshot, - || { + let Some(outcome) = self + .retry_state(WorkflowStateOperation::FailInvalidSourceSnapshot, || { self.state_db .workflows() .fail_workflow_run_if_source_snapshot_invalid(run_id) - }, - ) - .await? + }) + .await? else { return Ok(()); }; @@ -2651,12 +2673,14 @@ mod tests { ) -> anyhow::Result<()> { Ok(()) } + } struct ActivationTestFixture { _temp_dir: tempfile::TempDir, state_db: Arc, source_thread_id: ThreadId, workflow_record_id: String, + source_yaml_sha256: String, } async fn activation_test_fixture() -> ActivationTestFixture { @@ -2692,6 +2716,7 @@ mod tests { state_db, source_thread_id, workflow_record_id: spec.workflow_record_id, + source_yaml_sha256: spec.source_yaml_sha256, } } @@ -3371,8 +3396,7 @@ cleanup: fault_plan.push(WorkflowStateOperation::AdmitBranches, error); } let service = WorkflowActivationService::new(Arc::clone(&fixture.state_db)) - .with_state_fault_plan(fault_plan) - .with_state_retry_policy(tiny_state_retry_policy()); + .with_state_fault_plan(fault_plan); let activation_config = WorkflowActivationConfig { route_runtime: supported_route_runtime(), max_active_background_agent_runs: Some(10), @@ -3383,6 +3407,7 @@ cleanup: .start_workflow_run(WorkflowStartRequest { workflow_record_id: fixture.workflow_record_id.clone(), source_thread_id: fixture.source_thread_id, + expected_source_yaml_sha256: fixture.source_yaml_sha256.clone(), idempotency_key: Some(idempotency_key.to_string()), activation_config: activation_config.clone(), }) @@ -3422,6 +3447,7 @@ cleanup: .start_workflow_run(WorkflowStartRequest { workflow_record_id: fixture.workflow_record_id, source_thread_id: fixture.source_thread_id, + expected_source_yaml_sha256: fixture.source_yaml_sha256.clone(), idempotency_key: Some(idempotency_key.to_string()), activation_config: activation_config.clone(), }) @@ -3480,6 +3506,7 @@ cleanup: .create_workflow_run(WorkflowRunCreateParams { workflow_record_id: fixture.workflow_record_id, source_thread_id: Some(fixture.source_thread_id), + expected_source_yaml_sha256: fixture.source_yaml_sha256.clone(), idempotency_key: Some(idempotency_key.to_string()), }) .await @@ -3500,8 +3527,7 @@ cleanup: fault_plan.push(WorkflowStateOperation::AdmitBranches, error); } let service = WorkflowActivationService::new(Arc::clone(&fixture.state_db)) - .with_state_fault_plan(fault_plan) - .with_state_retry_policy(tiny_state_retry_policy()); + .with_state_fault_plan(fault_plan); let activation_config = WorkflowActivationConfig { route_runtime: supported_route_runtime(), max_active_background_agent_runs: Some(10), @@ -3609,10 +3635,11 @@ cleanup: } let service = WorkflowActivationService::new(Arc::clone(&fixture.state_db)) .with_state_fault_plan(fault_plan) - .with_state_retry_policy(tiny_state_retry_policy()); + .with_fault_retry_policy(tiny_state_retry_policy()); let request = WorkflowStartRequest { workflow_record_id: fixture.workflow_record_id, source_thread_id: fixture.source_thread_id, + expected_source_yaml_sha256: fixture.source_yaml_sha256, idempotency_key: Some("always-busy-initial-start".to_string()), activation_config: WorkflowActivationConfig { route_runtime: supported_route_runtime(),