use super::*;
impl RunQueue {
pub(crate) fn put_prev_task(&mut self, id: ThreadId) -> Result<SchedulingEntity, TaskError> {
if self.linked_current() != Some(id) {
return Err(TaskError::NotReady);
}
let class = self.membership_class(id).ok_or(TaskError::NotReady)?;
let policy = self
.queued_thread_including_current(id)
.ok_or(TaskError::NotReady)?
.policy;
let entity = SchedulerClass::for_policy(policy).put_prev_task(self, class, id)?;
self.refresh_class_pushable(id, None);
self.mark_publication_dirty();
Ok(entity)
}
#[inline(always)]
pub(crate) fn yield_realtime_current(
&mut self,
id: ThreadId,
) -> Result<LinkedRqTaskRef, TaskError> {
let Some(QueueMembershipClass::Realtime(key)) = self.membership_class(id) else {
return Err(TaskError::InvalidConfiguration);
};
self.rt.yield_current(key).ok_or(TaskError::NotReady)
}
#[inline(always)]
pub(crate) fn put_prev_realtime_task(&mut self, id: ThreadId, migration_capable: bool) {
if migration_capable {
self.refresh_class_pushable(id, None);
}
}
pub(crate) fn requeue_realtime_wakee_head(&mut self, id: ThreadId) -> bool {
let Some(QueueMembershipClass::Realtime(key)) = self.membership_class(id) else {
return false;
};
self.rt.requeue_head(key)
}
pub(crate) fn detach_for_transfer(
&mut self,
id: ThreadId,
current_fair: Option<FairEntity>,
timing_granularity_ns: u64,
) -> Option<QueuedThread> {
if self.linked_current() == Some(id) {
return None;
}
let class = self.membership_class(id)?;
let is_fair = matches!(class, QueueMembershipClass::Fair);
if is_fair {
self.update_fair_virtual_time(current_fair);
}
if class == QueueMembershipClass::DeadlineThrottled {
let thread = self.deadline.take_throttled(id)?;
self.unregister_membership(id);
return Some(thread);
}
let thread = SchedulerClass::for_policy(self.queued_thread_including_current(id)?.policy())
.migrate_task_rq(self, class, id, timing_granularity_ns)?;
self.nr_running = self
.nr_running
.checked_sub(1)
.expect("migration must detach one runnable entity");
self.fixed_placement_demand = self
.fixed_placement_demand
.saturating_sub(fixed_placement_demand(thread.active.policy()));
self.unregister_membership(thread.id);
self.refresh_class_pushable(thread.id, self.linked_current());
if is_fair {
self.update_fair_virtual_time(current_fair);
}
self.mark_publication_dirty();
Some(thread)
}
#[inline(always)]
pub(crate) fn pick_next_task(
&mut self,
rt_eligibility: RtEligibility,
skip_delayed: bool,
protected_fair_current: Option<ThreadId>,
) -> Option<PickTaskResult> {
for class in SchedulerClass::PICK_ORDER {
if let Some(picked) =
class.pick_task(self, rt_eligibility, skip_delayed, protected_fair_current)
{
return Some(picked);
}
}
None
}
#[inline(always)]
pub(crate) fn set_next_task(&mut self, picked: &PickedThread) {
SchedulerClass::for_policy(picked.policy()).set_next_task(self, picked);
if picked.metadata().affinity.is_migration_capable() {
self.refresh_class_pushable(picked.id(), Some(picked.id()));
}
}
#[inline(always)]
pub(crate) fn set_next_realtime_task(&mut self, picked: LinkedRqTaskRef) {
let thread = picked.thread();
debug_assert!(matches!(
thread.policy(),
SchedulePolicy::Fifo { .. } | SchedulePolicy::RoundRobin { .. }
));
if thread.metadata.affinity.is_migration_capable() {
self.refresh_class_pushable(thread.id, Some(thread.id));
}
}
}