use core::task::Poll;
use crate::task::{iter_task_regs, TaskReg, TaskRegistry};
use crate::waker;
#[cfg(not(loom))]
static CURRENT_TASK: crate::sync::atomic::AtomicU32 = crate::sync::atomic::AtomicU32::new(u32::MAX);
#[cfg(loom)]
loom::lazy_static! {
static ref CURRENT_TASK: crate::sync::atomic::AtomicU32 = crate::sync::atomic::AtomicU32::new(u32::MAX);
}
fn set_current(id: crate::task::TaskId) {
CURRENT_TASK.store(id.as_u16() as u32, crate::sync::atomic::Ordering::Release);
}
fn clear_current() {
CURRENT_TASK.store(u32::MAX, crate::sync::atomic::Ordering::Release);
}
#[doc(hidden)]
pub fn set_current_for_test(priority: u8, index: u8) {
set_current(crate::task::TaskId::new(priority, index));
}
#[doc(hidden)]
pub fn clear_current_for_test() {
clear_current();
}
pub fn current_task() -> Option<crate::task::TaskId> {
let v = CURRENT_TASK.load(crate::sync::atomic::Ordering::Acquire);
if v == u32::MAX {
None
} else {
Some(crate::task::TaskId::from_u16(v as u16))
}
}
#[cfg(not(loom))]
static LIVE_TASKS: crate::sync::atomic::AtomicUsize = crate::sync::atomic::AtomicUsize::new(0);
#[cfg(loom)]
loom::lazy_static! {
static ref LIVE_TASKS: crate::sync::atomic::AtomicUsize = crate::sync::atomic::AtomicUsize::new(0);
}
pub struct Executor {
registry: TaskRegistry,
}
impl Executor {
pub const fn new() -> Self {
Self {
registry: TaskRegistry::new(),
}
}
pub fn init(&mut self) {
let mut counts: [u8; 32] = [0; 32];
for reg in iter_task_regs() {
let prio = reg.priority as usize;
if prio > (crate::task::MAX_PRIORITY as usize) {
panic!(
"rivet: task priority {} exceeds MAX_PRIORITY {} \
(check #[rivet::task(priority = ...)])",
reg.priority,
crate::task::MAX_PRIORITY
);
}
let idx = counts[prio] as usize;
if idx >= crate::task::MAX_TASKS {
panic!(
"rivet: too many #[rivet::task]s at priority {} \
(limit MAX_TASKS = {} per priority)",
reg.priority,
crate::task::MAX_TASKS
);
}
self.registry.tasks[prio][idx] = Some(reg as *const TaskReg);
counts[prio] += 1;
self.registry.total += 1;
LIVE_TASKS.fetch_add(1, crate::sync::atomic::Ordering::Relaxed);
}
self.registry.count_per_priority[..=(crate::task::MAX_PRIORITY as usize)]
.copy_from_slice(&counts[..=(crate::task::MAX_PRIORITY as usize)]);
for (p, &count) in counts[..=(crate::task::MAX_PRIORITY as usize)]
.iter()
.enumerate()
{
for i in 0..(count as usize) {
crate::waker::mark_ready(crate::task::TaskId::new(p as u8, i as u8));
}
}
}
pub fn run(&self) -> ! {
loop {
waker::clear_pend();
while let Some(id) = waker::next_ready() {
let reg = match self.lookup_task(id.priority(), id.index()) {
Some(r) => r,
None => continue,
};
unsafe {
if (reg.completed_fn)(reg.user_data) {
continue;
}
}
let task_waker = waker::task_waker(id);
set_current(id);
let result = unsafe { (reg.poll_fn)(reg.user_data, &task_waker) };
clear_current();
if result == Poll::Ready(()) {
LIVE_TASKS.fetch_sub(1, crate::sync::atomic::Ordering::Relaxed);
}
}
if !waker::has_pending() {
crate::port::arch::idle();
}
}
}
fn lookup_task(&self, priority: u8, index: u8) -> Option<&'static TaskReg> {
let p = priority as usize;
let i = index as usize;
if p > (crate::task::MAX_PRIORITY as usize) || i >= crate::task::MAX_TASKS {
return None;
}
self.registry.tasks[p][i].map(|ptr| unsafe { &*ptr })
}
}
#[cfg(feature = "test-support")]
pub(crate) fn reset_for_test() {
clear_current();
}
impl Default for Executor {
fn default() -> Self {
Self::new()
}
}
pub static mut EXECUTOR: Executor = Executor::new();