use crate::cancel::progress_certificate::{DrainPhase, ProgressCertificate};
use crate::obligation::lyapunov::{
LyapunovGovernor, PotentialWeights, SchedulingSuggestion, StateSnapshot,
};
use crate::observability::spectral_health::{SpectralHealthMonitor, SpectralThresholds};
use crate::runtime::config::SchedulerPlacementMode;
use crate::runtime::io_driver::IoDriverHandle;
use crate::runtime::scheduler::global_injector::{GlobalInjector, PriorityTask};
use crate::runtime::scheduler::local_queue::{self, LocalQueue};
use crate::runtime::scheduler::priority::Scheduler as PriorityScheduler;
use crate::runtime::scheduler::swarm_evidence::{
SCHEDULER_EVIDENCE_SCHEMA_VERSION, SchedulerEvidenceArtifact, SchedulerEvidenceMetrics,
SchedulerKnobProfile, SchedulerTopologyDescriptor, SchedulerWorkloadClass,
};
use crate::runtime::scheduler::worker::Parker;
use crate::runtime::stored_task::AnyStoredTask;
use crate::runtime::{RuntimeState, TaskTable};
use crate::sync::ContendedMutex;
use crate::time::TimerDriverHandle;
use crate::tracing_compat::{error, trace};
use crate::types::{CxInner, TaskId, Time};
use crate::util::{CachePadded, DetHashMap, DetHasher, DetRng};
use parking_lot::Mutex;
use parking_lot::RwLock;
use smallvec::SmallVec;
use std::cell::RefCell;
use std::collections::{BTreeMap, BTreeSet, HashSet, VecDeque};
use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
use std::sync::{Arc, Weak};
use std::task::{Context, Poll, Waker};
use std::time::Duration;
pub type WorkerId = usize;
const DEFAULT_CANCEL_STREAK_LIMIT: usize = 16;
const DEFAULT_BROWSER_READY_HANDOFF_LIMIT: usize = 0;
const DEFAULT_STEAL_BATCH_SIZE: usize = 4;
const GLOBAL_READY_BATCH_DRAIN_MIN_DEPTH: usize = 8;
const DEFAULT_ENABLE_PARKING: bool = true;
const LOCAL_SCHEDULER_BURST_BUDGET: usize = 2048;
const LOCAL_SCHEDULER_MIN_CAPACITY: usize = 128;
const LOCAL_SCHEDULER_MAX_CAPACITY: usize = 1024;
const ADAPTIVE_STREAK_ARMS: [usize; 5] = [4, 8, 16, 32, 64];
const ADAPTIVE_UCB_DISCOUNT: f64 = 0.95;
const ADAPTIVE_UCB_CONFIDENCE: f64 = 2.0;
const ADAPTIVE_EPROCESS_LAMBDA: f64 = 0.5;
const SPIN_LIMIT: u32 = 8;
const YIELD_LIMIT: u32 = 2;
const EMPTY_BACKOFF_PARK_THRESHOLD: u32 = SPIN_LIMIT + YIELD_LIMIT;
const STALE_DUE_DEADLINE_PARK_NANOS: u64 = 1;
const SHORT_WAIT_LE_5MS_NANOS: u64 = 5_000_000;
const IDLE_IO_POLL_MAX_TIMEOUT: Duration = Duration::from_millis(250);
#[cfg(any(test, feature = "test-internals"))]
const DEFAULT_SCHEDULER_EVIDENCE_MAX_INFLIGHT_MULTIPLIER: usize = 4;
#[derive(Debug)]
pub(crate) struct LocalReadyQueueInner {
ready: VecDeque<TaskId>,
present: HashSet<TaskId>,
tombstones: HashSet<TaskId>,
}
impl LocalReadyQueueInner {
pub(crate) fn new(ready: VecDeque<TaskId>) -> Self {
let present: HashSet<TaskId> = ready.iter().copied().collect();
Self {
ready,
present,
tombstones: HashSet::new(),
}
}
fn push_back(&mut self, task: TaskId) {
if self.present.insert(task) {
self.ready.push_back(task);
} else {
self.tombstones.remove(&task);
}
}
pub(crate) fn pop_front(&mut self) -> Option<TaskId> {
while let Some(task) = self.ready.pop_front() {
self.present.remove(&task);
if self.tombstones.remove(&task) {
continue;
}
return Some(task);
}
None
}
fn tombstone(&mut self, task: TaskId) {
if self.present.contains(&task) {
self.tombstones.insert(task);
}
}
#[cfg(test)]
fn contains(&self, task: &TaskId) -> bool {
self.present.contains(task) && !self.tombstones.contains(task)
}
fn iter(&self) -> impl Iterator<Item = &TaskId> {
self.ready
.iter()
.filter(|task| !self.tombstones.contains(task))
}
#[cfg(test)]
fn drain(&mut self, _range: std::ops::RangeFull) -> std::vec::IntoIter<TaskId> {
let mut live = Vec::with_capacity(self.len());
while let Some(task) = self.pop_front() {
live.push(task);
}
live.into_iter()
}
fn len(&self) -> usize {
self.present.len().saturating_sub(self.tombstones.len())
}
pub(crate) fn is_empty(&self) -> bool {
self.present.len() == self.tombstones.len()
}
fn snapshot(&self) -> Vec<TaskId> {
self.ready
.iter()
.copied()
.filter(|t| !self.tombstones.contains(t))
.collect()
}
}
impl Extend<TaskId> for LocalReadyQueueInner {
fn extend<T>(&mut self, iter: T)
where
T: IntoIterator<Item = TaskId>,
{
for task in iter {
self.push_back(task);
}
}
}
impl std::ops::Index<usize> for LocalReadyQueueInner {
type Output = TaskId;
fn index(&self, index: usize) -> &Self::Output {
self.iter()
.nth(index)
.expect("local-ready live index out of bounds")
}
}
type LocalReadyQueue = Mutex<LocalReadyQueueInner>;
fn local_ready_queue(initial: VecDeque<TaskId>) -> LocalReadyQueue {
Mutex::new(LocalReadyQueueInner::new(initial))
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum IoPhaseOutcome {
Progress,
Follower,
NoProgress,
}
#[inline]
fn select_io_poll_timeout(
idle_timeout: Option<Duration>,
fast_queue_empty: bool,
spawn_mailbox_has_work: bool,
) -> Option<Duration> {
if fast_queue_empty && !spawn_mailbox_has_work {
idle_timeout
} else {
Some(Duration::ZERO)
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum BackoffTimeoutDecision {
ParkTimeout { nanos: u64 },
DeadlineDue,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum EmptyBackoffAction {
Spin,
Yield,
Park,
}
#[inline]
fn select_backoff_deadline(
io_phase: IoPhaseOutcome,
timer_deadline: Option<Time>,
local_deadline: Option<Time>,
global_deadline: Option<Time>,
) -> Option<Time> {
if matches!(io_phase, IoPhaseOutcome::Follower) {
local_deadline
} else {
[timer_deadline, local_deadline, global_deadline]
.into_iter()
.flatten()
.min()
}
}
#[inline]
fn record_backoff_deadline_selection(
metrics: &mut PreemptionMetrics,
io_phase: IoPhaseOutcome,
timer_deadline: Option<Time>,
global_deadline: Option<Time>,
) {
if matches!(io_phase, IoPhaseOutcome::Follower)
&& (timer_deadline.is_some() || global_deadline.is_some())
{
metrics.follower_shared_deadline_ignored += 1;
}
}
#[inline]
fn record_backoff_timeout_park(
metrics: &mut PreemptionMetrics,
io_phase: IoPhaseOutcome,
nanos: u64,
) {
metrics.backoff_parks_total += 1;
metrics.backoff_timeout_parks_total += 1;
metrics.backoff_timeout_nanos_total = metrics.backoff_timeout_nanos_total.saturating_add(nanos);
if nanos <= SHORT_WAIT_LE_5MS_NANOS {
metrics.short_wait_le_5ms += 1;
}
if matches!(io_phase, IoPhaseOutcome::Follower) {
metrics.follower_timeout_parks += 1;
}
}
#[inline]
fn classify_backoff_timeout_decision(
_io_phase: IoPhaseOutcome,
next_deadline: Time,
now: Time,
) -> BackoffTimeoutDecision {
if next_deadline <= now {
BackoffTimeoutDecision::DeadlineDue
} else {
let nanos = next_deadline.duration_since(now);
BackoffTimeoutDecision::ParkTimeout { nanos }
}
}
#[inline]
fn record_backoff_indefinite_park(metrics: &mut PreemptionMetrics, io_phase: IoPhaseOutcome) {
metrics.backoff_parks_total += 1;
metrics.backoff_indefinite_parks += 1;
if matches!(io_phase, IoPhaseOutcome::Follower) {
metrics.follower_indefinite_parks += 1;
}
}
#[inline]
#[allow(clippy::cast_precision_loss)]
#[allow(dead_code)]
fn usize_to_f64(value: usize) -> f64 {
value as f64
}
#[inline]
#[allow(clippy::cast_precision_loss)]
fn u64_to_f64(value: u64) -> f64 {
value as f64
}
#[inline]
#[allow(clippy::cast_precision_loss)]
fn normalized_entropy(probs: &[f64]) -> f64 {
if probs.len() <= 1 {
return 0.0;
}
let mut entropy = 0.0_f64;
for &p in probs {
if p > f64::EPSILON {
entropy = p.mul_add(-p.ln(), entropy);
}
}
let max_entropy = (probs.len() as f64).ln();
if max_entropy <= f64::EPSILON {
0.0
} else {
(entropy / max_entropy).clamp(0.0, 1.0)
}
}
#[derive(Debug, Clone, Copy)]
pub(crate) struct AdaptiveEpochSnapshot {
potential: f64,
deadline_pressure: f64,
effective_limit_exceedances: u64,
fallback_cancel_dispatches: u64,
}
impl AdaptiveEpochSnapshot {
fn reward_against(self, end: Self, epoch_steps: u32) -> f64 {
let denom = self.potential.abs() + 1.0;
let normalized_drop = ((self.potential - end.potential) / denom).clamp(-1.0, 1.0);
let deadline_penalty = ((end.deadline_pressure - self.deadline_pressure).max(0.0)
/ (self.deadline_pressure.abs() + 1.0))
.clamp(0.0, 1.0);
let eps = f64::from(epoch_steps.max(1));
let effective_exceedances = u64_to_f64(
end.effective_limit_exceedances
.saturating_sub(self.effective_limit_exceedances),
);
let fairness_penalty = effective_exceedances / eps;
let fallback_penalty = u64_to_f64(
end.fallback_cancel_dispatches
.saturating_sub(self.fallback_cancel_dispatches),
) / eps;
let reward = 0.5f64.mul_add(normalized_drop, 0.5);
let reward = (-0.2f64).mul_add(deadline_penalty, reward);
let reward = (-0.2f64).mul_add(fairness_penalty.clamp(0.0, 1.0), reward);
let reward = (-0.1f64).mul_add(fallback_penalty.clamp(0.0, 1.0), reward);
reward.clamp(0.0, 1.0)
}
}
#[derive(Debug, Clone)]
pub(crate) struct AdaptiveCancelStreakPolicy {
arms: [usize; ADAPTIVE_STREAK_ARMS.len()],
mean_rewards: [f64; ADAPTIVE_STREAK_ARMS.len()],
discounted_pulls: [f64; ADAPTIVE_STREAK_ARMS.len()],
pulls: [u64; ADAPTIVE_STREAK_ARMS.len()],
selected_arm: usize,
epoch_steps: u32,
steps_in_epoch: u32,
epoch_count: u64,
reward_ema: f64,
e_process_log: f64,
epoch_start: Option<AdaptiveEpochSnapshot>,
}
impl AdaptiveCancelStreakPolicy {
fn new(epoch_steps: u32) -> Self {
let arms = ADAPTIVE_STREAK_ARMS;
Self {
arms,
mean_rewards: [0.0; ADAPTIVE_STREAK_ARMS.len()],
discounted_pulls: [0.0; ADAPTIVE_STREAK_ARMS.len()],
pulls: [0; ADAPTIVE_STREAK_ARMS.len()],
selected_arm: 2, epoch_steps: epoch_steps.max(1),
steps_in_epoch: 0,
epoch_count: 0,
reward_ema: 0.5,
e_process_log: 0.0,
epoch_start: None,
}
}
fn set_epoch_steps(&mut self, epoch_steps: u32) {
let epoch_steps = epoch_steps.max(1);
if self.epoch_steps == epoch_steps {
return;
}
self.epoch_steps = epoch_steps;
self.steps_in_epoch = 0;
self.epoch_start = None;
}
fn abort_epoch(&mut self) {
self.steps_in_epoch = 0;
self.epoch_start = None;
}
fn reset_to_priors(&mut self) {
let epoch_steps = self.epoch_steps;
*self = Self::new(epoch_steps);
}
fn current_limit(&self) -> usize {
self.arms[self.selected_arm]
}
fn select_arm_ucb(&self) -> usize {
let total_discounted_pulls: f64 = self.discounted_pulls.iter().sum();
if total_discounted_pulls < f64::EPSILON {
return 2; }
for (i, &n_i) in self.discounted_pulls.iter().enumerate() {
if n_i < f64::EPSILON {
return i;
}
}
let exploration_scale = ADAPTIVE_UCB_CONFIDENCE * total_discounted_pulls.ln().sqrt();
let mut best_arm = 0;
let mut best_ucb = f64::NEG_INFINITY;
for i in 0..self.arms.len() {
let n_i = self.discounted_pulls[i];
let confidence_bound = exploration_scale / n_i.sqrt();
let ucb_value = self.mean_rewards[i] + confidence_bound;
if ucb_value > best_ucb {
best_ucb = ucb_value;
best_arm = i;
}
}
best_arm
}
fn begin_epoch(&mut self, snapshot: AdaptiveEpochSnapshot) {
self.epoch_start = Some(snapshot);
}
fn on_dispatch(&mut self) -> bool {
self.steps_in_epoch = self.steps_in_epoch.saturating_add(1);
self.steps_in_epoch >= self.epoch_steps
}
fn complete_epoch(&mut self, end: AdaptiveEpochSnapshot) -> Option<f64> {
let start = self.epoch_start?;
let reward = start.reward_against(end, self.epoch_steps);
let chosen = self.selected_arm;
for i in 0..self.arms.len() {
self.discounted_pulls[i] *= ADAPTIVE_UCB_DISCOUNT;
}
let old_n = self.discounted_pulls[chosen];
let new_n = old_n + 1.0;
let delta = reward - self.mean_rewards[chosen];
self.mean_rewards[chosen] += delta / new_n;
self.discounted_pulls[chosen] = new_n;
self.e_process_log += ADAPTIVE_EPROCESS_LAMBDA
.mul_add(reward - 0.5, -(ADAPTIVE_EPROCESS_LAMBDA.powi(2) / 8.0));
self.reward_ema = 0.9f64.mul_add(self.reward_ema, 0.1 * reward);
self.pulls[chosen] = self.pulls[chosen].saturating_add(1);
self.epoch_count = self.epoch_count.saturating_add(1);
self.steps_in_epoch = 0;
self.selected_arm = self.select_arm_ucb();
self.epoch_start = Some(end);
Some(reward)
}
fn e_value(&self) -> f64 {
self.e_process_log.clamp(-60.0, 60.0).exp()
}
}
#[cfg(feature = "test-internals")]
#[derive(Debug, Clone, Copy)]
pub struct AdaptivePolicyBenchSnapshot(AdaptiveEpochSnapshot);
#[cfg(feature = "test-internals")]
impl AdaptivePolicyBenchSnapshot {
#[must_use]
pub fn new(
potential: f64,
deadline_pressure: f64,
_base_limit_exceedances: u64,
effective_limit_exceedances: u64,
fallback_cancel_dispatches: u64,
) -> Self {
Self(AdaptiveEpochSnapshot {
potential,
deadline_pressure,
effective_limit_exceedances,
fallback_cancel_dispatches,
})
}
}
#[cfg(feature = "test-internals")]
#[derive(Debug, Clone)]
pub struct AdaptiveCancelStreakPolicyBench {
policy: AdaptiveCancelStreakPolicy,
}
#[cfg(feature = "test-internals")]
impl AdaptiveCancelStreakPolicyBench {
#[must_use]
pub fn new(epoch_steps: u32) -> Self {
Self {
policy: AdaptiveCancelStreakPolicy::new(epoch_steps),
}
}
#[must_use]
pub fn arm_count(&self) -> usize {
self.policy.arms.len()
}
pub fn force_selected_arm(&mut self, arm_index: usize) {
assert!(arm_index < self.policy.arms.len(), "arm index out of range");
self.policy.selected_arm = arm_index;
}
pub fn seed_history(
&mut self,
mean_rewards: [f64; ADAPTIVE_STREAK_ARMS.len()],
discounted_pulls: [f64; ADAPTIVE_STREAK_ARMS.len()],
) {
self.policy.mean_rewards = mean_rewards;
self.policy.discounted_pulls = discounted_pulls;
}
pub fn begin_epoch(&mut self, snapshot: AdaptivePolicyBenchSnapshot) {
self.policy.begin_epoch(snapshot.0);
}
pub fn complete_epoch(&mut self, end: AdaptivePolicyBenchSnapshot) -> Option<f64> {
self.policy.complete_epoch(end.0)
}
#[must_use]
pub fn discounted_pulls(&self) -> [f64; ADAPTIVE_STREAK_ARMS.len()] {
self.policy.discounted_pulls
}
#[must_use]
pub fn mean_rewards(&self) -> [f64; ADAPTIVE_STREAK_ARMS.len()] {
self.policy.mean_rewards
}
#[must_use]
pub fn e_value(&self) -> f64 {
self.policy.e_value()
}
#[must_use]
pub fn select_arm_ucb(&self) -> usize {
self.policy.select_arm_ucb()
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
pub enum AdaptiveBatchDecisionReason {
Disabled,
FixedFallback,
ReadyContentionScaleUp,
CancelDebtFloor,
CooldownHold,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct AdaptiveBatchSizingProfile {
pub enabled: bool,
pub min_batch_size: usize,
pub max_batch_size: usize,
pub scale_up_ready_depth: usize,
pub scale_up_in_flight: usize,
pub scale_up_claim_failures: usize,
pub cancel_debt_floor: usize,
pub cooldown_steps: usize,
}
impl AdaptiveBatchSizingProfile {
#[inline]
fn normalized(self, fixed_batch_size: usize) -> Self {
let fixed_batch_size = fixed_batch_size.max(1);
let min_batch_size = self.min_batch_size.max(1);
let max_batch_size = self
.max_batch_size
.max(min_batch_size)
.max(fixed_batch_size);
Self {
enabled: self.enabled,
min_batch_size,
max_batch_size,
scale_up_ready_depth: self.scale_up_ready_depth,
scale_up_in_flight: self.scale_up_in_flight,
scale_up_claim_failures: self.scale_up_claim_failures,
cancel_debt_floor: self.cancel_debt_floor,
cooldown_steps: self.cooldown_steps,
}
}
#[inline]
fn contention_scale_up_batch_size(self, fixed_batch_size: usize) -> usize {
let fixed_batch_size = fixed_batch_size.max(1).min(self.max_batch_size);
if self.max_batch_size <= fixed_batch_size {
return fixed_batch_size;
}
let headroom = self.max_batch_size.saturating_sub(fixed_batch_size);
fixed_batch_size
.saturating_add((headroom / 2).max(1))
.clamp(self.min_batch_size, self.max_batch_size)
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct AdaptiveBatchDecisionSnapshot {
pub selected_batch_size: usize,
pub fixed_batch_size: usize,
pub ready_depth: usize,
pub cancel_debt: usize,
pub combiner_in_flight: usize,
pub combiner_claim_failures_delta: usize,
pub reason: AdaptiveBatchDecisionReason,
}
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
struct AdaptiveBatchRuntimeState {
active_batch_size: usize,
cooldown_remaining: usize,
last_combiner_claim_failures: usize,
last_snapshot: Option<AdaptiveBatchDecisionSnapshot>,
}
#[derive(Debug)]
pub(crate) struct WorkerCoordinator {
parkers: SmallVec<[Parker; 16]>,
next_wake: CachePadded<AtomicUsize>,
mask: Option<usize>,
io_driver: Option<IoDriverHandle>,
}
impl WorkerCoordinator {
pub(crate) fn new(parkers: SmallVec<[Parker; 16]>, io_driver: Option<IoDriverHandle>) -> Self {
let count = parkers.len();
let mask = if count > 0 && count.is_power_of_two() {
Some(count - 1)
} else {
None
};
Self {
parkers,
next_wake: CachePadded::new(AtomicUsize::new(0)),
mask,
io_driver,
}
}
#[inline]
fn wake_one_parker_prefer_waiter(&self) -> bool {
let count = self.parkers.len();
if count == 0 {
return false;
}
let start = self.next_wake.fetch_add(1, Ordering::AcqRel);
let slot_for = |index: usize| {
self.mask.map_or_else(|| index % count, |mask| index & mask)
};
for offset in 0..count {
let slot = slot_for(start.wrapping_add(offset));
if self.parkers[slot].unpark_if_waiting() {
return true;
}
}
self.parkers[slot_for(start)].unpark();
true
}
#[inline]
pub(crate) fn wake_one(&self) {
if !self.wake_one_parker_prefer_waiter() {
return;
}
if let Some(io) = &self.io_driver {
let _ = io.wake();
}
}
#[inline]
pub(crate) fn wake_one_parker(&self) {
self.wake_one_parker_prefer_waiter();
}
#[inline]
pub(crate) fn wake_many(&self, num_wakes: usize) {
let count = self.parkers.len();
if count == 0 || num_wakes == 0 {
return;
}
if num_wakes >= count {
self.wake_all();
return;
}
for _ in 0..num_wakes {
self.wake_one_parker_prefer_waiter();
}
if let Some(io) = &self.io_driver {
let _ = io.wake();
}
}
#[inline]
pub(crate) fn wake_worker(&self, worker_id: WorkerId) {
if let Some(parker) = self.parkers.get(worker_id) {
parker.unpark();
}
if let Some(io) = &self.io_driver {
let _ = io.wake();
}
}
#[inline]
pub(crate) fn wake_all(&self) {
for parker in &self.parkers {
parker.unpark();
}
if let Some(io) = &self.io_driver {
let _ = io.wake();
}
}
}
thread_local! {
static CURRENT_LOCAL: RefCell<Option<Arc<Mutex<PriorityScheduler>>>> =
const { RefCell::new(None) };
static CURRENT_LOCAL_READY: RefCell<Option<Arc<LocalReadyQueue>>> =
const { RefCell::new(None) };
static CURRENT_WORKER_ID: RefCell<Option<WorkerId>> = const { RefCell::new(None) };
}
#[derive(Debug)]
pub(crate) struct ScopedLocalScheduler {
prev: Option<Arc<Mutex<PriorityScheduler>>>,
}
impl ScopedLocalScheduler {
pub(crate) fn new(local: Arc<Mutex<PriorityScheduler>>) -> Self {
let prev = CURRENT_LOCAL.with(|cell| cell.replace(Some(local)));
Self { prev }
}
}
impl Drop for ScopedLocalScheduler {
fn drop(&mut self) {
let prev = self.prev.take();
CURRENT_LOCAL.with(|cell| {
*cell.borrow_mut() = prev;
});
}
}
pub(crate) struct ScopedWorkerId {
prev: Option<WorkerId>,
}
impl ScopedWorkerId {
pub(crate) fn new(id: WorkerId) -> Self {
let prev = CURRENT_WORKER_ID.with(|cell| cell.replace(Some(id)));
Self { prev }
}
}
impl Drop for ScopedWorkerId {
fn drop(&mut self) {
let prev = self.prev.take();
let _ = CURRENT_WORKER_ID.try_with(|cell| {
*cell.borrow_mut() = prev;
});
}
}
pub(crate) struct ScopedLocalReady {
prev: Option<Arc<LocalReadyQueue>>,
}
impl ScopedLocalReady {
pub(crate) fn new(queue: Arc<LocalReadyQueue>) -> Self {
let prev = CURRENT_LOCAL_READY.with(|cell| cell.replace(Some(queue)));
Self { prev }
}
}
impl Drop for ScopedLocalReady {
fn drop(&mut self) {
CURRENT_LOCAL_READY.with(|cell| {
*cell.borrow_mut() = self.prev.take();
});
}
}
#[inline]
pub(crate) fn schedule_local_task(task: TaskId) -> bool {
CURRENT_LOCAL_READY.with(|cell| {
cell.borrow().as_ref().is_some_and(|queue| {
queue.lock().push_back(task);
true
})
})
}
#[inline]
pub(crate) fn current_worker_id() -> Option<WorkerId> {
CURRENT_WORKER_ID.with(|cell| *cell.borrow())
}
fn trapped_scc_with_edge_observer<F>(
adjacency: &[Vec<usize>],
mut observe_edge: F,
) -> Option<Vec<usize>>
where
F: FnMut(usize, usize),
{
struct Tarjan<'a, F> {
adjacency: &'a [Vec<usize>],
observe_edge: &'a mut F,
index: usize,
stack: Vec<usize>,
on_stack: Vec<bool>,
indices: Vec<Option<usize>>,
lowlink: Vec<usize>,
trapped: Option<Vec<usize>>,
}
impl<F: FnMut(usize, usize)> Tarjan<'_, F> {
fn strongconnect(&mut self, v: usize) {
if self.trapped.is_some() {
return;
}
self.indices[v] = Some(self.index);
self.lowlink[v] = self.index;
self.index += 1;
self.stack.push(v);
self.on_stack[v] = true;
for &w in &self.adjacency[v] {
if self.trapped.is_some() {
return;
}
(self.observe_edge)(v, w);
if self.indices[w].is_none() {
self.strongconnect(w);
if self.trapped.is_some() {
return;
}
self.lowlink[v] = self.lowlink[v].min(self.lowlink[w]);
} else if self.on_stack[w] {
self.lowlink[v] = self.lowlink[v].min(self.indices[w].unwrap_or(usize::MAX));
}
}
if self.lowlink[v] == self.indices[v].unwrap_or(usize::MAX) {
let mut component = Vec::new();
while let Some(w) = self.stack.pop() {
self.on_stack[w] = false;
component.push(w);
if w == v {
break;
}
}
let cyclic = component.len() > 1
|| component
.first()
.is_some_and(|n| self.adjacency[*n].contains(n));
if cyclic {
let component_set: BTreeSet<usize> = component.iter().copied().collect();
let mut has_egress = false;
for &u in &component {
if self.adjacency[u].iter().any(|v| !component_set.contains(v)) {
has_egress = true;
break;
}
}
if !has_egress {
component.sort_unstable();
self.trapped = Some(component);
}
}
}
}
}
let n = adjacency.len();
let mut tarjan = Tarjan {
adjacency,
observe_edge: &mut observe_edge,
index: 0,
stack: Vec::new(),
on_stack: vec![false; n],
indices: vec![None; n],
lowlink: vec![0; n],
trapped: None,
};
for v in 0..n {
if tarjan.indices[v].is_none() {
tarjan.strongconnect(v);
if tarjan.trapped.is_some() {
return tarjan.trapped;
}
}
}
None
}
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, PartialOrd, Ord, serde::Serialize)]
#[serde(rename_all = "snake_case")]
#[allow(dead_code)]
enum WaitCause {
Lock,
Channel,
Notify,
Join,
#[default]
Unknown,
}
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, PartialOrd, Ord, serde::Serialize)]
struct WaitLocation {
file: Option<&'static str>,
line: Option<u32>,
label: Option<&'static str>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Serialize)]
struct WaitGraphEdgeSnapshot {
waiter: TaskId,
cause: WaitCause,
location: WaitLocation,
}
#[derive(Debug, Clone)]
struct WaitGraphTaskSnapshot {
id: TaskId,
waiters: Vec<TaskId>,
wait_edges: Vec<WaitGraphEdgeSnapshot>,
}
#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize)]
struct DeadlockWaitEdgeReport {
waiter: TaskId,
blocked_on: TaskId,
cause: WaitCause,
location: WaitLocation,
}
#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize)]
struct DeadlockCycleReport {
tasks: Vec<TaskId>,
edges: Vec<DeadlockWaitEdgeReport>,
}
#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize)]
struct WaitGraphSignalReport {
node_count: usize,
undirected_edges: Vec<(usize, usize)>,
trapped_wait_cycle: bool,
trapped_cycle: Option<DeadlockCycleReport>,
}
fn wait_graph_snapshot_from_state(state: &RuntimeState) -> Vec<WaitGraphTaskSnapshot> {
wait_graph_snapshot_from_tasks(&state.tasks)
}
fn wait_graph_snapshot_from_tasks(tasks: &TaskTable) -> Vec<WaitGraphTaskSnapshot> {
let mut snapshots = Vec::new();
for (_, task) in tasks.iter() {
if !task.state.is_terminal() {
let wait_edges = task
.waiters
.iter()
.copied()
.map(|waiter| WaitGraphEdgeSnapshot {
waiter,
cause: WaitCause::Unknown,
location: WaitLocation::default(),
})
.collect();
snapshots.push(WaitGraphTaskSnapshot {
id: task.id,
waiters: task.waiters.to_vec(),
wait_edges,
});
}
}
snapshots
}
fn wait_graph_signal_report_from_snapshot(
tasks: &[WaitGraphTaskSnapshot],
) -> WaitGraphSignalReport {
let mut live_tasks: Vec<TaskId> = tasks.iter().map(|task| task.id).collect();
live_tasks.sort();
let index_by_task: BTreeMap<TaskId, usize> = live_tasks
.iter()
.enumerate()
.map(|(idx, id)| (*id, idx))
.collect();
let mut undirected_edges: BTreeSet<(usize, usize)> = BTreeSet::new();
let mut adjacency = vec![Vec::new(); live_tasks.len()];
for task in tasks {
let Some(&task_idx) = index_by_task.get(&task.id) else {
continue;
};
for edge in &task.wait_edges {
if let Some(&waiter_idx) = index_by_task.get(&edge.waiter) {
adjacency[waiter_idx].push(task_idx);
if waiter_idx == task_idx {
continue;
}
undirected_edges.insert(if waiter_idx < task_idx {
(waiter_idx, task_idx)
} else {
(task_idx, waiter_idx)
});
}
}
if task.wait_edges.is_empty() {
for waiter in &task.waiters {
if let Some(&waiter_idx) = index_by_task.get(waiter) {
adjacency[waiter_idx].push(task_idx);
if waiter_idx == task_idx {
continue;
}
undirected_edges.insert(if waiter_idx < task_idx {
(waiter_idx, task_idx)
} else {
(task_idx, waiter_idx)
});
}
}
}
}
for edges in &mut adjacency {
edges.sort_unstable();
edges.dedup();
}
let trapped_scc = trapped_scc_with_edge_observer(&adjacency, |_, _| {});
let trapped_cycle = trapped_scc.as_ref().map(|component| {
let component_set: BTreeSet<usize> = component.iter().copied().collect();
let cycle_tasks: Vec<TaskId> = component.iter().map(|idx| live_tasks[*idx]).collect();
let mut edges = Vec::new();
for snapshot in tasks {
let Some(&task_idx) = index_by_task.get(&snapshot.id) else {
continue;
};
if !component_set.contains(&task_idx) {
continue;
}
for edge in &snapshot.wait_edges {
let Some(&waiter_idx) = index_by_task.get(&edge.waiter) else {
continue;
};
if component_set.contains(&waiter_idx) {
edges.push(DeadlockWaitEdgeReport {
waiter: edge.waiter,
blocked_on: snapshot.id,
cause: edge.cause,
location: edge.location,
});
}
}
if snapshot.wait_edges.is_empty() {
for waiter in &snapshot.waiters {
let Some(&waiter_idx) = index_by_task.get(waiter) else {
continue;
};
if component_set.contains(&waiter_idx) {
edges.push(DeadlockWaitEdgeReport {
waiter: *waiter,
blocked_on: snapshot.id,
cause: WaitCause::Unknown,
location: WaitLocation::default(),
});
}
}
}
}
edges.sort_by_key(|edge| (edge.waiter, edge.blocked_on, edge.cause, edge.location));
DeadlockCycleReport {
tasks: cycle_tasks,
edges,
}
});
WaitGraphSignalReport {
node_count: live_tasks.len(),
undirected_edges: undirected_edges.into_iter().collect(),
trapped_wait_cycle: trapped_cycle.is_some(),
trapped_cycle,
}
}
fn wait_graph_signals_from_snapshot(
tasks: &[WaitGraphTaskSnapshot],
) -> (usize, Vec<(usize, usize)>, bool) {
let report = wait_graph_signal_report_from_snapshot(tasks);
(
report.node_count,
report.undirected_edges,
report.trapped_wait_cycle,
)
}
#[cfg(test)]
fn wait_graph_signals_from_state(state: &RuntimeState) -> (usize, Vec<(usize, usize)>, bool) {
let snapshot = wait_graph_snapshot_from_state(state);
wait_graph_signals_from_snapshot(&snapshot)
}
#[inline]
pub(crate) fn schedule_on_current_local(task: TaskId, priority: u8) -> bool {
if LocalQueue::schedule_local(task) {
return true;
}
CURRENT_LOCAL.with(|cell| {
if let Some(local) = cell.borrow().as_ref() {
local.lock().schedule(task, priority);
return true;
}
false
})
}
#[inline]
fn move_local_ready_task_to_cancel_lane(
local: &Mutex<PriorityScheduler>,
local_ready: &LocalReadyQueue,
task: TaskId,
priority: u8,
) {
let mut local_guard = local.lock();
local_ready.lock().tombstone(task);
local_guard.move_to_cancel_lane(task, priority);
}
#[inline]
pub(crate) fn schedule_cancel_on_current_local(task: TaskId, priority: u8) -> bool {
CURRENT_LOCAL.with(|cell| {
let borrow = cell.borrow();
let Some(local) = borrow.as_ref() else {
return false;
};
let mut local_guard = local.lock();
CURRENT_LOCAL_READY.with(|lr_cell| {
if let Some(queue) = lr_cell.borrow().as_ref() {
queue.lock().tombstone(task);
}
});
local_guard.move_to_cancel_lane(task, priority);
drop(local_guard);
true
})
}
#[derive(Debug)]
pub struct ThreeLaneScheduler {
global: Arc<GlobalInjector>,
local_schedulers: Vec<Arc<Mutex<PriorityScheduler>>>,
local_ready: SmallVec<[Arc<LocalReadyQueue>; 16]>,
parkers: SmallVec<[Parker; 16]>,
workers: SmallVec<[ThreeLaneWorker; 16]>,
shutdown: Arc<AtomicBool>,
coordinator: Arc<WorkerCoordinator>,
browser_ready_handoff_limit: usize,
steal_batch_size: usize,
enable_parking: bool,
#[allow(dead_code)] timer_driver: Option<TimerDriverHandle>,
state: Arc<ContendedMutex<RuntimeState>>,
task_table: Option<Arc<ContendedMutex<TaskTable>>>,
global_queue_limit: usize,
scheduler_evidence: Option<Arc<Mutex<SchedulerEvidenceCollector>>>,
placement_mode: SchedulerPlacementMode,
worker_cohort_map: Option<Vec<usize>>,
cohort_count: usize,
spawn_mailbox: Option<Arc<crate::runtime::spawn_mailbox::SpawnMailbox>>,
}
#[derive(Copy, Clone)]
enum ScheduleIntent {
Spawn,
Wake,
}
impl ScheduleIntent {
fn local_route_failure_assert(self, task: TaskId) -> String {
match self {
Self::Spawn => format!(
"Attempted to spawn local task {task:?} from non-owner thread or outside worker context"
),
Self::Wake => format!(
"Attempted to wake local task {task:?} via scheduler from non-owner thread. Use Waker instead."
),
}
}
fn local_route_failure_log(self) -> &'static str {
match self {
Self::Spawn => {
"spawn: local task cannot be scheduled from non-owner thread, spawn skipped"
}
Self::Wake => "wake: local task cannot be woken from non-owner thread, wake skipped",
}
}
}
#[derive(Clone)]
pub struct SchedulerConstructionHandles {
pub io_driver: Option<IoDriverHandle>,
pub timer_driver: Option<TimerDriverHandle>,
pub pending_cancel_dispatch_ready: Arc<AtomicBool>,
}
impl std::fmt::Debug for SchedulerConstructionHandles {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("SchedulerConstructionHandles")
.field("has_io_driver", &self.io_driver.is_some())
.field("has_timer_driver", &self.timer_driver.is_some())
.field(
"pending_cancel_dispatch_ready",
&self
.pending_cancel_dispatch_ready
.load(std::sync::atomic::Ordering::Relaxed),
)
.finish()
}
}
impl SchedulerConstructionHandles {
#[must_use]
pub fn extract_from_unified(state: &Arc<ContendedMutex<RuntimeState>>) -> Self {
let guard = state
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
Self {
io_driver: guard.io_driver_handle(),
timer_driver: guard.timer_driver_handle(),
pending_cancel_dispatch_ready: guard.pending_cancel_dispatch_ready_handle(),
}
}
#[must_use]
pub fn extract_from_sharded(
shards: &crate::runtime::sharded_state::ShardedState,
state: &Arc<ContendedMutex<RuntimeState>>,
) -> Self {
let pending_cancel_dispatch_ready = {
let guard = state
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
guard.pending_cancel_dispatch_ready_handle()
};
Self {
io_driver: shards.io_driver_handle(),
timer_driver: shards.timer_driver_handle(),
pending_cancel_dispatch_ready,
}
}
}
impl ThreeLaneScheduler {
#[inline]
fn initial_local_scheduler_capacity(worker_count: usize) -> usize {
let workers = worker_count.max(1);
let per_worker = LOCAL_SCHEDULER_BURST_BUDGET.div_ceil(workers);
per_worker.clamp(LOCAL_SCHEDULER_MIN_CAPACITY, LOCAL_SCHEDULER_MAX_CAPACITY)
}
pub fn new(worker_count: usize, state: &Arc<ContendedMutex<RuntimeState>>) -> Self {
Self::new_with_options(worker_count, state, DEFAULT_CANCEL_STREAK_LIMIT, false, 32)
}
pub fn try_new(
worker_count: usize,
state: &Arc<ContendedMutex<RuntimeState>>,
) -> Result<Self, crate::error::Error> {
Self::try_new_with_options_and_task_table(
worker_count,
state,
None,
DEFAULT_CANCEL_STREAK_LIMIT,
false,
32,
)
}
pub fn new_with_cancel_limit(
worker_count: usize,
state: &Arc<ContendedMutex<RuntimeState>>,
cancel_streak_limit: usize,
) -> Self {
Self::new_with_options(worker_count, state, cancel_streak_limit, false, 32)
}
pub fn new_with_options(
worker_count: usize,
state: &Arc<ContendedMutex<RuntimeState>>,
cancel_streak_limit: usize,
enable_governor: bool,
governor_interval: u32,
) -> Self {
Self::new_with_options_and_task_table(
worker_count,
state,
None,
cancel_streak_limit,
enable_governor,
governor_interval,
)
}
pub fn new_with_options_and_task_table(
worker_count: usize,
state: &Arc<ContendedMutex<RuntimeState>>,
task_table: Option<Arc<ContendedMutex<TaskTable>>>,
cancel_streak_limit: usize,
enable_governor: bool,
governor_interval: u32,
) -> Self {
let handles = SchedulerConstructionHandles::extract_from_unified(state);
let scheduler = Self::new_with_options_task_table_and_handles(
worker_count,
state,
task_table,
handles,
cancel_streak_limit,
enable_governor,
governor_interval,
);
scheduler.install_pending_cancel_dispatch_coordinator(state);
scheduler
}
#[must_use]
pub(crate) fn dispatch_task_table(&self) -> Option<Arc<ContendedMutex<TaskTable>>> {
self.task_table.clone()
}
pub fn new_with_sharded_state(
worker_count: usize,
state: &Arc<ContendedMutex<RuntimeState>>,
shards: &crate::runtime::sharded_state::ShardedState,
cancel_streak_limit: usize,
enable_governor: bool,
governor_interval: u32,
) -> Self {
let handles = SchedulerConstructionHandles::extract_from_sharded(shards, state);
let scheduler = Self::new_with_options_task_table_and_handles(
worker_count,
state,
Some(shards.task_shard_handle()),
handles,
cancel_streak_limit,
enable_governor,
governor_interval,
);
scheduler.install_pending_cancel_dispatch_coordinator(state);
scheduler
}
#[allow(clippy::too_many_lines)]
pub fn new_with_options_task_table_and_handles(
worker_count: usize,
state: &Arc<ContendedMutex<RuntimeState>>,
task_table: Option<Arc<ContendedMutex<TaskTable>>>,
handles: SchedulerConstructionHandles,
cancel_streak_limit: usize,
enable_governor: bool,
governor_interval: u32,
) -> Self {
let worker_count = worker_count.max(1);
let cancel_streak_limit = cancel_streak_limit.max(1);
let browser_ready_handoff_limit = DEFAULT_BROWSER_READY_HANDOFF_LIMIT;
let governor_interval = governor_interval.max(1);
let steal_batch_size = DEFAULT_STEAL_BATCH_SIZE;
let enable_parking = DEFAULT_ENABLE_PARKING;
let global = Arc::new(GlobalInjector::new());
let scheduler_evidence = None;
let shutdown = Arc::new(AtomicBool::new(false));
let mut workers = SmallVec::<[ThreeLaneWorker; 16]>::with_capacity(worker_count);
let mut parkers = SmallVec::<[Parker; 16]>::with_capacity(worker_count);
let mut local_schedulers: Vec<Arc<Mutex<PriorityScheduler>>> =
Vec::with_capacity(worker_count);
let mut local_ready = SmallVec::<[Arc<LocalReadyQueue>; 16]>::with_capacity(worker_count);
let local_scheduler_capacity = Self::initial_local_scheduler_capacity(worker_count);
let SchedulerConstructionHandles {
io_driver,
timer_driver,
pending_cancel_dispatch_ready,
} = handles;
for _ in 0..worker_count {
local_schedulers.push(Arc::new(Mutex::new(PriorityScheduler::with_capacity(
local_scheduler_capacity,
))));
}
for _ in 0..worker_count {
local_ready.push(Arc::new(local_ready_queue(VecDeque::with_capacity(32))));
}
for _ in 0..worker_count {
parkers.push(Parker::new());
}
let coordinator = Arc::new(WorkerCoordinator::new(parkers.clone(), io_driver.clone()));
let fast_queues: Vec<LocalQueue> = (0..worker_count)
.map(|_| {
task_table.as_ref().map_or_else(
|| LocalQueue::new(Arc::clone(state)),
|tt| LocalQueue::new_with_task_table(Arc::clone(tt)),
)
})
.collect();
for id in 0..worker_count {
let parker = parkers[id].clone();
let stealers: SmallVec<[Arc<Mutex<PriorityScheduler>>; 16]> = local_schedulers
.iter()
.enumerate()
.filter(|(i, _)| *i != id)
.map(|(_, sched)| Arc::clone(sched))
.collect();
let heap_stealer_locality: SmallVec<[StealerLocality; 16]> = (0..stealers.len())
.map(|_| StealerLocality::SameCohort)
.collect();
let fast_stealers: SmallVec<[local_queue::Stealer; 16]> = fast_queues
.iter()
.enumerate()
.filter(|(i, _)| *i != id)
.map(|(_, q)| q.stealer())
.collect();
let fast_stealer_locality: SmallVec<[StealerLocality; 16]> = (0..fast_stealers.len())
.map(|_| StealerLocality::SameCohort)
.collect();
workers.push(ThreeLaneWorker {
id,
local: Arc::clone(&local_schedulers[id]),
stealers,
preferred_heap_stealer_count: worker_count.saturating_sub(1),
heap_stealer_locality,
fast_queue: fast_queues[id].clone(),
global_ready_buffer: Vec::with_capacity(steal_batch_size),
fast_stealers,
preferred_fast_stealer_count: worker_count.saturating_sub(1),
fast_stealer_locality,
local_ready: Arc::clone(&local_ready[id]),
all_local_ready: local_ready.clone(),
all_local_schedulers: local_schedulers.iter().cloned().collect(),
global: Arc::clone(&global),
state: Arc::clone(state),
pending_cancel_dispatch_ready: Arc::clone(&pending_cancel_dispatch_ready),
task_table: task_table.clone(),
parker,
coordinator: Arc::clone(&coordinator),
spawn_mailbox: None,
rng: DetRng::new(id as u64),
shutdown: Arc::clone(&shutdown),
io_driver: io_driver.clone(),
timer_driver: timer_driver.clone(),
steal_buffer: Vec::new(),
steal_batch_size,
enable_parking,
empty_backoff: 0,
cancel_streak: 0,
ready_dispatch_streak: 0,
browser_ready_handoff_limit,
cancel_streak_limit,
governor: if enable_governor {
Some(LyapunovGovernor::with_defaults())
} else {
None
},
cached_suggestion: SchedulingSuggestion::NoPreference,
steps_since_snapshot: governor_interval.saturating_sub(1),
governor_interval,
preemption_metrics: PreemptionMetrics {
adaptive_current_limit: cancel_streak_limit,
adaptive_e_value: 1.0,
..PreemptionMetrics::default()
},
evidence_sink: None,
decision_contract: if enable_governor {
Some(super::decision_contract::SchedulerDecisionContract::new())
} else {
None
},
decision_posterior: if enable_governor {
Some(franken_decision::Posterior::uniform(
super::decision_contract::state::COUNT,
))
} else {
None
},
adaptive_cancel_policy: None,
spectral_monitor: if enable_governor {
Some(SpectralHealthMonitor::new(SpectralThresholds::default()))
} else {
None
},
drain_certificate: if enable_governor {
Some(ProgressCertificate::with_defaults())
} else {
None
},
decision_sequence: 0,
fairness_monitor: Mutex::new(FairnessMonitor::with_defaults()),
invariant_monitor: Mutex::new(
super::invariant_monitor::SchedulerInvariantMonitor::with_defaults(),
),
fast_queue_dispatch_streak: 0,
fast_queue_fairness_limit: 4, timed_dispatch_streak: 0,
timed_fairness_limit: 6, adaptive_batch_profile: None,
adaptive_batch_state: AdaptiveBatchRuntimeState::default(),
steal_locality_counters: StealLocalityCounters::default(),
scheduler_evidence: scheduler_evidence.clone(),
});
}
Self {
global,
local_schedulers,
local_ready,
parkers,
workers,
shutdown,
coordinator,
spawn_mailbox: None,
timer_driver,
state: Arc::clone(state),
task_table,
browser_ready_handoff_limit,
steal_batch_size,
enable_parking,
global_queue_limit: 0,
scheduler_evidence,
placement_mode: SchedulerPlacementMode::default(),
worker_cohort_map: None,
cohort_count: 1,
}
}
pub fn try_new_with_options_and_task_table(
worker_count: usize,
state: &Arc<ContendedMutex<RuntimeState>>,
task_table: Option<Arc<ContendedMutex<TaskTable>>>,
cancel_streak_limit: usize,
enable_governor: bool,
governor_interval: u32,
) -> Result<Self, crate::error::Error> {
if worker_count == 0 {
return Err(
crate::error::Error::new(crate::error::ErrorKind::ConfigError).with_message(
"ThreeLaneScheduler requires worker_count >= 1; \
a zero-worker scheduler cannot dispatch any task and \
silently hangs block_on. Use try_new_with_options_and_task_table \
to surface this as ConfigError; the infallible \
constructors clamp to 1 instead.",
),
);
}
Ok(Self::new_with_options_and_task_table(
worker_count,
state,
task_table,
cancel_streak_limit,
enable_governor,
governor_interval,
))
}
pub fn install_pending_cancel_dispatch_coordinator(
&self,
state: &Arc<ContendedMutex<RuntimeState>>,
) {
let mut guard = state
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
guard.set_pending_cancel_dispatch_coordinator(&self.coordinator);
}
pub fn set_steal_batch_size(&mut self, size: usize) {
let size = size.max(1);
self.steal_batch_size = size;
for worker in &mut self.workers {
worker.steal_batch_size = size;
if worker.steal_buffer.capacity() < size {
worker
.steal_buffer
.reserve(size - worker.steal_buffer.capacity());
}
if worker.global_ready_buffer.capacity() < size {
worker
.global_ready_buffer
.reserve(size - worker.global_ready_buffer.capacity());
}
worker.reset_adaptive_batch_state();
}
}
pub fn set_adaptive_batch_profile(&mut self, profile: Option<AdaptiveBatchSizingProfile>) {
for worker in &mut self.workers {
worker.adaptive_batch_profile = profile;
worker.reset_adaptive_batch_state();
worker.preemption_metrics.adaptive_batch_scale_up_events = 0;
worker.preemption_metrics.adaptive_batch_cancel_floor_hits = 0;
worker.preemption_metrics.adaptive_batch_cooldown_holds = 0;
worker.preemption_metrics.adaptive_batch_max_selected = worker.fixed_ready_batch_size();
}
}
#[doc(hidden)]
#[cfg(any(test, feature = "test-internals"))]
pub fn set_adaptive_batch_profile_for_test(
&mut self,
profile: Option<AdaptiveBatchSizingProfile>,
) {
self.set_adaptive_batch_profile(profile);
}
#[doc(hidden)]
#[cfg(any(test, feature = "test-internals"))]
pub fn seed_ready_combiner_pressure_for_test(
&self,
max_in_flight: usize,
combiner_claim_failures: usize,
) {
self.global
.seed_ready_combiner_pressure_for_test(max_in_flight, combiner_claim_failures);
}
fn ordered_steal_peers(
worker_id: usize,
worker_to_cohort: &[usize],
mode: SchedulerPlacementMode,
) -> Vec<usize> {
let worker_count = worker_to_cohort.len();
let my_cohort = worker_to_cohort[worker_id];
let mut peers = (0..worker_count)
.filter(|&peer_id| peer_id != worker_id)
.collect::<Vec<_>>();
match mode {
SchedulerPlacementMode::LocalityFirst => {
peers.sort_by_key(|&peer_id| (worker_to_cohort[peer_id] != my_cohort, peer_id));
}
SchedulerPlacementMode::LatencyFirst => {
peers.sort_by_key(|&peer_id| {
(
worker_to_cohort[peer_id] != my_cohort,
Self::worker_slot_distance(worker_id, peer_id, worker_count),
peer_id,
)
});
}
SchedulerPlacementMode::ThroughputFirst => {
peers.sort_unstable();
}
}
peers
}
#[inline]
fn preferred_stealer_count(
mode: SchedulerPlacementMode,
my_cohort: usize,
worker_to_cohort: &[usize],
ordered_peers: &[usize],
) -> usize {
if matches!(mode, SchedulerPlacementMode::ThroughputFirst) {
return ordered_peers.len();
}
ordered_peers
.iter()
.take_while(|&&peer_id| worker_to_cohort[peer_id] == my_cohort)
.count()
}
#[inline]
fn worker_slot_distance(lhs: usize, rhs: usize, worker_count: usize) -> usize {
let forward = if rhs >= lhs {
rhs - lhs
} else {
worker_count - (lhs - rhs)
};
forward.min(worker_count.saturating_sub(forward))
}
pub fn set_worker_cohort_map(
&mut self,
worker_to_cohort: &[usize],
) -> Result<(), crate::error::Error> {
let worker_count = self.workers.len();
if worker_count == 0 {
return Err(
crate::error::Error::new(crate::error::ErrorKind::ConfigError)
.with_message("worker cohort map requires at least one worker"),
);
}
if worker_to_cohort.len() != worker_count {
return Err(
crate::error::Error::new(crate::error::ErrorKind::ConfigError)
.with_message("worker cohort map length must match worker_threads".to_string()),
);
}
self.rebuild_worker_stealers(worker_to_cohort);
self.worker_cohort_map = Some(worker_to_cohort.to_vec());
self.cohort_count = worker_to_cohort
.iter()
.copied()
.max()
.map_or(1, |max_cohort| max_cohort.saturating_add(1));
Ok(())
}
pub fn set_scheduler_placement_mode(&mut self, mode: SchedulerPlacementMode) {
self.placement_mode = mode;
if let Some(worker_to_cohort) = self.worker_cohort_map.clone() {
self.rebuild_worker_stealers(&worker_to_cohort);
}
}
#[must_use]
pub const fn scheduler_placement_mode(&self) -> SchedulerPlacementMode {
self.placement_mode
}
fn rebuild_worker_stealers(&mut self, worker_to_cohort: &[usize]) {
let worker_count = self.workers.len();
let fast_queues: Vec<_> = self
.workers
.iter()
.map(|worker| worker.fast_queue.clone())
.collect();
let local_schedulers = self.local_schedulers.clone();
for (worker_id, worker) in self.workers.iter_mut().enumerate() {
let my_cohort = worker_to_cohort[worker_id];
let ordered_peers =
Self::ordered_steal_peers(worker_id, worker_to_cohort, self.placement_mode);
let preferred_count = Self::preferred_stealer_count(
self.placement_mode,
my_cohort,
worker_to_cohort,
&ordered_peers,
);
let mut fast_stealers = SmallVec::<[local_queue::Stealer; 16]>::new();
let mut fast_stealer_locality = SmallVec::<[StealerLocality; 16]>::new();
let mut heap_stealers = SmallVec::<[Arc<Mutex<PriorityScheduler>>; 16]>::new();
let mut heap_stealer_locality = SmallVec::<[StealerLocality; 16]>::new();
for peer_id in ordered_peers {
let locality =
StealerLocality::from_same_cohort(worker_to_cohort[peer_id] == my_cohort);
fast_stealers.push(fast_queues[peer_id].stealer());
fast_stealer_locality.push(locality);
heap_stealers.push(Arc::clone(&local_schedulers[peer_id]));
heap_stealer_locality.push(locality);
}
debug_assert_eq!(fast_stealers.len(), worker_count.saturating_sub(1));
debug_assert_eq!(heap_stealers.len(), worker_count.saturating_sub(1));
worker.fast_stealers = fast_stealers;
worker.preferred_fast_stealer_count = preferred_count;
worker.fast_stealer_locality = fast_stealer_locality;
worker.stealers = heap_stealers;
worker.preferred_heap_stealer_count = preferred_count;
worker.heap_stealer_locality = heap_stealer_locality;
worker.steal_locality_counters = StealLocalityCounters::default();
}
}
#[doc(hidden)]
#[cfg(feature = "test-internals")]
pub fn seed_worker_fast_ready_for_test(&mut self, worker_id: usize, task: TaskId) {
self.workers[worker_id].fast_queue.push(task);
}
#[doc(hidden)]
#[cfg(feature = "test-internals")]
pub fn seed_worker_priority_ready_for_test(
&mut self,
worker_id: usize,
task: TaskId,
priority: u8,
) {
self.workers[worker_id]
.local
.lock()
.schedule(task, priority);
}
pub fn set_enable_parking(&mut self, enable: bool) {
self.enable_parking = enable;
for worker in &mut self.workers {
worker.enable_parking = enable;
}
}
pub fn set_browser_ready_handoff_limit(&mut self, limit: usize) {
self.browser_ready_handoff_limit = limit;
for worker in &mut self.workers {
worker.browser_ready_handoff_limit = limit;
if limit == 0 {
worker.ready_dispatch_streak = 0;
}
}
}
pub fn set_adaptive_cancel_streak(&mut self, enable: bool, epoch_steps: u32) {
let epoch_steps = epoch_steps.max(1);
for worker in &mut self.workers {
if enable {
if let Some(policy) = worker.adaptive_cancel_policy.as_mut() {
policy.set_epoch_steps(epoch_steps);
} else {
worker.adaptive_cancel_policy =
Some(AdaptiveCancelStreakPolicy::new(epoch_steps));
}
if let Some(policy) = worker.adaptive_cancel_policy.as_ref() {
worker.preemption_metrics.adaptive_current_limit = policy.current_limit();
worker.preemption_metrics.adaptive_reward_ema = policy.reward_ema;
worker.preemption_metrics.adaptive_e_value = policy.e_value();
}
} else {
worker.adaptive_cancel_policy = None;
worker.preemption_metrics.adaptive_epochs = 0;
worker.preemption_metrics.adaptive_current_limit = worker.cancel_streak_limit;
worker.preemption_metrics.adaptive_reward_ema = 0.0;
worker.preemption_metrics.adaptive_e_value = 1.0;
}
}
}
pub fn set_global_queue_limit(&mut self, limit: usize) {
self.global_queue_limit = limit;
}
#[inline]
fn record_scheduler_evidence_enqueue(&self, task: TaskId) {
let Some(collector) = &self.scheduler_evidence else {
return;
};
let timestamp_ns = crate::time::wall_now().as_nanos();
collector.lock().record_task_enqueue(task, timestamp_ns);
}
#[inline]
fn contain_publication_effect(effect: impl FnOnce()) {
if let Err(payload) = std::panic::catch_unwind(std::panic::AssertUnwindSafe(effect)) {
std::mem::forget(payload);
}
}
#[inline]
fn finish_global_ready_publication(
&self,
task: TaskId,
priority: u8,
ready_count_before: usize,
) {
let _ = priority;
Self::contain_publication_effect(|| self.record_scheduler_evidence_enqueue(task));
if self.global_queue_limit > 0 && ready_count_before >= self.global_queue_limit {
Self::contain_publication_effect(|| {
crate::tracing_compat::warn!(
?task,
priority,
limit = self.global_queue_limit,
current = ready_count_before,
"inject_ready: global ready queue at capacity, scheduling anyway"
);
});
}
Self::contain_publication_effect(|| self.wake_one());
}
fn scheduler_evidence_remote_steal_ratio_pct(&self) -> Option<u8> {
let (preferred, remote) =
self.workers
.iter()
.fold((0_u64, 0_u64), |(preferred, remote), worker| {
let counters = worker.steal_locality_counters;
(
preferred
.saturating_add(counters.preferred_fast_steals)
.saturating_add(counters.preferred_heap_steals),
remote
.saturating_add(counters.remote_fast_steals)
.saturating_add(counters.remote_heap_steals),
)
});
let total = preferred.saturating_add(remote);
if total == 0 {
return None;
}
let pct = remote.saturating_mul(100).saturating_add(total / 2) / total;
Some(u8::try_from(pct.min(100)).expect("remote steal ratio should fit in u8"))
}
#[cfg(any(test, feature = "test-internals"))]
pub fn set_scheduler_evidence_window(&mut self, sample_window: usize) {
let collector = (sample_window > 0)
.then(|| Arc::new(Mutex::new(SchedulerEvidenceCollector::new(sample_window))));
self.scheduler_evidence.clone_from(&collector);
for worker in &mut self.workers {
worker.scheduler_evidence.clone_from(&collector);
}
}
#[must_use]
pub fn scheduler_evidence_artifact(
&self,
run_label: &str,
workload_class: SchedulerWorkloadClass,
memory_budget_gib: usize,
) -> Option<SchedulerEvidenceArtifact> {
if self.workers.is_empty() {
return None;
}
let collector = self.scheduler_evidence.as_ref()?;
let remote_steal_ratio_pct = self.scheduler_evidence_remote_steal_ratio_pct();
let collector = collector.lock();
let sample_window = collector.sample_window();
let (wake_to_run_samples, queue_residency_samples, ready_backlog_samples, cancel_samples) =
collector.sample_counts();
let metrics = collector.snapshot_metrics(remote_steal_ratio_pct);
drop(collector);
let cancel_streak_limit = self
.workers
.first()
.map_or(DEFAULT_CANCEL_STREAK_LIMIT, |worker| {
worker.cancel_streak_limit
});
Some(SchedulerEvidenceArtifact {
schema_version: SCHEDULER_EVIDENCE_SCHEMA_VERSION.to_string(),
run_label: run_label.to_string(),
workload_class,
topology: SchedulerTopologyDescriptor {
worker_threads: self.workers.len(),
cohort_count: self.cohort_count.max(1),
memory_budget_gib,
},
current_knobs: SchedulerKnobProfile {
worker_threads: self.workers.len(),
steal_batch_size: self.steal_batch_size,
cancel_streak_limit,
global_queue_limit: self.global_queue_limit,
parking_enabled: self.enable_parking,
},
metrics,
notes: vec![
"runtime_capture".to_string(),
format!("placement_mode={}", self.placement_mode.as_str()),
format!("sample_window={sample_window}"),
format!(
"sample_counts=wake_to_run:{wake_to_run_samples},queue_residency:{queue_residency_samples},ready_backlog:{ready_backlog_samples},cancel_debt:{cancel_samples}"
),
],
})
}
#[doc(hidden)]
#[cfg(any(test, feature = "test-internals"))]
pub fn worker_mut_for_test(&mut self, worker_id: usize) -> &mut ThreeLaneWorker {
&mut self.workers[worker_id]
}
#[must_use]
pub fn global_injector(&self) -> Arc<GlobalInjector> {
self.global.clone()
}
#[inline]
fn with_task_table_ref<R, F: FnOnce(&TaskTable) -> R>(&self, f: F) -> R {
if let Some(tt) = &self.task_table {
let guard = tt.lock().unwrap_or_else(std::sync::PoisonError::into_inner);
f(&guard)
} else {
let state = self
.state
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
f(&state.tasks)
}
}
#[inline]
fn clear_task_wake_state(&self, task: TaskId) {
self.with_task_table_ref(|tt| {
if let Some(record) = tt.task(task) {
record.wake_state.clear();
}
});
}
pub fn inject_cancel(&self, task: TaskId, priority: u8) {
let (is_local, pinned_worker) = self.with_task_table_ref(|tt| {
tt.task(task).map_or((false, None), |record| {
if record.is_local() {
record.wake_state.notify();
}
(record.is_local(), record.pinned_worker())
})
});
if is_local {
if let Some(worker_id) = pinned_worker {
if let Some(local) = self.local_schedulers.get(worker_id) {
let mut local_guard = local.lock();
if let Some(local_ready) = self.local_ready.get(worker_id) {
local_ready.lock().tombstone(task);
}
local_guard.move_to_cancel_lane(task, priority);
drop(local_guard);
Self::contain_publication_effect(|| {
self.record_scheduler_evidence_enqueue(task);
});
if let Some(parker) = self.parkers.get(worker_id) {
Self::contain_publication_effect(|| parker.unpark());
}
return;
}
}
if schedule_cancel_on_current_local(task, priority) {
Self::contain_publication_effect(|| {
self.record_scheduler_evidence_enqueue(task);
});
return;
}
self.clear_task_wake_state(task);
debug_assert!(
false,
"Attempted to inject_cancel local task {task:?} without owner worker"
);
Self::contain_publication_effect(|| {
error!(
?task,
"inject_cancel: cannot route local task to owner worker, cancel skipped"
);
});
return;
}
self.with_task_table_ref(|tt| {
if let Some(record) = tt.task(task) {
record.wake_state.notify();
}
self.global.remove_timed(task);
self.global.inject_cancel(task, priority);
});
Self::contain_publication_effect(|| self.record_scheduler_evidence_enqueue(task));
Self::contain_publication_effect(|| self.wake_one());
}
pub fn inject_timed(&self, task: TaskId, deadline: Time) {
let injected = self.with_task_table_ref(|tt| {
match tt.task(task) {
Some(record) => {
if record.wake_state.notify() {
self.global.inject_timed(task, deadline);
true
} else {
false
}
}
None => {
self.global.inject_timed(task, deadline);
true
}
}
});
if injected {
self.record_scheduler_evidence_enqueue(task);
self.wake_one();
}
}
#[inline]
fn inject_global_ready_checked(&self, task: TaskId, priority: u8) {
let ready_count_before = self.global.ready_count();
self.global.inject_ready(task, priority);
self.finish_global_ready_publication(task, priority, ready_count_before);
}
pub fn inject_ready(&self, task: TaskId, priority: u8) {
let (injected, is_local, ready_count_before) = self.with_task_table_ref(|tt| {
match tt.task(task) {
Some(record) => {
let is_local = record.is_local();
if is_local {
(false, true, 0)
} else if record.wake_state.notify() || self.global.remove_timed(task) {
let ready_count_before = self.global.ready_count();
self.global.inject_ready(task, priority);
(true, false, ready_count_before)
} else {
(false, false, 0)
}
}
None => {
let ready_count_before = self.global.ready_count();
self.global.inject_ready(task, priority);
(true, false, ready_count_before)
}
}
});
debug_assert!(
!is_local,
"Attempted to globally inject local task {task:?}. Local tasks must be scheduled on their owner thread."
);
if is_local {
error!(
?task,
"inject_ready: refusing to globally inject local task, scheduling skipped"
);
return;
}
if injected {
self.finish_global_ready_publication(task, priority, ready_count_before);
Self::contain_publication_effect(|| {
trace!(
?task,
priority, "inject_ready: task injected into global ready queue"
);
});
} else {
Self::contain_publication_effect(|| {
trace!(
?task,
priority, "inject_ready: task NOT scheduled (should_schedule=false)"
);
});
}
}
#[inline]
pub fn spawn(&self, task: TaskId, priority: u8) {
self.schedule_internal(task, priority, ScheduleIntent::Spawn);
}
#[inline]
pub fn wake(&self, task: TaskId, priority: u8) {
self.schedule_internal(task, priority, ScheduleIntent::Wake);
}
fn schedule_internal(&self, task: TaskId, priority: u8, intent: ScheduleIntent) {
let (should_schedule, is_local, pinned_worker) = self.with_task_table_ref(|tt| {
tt.task(task).map_or((true, false, None), |record| {
(
record.wake_state.notify(),
record.is_local(),
record.pinned_worker(),
)
})
});
if !should_schedule {
return;
}
if is_local {
let current_worker = current_worker_id();
let is_pinned_here = match (pinned_worker, current_worker) {
(Some(pw), Some(cw)) => pw == cw,
(None, Some(_)) => true,
_ => false,
};
if is_pinned_here && schedule_local_task(task) {
self.record_scheduler_evidence_enqueue(task);
return;
}
if let Some(worker_id) = pinned_worker {
if let Some(queue) = self.local_ready.get(worker_id) {
queue.lock().push_back(task);
self.record_scheduler_evidence_enqueue(task);
self.coordinator.wake_worker(worker_id);
return;
}
}
let assert_msg = intent.local_route_failure_assert(task);
let _error_msg = intent.local_route_failure_log();
self.clear_task_wake_state(task);
debug_assert!(false, "{}", assert_msg);
error!(?task, "{}", _error_msg);
return;
}
if schedule_on_current_local(task, priority) {
self.record_scheduler_evidence_enqueue(task);
return;
}
self.inject_global_ready_checked(task, priority);
}
#[inline]
fn wake_one(&self) {
self.coordinator.wake_one();
}
pub fn wake_all(&self) {
self.coordinator.wake_all();
}
pub fn attach_spawn_mailbox(
&mut self,
mailbox: Arc<crate::runtime::spawn_mailbox::SpawnMailbox>,
) {
for worker in &mut self.workers {
worker.spawn_mailbox = Some(Arc::clone(&mailbox));
}
self.spawn_mailbox = Some(mailbox);
}
pub fn notify_spawn_enqueued(&self) {
self.coordinator.wake_one();
}
#[must_use]
pub fn spawn_enqueued_notifier(&self) -> Arc<dyn Fn() + Send + Sync> {
let coordinator = Arc::clone(&self.coordinator);
Arc::new(move || coordinator.wake_one())
}
pub fn take_workers(&mut self) -> Vec<ThreeLaneWorker> {
std::mem::take(&mut self.workers).into_vec()
}
pub fn shutdown(&self) {
self.shutdown.store(true, Ordering::Release);
self.wake_all();
}
#[must_use]
pub fn is_shutdown(&self) -> bool {
self.shutdown.load(Ordering::Acquire)
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum StealerLocality {
SameCohort,
CrossCohort,
}
impl StealerLocality {
#[inline]
const fn from_same_cohort(same_cohort: bool) -> Self {
if same_cohort {
Self::SameCohort
} else {
Self::CrossCohort
}
}
#[inline]
const fn is_same_cohort(self) -> bool {
matches!(self, Self::SameCohort)
}
}
#[derive(Debug)]
pub struct ThreeLaneWorker {
pub id: WorkerId,
pub local: Arc<Mutex<PriorityScheduler>>,
pub stealers: SmallVec<[Arc<Mutex<PriorityScheduler>>; 16]>,
preferred_heap_stealer_count: usize,
heap_stealer_locality: SmallVec<[StealerLocality; 16]>,
pub fast_queue: LocalQueue,
global_ready_buffer: Vec<PriorityTask>,
fast_stealers: SmallVec<[local_queue::Stealer; 16]>,
preferred_fast_stealer_count: usize,
fast_stealer_locality: SmallVec<[StealerLocality; 16]>,
local_ready: Arc<LocalReadyQueue>,
all_local_ready: SmallVec<[Arc<LocalReadyQueue>; 16]>,
all_local_schedulers: SmallVec<[Arc<Mutex<PriorityScheduler>>; 16]>,
pub global: Arc<GlobalInjector>,
pub state: Arc<ContendedMutex<RuntimeState>>,
pending_cancel_dispatch_ready: Arc<AtomicBool>,
pub task_table: Option<Arc<ContendedMutex<TaskTable>>>,
pub parker: Parker,
pub(crate) coordinator: Arc<WorkerCoordinator>,
pub(crate) spawn_mailbox: Option<Arc<crate::runtime::spawn_mailbox::SpawnMailbox>>,
pub rng: DetRng,
pub shutdown: Arc<AtomicBool>,
pub io_driver: Option<IoDriverHandle>,
pub timer_driver: Option<TimerDriverHandle>,
steal_buffer: Vec<(TaskId, u8)>,
steal_batch_size: usize,
enable_parking: bool,
empty_backoff: u32,
cancel_streak: usize,
ready_dispatch_streak: usize,
browser_ready_handoff_limit: usize,
cancel_streak_limit: usize,
governor: Option<LyapunovGovernor>,
cached_suggestion: SchedulingSuggestion,
steps_since_snapshot: u32,
governor_interval: u32,
preemption_metrics: PreemptionMetrics,
evidence_sink: Option<Arc<dyn crate::evidence_sink::EvidenceSink>>,
decision_contract: Option<super::decision_contract::SchedulerDecisionContract>,
decision_posterior: Option<franken_decision::Posterior>,
adaptive_cancel_policy: Option<AdaptiveCancelStreakPolicy>,
spectral_monitor: Option<SpectralHealthMonitor>,
drain_certificate: Option<ProgressCertificate>,
decision_sequence: u64,
fairness_monitor: Mutex<FairnessMonitor>,
invariant_monitor: Mutex<super::invariant_monitor::SchedulerInvariantMonitor>,
fast_queue_dispatch_streak: usize,
fast_queue_fairness_limit: usize,
timed_dispatch_streak: usize,
timed_fairness_limit: usize,
adaptive_batch_profile: Option<AdaptiveBatchSizingProfile>,
adaptive_batch_state: AdaptiveBatchRuntimeState,
steal_locality_counters: StealLocalityCounters,
scheduler_evidence: Option<Arc<Mutex<SchedulerEvidenceCollector>>>,
}
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
pub struct StealLocalityCounters {
pub preferred_fast_steals: u64,
pub remote_fast_steals: u64,
pub preferred_heap_steals: u64,
pub remote_heap_steals: u64,
}
#[derive(Debug)]
struct SchedulerEvidenceCollector {
sample_window: usize,
max_inflight: usize,
next_sequence: u64,
pending_enqueue: DetHashMap<TaskId, (u64, u64)>,
pending_wake: DetHashMap<TaskId, (u64, u64)>,
wake_order: VecDeque<(TaskId, u64)>,
enqueue_order: VecDeque<(TaskId, u64)>,
wake_to_run_samples_ns: VecDeque<u64>,
queue_residency_samples_ns: VecDeque<u64>,
ready_backlog_samples: VecDeque<usize>,
cancel_debt_samples: VecDeque<usize>,
}
impl SchedulerEvidenceCollector {
#[cfg(any(test, feature = "test-internals"))]
fn new(sample_window: usize) -> Self {
let sample_window = sample_window.max(1);
Self {
sample_window,
max_inflight: sample_window
.saturating_mul(DEFAULT_SCHEDULER_EVIDENCE_MAX_INFLIGHT_MULTIPLIER)
.max(sample_window),
next_sequence: 0,
pending_enqueue: DetHashMap::default(),
pending_wake: DetHashMap::default(),
wake_order: VecDeque::with_capacity(sample_window),
enqueue_order: VecDeque::with_capacity(sample_window),
wake_to_run_samples_ns: VecDeque::with_capacity(sample_window),
queue_residency_samples_ns: VecDeque::with_capacity(sample_window),
ready_backlog_samples: VecDeque::with_capacity(sample_window),
cancel_debt_samples: VecDeque::with_capacity(sample_window),
}
}
fn record_task_enqueue(&mut self, task_id: TaskId, timestamp_ns: u64) {
self.next_sequence = self.next_sequence.saturating_add(1);
let sequence = self.next_sequence;
self.pending_enqueue
.insert(task_id, (timestamp_ns, sequence));
self.enqueue_order.push_back((task_id, sequence));
self.pending_wake.insert(task_id, (timestamp_ns, sequence));
self.wake_order.push_back((task_id, sequence));
self.trim_pending();
}
fn record_task_dispatch(
&mut self,
task_id: TaskId,
dispatch_time_ns: u64,
ready_backlog: usize,
cancel_debt: usize,
) {
let sample_window = self.sample_window;
if let Some((enqueue_time_ns, _)) = self.pending_enqueue.remove(&task_id) {
Self::push_u64_sample(
&mut self.queue_residency_samples_ns,
dispatch_time_ns.saturating_sub(enqueue_time_ns),
sample_window,
);
}
if let Some((wake_time_ns, _)) = self.pending_wake.remove(&task_id) {
Self::push_u64_sample(
&mut self.wake_to_run_samples_ns,
dispatch_time_ns.saturating_sub(wake_time_ns),
sample_window,
);
}
Self::push_usize_sample(
&mut self.ready_backlog_samples,
ready_backlog,
sample_window,
);
Self::push_usize_sample(&mut self.cancel_debt_samples, cancel_debt, sample_window);
}
fn sample_window(&self) -> usize {
self.sample_window
}
fn sample_counts(&self) -> (usize, usize, usize, usize) {
(
self.wake_to_run_samples_ns.len(),
self.queue_residency_samples_ns.len(),
self.ready_backlog_samples.len(),
self.cancel_debt_samples.len(),
)
}
fn snapshot_metrics(&self, remote_steal_ratio_pct: Option<u8>) -> SchedulerEvidenceMetrics {
SchedulerEvidenceMetrics {
wake_to_run_p50_ns: percentile_u64(&self.wake_to_run_samples_ns, 50),
wake_to_run_p95_ns: percentile_u64(&self.wake_to_run_samples_ns, 95),
wake_to_run_p99_ns: percentile_u64(&self.wake_to_run_samples_ns, 99),
queue_residency_p50_ns: percentile_u64(&self.queue_residency_samples_ns, 50),
queue_residency_p95_ns: percentile_u64(&self.queue_residency_samples_ns, 95),
queue_residency_p99_ns: percentile_u64(&self.queue_residency_samples_ns, 99),
ready_backlog_p95: percentile_usize(&self.ready_backlog_samples, 95),
ready_backlog_p99: percentile_usize(&self.ready_backlog_samples, 99),
cancel_debt_p95: percentile_usize(&self.cancel_debt_samples, 95),
cancel_debt_p99: percentile_usize(&self.cancel_debt_samples, 99),
remote_steal_ratio_pct,
cross_cohort_wake_p99_ns: None,
}
}
fn trim_pending(&mut self) {
while self.pending_enqueue.len() > self.max_inflight {
let Some((task_id, sequence)) = self.enqueue_order.pop_front() else {
break;
};
if self
.pending_enqueue
.get(&task_id)
.is_some_and(|(_, current_sequence)| *current_sequence == sequence)
{
self.pending_enqueue.remove(&task_id);
}
}
while self.pending_wake.len() > self.max_inflight {
let Some((task_id, sequence)) = self.wake_order.pop_front() else {
break;
};
if self
.pending_wake
.get(&task_id)
.is_some_and(|(_, current_sequence)| *current_sequence == sequence)
{
self.pending_wake.remove(&task_id);
}
}
}
fn push_u64_sample(samples: &mut VecDeque<u64>, value: u64, sample_window: usize) {
if samples.len() == sample_window {
samples.pop_front();
}
samples.push_back(value);
}
fn push_usize_sample(samples: &mut VecDeque<usize>, value: usize, sample_window: usize) {
if samples.len() == sample_window {
samples.pop_front();
}
samples.push_back(value);
}
}
fn percentile_index(len: usize, percentile: usize) -> usize {
debug_assert!(len > 0);
let percentile = percentile.clamp(1, 100);
percentile
.saturating_mul(len)
.div_ceil(100)
.saturating_sub(1)
.min(len.saturating_sub(1))
}
fn percentile_u64(samples: &VecDeque<u64>, percentile: usize) -> u64 {
if samples.is_empty() {
return 0;
}
let mut values = samples.iter().copied().collect::<Vec<_>>();
values.sort_unstable();
values[percentile_index(values.len(), percentile)]
}
fn percentile_usize(samples: &VecDeque<usize>, percentile: usize) -> usize {
if samples.is_empty() {
return 0;
}
let mut values = samples.iter().copied().collect::<Vec<_>>();
values.sort_unstable();
values[percentile_index(values.len(), percentile)]
}
#[derive(Debug, Clone)]
struct WaiterWakeMetadata {
priority: u8,
is_local: bool,
pinned_worker: Option<WorkerId>,
wake_state: Arc<crate::record::task::TaskWakeState>,
notified: bool,
}
#[derive(Debug, Clone, Default)]
pub struct PreemptionMetrics {
pub cancel_dispatches: u64,
pub timed_dispatches: u64,
pub ready_dispatches: u64,
pub browser_ready_handoff_yields: u64,
pub fairness_yields: u64,
pub max_ready_dispatch_stall: usize,
pub max_timed_dispatch_stall: usize,
pub ready_priority_inversions: u64,
pub max_ready_priority_inversion_gap: u8,
pub max_cancel_streak: usize,
pub fallback_cancel_dispatches: u64,
pub base_limit_exceedances: u64,
pub effective_limit_exceedances: u64,
pub max_effective_limit_observed: usize,
pub adaptive_epochs: u64,
pub adaptive_current_limit: usize,
pub adaptive_reward_ema: f64,
pub adaptive_e_value: f64,
pub backoff_parks_total: u64,
pub backoff_timeout_parks_total: u64,
pub backoff_indefinite_parks: u64,
pub backoff_timeout_nanos_total: u64,
pub short_wait_le_5ms: u64,
pub follower_shared_deadline_ignored: u64,
pub follower_timeout_parks: u64,
pub follower_indefinite_parks: u64,
pub follower_short_wait_skip_le_5ms: u64,
pub global_ready_batch_drains: u64,
pub global_ready_batch_tasks: u64,
pub adaptive_batch_scale_up_events: u64,
pub adaptive_batch_cancel_floor_hits: u64,
pub adaptive_batch_cooldown_holds: u64,
pub adaptive_batch_max_selected: usize,
}
impl PreemptionMetrics {
const RATIO_BPS_SCALE: u64 = 10_000;
#[inline]
fn ratio_bps(numerator: u64, denominator: u64) -> u16 {
if denominator == 0 {
return 0;
}
let raw = numerator
.saturating_mul(Self::RATIO_BPS_SCALE)
.saturating_div(denominator)
.min(Self::RATIO_BPS_SCALE);
raw as u16
}
#[must_use]
pub fn avg_timeout_park_nanos(&self) -> u64 {
if self.backoff_timeout_parks_total == 0 {
return 0;
}
self.backoff_timeout_nanos_total
.saturating_div(self.backoff_timeout_parks_total)
}
#[must_use]
pub fn short_wait_ratio_bps(&self) -> u16 {
Self::ratio_bps(self.short_wait_le_5ms, self.backoff_timeout_parks_total)
}
#[must_use]
pub fn follower_short_wait_avoidance_bps(&self) -> u16 {
let opportunities = self
.follower_short_wait_skip_le_5ms
.saturating_add(self.follower_timeout_parks);
Self::ratio_bps(self.follower_short_wait_skip_le_5ms, opportunities)
}
#[must_use]
pub fn max_non_cancel_dispatch_stall(&self) -> usize {
self.max_ready_dispatch_stall
.max(self.max_timed_dispatch_stall)
}
}
#[derive(Debug, Clone)]
pub struct FairnessConfig {
pub starvation_threshold_ns: u64,
pub analysis_window_size: usize,
pub priority_inversion_threshold: u8,
pub max_tracked_tasks: usize,
pub enable_per_task_tracking: bool,
}
impl Default for FairnessConfig {
fn default() -> Self {
Self {
starvation_threshold_ns: 100_000_000, analysis_window_size: 1000,
priority_inversion_threshold: 5,
max_tracked_tasks: 10_000,
enable_per_task_tracking: true,
}
}
}
#[derive(Debug, Clone)]
struct TaskStarvationInfo {
task_id: TaskId,
priority: u8,
enqueue_time_ns: u64,
skip_count: u32,
last_skip_time_ns: u64,
current_lane: u8,
total_wait_time_ns: u64,
}
impl TaskStarvationInfo {
fn new(task_id: TaskId, priority: u8, current_time_ns: u64, lane: u8) -> Self {
Self {
task_id,
priority,
enqueue_time_ns: current_time_ns,
skip_count: 0,
last_skip_time_ns: 0,
current_lane: lane,
total_wait_time_ns: 0,
}
}
fn refresh_queue_membership(&mut self, priority: u8, current_time_ns: u64, lane: u8) {
self.priority = priority;
self.current_lane = lane;
self.total_wait_time_ns = self
.total_wait_time_ns
.max(self.current_wait_time_ns(current_time_ns));
}
fn record_skip(&mut self, current_time_ns: u64) {
self.skip_count = self.skip_count.saturating_add(1);
self.last_skip_time_ns = current_time_ns;
self.total_wait_time_ns = self.current_wait_time_ns(current_time_ns);
}
fn current_wait_time_ns(&self, current_time_ns: u64) -> u64 {
current_time_ns.saturating_sub(self.enqueue_time_ns)
}
fn is_starved(&self, threshold_ns: u64, current_time_ns: u64) -> bool {
self.current_wait_time_ns(current_time_ns) >= threshold_ns
}
}
#[derive(Debug, Clone)]
struct PriorityInversionEvent {
blocked_task_id: TaskId,
blocked_priority: u8,
executing_task_id: TaskId,
executing_priority: u8,
timestamp_ns: u64,
duration_ns: u64,
}
#[derive(Debug, Clone)]
struct StarvationAnalysisWindow {
events: Vec<u64>,
write_pos: usize,
total_events: u64,
size: usize,
}
impl StarvationAnalysisWindow {
fn new(size: usize) -> Self {
Self {
events: vec![0; size.max(1)],
write_pos: 0,
size: size.max(1),
total_events: 0,
}
}
fn record_event(&mut self, timestamp_ns: u64) {
self.events[self.write_pos] = timestamp_ns;
self.write_pos = (self.write_pos + 1) % self.size;
self.total_events = self.total_events.saturating_add(1);
}
fn events_in_window(&self, window_duration_ns: u64, current_time_ns: u64) -> u32 {
let threshold_time = current_time_ns.saturating_sub(window_duration_ns);
let mut count = 0;
let recorded_events = usize::try_from(self.total_events)
.unwrap_or(usize::MAX)
.min(self.size);
for &event_time in self.events.iter().take(recorded_events) {
if event_time >= threshold_time && event_time <= current_time_ns {
count += 1;
}
}
count
}
fn is_pattern_detected(
&self,
min_events: u32,
window_duration_ns: u64,
current_time_ns: u64,
) -> bool {
self.events_in_window(window_duration_ns, current_time_ns) >= min_events
}
}
#[derive(Debug)]
pub struct FairnessMonitor {
config: FairnessConfig,
tracked_tasks: BTreeMap<TaskId, TaskStarvationInfo>,
priority_inversions: Vec<PriorityInversionEvent>,
starvation_window: StarvationAnalysisWindow,
total_starvation_events: u64,
total_priority_inversions: u64,
max_task_wait_time_ns: u64,
last_cleanup_time_ns: u64,
}
impl FairnessMonitor {
#[must_use]
pub fn new(config: FairnessConfig) -> Self {
let window_size = config.analysis_window_size;
Self {
config,
tracked_tasks: BTreeMap::new(),
priority_inversions: Vec::new(),
starvation_window: StarvationAnalysisWindow::new(window_size),
total_starvation_events: 0,
total_priority_inversions: 0,
max_task_wait_time_ns: 0,
last_cleanup_time_ns: 0,
}
}
#[must_use]
pub fn with_defaults() -> Self {
Self::new(FairnessConfig::default())
}
pub fn record_task_enqueue(
&mut self,
task_id: TaskId,
priority: u8,
current_time_ns: u64,
lane: u8,
) {
if !self.config.enable_per_task_tracking {
return;
}
if let Some(info) = self.tracked_tasks.get_mut(&task_id) {
info.refresh_queue_membership(priority, current_time_ns, lane);
return;
}
self.cleanup_if_needed(current_time_ns);
if self.tracked_tasks.len() >= self.config.max_tracked_tasks {
if let Some((oldest_task_id, _)) = self
.tracked_tasks
.iter()
.min_by_key(|(id, info)| (info.enqueue_time_ns, **id))
.map(|(id, info)| (*id, info.clone()))
{
self.tracked_tasks.remove(&oldest_task_id);
}
}
let info = TaskStarvationInfo::new(task_id, priority, current_time_ns, lane);
self.tracked_tasks.insert(task_id, info);
}
pub fn record_task_dispatch(&mut self, task_id: TaskId, current_time_ns: u64) -> Option<u64> {
if let Some(info) = self.tracked_tasks.remove(&task_id) {
let wait_time = info.current_wait_time_ns(current_time_ns);
if wait_time > self.max_task_wait_time_ns {
self.max_task_wait_time_ns = wait_time;
}
Some(wait_time)
} else {
None
}
}
pub fn record_task_skip(
&mut self,
skipped_task_id: TaskId,
executing_task_id: TaskId,
executing_priority: u8,
current_time_ns: u64,
) {
let (should_record_starvation, should_record_inversion, blocked_priority) = {
if let Some(info) = self.tracked_tasks.get_mut(&skipped_task_id) {
info.record_skip(current_time_ns);
let is_starved =
info.is_starved(self.config.starvation_threshold_ns, current_time_ns);
let is_inversion = info.priority > executing_priority;
let priority = info.priority;
(is_starved, is_inversion, priority)
} else {
(false, false, 0)
}
};
if should_record_starvation {
self.record_starvation_event(current_time_ns);
}
if should_record_inversion {
self.record_priority_inversion(
skipped_task_id,
blocked_priority,
executing_task_id,
executing_priority,
current_time_ns,
);
}
}
fn record_starvation_event(&mut self, timestamp_ns: u64) {
self.total_starvation_events = self.total_starvation_events.saturating_add(1);
self.starvation_window.record_event(timestamp_ns);
}
fn record_priority_inversion(
&mut self,
blocked_task: TaskId,
blocked_priority: u8,
executing_task: TaskId,
executing_priority: u8,
timestamp_ns: u64,
) {
self.total_priority_inversions = self.total_priority_inversions.saturating_add(1);
let inversion = PriorityInversionEvent {
blocked_task_id: blocked_task,
blocked_priority,
executing_task_id: executing_task,
executing_priority,
timestamp_ns,
duration_ns: 0, };
self.priority_inversions.push(inversion);
const MAX_TRACKED_INVERSIONS: usize = 1000;
if self.priority_inversions.len() > MAX_TRACKED_INVERSIONS {
self.priority_inversions
.drain(0..self.priority_inversions.len() - MAX_TRACKED_INVERSIONS);
}
}
#[must_use]
pub fn detect_starvation_pattern(&self, current_time_ns: u64) -> bool {
const PATTERN_WINDOW_NS: u64 = 1_000_000_000; const MIN_EVENTS_FOR_PATTERN: u32 = 10;
self.starvation_window.is_pattern_detected(
MIN_EVENTS_FOR_PATTERN,
PATTERN_WINDOW_NS,
current_time_ns,
)
}
#[must_use]
pub fn count_starved_tasks(&self, current_time_ns: u64) -> u32 {
self.tracked_tasks
.values()
.filter(|info| info.is_starved(self.config.starvation_threshold_ns, current_time_ns))
.count() as u32
}
#[must_use]
pub fn starvation_stats(&self, current_time_ns: u64) -> StarvationStats {
let currently_starved = self.count_starved_tasks(current_time_ns);
let total_tracked_wait_time_ns = self
.tracked_tasks
.values()
.map(|info| {
info.total_wait_time_ns
.max(info.current_wait_time_ns(current_time_ns))
})
.sum::<u64>();
let avg_wait_time_ns = if self.tracked_tasks.is_empty() {
0
} else {
total_tracked_wait_time_ns / self.tracked_tasks.len() as u64
};
let oldest_tracked_task = self
.tracked_tasks
.values()
.max_by_key(|info| info.current_wait_time_ns(current_time_ns))
.map(|info| StarvedTaskSummary {
task_id: info.task_id,
priority: info.priority,
current_lane: info.current_lane,
skip_count: info.skip_count,
wait_time_ns: info.current_wait_time_ns(current_time_ns),
total_wait_time_ns: info
.total_wait_time_ns
.max(info.current_wait_time_ns(current_time_ns)),
});
let latest_priority_inversion =
self.priority_inversions
.last()
.map(|event| PriorityInversionSummary {
blocked_task_id: event.blocked_task_id,
blocked_priority: event.blocked_priority,
executing_task_id: event.executing_task_id,
executing_priority: event.executing_priority,
priority_gap: event
.blocked_priority
.saturating_sub(event.executing_priority),
timestamp_ns: event.timestamp_ns,
duration_ns: event.duration_ns,
});
let max_priority_inversion_gap = self
.priority_inversions
.iter()
.map(|event| {
event
.blocked_priority
.saturating_sub(event.executing_priority)
})
.max()
.unwrap_or(0);
StarvationStats {
total_starvation_events: self.total_starvation_events,
currently_starved_tasks: currently_starved,
max_task_wait_time_ns: self.max_task_wait_time_ns,
avg_task_wait_time_ns: avg_wait_time_ns,
total_priority_inversions: self.total_priority_inversions,
tracked_tasks_count: self.tracked_tasks.len() as u32,
pattern_detected: self.detect_starvation_pattern(current_time_ns),
total_tracked_wait_time_ns,
oldest_tracked_task,
max_priority_inversion_gap,
latest_priority_inversion,
}
}
fn cleanup_if_needed(&mut self, current_time_ns: u64) {
const CLEANUP_INTERVAL_NS: u64 = 60_000_000_000; const MAX_TASK_AGE_NS: u64 = 300_000_000_000;
if current_time_ns.saturating_sub(self.last_cleanup_time_ns) < CLEANUP_INTERVAL_NS {
return;
}
self.last_cleanup_time_ns = current_time_ns;
let cutoff_time = current_time_ns.saturating_sub(MAX_TASK_AGE_NS);
self.tracked_tasks
.retain(|_, info| info.enqueue_time_ns >= cutoff_time);
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct StarvedTaskSummary {
pub task_id: TaskId,
pub priority: u8,
pub current_lane: u8,
pub skip_count: u32,
pub wait_time_ns: u64,
pub total_wait_time_ns: u64,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct PriorityInversionSummary {
pub blocked_task_id: TaskId,
pub blocked_priority: u8,
pub executing_task_id: TaskId,
pub executing_priority: u8,
pub priority_gap: u8,
pub timestamp_ns: u64,
pub duration_ns: u64,
}
#[derive(Debug, Clone, Default)]
pub struct StarvationStats {
pub total_starvation_events: u64,
pub currently_starved_tasks: u32,
pub max_task_wait_time_ns: u64,
pub avg_task_wait_time_ns: u64,
pub total_priority_inversions: u64,
pub tracked_tasks_count: u32,
pub pattern_detected: bool,
pub total_tracked_wait_time_ns: u64,
pub oldest_tracked_task: Option<StarvedTaskSummary>,
pub max_priority_inversion_gap: u8,
pub latest_priority_inversion: Option<PriorityInversionSummary>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct PreemptionFairnessCertificate {
pub base_limit: usize,
pub effective_limit: usize,
pub observed_max_cancel_streak: usize,
pub cancel_dispatches: u64,
pub timed_dispatches: u64,
pub ready_dispatches: u64,
pub fairness_yields: u64,
pub observed_max_ready_stall_steps: usize,
pub observed_max_timed_stall_steps: usize,
pub ready_priority_inversions: u64,
pub max_ready_priority_inversion_gap: u8,
pub fallback_cancel_dispatches: u64,
pub base_limit_exceedances: u64,
pub effective_limit_exceedances: u64,
pub adaptive_enabled: bool,
pub adaptive_current_limit: usize,
}
impl PreemptionFairnessCertificate {
#[must_use]
pub fn ready_stall_bound_steps(&self) -> usize {
self.effective_limit.saturating_add(1)
}
#[must_use]
pub fn observed_non_cancel_stall_steps(&self) -> usize {
self.observed_max_ready_stall_steps
.max(self.observed_max_timed_stall_steps)
}
#[must_use]
pub fn invariant_holds(&self) -> bool {
self.effective_limit_exceedances == 0
&& self.observed_max_cancel_streak <= self.effective_limit
&& self.ready_priority_inversions == 0
}
#[must_use]
pub fn witness_hash(&self) -> u64 {
use std::hash::{Hash, Hasher};
let mut h = DetHasher::default();
self.base_limit.hash(&mut h);
self.effective_limit.hash(&mut h);
self.observed_max_cancel_streak.hash(&mut h);
self.cancel_dispatches.hash(&mut h);
self.timed_dispatches.hash(&mut h);
self.ready_dispatches.hash(&mut h);
self.fairness_yields.hash(&mut h);
self.observed_max_ready_stall_steps.hash(&mut h);
self.observed_max_timed_stall_steps.hash(&mut h);
self.ready_priority_inversions.hash(&mut h);
self.max_ready_priority_inversion_gap.hash(&mut h);
self.fallback_cancel_dispatches.hash(&mut h);
self.base_limit_exceedances.hash(&mut h);
self.effective_limit_exceedances.hash(&mut h);
self.adaptive_enabled.hash(&mut h);
self.adaptive_current_limit.hash(&mut h);
h.finish()
}
}
static THREE_LANE_TIME_FALLBACK_WARNED: AtomicBool = AtomicBool::new(false);
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
struct DeferredCancelLanePublication {
priority: u8,
wake_target: DeferredCancelWakeTarget,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum DeferredCancelWakeTarget {
AnyWorker,
PinnedWorker(usize),
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum DeferredCancelLaneError {
MissingPinnedWorker,
PinnedWorkerUnavailable(usize),
}
type ReadyFinalizerPublication = (
smallvec::SmallVec<[TaskId; 2]>,
smallvec::SmallVec<[crate::runtime::state::TaskSpawnEffects; 2]>,
);
enum PolledCompletion {
Ready(crate::types::Outcome<(), crate::error::Error>),
Panicked(crate::types::Outcome<(), crate::error::Error>),
}
struct PolledCompletionArtifacts {
completion_observer: crate::runtime::state::TaskCompletionObserver,
cancel_wakes: crate::types::task_context::CancelWakeEffects,
detached_record: Option<crate::record::task::TaskRecord>,
finalizer_publication: ReadyFinalizerPublication,
}
impl PolledCompletionArtifacts {
fn dispatch_post_lock(self, worker: &ThreeLaneWorker) {
ThreeLaneWorker::retire_detached_task_record(self.detached_record);
worker.finish_ready_finalizer_publication(self.finalizer_publication);
self.completion_observer.dispatch();
self.cancel_wakes.dispatch();
}
}
struct UnwindCompletionArtifacts {
cancel_waker_retirements: crate::runtime::state::TaskCompletionRetirements,
detached_record: Option<crate::record::task::TaskRecord>,
finalizer_publication: ReadyFinalizerPublication,
}
impl UnwindCompletionArtifacts {
fn dispatch_post_lock(self, worker: &ThreeLaneWorker) {
worker.finish_ready_finalizer_publication(self.finalizer_publication);
self.cancel_waker_retirements.retire();
ThreeLaneWorker::retire_detached_task_record(self.detached_record);
}
}
impl ThreeLaneWorker {
#[inline]
fn current_time_ns(&self) -> u64 {
if let Some(timer) = self.timer_driver.as_ref() {
return timer.now().as_nanos();
}
if !THREE_LANE_TIME_FALLBACK_WARNED.swap(true, Ordering::Relaxed) {
crate::tracing_compat::warn!(
target: "asupersync::runtime::scheduler::three_lane",
"br-asupersync-9nn568: ThreeLaneWorker has no TimerDriverHandle attached; \
FairnessMonitor falling back to wall_now() for current_time_ns. Replay \
determinism in the lab runtime requires a timer driver."
);
}
crate::time::wall_now().as_nanos()
}
#[inline]
fn record_scheduler_evidence_enqueue_at(&self, task: TaskId, timestamp_ns: u64) {
let Some(collector) = &self.scheduler_evidence else {
return;
};
collector.lock().record_task_enqueue(task, timestamp_ns);
}
#[inline]
fn record_scheduler_evidence_enqueue(&self, task: TaskId) {
self.record_scheduler_evidence_enqueue_at(task, self.current_time_ns());
}
pub fn with_fairness_monitor<T>(&self, f: impl FnOnce(&FairnessMonitor) -> T) -> T {
f(&self.fairness_monitor.lock())
}
#[must_use]
pub fn starvation_stats(&self) -> StarvationStats {
let current_time = self.current_time_ns();
self.fairness_monitor.lock().starvation_stats(current_time)
}
#[must_use]
pub fn invariant_stats(&self) -> super::invariant_monitor::InvariantStats {
self.invariant_monitor.lock().stats()
}
#[must_use]
pub fn invariant_violations(
&self,
) -> std::collections::VecDeque<super::invariant_monitor::InvariantViolation> {
self.invariant_monitor.lock().violations().clone()
}
pub fn verify_scheduler_invariants(&mut self) {
if !self.invariant_monitor.lock().is_enabled() {
return;
}
let current_time = Time::from_nanos(self.current_time_ns());
{
let local_ready_guard = self.local_ready.lock();
let local_ready_tasks: Vec<_> = local_ready_guard.snapshot();
let ready_snapshot = super::invariant_monitor::QueueSnapshot {
name: "local_ready_queue".to_string(),
reported_depth: local_ready_tasks.len(),
actual_tasks: local_ready_tasks,
priority_range: if local_ready_guard.is_empty() {
None
} else {
Some((0, 255)) },
time_range: Some((current_time, current_time)), };
drop(local_ready_guard);
self.invariant_monitor
.lock()
.verify_queue_consistency(&ready_snapshot, current_time);
}
let fast_queue_tasks = self.fast_queue.snapshot_tasks();
let fast_snapshot = super::invariant_monitor::QueueSnapshot {
name: "fast_queue".to_string(),
reported_depth: fast_queue_tasks.len(),
actual_tasks: fast_queue_tasks.to_vec(),
priority_range: None,
time_range: Some((current_time, current_time)),
};
self.invariant_monitor
.lock()
.verify_queue_consistency(&fast_snapshot, current_time);
}
pub fn record_task_completion(&mut self, task: TaskId) {
if !self.invariant_monitor.lock().is_enabled() {
return;
}
let current_time = Time::from_nanos(self.current_time_ns());
self.invariant_monitor
.lock()
.record_task_complete(task, self.id, current_time);
}
pub fn record_task_cancellation(&mut self, task: TaskId) {
if !self.invariant_monitor.lock().is_enabled() {
return;
}
let current_time = Time::from_nanos(self.current_time_ns());
self.invariant_monitor
.lock()
.record_task_cancel(task, current_time);
}
#[inline]
fn with_task_table<R, F: FnOnce(&mut TaskTable) -> R>(&self, f: F) -> R {
if let Some(tt) = &self.task_table {
let mut guard = tt.lock().unwrap_or_else(std::sync::PoisonError::into_inner);
f(&mut guard)
} else {
let mut state = self
.state
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
f(&mut state.tasks)
}
}
#[inline]
fn with_task_table_ref<R, F: FnOnce(&TaskTable) -> R>(&self, f: F) -> R {
if let Some(tt) = &self.task_table {
let guard = tt.lock().unwrap_or_else(std::sync::PoisonError::into_inner);
f(&guard)
} else {
let state = self
.state
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
f(&state.tasks)
}
}
#[must_use]
pub fn preemption_metrics(&self) -> &PreemptionMetrics {
&self.preemption_metrics
}
#[must_use]
pub fn steal_locality_counters(&self) -> StealLocalityCounters {
self.steal_locality_counters
}
#[must_use]
pub fn apply_changepoint_detection_to_adaptive_cancel_streak(
&mut self,
detection: crate::runtime::changepoint::ChangePointDetection,
) -> bool {
if matches!(
detection.series,
crate::runtime::changepoint::RuntimeMetricSeries::Custom(_)
) {
return false;
}
let Some(policy) = self.adaptive_cancel_policy.as_mut() else {
return false;
};
policy.reset_to_priors();
self.preemption_metrics.adaptive_epochs = policy.epoch_count;
self.preemption_metrics.adaptive_current_limit = policy.current_limit();
self.preemption_metrics.adaptive_reward_ema = policy.reward_ema;
self.preemption_metrics.adaptive_e_value = policy.e_value();
trace!(
worker_id = self.id,
series = detection.series.as_str(),
detector = ?detection.detector,
direction = ?detection.direction,
sample_index = detection.sample_index,
adaptive_limit = self.preemption_metrics.adaptive_current_limit,
"changepoint reset adaptive cancel-streak policy to priors"
);
true
}
#[must_use]
pub fn preemption_fairness_certificate(&self) -> PreemptionFairnessCertificate {
let adaptive_current_limit = self.adaptive_cancel_policy.as_ref().map_or(
self.cancel_streak_limit,
AdaptiveCancelStreakPolicy::current_limit,
);
let effective_limit = self
.preemption_metrics
.max_effective_limit_observed
.max(adaptive_current_limit)
.max(1);
PreemptionFairnessCertificate {
base_limit: adaptive_current_limit,
effective_limit,
observed_max_cancel_streak: self.preemption_metrics.max_cancel_streak,
cancel_dispatches: self.preemption_metrics.cancel_dispatches,
timed_dispatches: self.preemption_metrics.timed_dispatches,
ready_dispatches: self.preemption_metrics.ready_dispatches,
fairness_yields: self.preemption_metrics.fairness_yields,
observed_max_ready_stall_steps: self.preemption_metrics.max_ready_dispatch_stall,
observed_max_timed_stall_steps: self.preemption_metrics.max_timed_dispatch_stall,
ready_priority_inversions: self.preemption_metrics.ready_priority_inversions,
max_ready_priority_inversion_gap: self
.preemption_metrics
.max_ready_priority_inversion_gap,
fallback_cancel_dispatches: self.preemption_metrics.fallback_cancel_dispatches,
base_limit_exceedances: self.preemption_metrics.base_limit_exceedances,
effective_limit_exceedances: self.preemption_metrics.effective_limit_exceedances,
adaptive_enabled: self.adaptive_cancel_policy.is_some(),
adaptive_current_limit,
}
}
pub fn set_evidence_sink(&mut self, sink: Arc<dyn crate::evidence_sink::EvidenceSink>) {
self.evidence_sink = Some(sink);
}
#[cfg(any(test, feature = "test-internals"))]
pub fn set_cached_suggestion(&mut self, suggestion: SchedulingSuggestion) {
self.cached_suggestion = suggestion;
}
#[cfg(any(test, feature = "test-internals"))]
pub fn disable_decision_contract_for_test(&mut self) {
self.decision_contract = None;
self.decision_posterior = None;
}
fn emit_scheduler_evidence_for_suggestion(&self, suggestion: SchedulingSuggestion) {
let Some(ref sink) = self.evidence_sink else {
return;
};
let snapshot = {
let state = self
.state
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
self.lyapunov_snapshot_locked(&state)
};
let ready_queue_depth = self.ready_queue_depth_signal();
#[allow(clippy::cast_possible_truncation)]
let ready_queue_depth = ready_queue_depth as u32;
let suggestion_str = match suggestion {
SchedulingSuggestion::MeetDeadlines => "meet_deadlines",
SchedulingSuggestion::DrainObligations => "drain_obligations",
SchedulingSuggestion::DrainRegions => "drain_regions",
SchedulingSuggestion::NoPreference => "no_preference",
};
let cancel_depth =
snapshot.cancel_requested_tasks + snapshot.cancelling_tasks + snapshot.finalizing_tasks;
crate::evidence_sink::emit_scheduler_evidence(
sink.as_ref(),
suggestion_str,
cancel_depth,
snapshot.draining_regions,
ready_queue_depth,
self.decision_contract
.as_ref()
.is_some_and(|_| self.decision_posterior.is_some()),
);
}
#[inline]
fn current_base_cancel_limit(&self) -> usize {
self.adaptive_cancel_policy
.as_ref()
.map_or(
self.cancel_streak_limit,
AdaptiveCancelStreakPolicy::current_limit,
)
.max(1)
}
fn potential_from_snapshot(snapshot: &StateSnapshot) -> f64 {
let w = PotentialWeights::default();
let task_component = w.w_tasks * f64::from(snapshot.live_tasks);
#[allow(clippy::cast_precision_loss)]
let obligation_age_seconds = snapshot.obligation_age_sum_ns as f64 / 1_000_000_000.0;
let obligation_component = w.w_obligation_age * obligation_age_seconds;
let region_component = w.w_draining_regions * f64::from(snapshot.draining_regions);
let deadline_component = w.w_deadline_pressure * snapshot.deadline_pressure;
task_component + obligation_component + region_component + deadline_component
}
fn lyapunov_snapshot_locked(&self, state: &RuntimeState) -> StateSnapshot {
match &self.task_table {
Some(tt) => {
let table = tt.lock().unwrap_or_else(std::sync::PoisonError::into_inner);
StateSnapshot::from_runtime_state_with_tasks(state, &table)
}
None => StateSnapshot::from_runtime_state(state),
}
}
fn capture_adaptive_snapshot(&self) -> AdaptiveEpochSnapshot {
let snapshot = {
let state = self
.state
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
self.lyapunov_snapshot_locked(&state)
};
AdaptiveEpochSnapshot {
potential: Self::potential_from_snapshot(&snapshot),
deadline_pressure: snapshot.deadline_pressure,
effective_limit_exceedances: self.preemption_metrics.effective_limit_exceedances,
fallback_cancel_dispatches: self.preemption_metrics.fallback_cancel_dispatches,
}
}
fn ensure_adaptive_epoch_started(&mut self) {
if self
.adaptive_cancel_policy
.as_ref()
.is_none_or(|p| p.epoch_start.is_some())
{
return;
}
let snap = self.capture_adaptive_snapshot();
if let Some(policy) = self.adaptive_cancel_policy.as_mut() {
policy.begin_epoch(snap);
}
}
fn adaptive_on_dispatch(&mut self) {
self.ensure_adaptive_epoch_started();
let should_close_epoch = self
.adaptive_cancel_policy
.as_mut()
.is_some_and(AdaptiveCancelStreakPolicy::on_dispatch);
if !should_close_epoch {
return;
}
let snapshot_end = self.capture_adaptive_snapshot();
let reward = self
.adaptive_cancel_policy
.as_mut()
.and_then(|p| p.complete_epoch(snapshot_end));
if let Some(policy) = self.adaptive_cancel_policy.as_ref() {
self.preemption_metrics.adaptive_epochs = policy.epoch_count;
self.preemption_metrics.adaptive_current_limit = policy.current_limit();
self.preemption_metrics.adaptive_reward_ema = policy.reward_ema;
self.preemption_metrics.adaptive_e_value = policy.e_value();
}
if let Some(reward_value) = reward {
let _ = reward_value;
trace!(
worker_id = self.id,
reward = reward_value,
adaptive_limit = self.preemption_metrics.adaptive_current_limit,
adaptive_epochs = self.preemption_metrics.adaptive_epochs,
adaptive_e_value = self.preemption_metrics.adaptive_e_value,
"adaptive cancel-streak epoch update"
);
}
}
fn abort_adaptive_epoch(&mut self) {
if let Some(policy) = self.adaptive_cancel_policy.as_mut() {
policy.abort_epoch();
}
}
fn drive_io_phase(&self) -> IoPhaseOutcome {
let Some(io) = &self.io_driver else {
return IoPhaseOutcome::NoProgress;
};
let now = self.current_scheduler_time();
let local_deadline = self.local.lock().next_deadline();
let timer_deadline = self
.timer_driver
.as_ref()
.and_then(TimerDriverHandle::next_deadline);
let global_deadline = self.global.peek_earliest_deadline();
let next_deadline = [timer_deadline, local_deadline, global_deadline]
.into_iter()
.flatten()
.min();
let timeout = next_deadline
.map(|deadline| {
if deadline > now {
Duration::from_nanos(deadline.duration_since(now))
} else {
Duration::ZERO
}
})
.or(Some(IDLE_IO_POLL_MAX_TIMEOUT));
let io_timeout = select_io_poll_timeout(
timeout,
self.fast_queue.is_empty(),
self.pending_cancel_dispatch_ready.load(Ordering::Acquire)
|| self
.spawn_mailbox
.as_ref()
.is_some_and(|mailbox| !mailbox.is_empty()),
);
if self.shutdown.load(Ordering::Acquire) {
return IoPhaseOutcome::NoProgress;
}
match io.try_turn_with(io_timeout, |_, _| {}) {
Ok(Some(n)) => {
if n > 0 || io_timeout != Some(Duration::ZERO) {
IoPhaseOutcome::Progress
} else {
IoPhaseOutcome::NoProgress
}
}
Ok(None) | Err(_) => {
IoPhaseOutcome::Follower
}
}
}
#[inline]
fn reset_empty_backoff(&mut self) {
self.empty_backoff = 0;
}
#[inline]
fn advance_empty_backoff(&mut self) -> EmptyBackoffAction {
if self.empty_backoff < SPIN_LIMIT {
self.empty_backoff += 1;
EmptyBackoffAction::Spin
} else if self.empty_backoff < EMPTY_BACKOFF_PARK_THRESHOLD {
self.empty_backoff += 1;
EmptyBackoffAction::Yield
} else {
EmptyBackoffAction::Park
}
}
pub fn run_loop(&mut self) {
let _guard = ScopedLocalScheduler::new(Arc::clone(&self.local));
let _queue_guard = LocalQueue::set_current(self.fast_queue.clone());
let _local_ready_guard = ScopedLocalReady::new(Arc::clone(&self.local_ready));
let _worker_guard = ScopedWorkerId::new(self.id);
while !self.shutdown.load(Ordering::Relaxed) {
if let Some(task) = self.next_task() {
self.reset_empty_backoff();
self.execute(task);
continue;
}
if self.schedule_ready_finalizers() {
continue;
}
let io_phase = self.drive_io_phase();
if matches!(io_phase, IoPhaseOutcome::Progress) {
continue;
}
if !self.fast_queue.is_empty() {
continue;
}
loop {
if self.shutdown.load(Ordering::Relaxed) {
break;
}
let now = self.current_scheduler_time();
if !self.fast_queue.is_empty()
|| self.global.has_cancel_work()
|| self.global.has_ready_work()
|| self.pending_cancel_dispatch_ready.load(Ordering::Acquire)
|| self
.spawn_mailbox
.as_ref()
.is_some_and(|mailbox| !mailbox.is_empty())
{
break;
}
if self.global.has_runnable_work(now) {
match self.advance_empty_backoff() {
EmptyBackoffAction::Spin => {
crate::runtime::metrics::record_worker_spin();
std::hint::spin_loop();
break;
}
EmptyBackoffAction::Yield => {
crate::runtime::metrics::record_sched_yield();
std::thread::yield_now();
break;
}
EmptyBackoffAction::Park => {}
}
}
match self.advance_empty_backoff() {
EmptyBackoffAction::Spin => {
crate::runtime::metrics::record_worker_spin();
std::hint::spin_loop();
}
EmptyBackoffAction::Yield => {
crate::runtime::metrics::record_sched_yield();
std::thread::yield_now();
}
EmptyBackoffAction::Park if self.enable_parking => {
let (local_has_runnable, local_deadline) = {
let mut local = self.local.lock();
(local.has_runnable_work(now), local.next_deadline())
};
let local_ready_has_work = !self.local_ready.lock().is_empty();
let spawn_mailbox_has_work = self
.spawn_mailbox
.as_ref()
.is_some_and(|mailbox| !mailbox.is_empty());
let local_spawn_lane_has_work =
!crate::runtime::spawn_mailbox::local_spawn_lane_is_empty();
if local_has_runnable
|| local_ready_has_work
|| self.pending_cancel_dispatch_ready.load(Ordering::Acquire)
|| spawn_mailbox_has_work
|| local_spawn_lane_has_work
{
break;
}
let timer_deadline = self
.timer_driver
.as_ref()
.and_then(TimerDriverHandle::next_deadline);
let global_deadline = self.global.peek_earliest_deadline();
record_backoff_deadline_selection(
&mut self.preemption_metrics,
io_phase,
timer_deadline,
global_deadline,
);
let next_deadline = select_backoff_deadline(
io_phase,
timer_deadline,
local_deadline,
global_deadline,
);
if let Some(next_deadline) = next_deadline {
let now = self.current_scheduler_time();
match classify_backoff_timeout_decision(io_phase, next_deadline, now) {
BackoffTimeoutDecision::ParkTimeout { nanos } => {
record_backoff_timeout_park(
&mut self.preemption_metrics,
io_phase,
nanos,
);
self.parker.park_timeout(Duration::from_nanos(nanos));
self.reset_empty_backoff();
break;
}
BackoffTimeoutDecision::DeadlineDue => {
record_backoff_timeout_park(
&mut self.preemption_metrics,
io_phase,
STALE_DUE_DEADLINE_PARK_NANOS,
);
self.parker.park_timeout(Duration::from_nanos(
STALE_DUE_DEADLINE_PARK_NANOS,
));
let wheel_due = self
.timer_driver
.as_ref()
.and_then(TimerDriverHandle::next_deadline)
.is_some_and(|deadline| {
deadline <= self.current_scheduler_time()
});
if wheel_due {
self.reset_empty_backoff();
break;
}
}
}
} else {
record_backoff_indefinite_park(&mut self.preemption_metrics, io_phase);
self.parker.park();
}
self.reset_empty_backoff();
}
EmptyBackoffAction::Park => {
self.reset_empty_backoff();
break;
}
}
}
self.cancel_streak = 0;
self.ready_dispatch_streak = 0;
}
}
#[inline]
fn fixed_ready_batch_size(&self) -> usize {
self.steal_batch_size.max(1)
}
#[inline]
fn reset_adaptive_batch_state(&mut self) {
let fixed_batch_size = self.fixed_ready_batch_size();
let last_combiner_claim_failures = self
.global
.ready_combiner_snapshot()
.combiner_claim_failures;
self.adaptive_batch_state = AdaptiveBatchRuntimeState {
active_batch_size: fixed_batch_size,
cooldown_remaining: 0,
last_combiner_claim_failures,
last_snapshot: None,
};
}
#[doc(hidden)]
#[cfg(any(test, feature = "test-internals"))]
pub fn adaptive_batch_snapshot_for_test(&self) -> Option<AdaptiveBatchDecisionSnapshot> {
self.adaptive_batch_state.last_snapshot
}
#[inline]
fn select_ready_batch_decision(&mut self) -> AdaptiveBatchDecisionSnapshot {
let fixed_batch_size = self.fixed_ready_batch_size();
let ready_depth = self.global.ready_count();
let combiner = self.global.ready_combiner_snapshot();
let cancel_debt = self.cancel_debt_signal();
let claim_failures_delta = combiner
.combiner_claim_failures
.saturating_sub(self.adaptive_batch_state.last_combiner_claim_failures);
self.adaptive_batch_state.last_combiner_claim_failures = combiner.combiner_claim_failures;
let mut selected_batch_size = fixed_batch_size;
let mut reason = AdaptiveBatchDecisionReason::Disabled;
if let Some(profile) = self.adaptive_batch_profile {
let profile = profile.normalized(fixed_batch_size);
if profile.enabled {
if self.adaptive_batch_state.cooldown_remaining > 0 {
selected_batch_size = self
.adaptive_batch_state
.active_batch_size
.max(fixed_batch_size)
.clamp(profile.min_batch_size, profile.max_batch_size);
self.adaptive_batch_state.cooldown_remaining = self
.adaptive_batch_state
.cooldown_remaining
.saturating_sub(1);
self.preemption_metrics.adaptive_batch_cooldown_holds += 1;
reason = AdaptiveBatchDecisionReason::CooldownHold;
} else if cancel_debt >= profile.cancel_debt_floor
&& fixed_batch_size > profile.min_batch_size
{
selected_batch_size = profile.min_batch_size;
self.adaptive_batch_state.active_batch_size = selected_batch_size;
self.preemption_metrics.adaptive_batch_cancel_floor_hits += 1;
reason = AdaptiveBatchDecisionReason::CancelDebtFloor;
} else {
let combiner_ready = combiner.max_in_flight >= profile.scale_up_in_flight
|| combiner.current_in_flight >= profile.scale_up_in_flight;
let claim_ready = claim_failures_delta >= profile.scale_up_claim_failures;
if ready_depth >= profile.scale_up_ready_depth
&& combiner_ready
&& claim_ready
&& profile.max_batch_size > fixed_batch_size
{
selected_batch_size =
profile.contention_scale_up_batch_size(fixed_batch_size);
self.adaptive_batch_state.active_batch_size = selected_batch_size;
self.adaptive_batch_state.cooldown_remaining = profile.cooldown_steps;
self.preemption_metrics.adaptive_batch_scale_up_events += 1;
reason = AdaptiveBatchDecisionReason::ReadyContentionScaleUp;
} else {
selected_batch_size = fixed_batch_size;
self.adaptive_batch_state.active_batch_size = selected_batch_size;
reason = AdaptiveBatchDecisionReason::FixedFallback;
}
}
}
}
self.preemption_metrics.adaptive_batch_max_selected = self
.preemption_metrics
.adaptive_batch_max_selected
.max(selected_batch_size);
let snapshot = AdaptiveBatchDecisionSnapshot {
selected_batch_size,
fixed_batch_size,
ready_depth,
cancel_debt,
combiner_in_flight: combiner.max_in_flight.max(combiner.current_in_flight),
combiner_claim_failures_delta: claim_failures_delta,
reason,
};
self.adaptive_batch_state.last_snapshot = Some(snapshot);
snapshot
}
fn insert_deferred_cancel_lane_without_wake(
&self,
task_id: TaskId,
priority: u8,
is_local: bool,
pinned_worker: Option<usize>,
) -> Result<DeferredCancelLanePublication, DeferredCancelLaneError> {
if is_local {
let Some(worker_id) = pinned_worker else {
return Err(DeferredCancelLaneError::MissingPinnedWorker);
};
let Some(local) = self.all_local_schedulers.get(worker_id) else {
return Err(DeferredCancelLaneError::PinnedWorkerUnavailable(worker_id));
};
let mut local = local.lock();
if let Some(local_ready) = self.all_local_ready.get(worker_id) {
local_ready.lock().tombstone(task_id);
}
local.move_to_cancel_lane(task_id, priority);
drop(local);
return Ok(DeferredCancelLanePublication {
priority,
wake_target: DeferredCancelWakeTarget::PinnedWorker(worker_id),
});
}
self.global.remove_timed(task_id);
self.global.inject_cancel(task_id, priority);
Ok(DeferredCancelLanePublication {
priority,
wake_target: DeferredCancelWakeTarget::AnyWorker,
})
}
fn finish_deferred_cancel_lane_publication(
&self,
task_id: TaskId,
publication: DeferredCancelLanePublication,
) {
if let Err(payload) = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
self.record_scheduler_evidence_enqueue(task_id);
})) {
std::mem::forget(payload);
}
match publication.wake_target {
DeferredCancelWakeTarget::AnyWorker => self.coordinator.wake_one(),
DeferredCancelWakeTarget::PinnedWorker(worker_id) => {
self.coordinator.wake_worker(worker_id);
}
}
}
fn emit_cancel_diagnostic(diagnostic: impl FnOnce()) {
if let Err(payload) = std::panic::catch_unwind(std::panic::AssertUnwindSafe(diagnostic)) {
std::mem::forget(payload);
}
}
fn publish_deferred_cancel_task(&self, task_id: TaskId, priority: u8) -> bool {
let task_location = self.with_task_table_ref(|tasks| {
tasks.task(task_id).map(|record| {
record.wake_state.notify();
(record.is_local(), record.pinned_worker())
})
});
let Some((is_local, pinned_worker)) = task_location else {
Self::emit_cancel_diagnostic(|| {
error!(
?task_id,
"deferred cancellation task is absent from the task table"
);
});
return false;
};
match self.insert_deferred_cancel_lane_without_wake(
task_id,
priority,
is_local,
pinned_worker,
) {
Ok(publication) => {
self.finish_deferred_cancel_lane_publication(task_id, publication);
true
}
Err(error) => {
Self::emit_cancel_diagnostic(|| {
let _ = &error;
error!(
?task_id,
?error,
"deferred cancellation lane insertion failed"
);
});
false
}
}
}
fn drain_handle_cancel_requests(&self) {
const HANDLE_CANCEL_BATCH: usize = 16;
let Some(mailbox) = self.spawn_mailbox.as_ref() else {
return;
};
if mailbox.handle_cancels_are_empty() {
return;
}
let mut requests = Vec::with_capacity(HANDLE_CANCEL_BATCH);
if mailbox.dequeue_handle_cancels_into(HANDLE_CANCEL_BATCH, &mut requests) == 0 {
return;
}
let requests = crate::runtime::spawn_mailbox::coalesce_handle_cancel_requests(requests);
let (tasks, delegated, immediate_wakes, immediate_admitted_slots) = if self
.task_table
.is_some()
{
let (tasks, delegated, mut immediate_wakes, immediate_admitted_slots, new_requests) =
self.with_task_table(|tt| {
let mut tasks = Vec::with_capacity(requests.len());
let mut delegated = Vec::new();
let mut immediate_wakes = Vec::new();
let mut immediate_admitted_slots = Vec::new();
let mut new_requests = Vec::new();
for request in requests {
let task_id = request.task_id;
let reason = request.reason;
let admitted_slot = request.admitted_slot;
let Some((effects, region_id, spawned_at)) =
tt.update_task(task_id, |record| {
let effects = record.request_cancel_for_handle(&reason);
(effects, record.owner, record.created_at)
})
else {
continue;
};
let (update, task_wakes) = effects.into_parts();
if update.newly_cancelled {
new_requests.push((task_id, region_id, spawned_at));
}
match update.route {
Some(route) if route.delegated_initial => delegated.push((
task_id,
route.priority,
reason,
task_wakes,
admitted_slot,
)),
Some(route) => {
tasks.push((task_id, route.priority, task_wakes, admitted_slot));
}
None => {
immediate_wakes.push(task_wakes);
if let Some(admitted_slot) = admitted_slot {
immediate_admitted_slots.push(admitted_slot);
}
}
}
}
(
tasks,
delegated,
immediate_wakes,
immediate_admitted_slots,
new_requests,
)
});
if !new_requests.is_empty() {
let new_requests = new_requests
.into_iter()
.map(|(task_id, region_id, spawned_at)| {
let task_still_live =
self.with_task_table_ref(|tt| tt.task(task_id).is_some());
(task_id, region_id, spawned_at, !task_still_live)
})
.collect::<Vec<_>>();
let state = self
.state
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
for (task_id, region_id, spawned_at, allow_retired_noop) in new_requests {
if let Some(validation_result) = state.external_handle_cancel_request_violation(
task_id,
region_id,
spawned_at,
allow_retired_noop,
) {
let mut diagnostic = crate::types::task_context::CancelWakeEffects::empty();
diagnostic.push_cancel_protocol_violation(
"external-shard task-handle cancellation",
validation_result,
);
immediate_wakes.push(diagnostic);
}
}
}
(tasks, delegated, immediate_wakes, immediate_admitted_slots)
} else {
let mut state = self
.state
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
let mut tasks = Vec::with_capacity(requests.len());
let mut delegated = Vec::new();
let mut immediate_wakes = Vec::new();
let mut immediate_admitted_slots = Vec::new();
for request in requests {
let task_id = request.task_id;
let reason = request.reason;
let admitted_slot = request.admitted_slot;
let task_exists = state.task(task_id).is_some();
let effects = state.cancel_task_for_handle(task_id, &reason);
let (route, task_wakes) = effects.into_parts();
match route {
Some(route) if route.delegated_initial => {
delegated.push((task_id, route.priority, reason, task_wakes, admitted_slot))
}
Some(route) => {
tasks.push((task_id, route.priority, task_wakes, admitted_slot));
}
None => {
immediate_wakes.push(task_wakes);
if task_exists && let Some(admitted_slot) = admitted_slot {
immediate_admitted_slots.push(admitted_slot);
}
}
}
}
drop(state);
(tasks, delegated, immediate_wakes, immediate_admitted_slots)
};
let mut wakes_to_dispatch = immediate_wakes;
let mut spawn_effects_to_dispatch =
Vec::with_capacity(immediate_admitted_slots.len() + tasks.len() + delegated.len());
for admitted_slot in immediate_admitted_slots {
if let Some(effects) = admitted_slot.take_spawn_effects_if_lane_published() {
spawn_effects_to_dispatch.push(effects);
}
}
for (task_id, priority, task_wakes, admitted_slot) in tasks {
if self.publish_deferred_cancel_task(task_id, priority) {
wakes_to_dispatch.push(task_wakes);
if let Some(admitted_slot) = admitted_slot
&& let Some(effects) = admitted_slot.publish_spawn_lane_and_take_effects()
{
spawn_effects_to_dispatch.push(effects);
}
} else {
Self::emit_cancel_diagnostic(|| {
error!(
?task_id,
priority,
"handle cancellation promotion failed; suppressing only this task's \
Wakers fail-closed"
);
});
task_wakes.suppress();
}
}
for (task_id, requested_priority, reason, mut task_wakes, admitted_slot) in delegated {
let mut lane_error = None;
let effects = if self.task_table.is_some() {
self.with_task_table(|tt| {
tt.update_task(task_id, |record| {
record.publish_delegated_cancel_lane(|priority, is_local, pinned_worker| {
match self.insert_deferred_cancel_lane_without_wake(
task_id,
priority,
is_local,
pinned_worker,
) {
Ok(publication) => Some(publication),
Err(error) => {
lane_error = Some(error);
None
}
}
})
})
.unwrap_or_else(|| crate::types::task_context::CancellationEffects::ready(None))
})
} else {
let mut state = self
.state
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
state.publish_handle_cancel_lane(task_id, |priority, is_local, pinned_worker| {
match self.insert_deferred_cancel_lane_without_wake(
task_id,
priority,
is_local,
pinned_worker,
) {
Ok(publication) => Some(publication),
Err(error) => {
lane_error = Some(error);
None
}
}
})
};
let (publication, publication_wakes) = effects.into_parts();
if let Some(publication) = publication {
self.finish_deferred_cancel_lane_publication(task_id, publication);
let spawn_effects = admitted_slot
.as_ref()
.and_then(|slot| slot.publish_spawn_lane_and_take_effects());
debug_assert!(publication.priority >= requested_priority);
task_wakes.merge(publication_wakes);
wakes_to_dispatch.push(task_wakes);
if let Some(effects) = spawn_effects {
spawn_effects_to_dispatch.push(effects);
}
continue;
}
if let Some(error) = lane_error {
Self::emit_cancel_diagnostic(|| {
let _ = &error;
error!(
?task_id,
requested_priority,
?error,
"delegated handle cancellation lane insertion failed; suppressing only \
this attempt's Wakers without an internal retry loop"
);
});
}
drop(reason);
task_wakes.suppress();
publication_wakes.retire_without_dispatch();
}
for effects in spawn_effects_to_dispatch {
effects.dispatch();
}
for wakes in wakes_to_dispatch {
wakes.dispatch();
}
}
fn drain_deferred_cancel_dispatches(&self) {
if !self.pending_cancel_dispatch_ready.load(Ordering::Acquire) {
return;
}
let batches = {
let mut state = self
.state
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
state.take_deferred_cancel_dispatches()
};
let mut wakes = Vec::with_capacity(batches.len());
for batch in batches {
let (tasks, batch_wakes) = batch.into_parts();
let mut batch_published = true;
for (task_id, priority) in tasks {
batch_published &= self.publish_deferred_cancel_task(task_id, priority);
}
wakes.push((batch_published, batch_wakes));
}
for (batch_published, batch_wakes) in wakes {
if batch_published {
batch_wakes.dispatch();
} else {
Self::emit_cancel_diagnostic(|| {
error!(
"deferred cancellation batch publication was incomplete; suppressing \
only that batch's Wakers fail-closed"
);
});
batch_wakes.suppress();
}
}
}
#[allow(clippy::too_many_lines)]
fn drain_spawn_admissions(&mut self) {
const SPAWN_ADMISSION_BATCH: usize = 1;
let Some(mailbox) = self.spawn_mailbox.as_ref() else {
return;
};
if mailbox.spawn_requests_are_empty() {
return;
}
let mailbox = Arc::clone(mailbox);
let mut requests = Vec::with_capacity(SPAWN_ADMISSION_BATCH);
if mailbox.dequeue_batch_into(SPAWN_ADMISSION_BATCH, &mut requests) == 0 {
return;
}
let mut admitted: SmallVec<
[(
TaskId,
u8,
crate::runtime::spawn_mailbox::AdmissionPublication,
crate::runtime::state::TaskSpawnEffects,
); 16],
> = SmallVec::new();
let mut denied: SmallVec<
[(
crate::runtime::spawn_mailbox::SpawnRequestParts,
crate::runtime::state::SpawnError,
); 4],
> = SmallVec::new();
{
let mut state = self
.state
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
let mut admission_tasks = self.task_table.as_ref().map_or(
crate::runtime::state::AdmissionTaskTarget::Embedded,
|tt| {
crate::runtime::state::AdmissionTaskTarget::External(
tt.lock().unwrap_or_else(std::sync::PoisonError::into_inner),
)
},
);
for request in requests {
match state.admit_spawn_request_in(
request.into_parts(),
&mut admission_tasks,
&crate::runtime::state::AdmissionRegionTarget::Embedded,
) {
crate::runtime::state::SpawnAdmission::Admitted {
task_id,
priority,
cancel_publication,
spawn_effects,
} => {
admitted.push((task_id, priority, cancel_publication, spawn_effects));
}
crate::runtime::state::SpawnAdmission::Denied { parts, error } => {
denied.push((parts, error));
}
}
}
}
let mut cancel_wakes = SmallVec::<
[crate::types::task_context::CancelWakeEffects; SPAWN_ADMISSION_BATCH],
>::new();
let mut spawn_effects =
SmallVec::<[crate::runtime::state::TaskSpawnEffects; SPAWN_ADMISSION_BATCH]>::new();
for (task_id, priority, publication, effects) in admitted {
let (wakes, effects) =
publication.publish_with_spawn_effects(effects, |cancel_priority| {
if let Some(cancel_priority) = cancel_priority {
self.global.inject_cancel(task_id, cancel_priority);
} else {
self.global.inject_ready(task_id, priority);
}
});
cancel_wakes.push(wakes);
if let Some(effects) = effects {
spawn_effects.push(effects);
}
}
for effects in spawn_effects {
effects.dispatch();
}
for (parts, error) in denied {
match error {
crate::runtime::state::SpawnError::RegionClosed(_)
| crate::runtime::state::SpawnError::RegionNotFound(_) => {
parts.resolve_cancelled(crate::types::CancelReason::new(
crate::types::CancelKind::ParentCancelled,
));
}
other => parts.resolve_failed(other),
}
}
for wakes in cancel_wakes {
wakes.dispatch();
}
}
fn drain_local_spawn_admissions(&mut self) {
const LOCAL_SPAWN_ADMISSION_BATCH: usize = 16;
if crate::runtime::spawn_mailbox::local_spawn_lane_is_empty() {
return;
}
let mut requests = Vec::with_capacity(LOCAL_SPAWN_ADMISSION_BATCH);
if crate::runtime::spawn_mailbox::drain_local_spawn_lane(
LOCAL_SPAWN_ADMISSION_BATCH,
&mut requests,
) == 0
{
return;
}
let mut admitted: Vec<(
TaskId,
u8,
crate::runtime::stored_task::LocalStoredTask,
crate::runtime::spawn_mailbox::AdmissionPublication,
crate::runtime::state::TaskSpawnEffects,
)> = Vec::with_capacity(requests.len());
let mut denied: SmallVec<
[(
crate::runtime::spawn_mailbox::LocalSpawnRequest,
crate::runtime::state::SpawnError,
); 4],
> = SmallVec::new();
{
let mut state = self
.state
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
let mut admission_tasks = self.task_table.as_ref().map_or(
crate::runtime::state::AdmissionTaskTarget::Embedded,
|tt| {
crate::runtime::state::AdmissionTaskTarget::External(
tt.lock().unwrap_or_else(std::sync::PoisonError::into_inner),
)
},
);
for request in requests {
match state.admit_local_spawn_request_in(
request,
&mut admission_tasks,
&crate::runtime::state::AdmissionRegionTarget::Embedded,
) {
crate::runtime::state::LocalSpawnAdmission::Admitted {
task_id,
priority,
stored,
cancel_publication,
spawn_effects,
} => {
match &mut admission_tasks {
crate::runtime::state::AdmissionTaskTarget::Embedded => {
if let Some(record) = state.task_mut(task_id) {
record.wake_state.notify();
}
}
crate::runtime::state::AdmissionTaskTarget::External(tt) => {
if let Some(record) = tt.task_mut(task_id) {
record.wake_state.notify();
}
}
}
admitted.push((
task_id,
priority,
stored,
cancel_publication,
spawn_effects,
));
}
crate::runtime::state::LocalSpawnAdmission::Denied { request, error } => {
denied.push((request, error));
}
}
}
}
let mut cancel_wakes = SmallVec::<
[crate::types::task_context::CancelWakeEffects; LOCAL_SPAWN_ADMISSION_BATCH],
>::new();
let mut spawn_effects = SmallVec::<
[crate::runtime::state::TaskSpawnEffects; LOCAL_SPAWN_ADMISSION_BATCH],
>::new();
for (task_id, _priority, stored, publication, effects) in admitted {
crate::runtime::local::store_local_task(task_id, stored);
let (wakes, effects) =
publication.publish_with_spawn_effects(effects, |cancel_priority| {
if let Some(cancel_priority) = cancel_priority {
self.local.lock().schedule_cancel(task_id, cancel_priority);
} else {
self.local_ready.lock().push_back(task_id);
}
});
cancel_wakes.push(wakes);
if let Some(effects) = effects {
spawn_effects.push(effects);
}
}
for effects in spawn_effects {
effects.dispatch();
}
for (request, error) in denied {
match error {
crate::runtime::state::SpawnError::RegionClosed(_)
| crate::runtime::state::SpawnError::RegionNotFound(_) => {
request.resolve_cancelled(crate::types::CancelReason::new(
crate::types::CancelKind::ParentCancelled,
));
}
other => request.resolve_failed(other),
}
}
for wakes in cancel_wakes {
wakes.dispatch();
}
}
pub fn next_task(&mut self) -> Option<TaskId> {
self.drain_handle_cancel_requests();
self.drain_deferred_cancel_dispatches();
if let Some(timer) = &self.timer_driver {
let _ = timer.process_timers();
}
self.drain_spawn_admissions();
self.drain_local_spawn_admissions();
self.drain_handle_cancel_requests();
self.drain_deferred_cancel_dispatches();
if self.pending_cancel_dispatch_ready.load(Ordering::Acquire)
|| self
.spawn_mailbox
.as_ref()
.is_some_and(|mailbox| !mailbox.handle_cancels_are_empty())
{
return None;
}
let suggestion = self.governor_suggest();
let base_limit = self.current_base_cancel_limit();
self.preemption_metrics.adaptive_current_limit = base_limit;
let effective_limit = match suggestion {
SchedulingSuggestion::DrainObligations | SchedulingSuggestion::DrainRegions => {
base_limit.saturating_mul(2)
}
_ => base_limit,
};
if effective_limit > self.preemption_metrics.max_effective_limit_observed {
self.preemption_metrics.max_effective_limit_observed = effective_limit;
}
let check_cancel = self.cancel_streak < effective_limit;
if !check_cancel {
self.preemption_metrics.fairness_yields += 1;
}
let check_timed = self.timed_dispatch_streak < self.timed_fairness_limit;
if !check_timed && suggestion == SchedulingSuggestion::MeetDeadlines {
if let Some(task) = self.try_phase3_ready_work() {
self.timed_dispatch_streak = 0; return Some(task);
}
self.preemption_metrics.fairness_yields += 1;
}
let now = self.current_scheduler_time();
if suggestion == SchedulingSuggestion::MeetDeadlines && check_timed {
if let Some(tt) = self.global.pop_timed_if_due(now) {
self.record_timed_dispatch();
self.timed_dispatch_streak += 1; return Some(self.dispatch_with_adaptive_epoch(tt.task));
}
} else {
if check_cancel {
if let Some(pt) = self.global.pop_cancel() {
self.cancel_streak += 1;
self.ready_dispatch_streak = 0;
self.record_cancel_dispatch(base_limit, effective_limit);
return Some(self.dispatch_with_adaptive_epoch(pt.task));
}
}
}
let mut local = self.local.lock();
let rng_hint = self.rng.next_u64();
if suggestion == SchedulingSuggestion::MeetDeadlines && check_timed {
if let Some(task) = local.pop_timed_only_with_hint(rng_hint, now) {
drop(local);
self.record_timed_dispatch();
self.timed_dispatch_streak += 1; return Some(self.dispatch_with_adaptive_epoch(task));
}
if check_cancel {
if let Some(pt) = self.global.pop_cancel() {
drop(local);
self.cancel_streak += 1;
self.ready_dispatch_streak = 0;
self.record_cancel_dispatch(base_limit, effective_limit);
return Some(self.dispatch_with_adaptive_epoch(pt.task));
}
if let Some(task) = local.pop_cancel_only_with_hint(rng_hint) {
drop(local);
self.cancel_streak += 1;
self.ready_dispatch_streak = 0;
self.record_cancel_dispatch(base_limit, effective_limit);
return Some(self.dispatch_with_adaptive_epoch(task));
}
}
} else {
if check_cancel {
if let Some(task) = local.pop_cancel_only_with_hint(rng_hint) {
drop(local);
self.cancel_streak += 1;
self.ready_dispatch_streak = 0;
self.record_cancel_dispatch(base_limit, effective_limit);
return Some(self.dispatch_with_adaptive_epoch(task));
}
}
if let Some(tt) = self.global.pop_timed_if_due(now) {
drop(local);
self.record_timed_dispatch();
return Some(self.dispatch_with_adaptive_epoch(tt.task));
}
if let Some(task) = local.pop_timed_only_with_hint(rng_hint, now) {
drop(local);
self.record_timed_dispatch();
return Some(self.dispatch_with_adaptive_epoch(task));
}
}
drop(local);
if self.should_force_ready_handoff() {
self.preemption_metrics.browser_ready_handoff_yields += 1;
self.cancel_streak = 0;
self.ready_dispatch_streak = 0;
return None;
}
if let Some(task) = self.try_phase3_ready_work() {
return Some(task);
}
if let Some(task) = self.try_steal() {
self.record_ready_dispatch();
return Some(self.dispatch_with_adaptive_epoch(task));
}
if !check_cancel {
if let Some(task) = self.try_cancel_work() {
self.preemption_metrics.fallback_cancel_dispatches += 1;
self.cancel_streak = 1;
self.ready_dispatch_streak = 0;
self.record_cancel_dispatch(base_limit, effective_limit);
return Some(self.dispatch_with_adaptive_epoch(task));
}
self.cancel_streak = 0;
}
self.ready_dispatch_streak = 0;
None
}
#[inline]
fn should_force_ready_handoff(&self) -> bool {
let limit = self.browser_ready_handoff_limit;
if limit == 0 || self.ready_dispatch_streak < limit {
return false;
}
if !self.fast_queue.is_empty()
|| !self.global_ready_buffer.is_empty()
|| self.global.has_ready_work()
{
return true;
}
if self
.local_ready
.try_lock()
.is_some_and(|queue| !queue.is_empty())
{
return true;
}
self.local.lock().has_ready_work()
}
#[inline]
fn peek_blocked_local_ready_for_inversion(&self) -> Option<(TaskId, u8)> {
self.local
.try_lock()
.and_then(|mut local| local.peek_ready_task())
}
#[inline]
fn take_global_ready_task(&mut self) -> Option<PriorityTask> {
if let Some(prefetched) = self.global_ready_buffer.pop() {
return Some(prefetched);
}
let decision = self.select_ready_batch_decision();
let batch_size = decision.selected_batch_size.max(1);
let batch_threshold = batch_size
.saturating_mul(2)
.max(GLOBAL_READY_BATCH_DRAIN_MIN_DEPTH);
if batch_size > 1 && self.global.ready_count() >= batch_threshold {
self.global_ready_buffer.clear();
let drained = self
.global
.pop_ready_batch_into(batch_size, &mut self.global_ready_buffer);
if drained > 0 {
self.global_ready_buffer.reverse();
self.preemption_metrics.global_ready_batch_drains += 1;
self.preemption_metrics.global_ready_batch_tasks += drained as u64;
return self.global_ready_buffer.pop();
}
}
self.global.pop_ready()
}
fn try_phase3_ready_work(&mut self) -> Option<TaskId> {
let local_ready_task = self.local_ready.lock().pop_front();
if let Some(task) = local_ready_task {
self.record_ready_dispatch();
self.fast_queue_dispatch_streak = 0; return Some(self.dispatch_with_adaptive_epoch(task));
}
let should_prioritize_local =
self.fast_queue_dispatch_streak >= self.fast_queue_fairness_limit;
if should_prioritize_local {
let rng_hint = self.rng.next_u64();
let local_task = {
let mut local = self.local.lock();
local.pop_ready_only_with_hint(rng_hint)
};
if let Some(task) = local_task {
self.record_ready_dispatch();
self.fast_queue_dispatch_streak = 0; return Some(self.dispatch_with_adaptive_epoch(task));
}
}
if let Some(task) = self.fast_queue.pop() {
if let Some(blocked_local_task) = self.peek_blocked_local_ready_for_inversion() {
let dispatched_priority = self.task_sched_priority(task);
self.record_ready_priority_inversion(
Some(blocked_local_task),
task,
dispatched_priority,
);
}
self.record_ready_dispatch();
self.fast_queue_dispatch_streak += 1; return Some(self.dispatch_with_adaptive_epoch(task));
}
if let Some(pt) = self.take_global_ready_task() {
if let Some(blocked_local_task) = self.peek_blocked_local_ready_for_inversion() {
self.record_ready_priority_inversion(
Some(blocked_local_task),
pt.task,
Some(pt.priority),
);
}
self.record_ready_dispatch();
self.fast_queue_dispatch_streak = 0; return Some(self.dispatch_with_adaptive_epoch(pt.task));
}
if !should_prioritize_local {
let rng_hint = self.rng.next_u64();
let local_task = {
let mut local = self.local.lock();
local.pop_ready_only_with_hint(rng_hint)
};
if let Some(task) = local_task {
self.record_ready_dispatch();
self.fast_queue_dispatch_streak = 0; return Some(self.dispatch_with_adaptive_epoch(task));
}
}
None
}
#[inline]
fn record_cancel_dispatch(&mut self, base_limit: usize, effective_limit: usize) {
self.preemption_metrics.cancel_dispatches += 1;
if self.cancel_streak > self.preemption_metrics.max_cancel_streak {
self.preemption_metrics.max_cancel_streak = self.cancel_streak;
}
if self.cancel_streak > base_limit {
self.preemption_metrics.base_limit_exceedances += 1;
}
if self.cancel_streak > effective_limit {
self.preemption_metrics.effective_limit_exceedances += 1;
}
self.timed_dispatch_streak = 0;
}
#[inline]
fn record_timed_dispatch(&mut self) {
if self.cancel_streak > self.preemption_metrics.max_timed_dispatch_stall {
self.preemption_metrics.max_timed_dispatch_stall = self.cancel_streak;
}
self.cancel_streak = 0;
self.ready_dispatch_streak = 0;
self.preemption_metrics.timed_dispatches += 1;
}
#[inline]
fn record_ready_dispatch(&mut self) {
if self.cancel_streak > self.preemption_metrics.max_ready_dispatch_stall {
self.preemption_metrics.max_ready_dispatch_stall = self.cancel_streak;
}
self.cancel_streak = 0;
self.ready_dispatch_streak = self.ready_dispatch_streak.saturating_add(1);
self.timed_dispatch_streak = 0;
self.preemption_metrics.ready_dispatches += 1;
}
fn record_ready_priority_inversion(
&mut self,
blocked_task: Option<(TaskId, u8)>,
executing_task: TaskId,
executing_priority: Option<u8>,
) {
let Some((blocked_task, blocked_priority)) = blocked_task else {
return;
};
let Some(executing_priority) = executing_priority else {
return;
};
if blocked_priority <= executing_priority {
return;
}
let timestamp = Time::from_nanos(self.current_time_ns());
self.preemption_metrics.ready_priority_inversions += 1;
let gap = blocked_priority.saturating_sub(executing_priority);
if gap > self.preemption_metrics.max_ready_priority_inversion_gap {
self.preemption_metrics.max_ready_priority_inversion_gap = gap;
}
{
let mut invariant_monitor = self.invariant_monitor.lock();
invariant_monitor.record_task_requeue(
blocked_task,
"local_ready_heap",
blocked_priority,
timestamp,
);
invariant_monitor.verify_priority_ordering(
executing_task,
executing_priority,
blocked_task,
blocked_priority,
timestamp,
);
}
self.fairness_monitor.lock().record_priority_inversion(
blocked_task,
blocked_priority,
executing_task,
executing_priority,
timestamp.as_nanos(),
);
}
#[inline]
fn task_sched_priority(&self, task: TaskId) -> Option<u8> {
self.with_task_table_ref(|tt| tt.task(task).map(|record| record.sched_priority))
}
#[inline]
fn dispatch_with_adaptive_epoch(&mut self, task: TaskId) -> TaskId {
self.ensure_adaptive_epoch_started();
self.finish_dispatch(task)
}
#[inline]
fn finish_dispatch(&mut self, task: TaskId) -> TaskId {
let current_time = self.current_time_ns();
self.fairness_monitor
.lock()
.record_task_dispatch(task, current_time);
self.invariant_monitor
.lock()
.record_task_dispatch(task, Time::from_nanos(current_time));
if let Some(collector) = &self.scheduler_evidence {
let ready_backlog = self.ready_queue_depth_signal();
let cancel_debt = self.cancel_debt_signal();
collector
.lock()
.record_task_dispatch(task, current_time, ready_backlog, cancel_debt);
}
task
}
#[inline]
fn ready_queue_depth_signal(&self) -> usize {
let global_ready = self.global.ready_count();
let prefetched_global_ready = self.global_ready_buffer.len();
let fast_ready = self.fast_queue.len();
let pinned_local_ready = self.local_ready.lock().len();
let local_priority_ready = self.local.lock().approx_ready_len();
global_ready
.saturating_add(prefetched_global_ready)
.saturating_add(fast_ready)
.saturating_add(pinned_local_ready)
.saturating_add(local_priority_ready)
}
#[inline]
fn cancel_debt_signal(&self) -> usize {
let global_cancel = self.global.cancel_count();
let local_cancel = self.local.lock().approx_cancel_len();
global_cancel.saturating_add(local_cancel)
}
#[allow(clippy::too_many_lines)]
fn governor_suggest(&mut self) -> SchedulingSuggestion {
let Some(governor) = &self.governor else {
return SchedulingSuggestion::NoPreference;
};
self.steps_since_snapshot += 1;
if self.steps_since_snapshot < self.governor_interval {
self.emit_scheduler_evidence_for_suggestion(self.cached_suggestion);
return self.cached_suggestion;
}
self.steps_since_snapshot = 0;
let state = self
.state
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
let (snapshot, wait_graph_snapshot) = if let Some(tt) = &self.task_table {
let table = tt.lock().unwrap_or_else(std::sync::PoisonError::into_inner);
let snapshot = StateSnapshot::from_runtime_state_with_tasks(&state, &table);
let wait_graph_snapshot = if self.spectral_monitor.is_some() {
Some(wait_graph_snapshot_from_tasks(&table))
} else {
None
};
(snapshot, wait_graph_snapshot)
} else {
let snapshot = StateSnapshot::from_runtime_state(&state);
let wait_graph_snapshot = if self.spectral_monitor.is_some() {
Some(wait_graph_snapshot_from_state(&state))
} else {
None
};
(snapshot, wait_graph_snapshot)
};
drop(state);
let (wait_graph_nodes, wait_graph_edges, trapped_wait_cycle) = wait_graph_snapshot
.as_ref()
.map_or((0, Vec::new(), false), |snapshot| {
wait_graph_signals_from_snapshot(snapshot)
});
let queue_depth = self.ready_queue_depth_signal();
#[allow(clippy::cast_possible_truncation)]
let snapshot = snapshot.with_ready_queue_depth(queue_depth as u32);
let lyapunov_suggestion = governor.suggest(&snapshot);
let drain_verdict = self.drain_certificate.as_mut().and_then(|cert| {
let is_drain_phase = matches!(
lyapunov_suggestion,
SchedulingSuggestion::DrainObligations | SchedulingSuggestion::DrainRegions
);
if is_drain_phase {
cert.observe(governor.compute_record(&snapshot).total);
if cert.len() > 128 {
cert.compact(64);
}
Some(cert.verdict())
} else {
if !cert.is_empty() {
cert.reset();
}
None
}
});
let mut spectral_report = None;
if let Some(monitor) = self.spectral_monitor.as_mut() {
if trapped_wait_cycle || wait_graph_nodes > 1 {
spectral_report = Some(monitor.analyze_with_trapped_cycle(
wait_graph_nodes,
&wait_graph_edges,
trapped_wait_cycle,
));
}
}
let mut suggestion = if let (Some(contract), Some(posterior)) =
(&self.decision_contract, &mut self.decision_posterior)
{
let likelihoods =
super::decision_contract::SchedulerDecisionContract::snapshot_likelihoods(
&snapshot,
);
posterior.bayesian_update(&likelihoods);
let probs = posterior.probs();
#[allow(clippy::cast_precision_loss)]
let uniform = 1.0 / probs.len().max(1) as f64;
let max_prob = probs
.iter()
.copied()
.fold(0.0_f64, f64::max)
.clamp(0.0, 1.0);
let concentration = if probs.len() > 1 {
((max_prob - uniform) / (1.0 - uniform)).clamp(0.0, 1.0)
} else {
1.0
};
let entropy = normalized_entropy(probs);
let conformal_hit = spectral_report
.as_ref()
.and_then(|report| {
report.bifurcation.as_ref().and_then(|bw| {
bw.conformal_lower_bound_next
.map(|lb| u8::from(report.decomposition.fiedler_value >= lb))
})
})
.map_or(1.0, f64::from);
let uncertainty_penalty = 0.35f64.mul_add(1.0 - concentration, 0.15 * entropy);
let conformal_penalty = 0.5 * (1.0 - conformal_hit);
let calibration_score = (1.0 - uncertainty_penalty - conformal_penalty).clamp(0.0, 1.0);
let ci_width = 0.5f64
.mul_add(1.0 - concentration, 0.25 * entropy)
.clamp(0.0, 1.0);
let adaptive_e = self.preemption_metrics.adaptive_e_value.max(1.0);
let spectral_e = spectral_report
.as_ref()
.and_then(|report| {
report
.bifurcation
.as_ref()
.map(|bw| bw.deterioration_e_value.max(1.0))
})
.unwrap_or(1.0);
let e_process = adaptive_e.max(spectral_e);
let seq = self.decision_sequence;
self.decision_sequence = self.decision_sequence.saturating_add(1);
let now_ms = self
.timer_driver
.as_ref()
.map_or(seq, |td| td.now().as_millis());
let random_bits = ((self.id as u128) << 64) | u128::from(seq);
let ctx = franken_decision::EvalContext {
calibration_score,
e_process,
ci_width,
decision_id: franken_kernel::DecisionId::from_parts(now_ms, random_bits),
trace_id: franken_kernel::TraceId::from_parts(
now_ms,
random_bits ^ 0xA5A5_A5A5_A5A5_A5A5_A5A5,
),
ts_unix_ms: now_ms,
};
let outcome = match franken_decision::evaluate(contract, posterior, &ctx) {
Ok(o) => o,
Err(_) => return lyapunov_suggestion,
};
if let Some(ref sink) = self.evidence_sink {
let evidence = outcome.audit_entry.to_evidence_ledger();
sink.emit(&evidence);
}
match outcome.action_index {
super::decision_contract::action::AGGRESSIVE => SchedulingSuggestion::NoPreference,
super::decision_contract::action::CONSERVATIVE => {
SchedulingSuggestion::MeetDeadlines
}
_ => lyapunov_suggestion,
}
} else {
lyapunov_suggestion
};
if let Some(report) = spectral_report.as_ref() {
let override_suggestion = match report.classification {
crate::observability::spectral_health::HealthClassification::Deadlocked => {
Some(SchedulingSuggestion::DrainObligations)
}
crate::observability::spectral_health::HealthClassification::Critical {
approaching_disconnect: true,
..
} => Some(SchedulingSuggestion::DrainObligations),
_ => report.bifurcation.as_ref().and_then(|bw| {
(bw.trend
== crate::observability::spectral_health::SpectralTrend::Deteriorating
&& (bw.confidence >= 0.6 || bw.deterioration_e_value >= 2.0))
.then_some(SchedulingSuggestion::DrainRegions)
}),
};
if let Some(ovr) = override_suggestion {
suggestion = ovr;
}
}
if trapped_wait_cycle {
suggestion = SchedulingSuggestion::DrainObligations;
}
if !trapped_wait_cycle {
if let Some(ref verdict) = drain_verdict {
match verdict.drain_phase {
DrainPhase::Stalled if verdict.stall_detected => {
suggestion = SchedulingSuggestion::DrainObligations;
}
DrainPhase::Quiescent => {
suggestion = SchedulingSuggestion::NoPreference;
}
_ => {}
}
}
}
self.emit_scheduler_evidence_for_suggestion(suggestion);
self.cached_suggestion = suggestion;
suggestion
}
fn current_scheduler_time(&self) -> Time {
if let Some(timer_driver) = self.timer_driver.as_ref() {
return TimerDriverHandle::now(timer_driver);
}
self.state
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.now
}
pub fn run_once(&mut self) -> bool {
if self.shutdown.load(Ordering::Relaxed) {
return false;
}
if let Some(task) = self.next_task() {
self.execute(task);
return true;
}
false
}
pub(crate) fn try_cancel_work(&mut self) -> Option<TaskId> {
if let Some(pt) = self.global.pop_cancel() {
return Some(pt.task);
}
let mut local = self.local.lock();
let rng_hint = self.rng.next_u64();
local.pop_cancel_only_with_hint(rng_hint)
}
#[allow(dead_code)] pub(crate) fn try_timed_work(&mut self) -> Option<TaskId> {
let now = self.current_scheduler_time();
if let Some(tt) = self.global.pop_timed_if_due(now) {
return Some(tt.task);
}
let mut local = self.local.lock();
let rng_hint = self.rng.next_u64();
local.pop_timed_only_with_hint(rng_hint, now)
}
#[cfg(any(test, feature = "test-internals"))]
pub fn ready_count(&self) -> usize {
let local_ready = self.local_ready.try_lock().map_or(0, |q| q.len());
let fast = self.fast_queue.len();
let prefetched_global = self.global_ready_buffer.len();
let global = self.global.ready_count();
local_ready + fast + prefetched_global + global
}
#[doc(hidden)]
#[cfg(any(test, feature = "test-internals"))]
pub fn enqueue_pinned_local_for_test(&self, task: TaskId) {
self.local_ready.lock().push_back(task);
}
#[doc(hidden)]
#[cfg(any(test, feature = "test-internals"))]
pub fn contains_pinned_local_for_test(&self, task: TaskId) -> bool {
self.local_ready.lock().snapshot().contains(&task)
}
#[cfg(any(test, feature = "test-internals"))]
pub fn bench_fast_ready_queue(&self) -> LocalQueue {
self.fast_queue.clone()
}
#[cfg(any(test, feature = "test-internals"))]
pub fn bench_local_priority_scheduler(&self) -> Arc<Mutex<PriorityScheduler>> {
Arc::clone(&self.local)
}
#[cfg(any(test, feature = "test-internals"))]
pub fn bench_try_phase3_ready_work(&mut self) -> Option<TaskId> {
self.try_phase3_ready_work()
}
#[cfg(any(test, feature = "test-internals"))]
pub fn bench_try_steal(&mut self) -> Option<TaskId> {
self.try_steal()
}
#[allow(dead_code)] pub(crate) fn try_ready_work(&mut self) -> Option<TaskId> {
if let Some(task) = self.local_ready.lock().pop_front() {
return Some(task);
}
if let Some(task) = self.fast_queue.pop() {
return Some(task);
}
if let Some(pt) = self.take_global_ready_task() {
return Some(pt.task);
}
let mut local = self.local.lock();
let rng_hint = self.rng.next_u64();
local.pop_ready_only_with_hint(rng_hint)
}
pub(crate) fn try_steal(&mut self) -> Option<TaskId> {
if !self.fast_stealers.is_empty() {
let preferred_len = self
.preferred_fast_stealer_count
.min(self.fast_stealers.len());
if preferred_len > 0 {
let start = self.rng.next_usize(preferred_len);
for i in 0..preferred_len {
let idx = (start + i) % preferred_len;
if let Some(task) = self.fast_stealers[idx].steal() {
debug_assert!(
!self.with_task_table_ref(|tt| {
tt.task(task)
.is_some_and(crate::record::task::TaskRecord::is_local)
}),
"BUG: stole a local (!Send) task {task:?} from another worker's fast_queue"
);
if self.fast_stealer_locality[idx].is_same_cohort() {
self.steal_locality_counters.preferred_fast_steals += 1;
} else {
self.steal_locality_counters.remote_fast_steals += 1;
}
self.invariant_monitor
.lock()
.record_task_dispatch(task, Time::from_nanos(self.current_time_ns()));
return Some(task);
}
}
}
let remote_len = self.fast_stealers.len().saturating_sub(preferred_len);
if remote_len > 0 {
let start = self.rng.next_usize(remote_len);
for i in 0..remote_len {
let idx = preferred_len + (start + i) % remote_len;
if let Some(task) = self.fast_stealers[idx].steal() {
debug_assert!(
!self.with_task_table_ref(|tt| {
tt.task(task)
.is_some_and(crate::record::task::TaskRecord::is_local)
}),
"BUG: stole a local (!Send) task {task:?} from another worker's fast_queue"
);
if self.fast_stealer_locality[idx].is_same_cohort() {
self.steal_locality_counters.preferred_fast_steals += 1;
} else {
self.steal_locality_counters.remote_fast_steals += 1;
}
self.invariant_monitor
.lock()
.record_task_dispatch(task, Time::from_nanos(self.current_time_ns()));
return Some(task);
}
}
}
}
if self.stealers.is_empty() {
return None;
}
let preferred_len = self.preferred_heap_stealer_count.min(self.stealers.len());
for &(segment_start, segment_len) in &[
(0usize, preferred_len),
(
preferred_len,
self.stealers.len().saturating_sub(preferred_len),
),
] {
if segment_len == 0 {
continue;
}
let start = self.rng.next_usize(segment_len);
for i in 0..segment_len {
let idx = segment_start + (start + i) % segment_len;
let stealer = &self.stealers[idx];
if let Some(mut victim) = stealer.try_lock() {
let stolen_count = victim
.steal_ready_batch_into(self.steal_batch_size, &mut self.steal_buffer);
drop(victim);
if stolen_count > 0 {
#[cfg(debug_assertions)]
{
for &(task, _) in &self.steal_buffer[..stolen_count] {
let is_local = self.with_task_table_ref(|tt| {
tt.task(task)
.is_some_and(crate::record::task::TaskRecord::is_local)
});
debug_assert!(
!is_local,
"BUG: stole a local (!Send) task {task:?} from PriorityScheduler"
);
}
}
let (first_task, _) = self.steal_buffer[0];
if self.heap_stealer_locality[idx].is_same_cohort() {
self.steal_locality_counters.preferred_heap_steals += 1;
} else {
self.steal_locality_counters.remote_heap_steals += 1;
}
self.invariant_monitor.lock().record_task_dispatch(
first_task,
Time::from_nanos(self.current_time_ns()),
);
let steal_back_into_local_ready =
stolen_count > 1 && self.local.lock().peek_ready_priority().is_some();
if stolen_count > 1 {
if steal_back_into_local_ready {
let mut local = self.local.lock();
for &(task, priority) in &self.steal_buffer[1..stolen_count] {
local.schedule(task, priority);
self.invariant_monitor.lock().record_task_requeue(
task,
"local_ready_stolen",
priority,
Time::from_nanos(self.current_time_ns()),
);
}
} else {
for &(task, priority) in
self.steal_buffer[1..stolen_count].iter().rev()
{
self.fast_queue.push(task);
self.invariant_monitor.lock().record_task_requeue(
task,
"fast_queue_stolen",
priority,
Time::from_nanos(self.current_time_ns()),
);
}
}
}
return Some(first_task);
}
}
}
}
None
}
#[doc(hidden)]
#[cfg(feature = "test-internals")]
pub fn steal_once_for_test(&mut self) -> Option<TaskId> {
self.try_steal()
}
pub fn schedule_local(&self, task: TaskId, priority: u8) {
let should_schedule = self.with_task_table_ref(|tt| {
tt.task(task).is_none_or(|record| {
if record.is_local() {
error!(
?task,
"schedule_local: refusing to enqueue local task into PriorityScheduler"
);
return false;
}
record.wake_state.notify()
})
});
if should_schedule {
let mut local = self.local.lock();
local.schedule(task, priority);
let current_time = self.current_time_ns();
self.fairness_monitor.lock().record_task_enqueue(
task,
priority,
current_time,
2, );
self.invariant_monitor.lock().record_task_enqueue(
task,
"local_ready_heap",
priority,
Time::from_nanos(current_time),
);
self.record_scheduler_evidence_enqueue_at(task, current_time);
self.parker.unpark();
}
}
pub fn schedule_local_cancel(&self, task: TaskId, priority: u8) {
self.with_task_table_ref(|tt| {
if let Some(record) = tt.task(task) {
record.wake_state.notify();
}
});
move_local_ready_task_to_cancel_lane(&self.local, &self.local_ready, task, priority);
let current_time = self.current_time_ns();
self.fairness_monitor.lock().record_task_enqueue(
task,
priority,
current_time,
0, );
self.invariant_monitor.lock().record_task_requeue(
task,
"local_cancel_queue",
priority,
Time::from_nanos(current_time),
);
self.record_scheduler_evidence_enqueue_at(task, current_time);
self.parker.unpark();
}
pub fn schedule_local_timed(&self, task: TaskId, deadline: Time) {
let should_schedule = self.with_task_table_ref(|tt| {
tt.task(task).is_none_or(|record| {
if record.is_local() {
error!(
?task,
"schedule_local_timed: refusing to enqueue local task into timed lane"
);
return false;
}
record.wake_state.notify()
})
});
if should_schedule {
let mut local = self.local.lock();
local.schedule_timed(task, deadline);
let current_time = self.current_time_ns();
self.fairness_monitor.lock().record_task_enqueue(
task,
0, current_time,
1, );
self.invariant_monitor.lock().record_task_enqueue(
task,
"local_timed_queue",
0, Time::from_nanos(current_time),
);
self.record_scheduler_evidence_enqueue_at(task, current_time);
self.parker.unpark();
}
}
fn waiter_wake_metadata(
&self,
state: &RuntimeState,
waiter: TaskId,
) -> Option<WaiterWakeMetadata> {
if let Some(tt) = &self.task_table {
let guard = tt.lock().unwrap_or_else(std::sync::PoisonError::into_inner);
let record = guard.task(waiter)?;
Some(WaiterWakeMetadata {
priority: record.sched_priority,
is_local: record.is_local(),
pinned_worker: record.pinned_worker(),
wake_state: Arc::clone(&record.wake_state),
notified: record.wake_state.notify(),
})
} else {
let record = state.task(waiter)?;
Some(WaiterWakeMetadata {
priority: record.sched_priority,
is_local: record.is_local(),
pinned_worker: record.pinned_worker(),
wake_state: Arc::clone(&record.wake_state),
notified: record.wake_state.notify(),
})
}
}
fn wake_dependents_locked(
&self,
state: &RuntimeState,
waiters: impl IntoIterator<Item = TaskId>,
) {
let mut global_tasks = smallvec::SmallVec::<[(TaskId, u8); 16]>::new();
for waiter in waiters {
let Some(metadata) = self.waiter_wake_metadata(state, waiter) else {
continue;
};
if metadata.notified {
if metadata.is_local {
if let Some(worker_id) = metadata.pinned_worker {
if let Some(queue) = self.all_local_ready.get(worker_id) {
queue.lock().push_back(waiter);
self.record_scheduler_evidence_enqueue(waiter);
self.coordinator.wake_worker(worker_id);
} else {
metadata.wake_state.clear();
error!(
?waiter,
worker_id,
"execute: pinned local waiter has invalid worker id, wake skipped and wake_state cleared"
);
}
} else {
self.local_ready.lock().push_back(waiter);
self.record_scheduler_evidence_enqueue(waiter);
self.parker.unpark();
}
} else {
global_tasks.push((waiter, metadata.priority));
}
}
}
let global_wakes = global_tasks.len();
if global_wakes > 0 {
let mut reservation = self.global.reserve_ready_count(global_wakes);
for (task, priority) in global_tasks {
self.global.inject_ready_uncounted(task, priority);
self.record_scheduler_evidence_enqueue(task);
reservation.publish_one();
}
self.coordinator.wake_many(global_wakes);
}
}
#[allow(clippy::too_many_lines)]
pub(crate) fn execute(&mut self, task_id: TaskId) {
struct TaskExecutionGuard<'a> {
worker: &'a ThreeLaneWorker,
task_id: TaskId,
completed: bool,
}
impl Drop for TaskExecutionGuard<'_> {
#[allow(clippy::significant_drop_tightening)] fn drop(&mut self) {
if !self.completed && std::thread::panicking() {
let cleanup = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
let artifacts =
self.worker.complete_task_after_unwind_ordered(self.task_id);
artifacts.dispatch_post_lock(self.worker);
}));
if let Err(payload) = cleanup {
std::mem::forget(payload);
}
}
}
}
Self::emit_cancel_diagnostic(|| {
trace!(task_id = ?task_id, worker_id = self.id, "executing task");
});
let (
mut stored,
wake_state,
priority,
task_cx,
cx_inner,
cached_waker,
cached_cancel_waker,
) = {
let merged = self.with_task_table(|tt| {
let global_stored = tt.remove_stored_future(task_id)?;
tt.update_task(task_id, |record| {
record.start_running();
record.wake_state.begin_poll();
let priority = record.sched_priority;
let wake_state = Arc::clone(&record.wake_state);
let task_cx = record.cx.clone();
let cached_waker = record.cached_waker.take();
let cached_cancel_waker = record.cached_cancel_waker.take();
let both_cached = cached_waker.is_some()
&& cached_cancel_waker
.as_ref()
.is_some_and(|(_, p)| *p == priority);
let cx_inner = if both_cached {
None
} else {
record.cx_inner.clone()
};
(
AnyStoredTask::Global(global_stored),
wake_state,
priority,
task_cx,
cx_inner,
cached_waker,
cached_cancel_waker,
)
})
});
if let Some(result) = merged {
result
} else {
let local = crate::runtime::local::remove_local_task(task_id);
let Some(local) = local else {
return;
};
let record_info = self.with_task_table(|tt| {
tt.update_task(task_id, |record| {
record.start_running();
record.wake_state.begin_poll();
let priority = record.sched_priority;
let wake_state = Arc::clone(&record.wake_state);
let task_cx = record.cx.clone();
let cached_waker = record.cached_waker.take();
let cached_cancel_waker = record.cached_cancel_waker.take();
let both_cached = cached_waker.is_some()
&& cached_cancel_waker
.as_ref()
.is_some_and(|(_, p)| *p == priority);
let cx_inner = if both_cached {
None
} else {
record.cx_inner.clone()
};
(
wake_state,
priority,
task_cx,
cx_inner,
cached_waker,
cached_cancel_waker,
)
})
});
let Some((
wake_state,
priority,
task_cx,
cx_inner,
cached_waker,
cached_cancel_waker,
)) = record_info
else {
return;
};
(
AnyStoredTask::Local(local),
wake_state,
priority,
task_cx,
cx_inner,
cached_waker,
cached_cancel_waker,
)
}
};
let is_local = stored.is_local();
let waker = if let Some((w, _)) = cached_waker {
w
} else {
let inner = cx_inner.as_ref().expect("cx_inner missing");
let fast_cancel = Arc::clone(&inner.read().fast_cancel);
let weak_inner = Arc::downgrade(inner);
if is_local {
Waker::from(Arc::new(ThreeLaneLocalWaker {
task_id,
priority,
wake_state: Arc::clone(&wake_state),
local: Arc::clone(&self.local),
local_ready: Arc::clone(&self.local_ready),
parker: self.parker.clone(),
fast_cancel,
cx_inner: weak_inner,
scheduler_evidence: self.scheduler_evidence.clone(),
}))
} else {
Waker::from(Arc::new(ThreeLaneWaker {
task_id,
wake_state: Arc::clone(&wake_state),
global: Arc::clone(&self.global),
coordinator: Arc::clone(&self.coordinator),
priority,
fast_cancel,
cx_inner: weak_inner,
scheduler_evidence: self.scheduler_evidence.clone(),
}))
}
};
let cancel_waker_for_cache = if cached_cancel_waker
.as_ref()
.is_some_and(|(_, p)| *p == priority)
{
cached_cancel_waker.map(|(w, _)| (w, priority))
} else {
cx_inner.as_ref().map(|inner| {
let w = if is_local {
Waker::from(Arc::new(ThreeLaneLocalCancelWaker {
task_id,
default_priority: priority,
wake_state: Arc::clone(&wake_state),
local: Arc::clone(&self.local),
local_ready: Arc::clone(&self.local_ready),
parker: self.parker.clone(),
cx_inner: Arc::downgrade(inner),
scheduler_evidence: self.scheduler_evidence.clone(),
}))
} else {
Waker::from(Arc::new(CancelLaneWaker {
task_id,
default_priority: priority,
wake_state: Arc::clone(&wake_state),
global: Arc::clone(&self.global),
coordinator: Arc::clone(&self.coordinator),
cx_inner: Arc::downgrade(inner),
scheduler_evidence: self.scheduler_evidence.clone(),
}))
};
let mut incoming_waker = Some(Arc::new(
crate::types::task_context::CancelWaker::new(w.clone()),
));
let retired_waker = {
let mut guard = inner.write();
if guard.cancel_waker_registry_closed {
None
} else if !guard
.cancel_waker
.as_ref()
.is_some_and(|existing| existing.will_wake(&w))
{
std::mem::replace(&mut guard.cancel_waker, incoming_waker.take())
} else {
None
}
};
drop(retired_waker);
drop(incoming_waker);
(w, priority)
})
};
let _cx_guard = crate::cx::Cx::set_current(task_cx);
let mut guard = TaskExecutionGuard {
worker: self,
task_id,
completed: false,
};
let poll_result = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
let mut cx = Context::from_waker(&waker);
stored.poll(&mut cx)
}));
let mut credit_adaptive_epoch = true;
match poll_result {
Ok(Poll::Ready(outcome)) => {
if matches!(outcome, crate::types::Outcome::Panicked(_)) {
credit_adaptive_epoch = false;
}
let task_outcome = outcome
.map_err(|()| crate::error::Error::new(crate::error::ErrorKind::Internal));
let artifacts = self
.complete_polled_task_ordered(task_id, PolledCompletion::Ready(task_outcome));
guard.completed = true;
wake_state.clear();
artifacts.dispatch_post_lock(self);
}
Ok(Poll::Pending) => {
let cancel_effects = match stored {
AnyStoredTask::Global(t) => self.with_task_table(move |tt| {
tt.store_spawned_task(task_id, t);
tt.update_task(task_id, |record| {
record.cached_waker = Some((waker, priority));
record.cached_cancel_waker = cancel_waker_for_cache;
record.consume_checkpoint_cancel_ack()
})
.unwrap_or_else(|| {
crate::types::task_context::CancellationEffects::ready(None)
})
}),
AnyStoredTask::Local(t) => {
crate::runtime::local::store_local_task(task_id, t);
self.with_task_table(move |tt| {
tt.update_task(task_id, |record| {
record.cached_waker = Some((waker, priority));
record.cached_cancel_waker = cancel_waker_for_cache;
record.consume_checkpoint_cancel_ack()
})
.unwrap_or_else(|| {
crate::types::task_context::CancellationEffects::ready(None)
})
})
}
};
let (cancel_ack, mut cancel_wakes) = cancel_effects.into_parts();
if let Some(receipt) = cancel_ack.as_ref() {
let state = self
.state
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
if self.task_table.is_some() {
if let Some(validation_result) = state
.external_checkpoint_cancel_materialization_violation(
task_id, receipt, false,
)
{
cancel_wakes.push_cancel_protocol_violation(
"external-shard checkpoint cancellation materialization",
validation_result,
);
}
} else if let Some(validation_result) =
state.checkpoint_cancel_materialization_violation(task_id, receipt)
{
cancel_wakes.push_cancel_protocol_violation(
"checkpoint cancellation materialization",
validation_result,
);
}
}
let acknowledged_priority =
cancel_ack.as_ref().map(|receipt| receipt.cleanup_priority);
if wake_state.finish_poll() || cancel_ack.is_some() {
let mut cancel_priority = acknowledged_priority.unwrap_or(priority);
let mut schedule_cancel = acknowledged_priority.is_some();
let cx_inner_for_finish = if cx_inner.is_some() {
cx_inner
} else {
self.with_task_table(|tt| tt.task(task_id).and_then(|r| r.cx_inner.clone()))
};
if !schedule_cancel && let Some(inner) = cx_inner_for_finish.as_ref() {
let guard = inner.read();
if guard.cancel_requested {
schedule_cancel = true;
if let Some(reason) = guard.cancel_reason.as_ref() {
cancel_priority = reason.cleanup_budget().priority;
}
}
}
if is_local {
if schedule_cancel {
move_local_ready_task_to_cancel_lane(
&self.local,
&self.local_ready,
task_id,
cancel_priority,
);
self.record_scheduler_evidence_enqueue(task_id);
} else {
self.local_ready.lock().push_back(task_id);
self.record_scheduler_evidence_enqueue(task_id);
}
self.parker.unpark();
} else {
if schedule_cancel {
self.global.inject_cancel(task_id, cancel_priority);
} else {
self.global.inject_ready(task_id, priority);
}
self.record_scheduler_evidence_enqueue(task_id);
self.coordinator.wake_one();
}
}
guard.completed = true;
cancel_wakes.dispatch();
}
Err(payload) => {
credit_adaptive_epoch = false;
let panic_message = crate::cx::scope::payload_to_string(&payload);
std::mem::forget(payload);
let panic_payload = crate::types::outcome::PanicPayload::new(panic_message);
let panic_outcome = crate::types::Outcome::Panicked(panic_payload);
let artifacts = self.complete_polled_task_ordered(
task_id,
PolledCompletion::Panicked(panic_outcome),
);
guard.completed = true;
wake_state.clear();
artifacts.dispatch_post_lock(self);
}
}
drop(guard);
if credit_adaptive_epoch {
self.adaptive_on_dispatch();
} else {
self.abort_adaptive_epoch();
}
}
fn apply_polled_completion(
record: &mut crate::record::task::TaskRecord,
completion: PolledCompletion,
cancel_ack: bool,
) {
match completion {
PolledCompletion::Ready(outcome) => {
Self::complete_polled_record(record, outcome, cancel_ack);
}
PolledCompletion::Panicked(outcome) => {
if !record.state.is_terminal() {
record.complete(outcome);
}
}
}
}
fn drain_ready_finalizers_locked(
&self,
state: &mut RuntimeState,
) -> smallvec::SmallVec<[(TaskId, u8, crate::runtime::state::TaskSpawnEffects); 2]> {
if !state.has_finalizing_regions() {
return smallvec::SmallVec::new();
}
let mut finalizer_tasks = self.task_table.as_ref().map_or(
crate::runtime::state::AdmissionTaskTarget::Embedded,
|tt| {
crate::runtime::state::AdmissionTaskTarget::External(
tt.lock().unwrap_or_else(std::sync::PoisonError::into_inner),
)
},
);
state.drain_ready_async_finalizers_in(
&crate::runtime::state::AdmissionRegionTarget::Embedded,
&mut finalizer_tasks,
)
}
fn complete_polled_task_ordered(
&self,
task_id: TaskId,
completion: PolledCompletion,
) -> PolledCompletionArtifacts {
if self.task_table.is_some() {
let (cancel_ack, mut cancel_wakes, mut detached_record) = self.with_task_table(|tt| {
let effects = Self::consume_cancel_ack_from_table(tt, task_id);
let (cancel_ack, cancel_wakes) = effects.into_parts();
let ack_materialized = cancel_ack.is_some();
let _ = tt.update_task(task_id, |record| {
Self::apply_polled_completion(record, completion, ack_materialized);
});
let detached_record = tt.remove_task(task_id);
(cancel_ack, cancel_wakes, detached_record)
});
let mut state = self
.state
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
if let Some(receipt) = cancel_ack.as_ref()
&& let Some(validation_result) = state
.external_checkpoint_cancel_materialization_violation(task_id, receipt, true)
{
cancel_wakes.push_cancel_protocol_violation(
"external-shard checkpoint cancellation materialization",
validation_result,
);
}
let completion_effects = match detached_record.as_mut() {
Some(record) => state.task_completed_from_external_record(record),
None => state.task_completed(task_id),
};
let (waiters, completion_observer) = completion_effects.into_parts();
let finalizers = self.drain_ready_finalizers_locked(&mut state);
self.wake_dependents_locked(&state, waiters);
let finalizer_publication = self.publish_ready_finalizers(finalizers);
drop(state);
PolledCompletionArtifacts {
completion_observer,
cancel_wakes,
detached_record,
finalizer_publication,
}
} else {
let mut state = self
.state
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
let (cancel_ack, cancel_wakes) =
Self::consume_cancel_ack_locked(&mut state, task_id).into_parts();
let ack_materialized = cancel_ack.is_some();
let _ = state.update_task(task_id, |record| {
Self::apply_polled_completion(record, completion, ack_materialized);
});
let (waiters, completion_observer) = state.task_completed(task_id).into_parts();
let finalizers = self.drain_ready_finalizers_locked(&mut state);
self.wake_dependents_locked(&state, waiters);
let finalizer_publication = self.publish_ready_finalizers(finalizers);
drop(state);
PolledCompletionArtifacts {
completion_observer,
cancel_wakes,
detached_record: None,
finalizer_publication,
}
}
}
fn complete_task_after_unwind_ordered(&self, task_id: TaskId) -> UnwindCompletionArtifacts {
let panic_outcome = crate::types::Outcome::Panicked(
crate::types::outcome::PanicPayload::new("task panicked during scheduler bookkeeping"),
);
let mut detached_record = if self.task_table.is_some() {
self.with_task_table(|tt| {
let _ = tt.update_task(task_id, |record| {
Self::apply_polled_completion(
record,
PolledCompletion::Panicked(panic_outcome),
false,
);
});
tt.remove_task(task_id)
})
} else {
self.with_task_table(|tt| {
let _ = tt.update_task(task_id, |record| {
Self::apply_polled_completion(
record,
PolledCompletion::Panicked(panic_outcome),
false,
);
});
});
None
};
let mut state = self
.state
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
let completion = match detached_record.as_mut() {
Some(record) => state.task_completed_from_external_record(record),
None => state.task_completed(task_id),
};
let (waiters, cancel_waker_retirements) =
completion.into_waiters_and_retirements_without_observers();
let finalizers = self.drain_ready_finalizers_locked(&mut state);
self.wake_dependents_locked(&state, waiters);
let finalizer_publication = self.publish_ready_finalizers(finalizers);
drop(state);
UnwindCompletionArtifacts {
cancel_waker_retirements,
detached_record,
finalizer_publication,
}
}
fn publish_ready_finalizers(
&self,
finalizers: smallvec::SmallVec<[(TaskId, u8, crate::runtime::state::TaskSpawnEffects); 2]>,
) -> (
smallvec::SmallVec<[TaskId; 2]>,
smallvec::SmallVec<[crate::runtime::state::TaskSpawnEffects; 2]>,
) {
let finalizer_wakes = finalizers.len();
if finalizer_wakes == 0 {
return (smallvec::SmallVec::new(), smallvec::SmallVec::new());
}
let mut tasks = smallvec::SmallVec::new();
let mut spawn_effects = smallvec::SmallVec::new();
let mut reservation = self.global.reserve_ready_count(finalizer_wakes);
for (finalizer_task, priority, effects) in finalizers {
self.global.inject_ready_uncounted(finalizer_task, priority);
reservation.publish_one();
tasks.push(finalizer_task);
spawn_effects.push(effects);
}
(tasks, spawn_effects)
}
fn finish_ready_finalizer_publication(
&self,
(tasks, spawn_effects): (
smallvec::SmallVec<[TaskId; 2]>,
smallvec::SmallVec<[crate::runtime::state::TaskSpawnEffects; 2]>,
),
) {
for &task in &tasks {
Self::emit_cancel_diagnostic(|| self.record_scheduler_evidence_enqueue(task));
}
Self::emit_cancel_diagnostic(|| self.coordinator.wake_many(tasks.len()));
for effects in spawn_effects {
effects.dispatch();
}
}
fn retire_detached_task_record(record: Option<crate::record::task::TaskRecord>) {
let Some(record) = record else {
return;
};
if let Err(payload) = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| {
drop(record);
})) {
std::mem::forget(payload);
}
}
fn schedule_ready_finalizers(&self) -> bool {
let tasks = {
let mut state = self
.state
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
self.drain_ready_finalizers_locked(&mut state)
};
if tasks.is_empty() {
return false;
}
let publication = self.publish_ready_finalizers(tasks);
self.finish_ready_finalizer_publication(publication);
true
}
#[allow(dead_code)] fn consume_cancel_ack(
&self,
task_id: TaskId,
) -> crate::types::task_context::CancellationEffects<
Option<crate::record::task::CheckpointCancelAck>,
> {
let effects = self.with_task_table(|tt| Self::consume_cancel_ack_from_table(tt, task_id));
let (receipt, mut wakes) = effects.into_parts();
if let Some(receipt) = receipt.as_ref() {
let state = self
.state
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
if self.task_table.is_some() {
let task_still_live = self.with_task_table_ref(|tt| tt.task(task_id).is_some());
if let Some(validation_result) = state
.external_checkpoint_cancel_materialization_violation(
task_id,
receipt,
!task_still_live,
)
{
wakes.push_cancel_protocol_violation(
"external-shard checkpoint cancellation materialization",
validation_result,
);
}
} else if let Some(validation_result) =
state.checkpoint_cancel_materialization_violation(task_id, receipt)
{
wakes.push_cancel_protocol_violation(
"checkpoint cancellation materialization",
validation_result,
);
}
}
crate::types::task_context::CancellationEffects::new(receipt, wakes)
}
fn consume_cancel_ack_locked(
state: &mut RuntimeState,
task_id: TaskId,
) -> crate::types::task_context::CancellationEffects<
Option<crate::record::task::CheckpointCancelAck>,
> {
state.consume_task_checkpoint_cancel_ack(task_id)
}
fn consume_cancel_ack_from_table(
tt: &mut TaskTable,
task_id: TaskId,
) -> crate::types::task_context::CancellationEffects<
Option<crate::record::task::CheckpointCancelAck>,
> {
tt.update_task(
task_id,
crate::record::task::TaskRecord::consume_checkpoint_cancel_ack,
)
.unwrap_or_else(|| crate::types::task_context::CancellationEffects::ready(None))
}
fn complete_polled_record(
record: &mut crate::record::task::TaskRecord,
task_outcome: crate::types::Outcome<(), crate::error::Error>,
cancel_ack: bool,
) {
if record.state.is_terminal() {
return;
}
let mut completed_via_cancel = false;
if matches!(task_outcome, crate::types::Outcome::Ok(())) {
let should_cancel = matches!(
record.state,
crate::record::task::TaskState::Cancelling { .. }
| crate::record::task::TaskState::Finalizing { .. }
) || (cancel_ack
&& matches!(
record.state,
crate::record::task::TaskState::CancelRequested { .. }
));
if should_cancel {
if matches!(
record.state,
crate::record::task::TaskState::CancelRequested { .. }
) {
let _ = record.acknowledge_cancel();
}
if matches!(
record.state,
crate::record::task::TaskState::Cancelling { .. }
) {
record.cleanup_done();
}
if matches!(
record.state,
crate::record::task::TaskState::Finalizing { .. }
) {
record.finalize_done();
}
completed_via_cancel = matches!(
record.state,
crate::record::task::TaskState::Completed(crate::types::Outcome::Cancelled(_))
);
}
}
if !completed_via_cancel {
record.complete(task_outcome);
}
}
}
struct ThreeLaneWaker {
task_id: TaskId,
wake_state: Arc<crate::record::task::TaskWakeState>,
global: Arc<GlobalInjector>,
coordinator: Arc<WorkerCoordinator>,
priority: u8,
fast_cancel: std::sync::Arc<std::sync::atomic::AtomicBool>,
cx_inner: Weak<RwLock<CxInner>>,
scheduler_evidence: Option<Arc<Mutex<SchedulerEvidenceCollector>>>,
}
impl ThreeLaneWaker {
#[inline]
fn schedule(&self) {
if self.wake_state.notify() {
let mut priority = self.priority;
let is_cancelling = self.fast_cancel.load(Ordering::Acquire);
if is_cancelling {
if let Some(inner) = self.cx_inner.upgrade() {
let guard = inner.read();
if let Some(reason) = &guard.cancel_reason {
priority = reason.cleanup_budget().priority;
}
}
}
if is_cancelling {
self.global.inject_cancel(self.task_id, priority);
} else {
self.global.inject_ready(self.task_id, priority);
}
if let Some(collector) = &self.scheduler_evidence {
collector
.lock()
.record_task_enqueue(self.task_id, crate::time::wall_now().as_nanos());
}
self.coordinator.wake_one();
}
}
}
use std::task::Wake;
impl Wake for ThreeLaneWaker {
#[inline]
fn wake(self: Arc<Self>) {
self.schedule();
}
#[inline]
fn wake_by_ref(self: &Arc<Self>) {
self.schedule();
}
}
struct ThreeLaneLocalWaker {
task_id: TaskId,
priority: u8,
wake_state: Arc<crate::record::task::TaskWakeState>,
local: Arc<Mutex<PriorityScheduler>>,
local_ready: Arc<LocalReadyQueue>,
parker: Parker,
fast_cancel: std::sync::Arc<std::sync::atomic::AtomicBool>,
cx_inner: Weak<RwLock<CxInner>>,
scheduler_evidence: Option<Arc<Mutex<SchedulerEvidenceCollector>>>,
}
impl ThreeLaneLocalWaker {
#[inline]
fn schedule(&self) {
if self.wake_state.notify() {
let is_cancelling = self.fast_cancel.load(Ordering::Acquire);
if is_cancelling {
let mut priority = self.priority;
if let Some(inner) = self.cx_inner.upgrade() {
let guard = inner.read();
if let Some(reason) = &guard.cancel_reason {
priority = reason.cleanup_budget().priority;
}
}
move_local_ready_task_to_cancel_lane(
&self.local,
&self.local_ready,
self.task_id,
priority,
);
} else {
self.local_ready.lock().push_back(self.task_id);
}
if let Some(collector) = &self.scheduler_evidence {
collector
.lock()
.record_task_enqueue(self.task_id, crate::time::wall_now().as_nanos());
}
self.parker.unpark();
}
}
}
impl Wake for ThreeLaneLocalWaker {
#[inline]
fn wake(self: Arc<Self>) {
self.schedule();
}
#[inline]
fn wake_by_ref(self: &Arc<Self>) {
self.schedule();
}
}
struct CancelLaneWaker {
task_id: TaskId,
default_priority: u8,
wake_state: Arc<crate::record::task::TaskWakeState>,
global: Arc<GlobalInjector>,
coordinator: Arc<WorkerCoordinator>,
cx_inner: Weak<RwLock<CxInner>>,
scheduler_evidence: Option<Arc<Mutex<SchedulerEvidenceCollector>>>,
}
impl CancelLaneWaker {
#[inline]
fn schedule(&self) {
let Some(inner) = self.cx_inner.upgrade() else {
return;
};
let (cancel_requested, priority) = {
let guard = inner.read();
let priority = guard
.cancel_reason
.as_ref()
.map_or(self.default_priority, |reason| {
reason.cleanup_budget().priority
});
(guard.cancel_requested, priority)
};
if !cancel_requested {
return;
}
self.wake_state.notify();
self.global.inject_cancel(self.task_id, priority);
if let Some(collector) = &self.scheduler_evidence {
collector
.lock()
.record_task_enqueue(self.task_id, crate::time::wall_now().as_nanos());
}
self.coordinator.wake_one();
}
}
impl Wake for CancelLaneWaker {
#[inline]
fn wake(self: Arc<Self>) {
self.schedule();
}
#[inline]
fn wake_by_ref(self: &Arc<Self>) {
self.schedule();
}
}
struct ThreeLaneLocalCancelWaker {
task_id: TaskId,
default_priority: u8,
wake_state: Arc<crate::record::task::TaskWakeState>,
local: Arc<Mutex<PriorityScheduler>>,
local_ready: Arc<LocalReadyQueue>,
parker: Parker,
cx_inner: Weak<RwLock<CxInner>>,
scheduler_evidence: Option<Arc<Mutex<SchedulerEvidenceCollector>>>,
}
impl ThreeLaneLocalCancelWaker {
#[inline]
fn schedule(&self) {
let Some(inner) = self.cx_inner.upgrade() else {
return;
};
let (cancel_requested, priority) = {
let guard = inner.read();
let priority = guard
.cancel_reason
.as_ref()
.map_or(self.default_priority, |reason| {
reason.cleanup_budget().priority
});
(guard.cancel_requested, priority)
};
if !cancel_requested {
return;
}
self.wake_state.notify();
{
move_local_ready_task_to_cancel_lane(
&self.local,
&self.local_ready,
self.task_id,
priority,
);
}
if let Some(collector) = &self.scheduler_evidence {
collector
.lock()
.record_task_enqueue(self.task_id, crate::time::wall_now().as_nanos());
}
self.parker.unpark();
}
}
impl Wake for ThreeLaneLocalCancelWaker {
#[inline]
fn wake(self: Arc<Self>) {
self.schedule();
}
#[inline]
fn wake_by_ref(self: &Arc<Self>) {
self.schedule();
}
}
#[cfg(test)]
#[path = "three_lane_tests.rs"]
mod tests;
#[cfg(test)]
#[path = "three_lane_metamorphic.rs"]
mod three_lane_metamorphic;