refactor(sched): unify realtime runqueue model (#2244)

Rename the FIFO-specific scheduler and per-CPU queue to represent the shared realtime scheduling class while preserving the currently supported FIFO policy behavior.

Harden realtime queue invariants by validating internal priorities, deleting exactly one matching task, moving only queued tasks during yield, and checking bitmap, bucket, and running-count consistency in diagnostic builds.

Replace the unbounded FIFO demo with a feature-gated per-CPU smoke test that deterministically covers same-priority yield ordering, blocking, wakeup, task exit, and remote CPU runqueues without leaving runnable realtime workers behind.

Validation includes x86_64, riscv64, and loongarch64 kernel builds; one- and two-vCPU FIFO smoke tests; scheduler policy, affinity, tracepoint, fork/signal, and RCU dunitests.

Refs: #760

Signed-off-by: longjin <longjin@dragonos.org>
This commit is contained in:
LoGin
2026-09-01 10:45:34 +08:00
committed by GitHub
parent 3466f18cf2
commit 0210af8096
6 changed files with 425 additions and 213 deletions
+1 -1
View File
@@ -31,7 +31,7 @@ kstack_protect = []
# initram
initram = []
# fifo_demo: 起一个使用FIFO测例的内核线程,每5秒打印1条消息
# fifo_demo: run a bounded per-CPU FIFO/realtime-runqueue smoke test
fifo_demo = []
# 运行时依赖项
+2 -2
View File
@@ -179,7 +179,7 @@ impl ProcessManager {
if running {
match old_class {
SchedClass::Realtime => {
crate::sched::fifo::FifoScheduler::put_prev_task(rq, pcb.clone())
crate::sched::realtime::RealtimeScheduler::put_prev_task(rq, pcb.clone())
}
SchedClass::Fair => {
crate::sched::fair::CompletelyFairScheduler::put_prev_task(rq, pcb.clone())
@@ -217,7 +217,7 @@ impl ProcessManager {
// Matches Linux __sched_setscheduler: after a running task changes its
// policy, set_next_task is required.
if running {
crate::sched::fifo::FifoScheduler::set_next_task(rq, pcb.clone());
crate::sched::realtime::RealtimeScheduler::set_next_task(rq, pcb.clone());
}
// check_class_changed → preemption check.
-174
View File
@@ -1,174 +0,0 @@
use alloc::{collections::VecDeque, sync::Arc, vec::Vec};
use crate::{process::ProcessControlBlock, sched::prio::MAX_RT_PRIO};
use super::{CpuRunQueue, DequeueFlag, EnqueueFlag, PrioUtil, SchedClass, Scheduler, WakeupFlags};
#[derive(Debug)]
pub struct FifoRunQueue {
queues: Vec<VecDeque<Arc<ProcessControlBlock>>>,
active: u128,
nr_running: usize,
}
impl FifoRunQueue {
pub fn new() -> Self {
let mut queues = Vec::with_capacity(MAX_RT_PRIO as usize);
queues.resize_with(MAX_RT_PRIO as usize, VecDeque::new);
Self {
queues,
active: 0,
nr_running: 0,
}
}
#[inline]
pub fn nr_running(&self) -> usize {
self.nr_running
}
#[inline]
fn prio_index(pcb: &ProcessControlBlock) -> usize {
let prio = pcb.sched_info().prio();
let prio = prio.clamp(0, MAX_RT_PRIO - 1);
prio as usize
}
#[inline]
fn set_active(&mut self, prio: usize) {
self.active |= 1u128 << prio;
}
#[inline]
fn clear_active_if_empty(&mut self, prio: usize) {
if self.queues[prio].is_empty() {
self.active &= !(1u128 << prio);
}
}
pub fn enqueue(&mut self, pcb: Arc<ProcessControlBlock>) {
let prio = Self::prio_index(&pcb);
self.queues[prio].push_back(pcb);
self.set_active(prio);
self.nr_running += 1;
}
pub fn dequeue(&mut self, pcb: &Arc<ProcessControlBlock>) -> bool {
let prio = Self::prio_index(pcb);
let q = &mut self.queues[prio];
let before = q.len();
q.retain(|p| !Arc::ptr_eq(p, pcb));
let removed = before != q.len();
if removed {
self.nr_running -= 1;
self.clear_active_if_empty(prio);
}
removed
}
pub fn yield_current(&mut self, pcb: &Arc<ProcessControlBlock>) {
let prio = Self::prio_index(pcb);
let q = &mut self.queues[prio];
if q.len() <= 1 {
return;
}
q.retain(|p| !Arc::ptr_eq(p, pcb));
q.push_back(pcb.clone());
self.set_active(prio);
}
pub fn pick_next(&self) -> Option<Arc<ProcessControlBlock>> {
let prio = self.highest_prio()?;
self.queues[prio].front().cloned()
}
pub fn highest_prio(&self) -> Option<usize> {
if self.active == 0 {
return None;
}
Some(self.active.trailing_zeros() as usize)
}
}
pub struct FifoScheduler;
impl FifoScheduler {
#[inline]
fn rt_prio(pcb: &ProcessControlBlock) -> i32 {
pcb.sched_info().prio()
}
/// 对标 Linux set_next_task_rt:将 FIFO 任务设为当前运行任务。
///
/// FIFO 不像 CFS 那样维护 per-class current entity
/// 但需要在 running 任务修改策略/优先级后被调用以保持与 Linux 一致的流程。
pub fn set_next_task(
_rq: &mut super::CpuRunQueue,
_pcb: alloc::sync::Arc<crate::process::ProcessControlBlock>,
) {
// FIFO 不维护 per-class 调度实体状态,无需额外操作
}
}
impl Scheduler for FifoScheduler {
fn enqueue(rq: &mut CpuRunQueue, pcb: Arc<ProcessControlBlock>, _flags: EnqueueFlag) {
rq.fifo.enqueue(pcb);
rq.add_nr_running(1);
}
fn dequeue(rq: &mut CpuRunQueue, pcb: Arc<ProcessControlBlock>, _flags: DequeueFlag) {
if rq.fifo.dequeue(&pcb) {
rq.sub_nr_running(1);
}
}
fn yield_task(rq: &mut CpuRunQueue) {
let curr = rq.current();
debug_assert_eq!(curr.sched_info().sched_class(), SchedClass::Realtime);
rq.fifo.yield_current(&curr);
rq.resched_current();
}
fn check_preempt_current(
rq: &mut CpuRunQueue,
pcb: &Arc<ProcessControlBlock>,
_flags: WakeupFlags,
) {
let curr = rq.current();
debug_assert_eq!(curr.sched_info().sched_class(), SchedClass::Realtime);
let new_prio = Self::rt_prio(pcb);
let curr_prio = Self::rt_prio(&curr);
if PrioUtil::rt_prio(new_prio) && PrioUtil::rt_prio(curr_prio) && new_prio < curr_prio {
rq.resched_current();
}
}
fn pick_task(rq: &mut CpuRunQueue) -> Option<Arc<ProcessControlBlock>> {
FifoScheduler::pick_next_task(rq, None)
}
fn pick_next_task(
rq: &mut CpuRunQueue,
_pcb: Option<Arc<ProcessControlBlock>>,
) -> Option<Arc<ProcessControlBlock>> {
rq.fifo.pick_next()
}
fn tick(rq: &mut CpuRunQueue, pcb: Arc<ProcessControlBlock>, _queued: bool) {
debug_assert_eq!(pcb.sched_info().sched_class(), SchedClass::Realtime);
let Some(highest) = rq.fifo.highest_prio() else {
return;
};
let curr_prio = Self::rt_prio(&pcb);
if PrioUtil::rt_prio(curr_prio) && (highest as i32) < curr_prio {
rq.resched_current();
}
}
fn task_fork(_pcb: Arc<ProcessControlBlock>) {}
fn put_prev_task(_rq: &mut CpuRunQueue, _prev: Arc<ProcessControlBlock>) {}
}
+156 -22
View File
@@ -1,38 +1,172 @@
#![allow(dead_code)]
use alloc::{boxed::Box, string::String};
use core::sync::atomic::{AtomicUsize, Ordering};
use alloc::{boxed::Box, string::String, sync::Arc, vec::Vec};
use crate::{
process::{
kthread::{KernelThreadClosure, KernelThreadMechanism},
ProcessManager,
ProcessControlBlock, ProcessManager,
},
sched::prio::MAX_RT_PRIO,
smp::cpu::ProcessorId,
time::{sleep::nanosleep, Duration, PosixTimeSpec},
sched::{completion::Completion, prio::MAX_RT_PRIO, OnRq},
smp::cpu::{smp_cpu_manager, ProcessorId},
time::{sleep::nanosleep, PosixTimeSpec},
};
pub fn fifo_demo_init() {
const YIELD_ROUNDS: usize = 16;
const EVENT_COUNT: usize = YIELD_ROUNDS * 2;
const EVENT_UNSET: usize = usize::MAX;
const FIFO_PRIO: i32 = MAX_RT_PRIO - 50;
struct FifoPairState {
first_start: Completion,
finished: [Completion; 2],
next_event: AtomicUsize,
events: Vec<AtomicUsize>,
resumed: AtomicUsize,
results: [AtomicUsize; 2],
}
impl FifoPairState {
fn new() -> Arc<Self> {
Arc::new(Self {
first_start: Completion::new(),
finished: [Completion::new(), Completion::new()],
next_event: AtomicUsize::new(0),
events: (0..EVENT_COUNT)
.map(|_| AtomicUsize::new(EVENT_UNSET))
.collect(),
resumed: AtomicUsize::new(0),
results: [AtomicUsize::new(EVENT_UNSET), AtomicUsize::new(EVENT_UNSET)],
})
}
}
fn worker_closure(worker: usize, state: Arc<FifoPairState>) -> KernelThreadClosure {
let closure: Box<dyn Fn() -> i32 + Send + Sync> = Box::new(move || {
let pcb = ProcessManager::current_pcb();
let mut result = 0usize;
// 设置CPU亲和性为Core 0
pcb.sched_info().set_on_cpu(Some(ProcessorId::new(0)));
// 设置调度策略为FIFO,优先级为50
ProcessManager::set_fifo_policy(&pcb, MAX_RT_PRIO - 50).expect("Failed to set FIFO policy");
loop {
log::info!("fifo is running");
// 睡眠5秒
let sleep_time = PosixTimeSpec::from(Duration::from_secs(5));
let _ = nanosleep(sleep_time);
if worker == 0 {
if state.first_start.wait_for_completion().is_err() {
result = 1;
}
} else {
state.first_start.complete();
}
if result == 0 {
for _ in 0..YIELD_ROUNDS {
let slot = state.next_event.fetch_add(1, Ordering::AcqRel);
if slot >= EVENT_COUNT {
result = 2;
break;
}
state.events[slot].store(worker, Ordering::Release);
crate::sched::sched_yield();
}
}
if result == 0 && nanosleep(PosixTimeSpec::new(0, 5_000_000)).is_err() {
result = 3;
}
if result == 0 {
state.resumed.fetch_or(1usize << worker, Ordering::AcqRel);
}
state.results[worker].store(result, Ordering::Release);
state.finished[worker].complete();
result as i32
});
let _ = KernelThreadMechanism::create_and_run(
KernelThreadClosure::EmptyClosure((closure, ())),
String::from("fifo_demo"),
KernelThreadClosure::EmptyClosure((closure, ()))
}
fn create_fifo_worker(
cpu: ProcessorId,
worker: usize,
state: Arc<FifoPairState>,
) -> Arc<ProcessControlBlock> {
let name = String::from(if worker == 0 {
"fifo_demo_a"
} else {
"fifo_demo_b"
});
let pcb = KernelThreadMechanism::create_on_cpu(worker_closure(worker, state), name, cpu)
.expect("fifo demo failed to create worker");
// create_on_cpu publishes the PCB before the child has necessarily
// finished its initial blocked switch-out. Wait for stable off-rq state
// before changing its scheduling policy.
pcb.sched_info().wait_until_not_running();
assert_eq!(
*pcb.sched_info().on_rq.lock_irqsave(),
OnRq::None,
"fifo demo worker did not become off-rq"
);
ProcessManager::set_fifo_policy(&pcb, FIFO_PRIO).expect("fifo demo failed to set FIFO policy");
pcb
}
fn run_fifo_pair(cpu: ProcessorId) {
let state = FifoPairState::new();
let first = create_fifo_worker(cpu, 0, state.clone());
let second = create_fifo_worker(cpu, 1, state.clone());
ProcessManager::wakeup(&first).expect("fifo demo failed to wake first worker");
ProcessManager::wakeup(&second).expect("fifo demo failed to wake second worker");
for completion in &state.finished {
completion
.wait_for_completion()
.expect("fifo demo worker completion failed");
}
assert_eq!(
state.next_event.load(Ordering::Acquire),
EVENT_COUNT,
"fifo demo recorded the wrong event count"
);
for (slot, event) in state.events.iter().enumerate() {
let expected = if slot % 2 == 0 { 1 } else { 0 };
assert_eq!(
event.load(Ordering::Acquire),
expected,
"fifo demo same-priority yield order mismatch at event {slot}"
);
}
assert_eq!(
state.resumed.load(Ordering::Acquire),
0b11,
"fifo demo worker did not resume after sleeping"
);
for (worker, pcb) in [first, second].iter().enumerate() {
assert_eq!(
state.results[worker].load(Ordering::Acquire),
0,
"fifo demo worker reported failure"
);
assert_eq!(
KernelThreadMechanism::stop(pcb),
Ok(0),
"fifo demo worker exited with failure"
);
}
log::info!("fifo_demo status=ok cpu={}", cpu.data());
}
pub fn fifo_demo_init() {
let cpus: Vec<ProcessorId> = smp_cpu_manager()
.present_cpus()
.iter_cpu()
.filter(|&cpu| smp_cpu_manager().is_online_cpu(cpu))
.take(2)
.collect();
assert!(!cpus.is_empty(), "fifo demo found no online CPU");
for cpu in cpus {
run_fifo_pair(cpu);
}
}
+14 -14
View File
@@ -3,7 +3,6 @@ pub mod completion;
pub mod cputime;
pub mod fair;
pub mod fair_tree;
pub mod fifo;
#[cfg(feature = "fifo_demo")]
pub mod fifo_demo;
pub mod idle;
@@ -11,6 +10,7 @@ pub mod loadavg;
pub mod pelt;
pub mod policy;
pub mod prio;
pub mod realtime;
pub mod syscall;
use core::{
@@ -56,8 +56,8 @@ use self::{
clock::{ClockUpdataFlag, SchedClock},
cputime::{irq_time_read, CpuTimeFunc, IrqTime},
fair::{CfsRunQueue, CompletelyFairScheduler, FairSchedEntity},
fifo::FifoScheduler,
prio::{PrioUtil, MAX_RT_PRIO},
realtime::RealtimeScheduler,
};
pub use policy::{LinuxSchedPolicy, SchedClass};
@@ -471,7 +471,7 @@ pub struct CpuRunQueue {
/// CFS调度器
cfs: Arc<CfsRunQueue>,
fifo: fifo::FifoRunQueue,
rt: realtime::RealtimeRunQueue,
clock_pelt: u64,
lost_idle_time: u64,
@@ -510,7 +510,7 @@ impl CpuRunQueue {
calc_load_update: clock() + (5 * HZ + 1),
calc_load_active: 0,
cfs: Arc::new(CfsRunQueue::new()),
fifo: fifo::FifoRunQueue::new(),
rt: realtime::RealtimeRunQueue::new(),
clock_pelt: 0,
lost_idle_time: 0,
clock_idle: 0,
@@ -587,7 +587,7 @@ impl CpuRunQueue {
}
match pcb.sched_info().sched_class() {
SchedClass::Realtime => FifoScheduler::enqueue(self, pcb, flags),
SchedClass::Realtime => RealtimeScheduler::enqueue(self, pcb, flags),
SchedClass::Fair => CompletelyFairScheduler::enqueue(self, pcb, flags),
SchedClass::Idle => IdleScheduler::enqueue(self, pcb, flags),
}
@@ -617,7 +617,7 @@ impl CpuRunQueue {
}
match pcb.sched_info().sched_class() {
SchedClass::Realtime => FifoScheduler::dequeue(self, pcb, flags),
SchedClass::Realtime => RealtimeScheduler::dequeue(self, pcb, flags),
SchedClass::Fair => CompletelyFairScheduler::dequeue(self, pcb, flags),
SchedClass::Idle => IdleScheduler::dequeue(self, pcb, flags),
}
@@ -651,7 +651,7 @@ impl CpuRunQueue {
SchedClass::Fair => {
CompletelyFairScheduler::check_preempt_current(self, pcb, flags)
}
SchedClass::Realtime => FifoScheduler::check_preempt_current(self, pcb, flags),
SchedClass::Realtime => RealtimeScheduler::check_preempt_current(self, pcb, flags),
SchedClass::Idle => IdleScheduler::check_preempt_current(self, pcb, flags),
}
} else if waking_class.outranks(current_class) {
@@ -696,7 +696,7 @@ impl CpuRunQueue {
} else if next_class == current_class {
match current_class {
SchedClass::Fair => {}
SchedClass::Realtime => FifoScheduler::check_preempt_current(self, pcb, flags),
SchedClass::Realtime => RealtimeScheduler::check_preempt_current(self, pcb, flags),
SchedClass::Idle => IdleScheduler::check_preempt_current(self, pcb, flags),
}
}
@@ -873,8 +873,8 @@ impl CpuRunQueue {
let mut next: Option<Arc<ProcessControlBlock>> = None;
if self.fifo.nr_running() > 0 {
next = FifoScheduler::pick_next_task(self, Some(prev.clone()));
if self.rt.nr_running() > 0 {
next = RealtimeScheduler::pick_next_task(self, Some(prev.clone()));
}
if next.is_none() {
@@ -893,7 +893,7 @@ impl CpuRunQueue {
if !Arc::ptr_eq(&prev, &next) {
match prev.sched_info().sched_class() {
SchedClass::Realtime => FifoScheduler::put_prev_task(self, prev),
SchedClass::Realtime => RealtimeScheduler::put_prev_task(self, prev),
SchedClass::Fair => CompletelyFairScheduler::put_prev_task(self, prev),
SchedClass::Idle => IdleScheduler::put_prev_task(self, prev),
}
@@ -1031,7 +1031,7 @@ pub fn scheduler_tick() {
// 更新请求队列时钟
rq.update_rq_clock();
match current.sched_info().sched_class() {
SchedClass::Realtime => FifoScheduler::tick(rq, current, false),
SchedClass::Realtime => RealtimeScheduler::tick(rq, current, false),
SchedClass::Fair => CompletelyFairScheduler::tick(rq, current, false),
SchedClass::Idle => IdleScheduler::tick(rq, current, false),
}
@@ -1321,7 +1321,7 @@ pub fn sched_cgroup_fork(pcb: &Arc<ProcessControlBlock>) {
__set_task_cpu(pcb, fork_cpu);
match pcb.sched_info().sched_class() {
SchedClass::Realtime => FifoScheduler::task_fork(pcb.clone()),
SchedClass::Realtime => RealtimeScheduler::task_fork(pcb.clone()),
SchedClass::Fair => CompletelyFairScheduler::task_fork(pcb.clone()),
SchedClass::Idle => unreachable!("a fork child cannot use the idle scheduling class"),
}
@@ -1491,7 +1491,7 @@ pub fn sched_yield() {
// TODO: schedstat_inc(rq->yld_count);
match pcb.sched_info().sched_class() {
SchedClass::Realtime => FifoScheduler::yield_task(rq),
SchedClass::Realtime => RealtimeScheduler::yield_task(rq),
SchedClass::Fair => CompletelyFairScheduler::yield_task(rq),
SchedClass::Idle => {}
}
+252
View File
@@ -0,0 +1,252 @@
use alloc::{collections::VecDeque, sync::Arc, vec::Vec};
use crate::{process::ProcessControlBlock, sched::prio::MAX_RT_PRIO};
use super::{CpuRunQueue, DequeueFlag, EnqueueFlag, PrioUtil, SchedClass, Scheduler, WakeupFlags};
const _: () = {
assert!(MAX_RT_PRIO > 1);
assert!(MAX_RT_PRIO <= u128::BITS as i32);
};
/// Per-CPU runqueue shared by every policy in the realtime scheduling class.
///
/// Lower numeric priorities run first. Tasks at the same priority keep FIFO
/// order unless the class explicitly moves an existing task to the tail.
#[derive(Debug)]
pub struct RealtimeRunQueue {
queues: Vec<VecDeque<Arc<ProcessControlBlock>>>,
active: u128,
nr_running: usize,
}
impl RealtimeRunQueue {
pub fn new() -> Self {
let mut queues = Vec::with_capacity(MAX_RT_PRIO as usize);
queues.resize_with(MAX_RT_PRIO as usize, VecDeque::new);
Self {
queues,
active: 0,
nr_running: 0,
}
}
#[inline]
pub fn nr_running(&self) -> usize {
self.nr_running
}
#[inline]
fn prio_index(pcb: &ProcessControlBlock) -> usize {
let prio = pcb.sched_info().prio();
assert!(
(0..MAX_RT_PRIO - 1).contains(&prio),
"realtime task has invalid internal priority {prio}"
);
prio as usize
}
#[inline]
fn set_active(&mut self, prio: usize) {
self.active |= 1u128 << prio;
}
#[inline]
fn clear_active_if_empty(&mut self, prio: usize) {
if self.queues[prio].is_empty() {
self.active &= !(1u128 << prio);
}
}
#[inline]
fn assert_not_queued(&self, _pcb: &Arc<ProcessControlBlock>) {
#[cfg(any(debug_assertions, feature = "fifo_demo"))]
assert!(
!self
.queues
.iter()
.flatten()
.any(|queued| Arc::ptr_eq(queued, _pcb)),
"realtime task is already queued"
);
}
#[inline]
fn assert_consistent(&self) {
#[cfg(any(debug_assertions, feature = "fifo_demo"))]
{
let mut expected_active = 0u128;
let mut expected_running = 0usize;
for (prio, queue) in self.queues.iter().enumerate() {
if !queue.is_empty() {
expected_active |= 1u128 << prio;
}
expected_running += queue.len();
for task in queue {
assert_eq!(
Self::prio_index(task),
prio,
"realtime task is queued in the wrong priority bucket"
);
}
}
assert_eq!(self.active, expected_active, "realtime bitmap mismatch");
assert_eq!(
self.nr_running, expected_running,
"realtime running count mismatch"
);
}
}
pub fn enqueue_tail(&mut self, pcb: Arc<ProcessControlBlock>) {
let prio = Self::prio_index(&pcb);
self.assert_not_queued(&pcb);
self.queues[prio].push_back(pcb);
self.set_active(prio);
self.nr_running += 1;
self.assert_consistent();
}
pub fn dequeue(&mut self, pcb: &Arc<ProcessControlBlock>) -> bool {
self.assert_consistent();
let prio = Self::prio_index(pcb);
let position = self.queues[prio]
.iter()
.position(|queued| Arc::ptr_eq(queued, pcb));
#[cfg(any(debug_assertions, feature = "fifo_demo"))]
assert!(position.is_some(), "realtime task is not queued");
let Some(position) = position else {
return false;
};
self.queues[prio].remove(position);
self.nr_running -= 1;
self.clear_active_if_empty(prio);
self.assert_consistent();
true
}
/// Move an existing task to the tail of its current priority bucket.
pub fn requeue_to_tail(&mut self, pcb: &Arc<ProcessControlBlock>) -> bool {
self.assert_consistent();
let prio = Self::prio_index(pcb);
let position = self.queues[prio]
.iter()
.position(|queued| Arc::ptr_eq(queued, pcb));
#[cfg(any(debug_assertions, feature = "fifo_demo"))]
assert!(position.is_some(), "realtime task is not queued");
let Some(position) = position else {
return false;
};
if position + 1 != self.queues[prio].len() {
let task = self.queues[prio]
.remove(position)
.expect("realtime queue position disappeared");
self.queues[prio].push_back(task);
}
self.assert_consistent();
true
}
pub fn pick_next(&self) -> Option<Arc<ProcessControlBlock>> {
let prio = self.highest_prio()?;
self.queues[prio].front().cloned()
}
pub fn highest_prio(&self) -> Option<usize> {
if self.active == 0 {
return None;
}
Some(self.active.trailing_zeros() as usize)
}
}
pub struct RealtimeScheduler;
impl RealtimeScheduler {
#[inline]
fn rt_prio(pcb: &ProcessControlBlock) -> i32 {
pcb.sched_info().prio()
}
/// Set a realtime task as the current task for its class.
///
/// DragonOS does not yet maintain a separate current RT entity, so there
/// is no class-local state to update here.
pub fn set_next_task(
_rq: &mut super::CpuRunQueue,
_pcb: alloc::sync::Arc<crate::process::ProcessControlBlock>,
) {
}
}
impl Scheduler for RealtimeScheduler {
fn enqueue(rq: &mut CpuRunQueue, pcb: Arc<ProcessControlBlock>, _flags: EnqueueFlag) {
rq.rt.enqueue_tail(pcb);
rq.add_nr_running(1);
}
fn dequeue(rq: &mut CpuRunQueue, pcb: Arc<ProcessControlBlock>, _flags: DequeueFlag) {
if rq.rt.dequeue(&pcb) {
rq.sub_nr_running(1);
}
}
fn yield_task(rq: &mut CpuRunQueue) {
let curr = rq.current();
debug_assert_eq!(curr.sched_info().sched_class(), SchedClass::Realtime);
if rq.rt.requeue_to_tail(&curr) {
rq.resched_current();
}
}
fn check_preempt_current(
rq: &mut CpuRunQueue,
pcb: &Arc<ProcessControlBlock>,
_flags: WakeupFlags,
) {
let curr = rq.current();
debug_assert_eq!(curr.sched_info().sched_class(), SchedClass::Realtime);
let new_prio = Self::rt_prio(pcb);
let curr_prio = Self::rt_prio(&curr);
if PrioUtil::rt_prio(new_prio) && PrioUtil::rt_prio(curr_prio) && new_prio < curr_prio {
rq.resched_current();
}
}
fn pick_task(rq: &mut CpuRunQueue) -> Option<Arc<ProcessControlBlock>> {
RealtimeScheduler::pick_next_task(rq, None)
}
fn pick_next_task(
rq: &mut CpuRunQueue,
_pcb: Option<Arc<ProcessControlBlock>>,
) -> Option<Arc<ProcessControlBlock>> {
rq.rt.pick_next()
}
fn tick(rq: &mut CpuRunQueue, pcb: Arc<ProcessControlBlock>, _queued: bool) {
debug_assert_eq!(pcb.sched_info().sched_class(), SchedClass::Realtime);
let Some(highest) = rq.rt.highest_prio() else {
return;
};
let curr_prio = Self::rt_prio(&pcb);
if PrioUtil::rt_prio(curr_prio) && (highest as i32) < curr_prio {
rq.resched_current();
}
}
fn task_fork(_pcb: Arc<ProcessControlBlock>) {}
fn put_prev_task(_rq: &mut CpuRunQueue, _prev: Arc<ProcessControlBlock>) {}
}