use crate::modules::input::Token;
use crate::{
constants::SLEEP_TOLERANCE,
executor,
futures::{
kernel_wait::KernelWait,
sleep::SleepMode,
task::sealed,
task::{Nothing, Task},
},
modules::event_desc::EventDesc,
modules::{
int_check::IntCheck,
kevent::KEvent,
kqueue::{self, Waited},
waiter::Waiter,
wake_target::WakeTarget,
},
};
use std::{
ptr,
time::{Duration, Instant},
};
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum Slept {
Waited,
Cancelled,
Refused,
}
#[derive(Clone)]
#[must_use = "a task does nothing until it is run or spawned"]
pub struct SleepTask {
pub(crate) sleep_for: Duration,
pub(crate) created: Instant,
pub(crate) mode: SleepMode,
pub(crate) until: Option<Instant>,
}
impl SleepTask {
pub(crate) fn new(time: Duration) -> Self {
Self {
sleep_for: time,
created: Instant::now(),
mode: SleepMode::Precise,
until: None,
}
}
pub(crate) fn until(when: Instant) -> Self {
let created = Instant::now();
Self {
sleep_for: when.saturating_duration_since(created),
created,
mode: SleepMode::Precise,
until: Some(when),
}
}
pub fn mode(mut self, mode: SleepMode) -> Self {
self.mode = mode;
self
}
#[inline(always)]
pub(crate) fn precise(&self) -> bool {
self.mode == SleepMode::Precise
}
pub(crate) fn spinlock(&self, until: Instant) {
while Instant::now() < until {}
}
#[inline(always)]
fn wait_on_own(&self, queue: i32, task_id: usize) -> Slept {
let registered = unsafe {
KEvent::register(
queue,
task_id,
self.get_intptr_t_data(),
WakeTarget::None.encode(),
EventDesc::new_timer(),
)
}
.check();
if registered.is_err() {
return Slept::Refused;
}
if !executor::waiting_on(queue) {
let _ = unsafe {
KEvent::register(
queue,
task_id,
0,
ptr::null_mut(),
EventDesc::new_timer_delete(),
)
}
.check();
return Slept::Cancelled;
}
let waited = kqueue::wait_for(queue, task_id, libc::EVFILT_TIMER);
if !executor::stopped_waiting() {
return Slept::Cancelled;
}
match waited {
Waited::Failed => Slept::Refused,
Waited::Cancelled => Slept::Cancelled,
Waited::Arrived => Slept::Waited,
}
}
#[inline(always)]
fn wait_on_reactor(&self, reactor_id: i32, task_id: usize) -> Slept {
let waiter = Waiter::new();
let registered = unsafe {
KEvent::register(
reactor_id,
task_id,
self.get_intptr_t_data(),
WakeTarget::Parked(&waiter as *const Waiter as *mut Waiter).encode(),
EventDesc::new_timer(),
)
}
.check();
if registered.is_err() {
return Slept::Refused;
}
waiter.wait();
Slept::Waited
}
}
impl sealed::Sealed for SleepTask {}
impl Task for SleepTask {
type Output = Duration;
type Input = Nothing;
#[inline(always)]
fn execute(&self, _token: Token, reactor_id: i32, task_id: usize) -> Self::Output {
if !self.precise() || self.sleep_for > SLEEP_TOLERANCE {
return self.offload(kqueue::id().ok(), reactor_id, task_id);
}
let until = self.created + self.sleep_for;
self.spinlock(until);
self.created.elapsed()
}
#[inline(always)]
fn prepare(&mut self, _token: Token) {
self.created = Instant::now();
if let Some(when) = self.until {
self.sleep_for = when.saturating_duration_since(self.created);
}
}
#[inline(always)]
fn blocking(&self, _token: Token) -> bool {
!self.precise() || self.sleep_for > SLEEP_TOLERANCE
}
}
impl KernelWait for SleepTask {
#[inline(always)]
fn get_intptr_t_data(&self) -> libc::intptr_t {
let target = if self.precise() {
self.sleep_for.saturating_sub(SLEEP_TOLERANCE)
} else {
self.sleep_for
};
let nanos = target.saturating_sub(self.created.elapsed()).as_nanos();
nanos.min(libc::intptr_t::MAX as u128) as libc::intptr_t
}
#[inline(always)]
fn offload(&self, queue: Option<i32>, reactor_id: i32, task_id: usize) -> Self::Output {
let slept = match queue {
Some(queue) => self.wait_on_own(queue, task_id),
None => self.wait_on_reactor(reactor_id, task_id),
};
let spin = match slept {
Slept::Cancelled => false,
Slept::Waited => self.precise(),
Slept::Refused => true,
};
if spin {
let until = self.created + self.sleep_for;
self.spinlock(until);
}
self.created.elapsed()
}
}