diff --git a/nativelink-scheduler/src/simple_scheduler_state_manager.rs b/nativelink-scheduler/src/simple_scheduler_state_manager.rs index 664aa1eba4c..79489efc9fe 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 @@ -680,23 +706,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, + }); } } @@ -711,17 +755,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 @@ -729,27 +779,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), } } @@ -929,28 +967,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 351474edae6..2df93f1625e 100644 --- a/nativelink-scheduler/tests/simple_scheduler_test.rs +++ b/nativelink-scheduler/tests/simple_scheduler_test.rs @@ -2571,7 +2571,9 @@ async fn worker_disconnect_loop_caps_at_max_worker_loss_retries_test() -> Result #[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)); @@ -2614,12 +2616,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(())