From d4b4b15471516de1de083f216f502e54fba5c960 Mon Sep 17 00:00:00 2001 From: Johnny Date: Wed, 30 Sep 2026 00:18:32 +0000 Subject: [PATCH] Give a live worker a grace period to report its own timeout should_timeout_operation ended an Executing action once Action.timeout had passed since it was assigned. The worker enforces the same timeout, but its clock starts when the command starts, after input fetch, so the scheduler always fired first. The action was requeued onto the worker still running it, which refused with AlreadyExists; the retries ran out, and when the worker's own DEADLINE_EXCEEDED result arrived it was treated as a stray and the worker was disconnected, taking every other action on it down too. The scheduler now enforces Action.timeout in every liveness case, at Action.timeout plus no_event_action_timeout, measured from the later of the assignment and the worker's last update on the action. A live worker reports first, so no liveness special-casing is needed, and the scheduler still ends an action whose worker never reports: a hung input fetch, an unkillable child, a broken timeout_handled_externally wrapper, or an orphan no instance reaps. The client's DEADLINE_EXCEEDED message now names the deadline that fired, instead of always citing no_event_action_timeout. This is the interim: the scheduler cannot observe when the command starts. #2827 tracks the worker reporting command start as an operation transition, so the deadline can be anchored to that event. --- .../src/simple_scheduler_state_manager.rs | 127 ++++++---- .../simple_scheduler_state_manager_test.rs | 229 +++++++++++++++++- .../tests/simple_scheduler_test.rs | 23 +- 3 files changed, 328 insertions(+), 51 deletions(-) diff --git a/nativelink-scheduler/src/simple_scheduler_state_manager.rs b/nativelink-scheduler/src/simple_scheduler_state_manager.rs index 7d794e3dcc3..e8f75e9dfee 100644 --- a/nativelink-scheduler/src/simple_scheduler_state_manager.rs +++ b/nativelink-scheduler/src/simple_scheduler_state_manager.rs @@ -46,6 +46,34 @@ use super::awaited_action_db::{ }; use crate::worker_registry::{ORPHANED_ACTION_TIMEOUT, SharedWorkerRegistry, WorkerLiveness}; +/// Why the scheduler times out an executing action. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum TimeoutCause { + /// The action's `Action.timeout` passed, plus the grace a live worker + /// has to report its own result. + ActionTimeout { timeout: Duration, grace: Duration }, + /// The worker sent no update on the action for this long. + NoWorkerUpdate(Duration), +} + +impl core::fmt::Display for TimeoutCause { + fn fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result { + match self { + Self::ActionTimeout { timeout, grace } => write!( + f, + "exceeded its Action.timeout of {} seconds, and its worker had not reported within {} seconds after that", + timeout.as_secs_f32(), + grace.as_secs_f32(), + ), + Self::NoWorkerUpdate(ceiling) => write!( + f, + "timed out after {} seconds with no update from its worker", + ceiling.as_secs_f32(), + ), + } + } +} + /// Maximum number of times an update to the database /// can fail before giving up. const MAX_UPDATE_RETRIES: usize = 5; @@ -312,25 +340,23 @@ where .err_tip(|| format!("Failed to upgrade weak reference to SimpleSchedulerStateManager in MatchingEngineActionStateResult::changed at attempt: {timeout_attempts}"))?; // Check if worker is alive via registry before timing out. - let should_timeout = simple_scheduler_state_manager - .should_timeout_operation(&awaited_action) - .await; - - if !should_timeout { + let Some(cause) = simple_scheduler_state_manager + .timeout_cause(&awaited_action) + .await + else { // Worker is alive, continue waiting for updates trace!( operation_id = %awaited_action.operation_id(), "Operation timeout check passed, worker is alive" ); continue; - } + }; warn!( ?awaited_action, - "OperationId {} / {} timed out after {} seconds issuing a retry", + "OperationId {} / {} {cause}, issuing a retry", awaited_action.operation_id(), awaited_action.state().client_operation_id, - self.no_event_action_timeout.as_secs_f32(), ); simple_scheduler_state_manager @@ -633,23 +659,41 @@ where } pub async fn should_timeout_operation(&self, awaited_action: &AwaitedAction) -> bool { + self.timeout_cause(awaited_action).await.is_some() + } + + /// Why the scheduler should time out an executing action, or `None` if + /// it should not. + pub async fn timeout_cause(&self, awaited_action: &AwaitedAction) -> Option { if !matches!(awaited_action.state().stage, ActionStage::Executing) { - return false; + return None; } let now = (self.now_fn)().now(); // Honor the per-action `Action.timeout` from the RBE protocol as a - // backend wall-clock deadline. Without this, the only enforcement is - // the Bazel client's --test_timeout, which surfaces as TIMEOUT/NO - // STATUS instead of a backend signal pointing at the worker. + // backend wall-clock deadline, in every liveness case below. The + // worker enforces the same timeout, but its clock starts when the + // command starts, after the inputs are fetched, which the scheduler + // cannot observe. Enforcing it at `Action.timeout` from assignment + // would pre-empt every live worker's own DEADLINE_EXCEEDED result, so + // allow `no_event_action_timeout` more, measured from the later of the + // assignment and the worker's last update: a live worker then reports + // first, and a wedged one (hung input fetch, unkillable child, broken + // external timeout, orphan) is still ended here. let action_timeout = awaited_action.action_info().timeout; if action_timeout > Duration::ZERO { - let executing_started_at = awaited_action.state().last_transition_timestamp; - if let Ok(elapsed) = now.duration_since(executing_started_at) - && elapsed > action_timeout + let anchor = awaited_action + .state() + .last_transition_timestamp + .max(awaited_action.last_worker_updated_timestamp()); + if let Ok(elapsed) = now.duration_since(anchor) + && elapsed > action_timeout.saturating_add(self.no_event_action_timeout) { - return true; + return Some(TimeoutCause::ActionTimeout { + timeout: action_timeout, + grace: self.no_event_action_timeout, + }); } } @@ -664,17 +708,23 @@ where _ => WorkerLiveness::Stale, }; + let last_update = awaited_action.last_worker_updated_timestamp(); + let silent_for = |ceiling: Duration| { + now.duration_since(last_update) + .is_ok_and(|elapsed| elapsed > ceiling) + .then_some(TimeoutCause::NoWorkerUpdate(ceiling)) + }; + match liveness { - // Ours and heartbeating: only the stuck-but-alive ceiling applies, - // and disabling that means no ceiling on a live worker. + // Ours and heartbeating: besides Action.timeout above, only the + // stuck-but-alive ceiling applies, and disabling that leaves an + // action with no Action.timeout no ceiling on a live worker. WorkerLiveness::Alive => { if self.max_executing_timeout > Duration::ZERO { - let last_update = awaited_action.last_worker_updated_timestamp(); - if let Ok(elapsed) = now.duration_since(last_update) { - return elapsed > self.max_executing_timeout; - } + silent_for(self.max_executing_timeout) + } else { + None } - false } // Usually a peer's healthy worker, so worker_timeout_s must not @@ -682,27 +732,15 @@ where // max_action_executing_timeout_s defaults to disabled, so fall // back to a ceiling rather than never timing out. WorkerLiveness::Unknown => { - let ceiling = if self.max_executing_timeout > Duration::ZERO { - self.max_executing_timeout + if self.max_executing_timeout > Duration::ZERO { + silent_for(self.max_executing_timeout) } else { - ORPHANED_ACTION_TIMEOUT - }; - let last_update = awaited_action.last_worker_updated_timestamp(); - match now.duration_since(last_update) { - Ok(elapsed) => elapsed > ceiling, - Err(_) => false, + silent_for(ORPHANED_ACTION_TIMEOUT) } } // Registered here and gone quiet: ours, and it looks dead. - WorkerLiveness::Stale => { - let worker_should_update_before = awaited_action - .last_worker_updated_timestamp() - .checked_add(self.no_event_action_timeout) - .unwrap_or(now); - - worker_should_update_before < now - } + WorkerLiveness::Stale => silent_for(self.no_event_action_timeout), } } @@ -882,28 +920,25 @@ where // Re-check under the lock against freshly loaded state, and delegate // rather than re-deriving the rule: the two copies had drifted. - if !self.should_timeout_operation(&awaited_action).await { + let Some(cause) = self.timeout_cause(&awaited_action).await else { trace!( %operation_id, worker_id = ?awaited_action.worker_id(), "Operation no longer needs timing out, skipping" ); return Ok(()); - } + }; warn!( %operation_id, worker_id = ?awaited_action.worker_id(), + %cause, "Timing out operation" ); self.assign_operation( operation_id, - Err(make_err!( - Code::DeadlineExceeded, - "Operation timed out after {} seconds", - self.no_event_action_timeout.as_secs_f32(), - )), + Err(make_err!(Code::DeadlineExceeded, "Operation {cause}")), ) .await } diff --git a/nativelink-scheduler/tests/simple_scheduler_state_manager_test.rs b/nativelink-scheduler/tests/simple_scheduler_state_manager_test.rs index eb96633b6cd..371487a668b 100644 --- a/nativelink-scheduler/tests/simple_scheduler_state_manager_test.rs +++ b/nativelink-scheduler/tests/simple_scheduler_state_manager_test.rs @@ -19,7 +19,9 @@ use nativelink_scheduler::awaited_action_db::{ }; use nativelink_scheduler::default_scheduler_factory::memory_awaited_action_db_factory; use nativelink_scheduler::simple_scheduler::SimpleScheduler; -use nativelink_scheduler::simple_scheduler_state_manager::SimpleSchedulerStateManager; +use nativelink_scheduler::simple_scheduler_state_manager::{ + SimpleSchedulerStateManager, TimeoutCause, +}; use nativelink_scheduler::worker::Worker; use nativelink_scheduler::worker_registry::WorkerRegistry; use nativelink_scheduler::worker_scheduler::WorkerScheduler; @@ -262,6 +264,231 @@ async fn does_not_time_out_a_live_worker_mid_action() -> Result<(), Error> { Ok(()) } +/// `Action.timeout` for the deadline tests below. The grace is +/// `WORKER_TIMEOUT`, the state managers' `no_event_action_timeout`. +const ACTION_TIMEOUT: Duration = Duration::from_mins(1); + +/// An action with `ACTION_TIMEOUT`, assigned to `worker_id` and executing +/// since `started`. +fn action_with_timeout(worker_id: &WorkerId, started: SystemTime) -> AwaitedAction { + let mut info = action_info(started); + info.timeout = ACTION_TIMEOUT; + let action_digest = info.digest(); + let operation_id = OperationId::default(); + let mut action = AwaitedAction::new(operation_id.clone(), Arc::new(info), started); + action.worker_set_state( + Arc::new(ActionState { + stage: ActionStage::Executing, + client_operation_id: operation_id, + action_digest, + last_transition_timestamp: started, + }), + started, + ); + action.set_worker_id(Some(worker_id.clone()), started); + action +} + +const fn action_timeout_cause() -> TimeoutCause { + TimeoutCause::ActionTimeout { + timeout: ACTION_TIMEOUT, + grace: WORKER_TIMEOUT, + } +} + +/// Just short of `Action.timeout` plus the grace. +const WITHIN_GRACE: Duration = + Duration::from_secs(ACTION_TIMEOUT.as_secs() + WORKER_TIMEOUT.as_secs() - 1); + +/// Just past `Action.timeout` plus the grace. +const PAST_GRACE: Duration = + Duration::from_secs(ACTION_TIMEOUT.as_secs() + WORKER_TIMEOUT.as_secs() + 1); + +/// A heartbeating worker enforces `Action.timeout` itself, from when the +/// command starts, and reports a completed `DEADLINE_EXCEEDED` result. The +/// scheduler's clock starts at assignment, earlier, so ending the action at +/// `Action.timeout` would requeue every action that runs to its timeout onto +/// the worker still running it, and the worker's own result would then evict +/// it. The grace lets the live worker report first. +#[nativelink_test] +async fn gives_a_live_worker_the_grace_to_report_its_own_timeout() -> Result<(), Error> { + MockClock::set_time(Duration::from_secs(NOW_TIME)); + let worker_id = WorkerId::from(String::from("enforces-its-own-timeout")); + let action = action_with_timeout(&worker_id, make_system_time(0)); + + let registry = Arc::new(WorkerRegistry::new()); + registry + .update_worker_heartbeat(&worker_id, make_system_time(WITHIN_GRACE.as_secs())) + .await; + let state_mgr = state_manager_no_executing_ceiling(registry); + + MockClock::advance(WITHIN_GRACE); + assert!( + !state_mgr.should_timeout_operation(&action).await, + "a live worker's action past Action.timeout but within the grace is the worker's to end" + ); + Ok(()) +} + +/// A heartbeating worker that never reports (a hung input fetch, an +/// unkillable child, a broken `timeout_handled_externally` wrapper) is timed +/// out by `Action.timeout` once the grace has passed, even with +/// `max_action_executing_timeout_s` disabled. +#[nativelink_test] +async fn times_out_a_live_but_wedged_worker_past_the_grace() -> Result<(), Error> { + MockClock::set_time(Duration::from_secs(NOW_TIME)); + let worker_id = WorkerId::from(String::from("heartbeats-but-never-reports")); + let action = action_with_timeout(&worker_id, make_system_time(0)); + + let registry = Arc::new(WorkerRegistry::new()); + registry + .update_worker_heartbeat(&worker_id, make_system_time(PAST_GRACE.as_secs())) + .await; + let state_mgr = state_manager_no_executing_ceiling(registry); + + MockClock::advance(PAST_GRACE); + let cause = state_mgr.timeout_cause(&action).await; + assert_eq!( + cause, + Some(action_timeout_cause()), + "a live worker's action past Action.timeout plus the grace must time out" + ); + assert!( + cause + .unwrap() + .to_string() + .contains("Action.timeout of 60 seconds"), + "the client's error must name the deadline that fired" + ); + Ok(()) +} + +/// The grace runs from the worker's last update when that is later than the +/// assignment, so it only has to cover the time since the worker last showed +/// signs of life. +#[nativelink_test] +async fn measures_the_grace_from_the_workers_last_update() -> Result<(), Error> { + MockClock::set_time(Duration::from_secs(NOW_TIME)); + let worker_id = WorkerId::from(String::from("updated-after-assignment")); + let mut action = action_with_timeout(&worker_id, make_system_time(0)); + let updated = Duration::from_secs(30); + let state = action.state().clone(); + action.worker_set_state(state, make_system_time(updated.as_secs())); + + let registry = Arc::new(WorkerRegistry::new()); + registry + .update_worker_heartbeat( + &worker_id, + make_system_time((PAST_GRACE + updated).as_secs()), + ) + .await; + let state_mgr = state_manager_no_executing_ceiling(registry); + + MockClock::advance(PAST_GRACE); + assert!( + !state_mgr.should_timeout_operation(&action).await, + "the grace must run from the last worker update, not the assignment" + ); + + MockClock::advance(updated); + assert_eq!( + state_mgr.timeout_cause(&action).await, + Some(action_timeout_cause()), + "past Action.timeout plus the grace from the last update, the action must time out" + ); + Ok(()) +} + +/// A peer instance's live worker enforces `Action.timeout` itself too, so it +/// gets the same grace here. +#[nativelink_test] +async fn gives_a_peer_instances_worker_the_grace() -> Result<(), Error> { + MockClock::set_time(Duration::from_secs(NOW_TIME)); + // Deliberately not registered here: it belongs to another instance. + let action = action_with_timeout( + &WorkerId::from(String::from("owned-by-another-scheduler")), + make_system_time(0), + ); + let state_mgr = state_manager_no_executing_ceiling(Arc::new(WorkerRegistry::new())); + + MockClock::advance(WITHIN_GRACE); + assert!( + !state_mgr.should_timeout_operation(&action).await, + "a peer's worker's action within the grace is the worker's to end" + ); + Ok(()) +} + +/// An orphan nobody reaps, or a peer's worker that never reports, is timed +/// out by `Action.timeout` past the grace rather than waiting for the orphan +/// ceiling. +#[nativelink_test] +async fn times_out_an_orphan_on_action_timeout_past_the_grace() -> Result<(), Error> { + MockClock::set_time(Duration::from_secs(NOW_TIME)); + let action = action_with_timeout( + &WorkerId::from(String::from("owner-was-scaled-down")), + make_system_time(0), + ); + let state_mgr = state_manager_no_executing_ceiling(Arc::new(WorkerRegistry::new())); + + MockClock::advance(PAST_GRACE); + assert_eq!( + state_mgr.timeout_cause(&action).await, + Some(action_timeout_cause()), + "an orphan past Action.timeout plus the grace must time out, not wait an hour" + ); + Ok(()) +} + +/// Registry heartbeats refresh only on keepalives, so a worker whose +/// keepalives lag reads Stale while its action is still running and still +/// updating. `Action.timeout` must not pre-empt it there either. +#[nativelink_test] +async fn gives_a_stale_worker_with_a_recent_update_the_grace() -> Result<(), Error> { + MockClock::set_time(Duration::from_secs(NOW_TIME)); + let worker_id = WorkerId::from(String::from("keepalives-lag")); + let assigned = WORKER_TIMEOUT + Duration::from_secs(30); + let action = action_with_timeout(&worker_id, make_system_time(assigned.as_secs())); + + let registry = Arc::new(WorkerRegistry::new()); + registry + .register_worker(&worker_id, make_system_time(0)) + .await; + let state_mgr = state_manager_no_executing_ceiling(registry); + + // Past Action.timeout, and less than worker_timeout_s since the action's + // last update, so only Action.timeout could fire. + MockClock::advance(assigned + ACTION_TIMEOUT + Duration::from_secs(1)); + assert!( + !state_mgr.should_timeout_operation(&action).await, + "a stale worker's recently updated action must get the grace too" + ); + Ok(()) +} + +/// A Stale worker's action past `Action.timeout` plus the grace is timed out +/// for that reason, which is what the client is told. +#[nativelink_test] +async fn times_out_a_stale_worker_on_action_timeout_past_the_grace() -> Result<(), Error> { + MockClock::set_time(Duration::from_secs(NOW_TIME)); + let worker_id = WorkerId::from(String::from("connected-to-us")); + let action = action_with_timeout(&worker_id, make_system_time(0)); + + let registry = Arc::new(WorkerRegistry::new()); + registry + .register_worker(&worker_id, make_system_time(0)) + .await; + let state_mgr = state_manager_no_executing_ceiling(registry); + + MockClock::advance(PAST_GRACE); + assert_eq!( + state_mgr.timeout_cause(&action).await, + Some(action_timeout_cause()), + "a stale worker's action past Action.timeout plus the grace must name Action.timeout" + ); + Ok(()) +} + /// A heartbeating worker stuck on one action past the executing ceiling must /// still be timed out, which is what `max_action_executing_timeout_s` is for. #[nativelink_test] diff --git a/nativelink-scheduler/tests/simple_scheduler_test.rs b/nativelink-scheduler/tests/simple_scheduler_test.rs index 141d6f6088b..9a353ca90a2 100644 --- a/nativelink-scheduler/tests/simple_scheduler_test.rs +++ b/nativelink-scheduler/tests/simple_scheduler_test.rs @@ -2570,7 +2570,9 @@ async fn worker_disconnect_loop_caps_at_max_job_retries_test() -> Result<(), Err #[nativelink_test] async fn action_timeout_is_enforced_backend_side_test() -> Result<(), Error> { use nativelink_scheduler::awaited_action_db::AwaitedAction; - use nativelink_scheduler::simple_scheduler_state_manager::SimpleSchedulerStateManager; + use nativelink_scheduler::simple_scheduler_state_manager::{ + SimpleSchedulerStateManager, TimeoutCause, + }; // Anchor MockClock so MockInstantWrapped::now() == make_system_time(0). MockClock::set_time(Duration::from_secs(NOW_TIME)); @@ -2613,12 +2615,25 @@ async fn action_timeout_is_enforced_backend_side_test() -> Result<(), Error> { "Should not time out before Action.timeout elapses", ); - // Advance past the 2s per-action deadline. + // Past the 2s per-action deadline, but within the grace a live worker + // has to report its own DEADLINE_EXCEEDED result. MockClock::advance(Duration::from_secs(5)); assert!( - state_mgr.should_timeout_operation(&awaited_action).await, - "Scheduler must mark Executing action timed out once Action.timeout has elapsed", + !state_mgr.should_timeout_operation(&awaited_action).await, + "Should not time out within the grace after Action.timeout", + ); + + // Past Action.timeout plus the grace (no_event_action_timeout). + MockClock::advance(Duration::from_mins(1)); + + assert_eq!( + state_mgr.timeout_cause(&awaited_action).await, + Some(TimeoutCause::ActionTimeout { + timeout: Duration::from_secs(2), + grace: Duration::from_mins(1), + }), + "Scheduler must time out an Executing action once Action.timeout plus the grace has elapsed", ); Ok(())