use super::*;
use crate::{
runtime::task_runtime,
sched::{FairMode, SchedulePolicy, algorithm::SchedulingEntity, system::DispatchCharge},
thread::DeadlineEntity,
};
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub(crate) enum SchedulerClass {
Stop,
Deadline,
Realtime,
Fair,
}
pub(super) struct ClassEnqueue {
pub(super) membership: QueueMembershipClass,
pub(super) entity: SchedulingEntity,
pub(super) reason: EnqueueReason,
}
#[derive(Clone, Copy)]
pub(crate) struct ClassTick {
pub(crate) slice_expired: bool,
pub(crate) request_reschedule: bool,
}
impl SchedulerClass {
pub(super) const PICK_ORDER: [Self; 4] =
[Self::Stop, Self::Deadline, Self::Realtime, Self::Fair];
pub(crate) const fn for_policy(policy: SchedulePolicy) -> Self {
match policy {
SchedulePolicy::KernelStop => Self::Stop,
SchedulePolicy::Deadline(_) => Self::Deadline,
SchedulePolicy::Fifo { .. } | SchedulePolicy::RoundRobin { .. } => Self::Realtime,
SchedulePolicy::Fair { .. } => Self::Fair,
}
}
pub(crate) fn has_selectable_higher_class(
self,
run_queue: &RunQueue,
rt_eligibility: RtEligibility,
) -> bool {
let stop = run_queue.stop.is_some();
let deadline = run_queue.deadline.has_runnable();
let realtime =
matches!(rt_eligibility, RtEligibility::Runnable) && run_queue.rt.has_any_rt();
match self {
Self::Stop => false,
Self::Deadline => stop,
Self::Realtime => stop || deadline,
Self::Fair => stop || deadline || realtime,
}
}
pub(super) fn enqueue_task(
self,
run_queue: &mut RunQueue,
mut thread: QueuedThread,
reason: EnqueueReason,
current_fair: Option<FairEntity>,
) -> Result<ClassEnqueue, TaskError> {
if let SchedulingEntity::Fair(fair) = thread.active.entity_mut() {
let virtual_time = run_queue.virtual_time();
match reason {
EnqueueReason::Wake => {
let (queue_weight, current_weight) =
run_queue.fair_placement_weights(current_fair);
fair.place_after_activation(
virtual_time,
queue_weight.saturating_add(current_weight),
)?;
}
EnqueueReason::Preempted => {}
EnqueueReason::Yield => fair.yield_request(virtual_time),
EnqueueReason::Migrated | EnqueueReason::PolicyChanged => {
let (queue_weight, current_weight) =
run_queue.fair_placement_weights(current_fair);
fair.place_after_transfer(
virtual_time,
queue_weight.saturating_add(current_weight),
)?;
}
EnqueueReason::Replenished => fair.place_at_least(virtual_time),
}
if !matches!(reason, EnqueueReason::Wake | EnqueueReason::Yield)
&& fair.request_exhausted()
{
fair.renew_request();
}
}
let entity = thread.active.entity().clone();
let membership = match self {
Self::Stop => {
thread.migration_capable = false;
assert!(
run_queue.stop.replace(thread).is_none(),
"one CPU runqueue can own only one stopper task"
);
QueueMembershipClass::Stop
}
Self::Deadline => {
if thread.active.entity().deadline().is_none_or(|deadline| {
deadline.absolute_deadline_ns().is_none() || deadline.is_throttled()
}) {
return Err(TaskError::NotReady);
}
QueueMembershipClass::Deadline(run_queue.deadline.insert(thread))
}
Self::Realtime => QueueMembershipClass::Realtime(run_queue.rt.enqueue(thread, reason)),
Self::Fair => {
run_queue.fair.insert(thread);
QueueMembershipClass::Fair
}
};
Ok(ClassEnqueue {
membership,
entity,
reason,
})
}
pub(super) fn dequeue_task(
self,
run_queue: &mut RunQueue,
membership: QueueMembershipClass,
id: ThreadId,
) -> Option<QueuedThread> {
match (self, membership) {
(Self::Stop, QueueMembershipClass::Stop) => run_queue.stop.take(),
(Self::Deadline, QueueMembershipClass::Deadline(key)) => run_queue.deadline.remove(key),
(Self::Realtime, QueueMembershipClass::Realtime(key)) => run_queue.rt.remove(key),
(Self::Fair, QueueMembershipClass::Fair) => run_queue.fair.remove(id),
_ => task_runtime::fatal_invariant(0x5251_1001, id.as_u64() as usize),
}
}
pub(super) fn migrate_task_rq(
self,
run_queue: &mut RunQueue,
membership: QueueMembershipClass,
id: ThreadId,
timing_granularity_ns: u64,
) -> Option<QueuedThread> {
if self == Self::Stop {
return None;
}
let source_fair_context = run_queue
.queued_thread_including_current(id)
.and_then(|thread| thread.base_entity.fair())
.map(|fair| {
(
run_queue.virtual_time(),
run_queue
.max_fair_service_request_ns()
.unwrap_or(fair.service_request_ns())
.max(fair.service_request_ns()),
)
});
let mut thread = self.dequeue_task(run_queue, membership, id)?;
if let Some((source_virtual_time, rq_max_slice_ns)) = source_fair_context {
thread.active.base_entity_mut().capture_fair_migration(
source_virtual_time,
rq_max_slice_ns,
timing_granularity_ns,
);
}
Some(thread)
}
#[inline(always)]
pub(super) fn pick_task(
self,
run_queue: &mut RunQueue,
rt_eligibility: RtEligibility,
skip_delayed: bool,
protected_fair_current: Option<ThreadId>,
) -> Option<PickTaskResult> {
match self {
Self::Stop => {
let picked = run_queue.stop.take()?;
run_queue.mark_publication_dirty();
Some(PickTaskResult::Continue(PickedThread::Owned(picked)))
}
Self::Deadline => {
let picked = run_queue.deadline.select_first();
picked
.map(PickedThread::Linked)
.map(PickTaskResult::Continue)
}
Self::Realtime => {
let picked = matches!(rt_eligibility, RtEligibility::Runnable)
.then(|| run_queue.rt.select())
.flatten();
picked
.map(PickedThread::Linked)
.map(PickTaskResult::Continue)
}
Self::Fair => {
run_queue.update_fair_virtual_time(None);
let queue = &mut run_queue.fair;
let virtual_time = queue.virtual_time();
let (mut thread, starts_dispatch) = if let Some(thread) = protected_fair_current
.and_then(|current| queue.take_protected_current(current, virtual_time))
{
#[cfg(feature = "qperf-metrics")]
crate::diagnostics::counters::record_fair_pick_protected_current();
(thread, false)
} else {
let thread = match queue.pick_eligible(virtual_time, skip_delayed)? {
FairPick::Runnable(thread) => thread,
FairPick::Delayed(core) => {
run_queue.mark_publication_dirty();
return Some(PickTaskResult::Break(core));
}
};
(thread, true)
};
if starts_dispatch {
let shortest_competing_slice_ns = queue.min_service_request_ns();
let SchedulingEntity::Fair(fair) = thread.active.entity_mut() else {
unreachable!("FairRunQueue can select only Fair entities")
};
fair.set_slice_protection(shortest_competing_slice_ns);
}
run_queue.mark_publication_dirty();
Some(PickTaskResult::Continue(PickedThread::Owned(thread)))
}
}
}
pub(super) fn put_prev_task(
self,
run_queue: &mut RunQueue,
membership: QueueMembershipClass,
id: ThreadId,
) -> Result<SchedulingEntity, TaskError> {
match (self, membership) {
(Self::Deadline, QueueMembershipClass::Deadline(key)) => {
let (new_key, entity) = run_queue
.deadline
.put_prev_current(key)
.ok_or(TaskError::NotReady)?;
run_queue.replace_membership_class(id, QueueMembershipClass::Deadline(new_key));
Ok(entity)
}
(Self::Realtime, QueueMembershipClass::Realtime(key)) => run_queue
.rt
.put_prev_current(key)
.ok_or(TaskError::NotReady),
_ => Err(TaskError::InvalidConfiguration),
}
}
pub(super) fn set_next_task(self, run_queue: &mut RunQueue, picked: &PickedThread) {
match self {
Self::Deadline | Self::Realtime => {}
Self::Stop | Self::Fair => {
run_queue.unregister_membership(picked.id());
}
}
}
pub(crate) fn task_tick(
self,
run_queue: &mut RunQueue,
current: ThreadId,
policy: SchedulePolicy,
current_entity: &SchedulingEntity,
charge: DispatchCharge,
periodic_tick_ns: Option<u64>,
) -> ClassTick {
match self {
Self::Deadline => ClassTick {
slice_expired: charge.slice_expired,
request_reschedule: charge.slice_expired,
},
Self::Realtime => match policy {
SchedulePolicy::RoundRobin { .. } => {
let tick_ns = periodic_tick_ns.unwrap_or_else(|| {
task_runtime::fatal_invariant(0x5251_1012, current.as_u64() as usize)
});
let key = match run_queue.membership_class(current) {
Some(QueueMembershipClass::Realtime(key)) => key,
_ => task_runtime::fatal_invariant(0x5251_1010, current.as_u64() as usize),
};
let tick = run_queue
.rt
.task_tick_round_robin(key, policy, tick_ns)
.unwrap_or_else(|| {
task_runtime::fatal_invariant(0x5251_1010, current.as_u64() as usize)
});
ClassTick {
slice_expired: tick.quantum_expired,
request_reschedule: tick.request_reschedule,
}
}
SchedulePolicy::Fifo { .. } => ClassTick {
slice_expired: false,
request_reschedule: false,
},
_ => task_runtime::fatal_invariant(0x5251_1011, current.as_u64() as usize),
},
Self::Fair => ClassTick {
slice_expired: charge.slice_expired,
request_reschedule: fair_tick_requests_reschedule(
run_queue.has_fair(),
current_entity,
charge,
),
},
Self::Stop => ClassTick {
slice_expired: false,
request_reschedule: false,
},
}
}
pub(super) fn check_preempt_curr(
self,
current_policy: SchedulePolicy,
current_entity: &SchedulingEntity,
current_is_idle: bool,
wakee_policy: SchedulePolicy,
wakee_entity: &SchedulingEntity,
fair_virtual_time: u64,
) -> bool {
if current_is_idle {
return true;
}
match self {
Self::Stop => !matches!(current_policy, SchedulePolicy::KernelStop),
Self::Deadline => match current_policy {
SchedulePolicy::KernelStop => false,
SchedulePolicy::Deadline(_) => {
deadline_key(wakee_entity) < deadline_key(current_entity)
}
_ => true,
},
Self::Realtime => {
let wakee_priority = wakee_policy
.rt_priority()
.expect("RT wakee must carry a fixed priority");
match current_policy {
SchedulePolicy::KernelStop | SchedulePolicy::Deadline(_) => false,
SchedulePolicy::Fifo { priority: current }
| SchedulePolicy::RoundRobin {
priority: current, ..
} => wakee_priority > current,
SchedulePolicy::Fair { .. } => true,
}
}
Self::Fair => fair_wakeup_preempts(
current_policy,
current_entity,
wakee_policy,
wakee_entity,
fair_virtual_time,
),
}
}
}
pub(crate) fn wakeup_preempts(
current_policy: SchedulePolicy,
current_entity: &SchedulingEntity,
current_is_idle: bool,
wakee_policy: SchedulePolicy,
wakee_entity: &SchedulingEntity,
fair_virtual_time: u64,
) -> bool {
SchedulerClass::for_policy(wakee_policy).check_preempt_curr(
current_policy,
current_entity,
current_is_idle,
wakee_policy,
wakee_entity,
fair_virtual_time,
)
}
pub(crate) fn default_sync_wakeup_preempts(
current_policy: SchedulePolicy,
current_entity: &SchedulingEntity,
current_is_idle: bool,
wakee_policy: SchedulePolicy,
wakee_entity: &SchedulingEntity,
fair_virtual_time: u64,
) -> bool {
wakeup_preempts(
current_policy,
current_entity,
current_is_idle,
wakee_policy,
wakee_entity,
fair_virtual_time,
)
}
fn fair_wakeup_preempts(
current_policy: SchedulePolicy,
current_entity: &SchedulingEntity,
wakee_policy: SchedulePolicy,
wakee_entity: &SchedulingEntity,
fair_virtual_time: u64,
) -> bool {
match current_policy {
SchedulePolicy::KernelStop
| SchedulePolicy::Deadline(_)
| SchedulePolicy::Fifo { .. }
| SchedulePolicy::RoundRobin { .. } => false,
SchedulePolicy::Fair {
mode: current_mode, ..
} => {
let wakee_mode = match wakee_policy {
SchedulePolicy::Fair { mode, .. } => mode,
_ => unreachable!("fair scheduler class requires a fair policy"),
};
let wakee = wakee_entity
.fair()
.expect("fair policy must own a fair scheduling entity");
let current = current_entity
.fair()
.expect("fair policy must own a fair scheduling entity");
#[cfg(feature = "qperf-metrics")]
crate::diagnostics::counters::record_fair_wake_distances(
crate::sched::algorithm::virtual_delta(wakee.vruntime(), fair_virtual_time),
crate::sched::algorithm::virtual_delta(current.vruntime(), fair_virtual_time),
);
if wakee_mode == FairMode::Idle {
false
} else if current_mode == FairMode::Idle {
true
} else if wakee_mode == FairMode::Batch
|| wakee_entity
.fair()
.is_some_and(|fair| !fair.is_eligible(fair_virtual_time))
{
#[cfg(feature = "qperf-metrics")]
crate::diagnostics::counters::record_fair_wake_wakee_ineligible();
false
} else {
if !current.is_eligible(fair_virtual_time) {
#[cfg(feature = "qperf-metrics")]
crate::diagnostics::counters::record_fair_wake_current_ineligible();
true
} else if current.slice_is_protected() && !wakee.has_shorter_slice_than(current) {
#[cfg(feature = "qperf-metrics")]
crate::diagnostics::counters::record_fair_wake_current_protected();
false
} else {
let precedes = wakee.deadline_precedes(current);
#[cfg(feature = "qperf-metrics")]
crate::diagnostics::counters::record_fair_wake_deadline(precedes);
precedes
}
}
}
}
}
fn fair_tick_requests_reschedule(
has_queued_peer: bool,
current_entity: &SchedulingEntity,
charge: DispatchCharge,
) -> bool {
let fair = current_entity
.fair()
.expect("Fair task_tick requires a Fair current entity");
has_queued_peer && (charge.slice_expired || !fair.slice_is_protected())
}
fn deadline_key(entity: &SchedulingEntity) -> u64 {
entity
.deadline()
.and_then(DeadlineEntity::absolute_deadline_ns)
.expect("a runnable Deadline entity must own an absolute deadline")
}
#[cfg(test)]
mod tests {
use super::*;
use crate::sched::{Nice, algorithm::FairEntity};
fn fair(vruntime: u64, virtual_deadline: u64) -> SchedulingEntity {
SchedulingEntity::Fair(FairEntity::test_state(
Nice::ZERO,
FairMode::Normal,
vruntime,
virtual_deadline,
))
}
fn normal_fair_policy() -> SchedulePolicy {
SchedulePolicy::fair(Nice::ZERO, FairMode::Normal)
}
#[test]
fn fair_wakeup_obeys_linux_eevdf_eligibility_and_slice_protection() {
let current = fair(2_000, 3_000);
let wakee = fair(1_000, 1_500);
assert!(!fair_wakeup_preempts(
normal_fair_policy(),
¤t,
normal_fair_policy(),
&wakee,
2_000,
));
let current = fair(3_000, 3_100);
let wakee = fair(1_000, 3_500);
assert!(fair_wakeup_preempts(
normal_fair_policy(),
¤t,
normal_fair_policy(),
&wakee,
2_000,
));
let mut current = FairEntity::new(Nice::ZERO, FairMode::Normal, 100, 2_000);
current.set_slice_protection(None);
let wakee = FairEntity::new(Nice::ZERO, FairMode::Normal, 50, 1_000);
assert!(fair_wakeup_preempts(
normal_fair_policy(),
&SchedulingEntity::Fair(current),
normal_fair_policy(),
&SchedulingEntity::Fair(wakee),
2_000,
));
}
}