Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
127 changes: 81 additions & 46 deletions nativelink-scheduler/src/simple_scheduler_state_manager.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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<TimeoutCause> {
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,
});
}
}

Expand All @@ -711,45 +755,39 @@ 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
// apply. It can also be an orphan no instance will ever reap, and
// 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),
}
}

Expand Down Expand Up @@ -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
}
Expand Down
Loading
Loading