use alloc::sync::Arc;
use core::num::NonZeroU64;
use crate::{
runtime::{
TaskSystem,
context::{
RuntimeIrqGuard, RuntimeSchedulerFrameGuard, runtime_current_cpu_mut,
runtime_task_system,
},
cpu::CpuLocal,
lock::PreemptScope,
service::SchedulerTickMode,
switch::{RuntimeScheduleOrigin, RuntimeSchedulerEntry, dispatch::execute_switch_plan},
task_runtime,
},
sched::{
CpuId, SchedulePolicy,
system::{DeadlineBaseGuardSource, SchedulerDeadlineDerivationSource},
},
thread::{
ParkCommit, ParkPrepare, TaskError, ThreadCore, ThreadId, ThreadWakeHandle,
current::{BlockingPermit, acquire_blocking_permit, current_thread_core_arc},
},
time::{MonotonicDeadline, MonotonicInstant, queue::TaskDeadlineKind},
};
#[derive(Debug)]
pub enum CurrentParkStart {
Notified,
Prepared(PreparedCurrentPark),
}
#[must_use = "a prepared current-thread park must be committed or cancelled"]
#[derive(Debug)]
pub struct PreparedCurrentPark {
thread: Arc<ThreadCore>,
ticket: Option<crate::thread::ParkTicket>,
system: &'static TaskSystem,
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub struct CurrentParkResume {
generation: u64,
deadline_expired: bool,
disposition: CurrentParkDisposition,
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub enum CurrentParkDisposition {
NotifiedBeforeBlock,
BlockedAndResumed,
}
impl CurrentParkResume {
pub const fn generation(self) -> u64 {
self.generation
}
pub const fn deadline_expired(self) -> bool {
self.deadline_expired
}
pub const fn disposition(self) -> CurrentParkDisposition {
self.disposition
}
pub const fn was_notified_before_block(self) -> bool {
matches!(
self.disposition,
CurrentParkDisposition::NotifiedBeforeBlock
)
}
}
impl PreparedCurrentPark {
pub fn thread_id(&self) -> ThreadId {
self.thread.id()
}
pub fn wake_handle(&self) -> ThreadWakeHandle {
ThreadWakeHandle::from_core(Arc::clone(&self.thread))
}
pub fn generation(&self) -> u64 {
self.ticket()
.expect("prepared park ticket remains owned")
.generation()
}
pub fn arm_deadline(&mut self, deadline: MonotonicDeadline) -> Result<(), TaskError> {
let ticket = self
.ticket
.as_mut()
.expect("prepared park ticket remains owned");
arm_current_park_deadline(&self.thread, ticket, deadline)
}
pub fn commit(mut self) -> Result<CurrentParkResume, TaskError> {
let mut ticket = self
.ticket
.take()
.expect("prepared park ticket remains owned");
let generation = ticket.generation();
let deadline_armed = ticket.has_deadline();
let disposition =
match commit_current_park_with_system(self.system, &self.thread, &mut ticket) {
Ok(disposition) => disposition,
Err(error) => {
let deadline_result = cancel_current_park_deadline(&self.thread, &mut ticket);
if cancel_current_park(&self.thread, &mut ticket).is_err() {
task_runtime::fatal_invariant(
0x5041_0002,
self.thread.id().as_u64() as usize,
);
}
let _cancelled = deadline_result?;
return Err(error);
}
};
let deadline_cancelled = cancel_current_park_deadline(&self.thread, &mut ticket)?;
Ok(CurrentParkResume {
generation,
deadline_expired: deadline_armed && !deadline_cancelled,
disposition,
})
}
pub fn cancel(mut self) -> Result<(), TaskError> {
let mut ticket = self
.ticket
.take()
.expect("prepared park ticket remains owned");
let deadline_result = cancel_current_park_deadline(&self.thread, &mut ticket);
let park_result = cancel_current_park(&self.thread, &mut ticket);
let _cancelled = deadline_result?;
park_result
}
fn ticket(&self) -> Option<&crate::thread::ParkTicket> {
self.ticket.as_ref()
}
}
impl Drop for PreparedCurrentPark {
fn drop(&mut self) {
if self
.ticket
.as_ref()
.is_some_and(|ticket| !ticket.is_resolved())
{
task_runtime::fatal_invariant(0x5041_0003, self.thread.id().as_u64() as usize);
}
}
}
pub fn begin_current_park() -> Result<CurrentParkStart, TaskError> {
let permit = acquire_blocking_permit()?;
begin_current_park_with_permit(&permit)
}
pub(crate) fn begin_current_park_with_permit(
_permit: &BlockingPermit,
) -> Result<CurrentParkStart, TaskError> {
let system = runtime_task_system()?;
let _current_pin = PreemptScope::enter();
let thread = current_thread_core_arc()?;
let prepare = system.prepare_current_park(&thread);
match prepare? {
ParkPrepare::Notified => Ok(CurrentParkStart::Notified),
ParkPrepare::Prepared(ticket) => Ok(CurrentParkStart::Prepared(PreparedCurrentPark {
thread,
ticket: Some(ticket),
system,
})),
}
}
pub fn on_clock_event(
now: MonotonicInstant,
budget: usize,
scheduler_event: ClaimedSchedulerDeadlines,
) -> Result<TaskClockEventOutcome, TaskError> {
let system = runtime_task_system()?;
let mut irq = RuntimeIrqGuard::enter();
let mut cpu = runtime_current_cpu_mut(&mut irq)?;
let periodic_tick = scheduler_event.runs_periodic_task_tick();
if periodic_tick && cpu.promote_lazy_reschedule() {
let _self_serviced = task_runtime::publish_local_scheduler_work();
}
let (charge, clock, current, task_tick_rq_observation) = match scheduler_event.accounting_kind()
{
ClockAccountingKind::RuntimeOnly => {
system.charge_current_until_with_clock(cpu.as_mut(), 0)?
}
ClockAccountingKind::SchedulerDeadline => {
system.clock_event_current_until_with_clock(cpu.as_mut(), 0)?
}
ClockAccountingKind::PeriodicTick => system.task_tick_current_until_with_clock(
cpu.as_mut(),
0,
scheduler_event.periodic_tick_ns(),
)?,
ClockAccountingKind::PeriodicTickWithSchedulerDeadline => system
.task_tick_and_clock_event_current_until_with_clock(
cpu.as_mut(),
0,
scheduler_event.periodic_tick_ns(),
)?,
};
let rt_period_rescheduled = system.service_rt_period(&cpu, now);
let hard = system.service_due_hard_timers(cpu.as_mut(), now, budget)?;
let batch = hard.soft();
let rq_observation =
match clock_event_rq_observation_plan(rt_period_rescheduled, hard.processed()) {
ClockEventRqObservationPlan::ReuseAccounted => cpu
.as_mut()
.scheduler_work_due_from_rq_observation(now, task_tick_rq_observation),
ClockEventRqObservationPlan::RefreshAndPublish => cpu.as_mut().scheduler_work_due(now),
};
let runtime_deadline = cpu.scheduler_runtime_deadline_for_rq_observation(rq_observation);
let update = cpu
.as_mut()
.next_scheduler_deadline_update_from_rq_observation(
rq_observation,
SchedulerDeadlineDerivationSource::ClockEvent,
)?;
Ok(TaskClockEventOutcome {
slice_expired: charge.slice_expired(),
deadline_overrun: charge.deadline_overrun(),
expired: hard.processed().saturating_add(batch.expired()),
update,
runtime_deadline: Some(runtime_deadline),
scheduler_tick: SchedulerTickStamp {
cpu: cpu.owner(),
thread: current,
observed_ns: clock.task().as_nanos(),
},
})
}
const fn clock_event_rq_observation_reusable(
rt_period_rescheduled: bool,
hard_timers_processed: usize,
) -> bool {
!rt_period_rescheduled && hard_timers_processed == 0
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
enum ClockEventRqObservationPlan {
ReuseAccounted,
RefreshAndPublish,
}
const fn clock_event_rq_observation_plan(
rt_period_rescheduled: bool,
hard_timers_processed: usize,
) -> ClockEventRqObservationPlan {
if clock_event_rq_observation_reusable(rt_period_rescheduled, hard_timers_processed) {
ClockEventRqObservationPlan::ReuseAccounted
} else {
ClockEventRqObservationPlan::RefreshAndPublish
}
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub struct ClaimedSchedulerDeadlines {
periodic_tick_ns: Option<NonZeroU64>,
scheduler_deadline_elapsed: bool,
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
enum ClockAccountingKind {
RuntimeOnly,
SchedulerDeadline,
PeriodicTick,
PeriodicTickWithSchedulerDeadline,
}
impl ClaimedSchedulerDeadlines {
pub const fn new(
periodic_tick_ns: Option<NonZeroU64>,
scheduler_deadline_elapsed: bool,
) -> Self {
Self {
periodic_tick_ns,
scheduler_deadline_elapsed,
}
}
const fn runs_periodic_task_tick(self) -> bool {
self.periodic_tick_ns.is_some()
}
fn periodic_tick_ns(self) -> u64 {
match self.periodic_tick_ns {
Some(tick_ns) => tick_ns.get(),
None => task_runtime::fatal_invariant(0x5251_1013, 0),
}
}
const fn accounting_kind(self) -> ClockAccountingKind {
match (
self.periodic_tick_ns.is_some(),
self.scheduler_deadline_elapsed,
) {
(false, false) => ClockAccountingKind::RuntimeOnly,
(false, true) => ClockAccountingKind::SchedulerDeadline,
(true, false) => ClockAccountingKind::PeriodicTick,
(true, true) => ClockAccountingKind::PeriodicTickWithSchedulerDeadline,
}
}
}
pub fn publish_scheduler_tick(
stamp: SchedulerTickStamp,
mode: SchedulerTickMode,
tick_ns: u64,
) -> Result<(), TaskError> {
if tick_ns == 0 {
return Err(TaskError::InvalidConfiguration);
}
let system = runtime_task_system()?;
let mut irq = RuntimeIrqGuard::enter();
let cpu = runtime_current_cpu_mut(&mut irq)?;
if cpu.owner() != stamp.cpu {
return Err(TaskError::CpuOwnerMismatch {
expected: stamp.cpu.as_u32(),
actual: cpu.owner().as_u32(),
});
}
system.publish_current_scheduler_tick_work(&cpu, stamp.thread, stamp.observed_ns, mode, tick_ns)
}
pub(crate) fn commit_current_park(
current: &Arc<ThreadCore>,
ticket: &mut crate::thread::ParkTicket,
) -> Result<CurrentParkDisposition, TaskError> {
let system = runtime_task_system()?;
commit_current_park_with_system(system, current, ticket)
}
fn commit_current_park_with_system(
system: &'static TaskSystem,
current: &Arc<ThreadCore>,
ticket: &mut crate::thread::ParkTicket,
) -> Result<CurrentParkDisposition, TaskError> {
let mut scheduler_frame = RuntimeSchedulerFrameGuard::enter(
RuntimeScheduleOrigin::Block,
RuntimeSchedulerEntry::Task,
)?;
let commit = {
let mut cpu = runtime_current_cpu_mut(&mut scheduler_frame)?;
unsafe { system.commit_park_in_scheduler_frame(cpu.as_mut(), current, ticket)? }
};
match commit {
ParkCommit::Notified => Ok(CurrentParkDisposition::NotifiedBeforeBlock),
ParkCommit::Blocked(mut decision) => {
execute_switch_plan(&mut scheduler_frame, &mut decision);
Ok(CurrentParkDisposition::BlockedAndResumed)
}
}
}
pub(crate) fn cancel_current_park(
current: &ThreadCore,
ticket: &mut crate::thread::ParkTicket,
) -> Result<(), TaskError> {
let mut irq = RuntimeIrqGuard::enter();
let mut cpu = runtime_current_cpu_mut(&mut irq)?;
runtime_task_system()?.cancel_current_park(cpu.as_mut(), current, ticket)
}
pub(crate) fn arm_current_park_deadline(
thread: &Arc<ThreadCore>,
ticket: &mut crate::thread::ParkTicket,
deadline: MonotonicDeadline,
) -> Result<(), TaskError> {
let mut irq = RuntimeIrqGuard::enter();
let cpu = runtime_current_cpu_mut(&mut irq)?;
if ticket.thread() != thread.id()
|| ticket.is_resolved()
|| ticket.has_deadline()
|| cpu.current() != Some(thread.id())
{
return Err(TaskError::StaleThreadId);
}
let owner = cpu.owner();
let (registration, update) = {
let mut deadline_base = cpu
.remote()
.lock_deadline_activity(DeadlineBaseGuardSource::Registration);
let non_timer = deadline_base.non_timer;
let kind = TaskDeadlineKind::park_timeout(ticket.generation());
let registration = if matches!(
thread.base_policy_snapshot(),
SchedulePolicy::Fifo { .. }
| SchedulePolicy::RoundRobin { .. }
| SchedulePolicy::Deadline(_)
) {
deadline_base.queue.arm_hard_park(
thread.sleep_timer(),
deadline,
kind,
Arc::clone(thread),
)
} else {
deadline_base
.queue
.arm(thread.sleep_timer(), deadline, kind)
}
.map_err(|error| match error {
crate::time::queue::TaskDeadlineError::Capacity => TaskError::TimerCapacity,
crate::time::queue::TaskDeadlineError::GenerationExhausted
| crate::time::queue::TaskDeadlineError::KindMismatch => {
TaskError::InvalidConfiguration
}
})?;
let token = registration.token();
thread.register_sleep_timer(owner, token.generation());
let update = match CpuLocal::update_scheduler_deadline_registration_publication(
&mut deadline_base,
non_timer,
) {
Ok(update) => update,
Err(error) => {
let removed = deadline_base.queue.cancel(®istration);
let completed = thread.complete_sleep_timer(token.generation());
if !removed || !completed {
task_runtime::fatal_invariant(0x5444_0005, thread.id().as_u64() as usize);
}
return Err(error);
}
};
(registration, update)
};
task_runtime::publish_scheduler_deadline(update);
if ticket.attach_deadline(registration).is_err() {
task_runtime::fatal_invariant(0x5444_0002, thread.id().as_u64() as usize);
}
Ok(())
}
pub(crate) fn cancel_current_park_deadline(
thread: &ThreadCore,
ticket: &mut crate::thread::ParkTicket,
) -> Result<bool, TaskError> {
if ticket.thread() != thread.id() {
return Err(TaskError::StaleThreadId);
}
let Some(token) = ticket.deadline().map(|registration| registration.token()) else {
return Ok(false);
};
let system = runtime_task_system()?;
let mut irq = RuntimeIrqGuard::enter();
let cpu = runtime_current_cpu_mut(&mut irq)?;
let actual = cpu.owner();
let Some(expected) = thread.sleep_timer_cpu_for(token.generation()) else {
if !ticket.clear_deadline(token) {
task_runtime::fatal_invariant(0x5444_0003, thread.id().as_u64() as usize);
}
return Ok(false);
};
if actual != expected {
let remote = system
.cpu_remote(expected)
.ok_or(TaskError::CpuOffline(expected.as_u32()))?;
let registration = ticket
.deadline()
.expect("the deadline registration remains owned until cancellation");
let (cancellation, expired) = {
let mut deadline_base =
remote.lock_deadline_activity(DeadlineBaseGuardSource::Registration);
let cancellation = deadline_base.queue.begin_cancel(registration);
let expired = if cancellation.is_none() {
deadline_base.cancel_expired_task_deadline(registration)
} else {
false
};
(cancellation, expired)
};
let cancelled = match (cancellation, expired) {
(Some(cancellation), _) => {
cancellation.commit();
true
}
(None, true) => false,
(None, false) if thread.sleep_timer_cpu_for(token.generation()).is_none() => {
if !ticket.clear_deadline(token) {
task_runtime::fatal_invariant(0x5444_0003, thread.id().as_u64() as usize);
}
return Ok(false);
}
(None, false) => {
task_runtime::fatal_invariant(0x5444_0006, thread.id().as_u64() as usize)
}
};
if !thread.complete_sleep_timer(token.generation()) || !ticket.clear_deadline(token) {
task_runtime::fatal_invariant(0x5444_0004, thread.id().as_u64() as usize);
}
return Ok(cancelled);
}
let (cancellation, update) = {
let registration = ticket
.deadline()
.expect("the deadline registration remains owned until cancellation");
let mut deadline_base = cpu
.remote()
.lock_deadline_activity(DeadlineBaseGuardSource::Registration);
let non_timer = deadline_base.non_timer;
let cancellation = deadline_base.queue.begin_cancel(registration);
let expired = if cancellation.is_none() {
deadline_base.cancel_expired_task_deadline(registration)
} else {
false
};
let cancellation = match (cancellation, expired) {
(Some(cancellation), _) => cancellation,
(None, true) => {
if !thread.complete_sleep_timer(token.generation()) || !ticket.clear_deadline(token)
{
task_runtime::fatal_invariant(0x5444_0004, thread.id().as_u64() as usize);
}
return Ok(false);
}
(None, false) if thread.sleep_timer_cpu_for(token.generation()).is_none() => {
if !ticket.clear_deadline(token) {
task_runtime::fatal_invariant(0x5444_0003, thread.id().as_u64() as usize);
}
return Ok(false);
}
(None, false) => {
task_runtime::fatal_invariant(0x5444_0006, thread.id().as_u64() as usize);
}
};
let update = match CpuLocal::update_scheduler_deadline_registration_publication(
&mut deadline_base,
non_timer,
) {
Ok(update) => update,
Err(error) => {
cancellation.rollback(&mut deadline_base.queue);
return Err(error);
}
};
(cancellation, update)
};
task_runtime::publish_scheduler_deadline(update);
cancellation.commit();
if !thread.complete_sleep_timer(token.generation()) || !ticket.clear_deadline(token) {
task_runtime::fatal_invariant(0x5444_0004, thread.id().as_u64() as usize);
}
Ok(true)
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub struct TaskClockEventOutcome {
slice_expired: bool,
deadline_overrun: bool,
expired: usize,
update: crate::runtime::cpu::SchedulerDeadlineUpdate,
runtime_deadline: Option<crate::runtime::cpu::SchedulerRuntimeDeadline>,
scheduler_tick: SchedulerTickStamp,
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub struct SchedulerTickStamp {
cpu: CpuId,
thread: ThreadId,
observed_ns: u64,
}
impl TaskClockEventOutcome {
pub const fn slice_expired(self) -> bool {
self.slice_expired
}
pub const fn deadline_overrun(self) -> bool {
self.deadline_overrun
}
pub const fn expired(self) -> usize {
self.expired
}
pub const fn update(self) -> crate::runtime::cpu::SchedulerDeadlineUpdate {
self.update
}
pub const fn runtime_deadline(self) -> Option<crate::runtime::cpu::SchedulerRuntimeDeadline> {
self.runtime_deadline
}
pub const fn scheduler_tick_stamp(self) -> SchedulerTickStamp {
self.scheduler_tick
}
pub const fn next_deadline(self) -> Option<MonotonicDeadline> {
self.update.deadline()
}
}
#[cfg(test)]
mod tests {
use core::num::NonZeroU64;
use super::{
ClaimedSchedulerDeadlines, ClockAccountingKind, ClockEventRqObservationPlan,
clock_event_rq_observation_plan, clock_event_rq_observation_reusable,
};
const _: () = assert!(matches!(
clock_event_rq_observation_plan(false, 1),
ClockEventRqObservationPlan::RefreshAndPublish
));
#[test]
fn only_periodic_clock_events_run_the_scheduler_tick() {
let tick_ns = NonZeroU64::new(10).unwrap();
assert!(!ClaimedSchedulerDeadlines::new(None, false).runs_periodic_task_tick());
assert!(!ClaimedSchedulerDeadlines::new(None, true).runs_periodic_task_tick());
assert!(ClaimedSchedulerDeadlines::new(Some(tick_ns), false).runs_periodic_task_tick());
assert!(ClaimedSchedulerDeadlines::new(Some(tick_ns), true).runs_periodic_task_tick());
}
#[test]
fn unrelated_physical_clockevent_only_accounts_runtime() {
let tick_ns = NonZeroU64::new(10).unwrap();
assert_eq!(
ClaimedSchedulerDeadlines::new(None, false).accounting_kind(),
ClockAccountingKind::RuntimeOnly
);
assert_eq!(
ClaimedSchedulerDeadlines::new(None, true).accounting_kind(),
ClockAccountingKind::SchedulerDeadline
);
assert_eq!(
ClaimedSchedulerDeadlines::new(Some(tick_ns), false).accounting_kind(),
ClockAccountingKind::PeriodicTick
);
assert_eq!(
ClaimedSchedulerDeadlines::new(Some(tick_ns), true).accounting_kind(),
ClockAccountingKind::PeriodicTickWithSchedulerDeadline
);
}
#[test]
fn linux_common_tick_reuses_the_task_tick_rq_observation() {
assert!(clock_event_rq_observation_reusable(false, 0));
assert!(!clock_event_rq_observation_reusable(true, 0));
assert!(!clock_event_rq_observation_reusable(false, 1));
}
#[test]
fn plain_hard_timer_rechecks_and_publishes_due_scheduler_work() {
assert_eq!(
clock_event_rq_observation_plan(false, 1),
ClockEventRqObservationPlan::RefreshAndPublish
);
}
}