use super::{
event_list_new, timer_rbtree_new, BoxedPollImpl, BoxedScheduler, Event, Poll, PollImpl,
Scheduler, SchedulerVTable, Timer, Waker, POLLIN,
};
use crate::{time, Result};
use core::mem::ManuallyDrop;
use core::ptr::{self, NonNull};
use core::time::Duration;
use hicollections::{List, RbTree};
use hioff::container_of_mut;
use hipool::{Allocator, Boxed};
struct StopEvent {
waker: Waker,
event: Event,
stopped: bool,
}
#[repr(C)]
pub struct SchedImpl<'a, A: Allocator> {
base: ManuallyDrop<Scheduler>,
poll: BoxedPollImpl<'a, A>,
events: List<Event>,
hi_pos: Option<NonNull<Event>>,
timers: RbTree<Timer>,
private_data: *const (),
now: Duration,
stop: StopEvent,
}
impl<'a, A: Allocator + 'a> SchedImpl<'a, A> {
const VTBL: SchedulerVTable = SchedulerVTable {
run: Self::vtbl_run,
stop: Self::vtbl_stop,
stopped: Self::vtbl_stopped,
now: Self::vtbl_now,
add_fd_event: Self::vtbl_add_fd_event,
mod_fd_event: Self::vtbl_mod_fd_event,
del_fd_event: Self::vtbl_del_fd_event,
add_event: Self::vtbl_add_event,
add_event_list: Self::vtbl_add_event_list,
del_event: Self::vtbl_del_event,
set_timer: Self::vtbl_set_timer,
del_timer: Self::vtbl_del_timer,
set_private_data: Self::vtbl_set_private_data,
private_data: Self::vtbl_private_data,
release: Self::vtbl_release,
};
fn from_ptr<'b>(this: *const ()) -> &'b mut Self {
unsafe { container_of_mut!(&*this.cast::<Scheduler>(), Self, base) }
}
fn vtbl_run(this: *const ()) {
Self::from_ptr(this).run()
}
fn vtbl_stop(this: *const ()) {
Self::from_ptr(this).stop()
}
fn vtbl_stopped(this: *const ()) -> bool {
Self::from_ptr(this).stopped()
}
fn vtbl_now(this: *const ()) -> Duration {
Self::from_ptr(this).now()
}
unsafe fn vtbl_add_fd_event(this: *const (), e: &Event, events: u32, fd: i32) -> Result<()> {
Self::from_ptr(this).add_fd_event(e, events, fd)
}
unsafe fn vtbl_mod_fd_event(this: *const (), e: &Event, events: u32, fd: i32) -> Result<()> {
Self::from_ptr(this).mod_fd_event(e, events, fd)
}
unsafe fn vtbl_del_fd_event(this: *const (), e: &Event, fd: i32) -> Result<()> {
Self::from_ptr(this).del_fd_event(e, fd)
}
unsafe fn vtbl_add_event(this: *const (), e: &Event, events: u32, priority: i32) {
Self::from_ptr(this).add_event(e, events, priority)
}
unsafe fn vtbl_add_event_list(this: *const (), events: &mut List<Event>, priority: i32) {
Self::from_ptr(this).add_event_list(events, priority)
}
unsafe fn vtbl_del_event(this: *const (), e: &Event) {
Self::from_ptr(this).del_event(e)
}
unsafe fn vtbl_set_timer(this: *const (), t: &Timer, msecs: u32) {
Self::from_ptr(this).set_timer(t, msecs)
}
unsafe fn vtbl_del_timer(this: *const (), t: &Timer) {
Self::from_ptr(this).del_timer(t)
}
fn vtbl_set_private_data(this: *const (), data: *const ()) {
Self::from_ptr(this).set_private_data(data)
}
fn vtbl_private_data(this: *const ()) -> *const () {
Self::from_ptr(this).private_data()
}
fn vtbl_release(this: *const ()) {
unsafe { core::ptr::drop_in_place(Self::from_ptr(this)) };
}
}
impl<'a, A: Allocator + 'a> SchedImpl<'a, A> {
pub fn new_in(pool: &'a A) -> Result<BoxedScheduler<'a, A>> {
let waker = Waker::new()?;
let mut boxed = Boxed::new_in(
pool,
Self {
base: ManuallyDrop::new(Scheduler::new(&Self::VTBL)),
stop: StopEvent {
event: Event::new(Self::stop_event_handle),
waker,
stopped: false,
},
poll: PollImpl::new_in(pool)?,
events: event_list_new(),
hi_pos: None,
timers: timer_rbtree_new(),
now: time::now(),
private_data: ptr::null(),
},
)?;
boxed.active_stop_event();
Ok(unsafe { boxed.cast_unchecked::<Scheduler>() })
}
}
impl<'a, A: Allocator + 'a> SchedImpl<'a, A> {
fn active_stop_event(&mut self) {
let _ = self.poll.add_event(
self.stop.waker.fd.fd(),
POLLIN,
ptr::addr_of!(self.stop.event) as u64,
);
}
fn stop_event_handle(event: &Event, _events: u32, _sched: &mut Scheduler) {
let this = unsafe { container_of_mut!(event, Self, stop.event) };
if this.stop.waker.awaken() {
this.stop.stopped = true;
}
}
fn wait_timeout(&self) -> i32 {
if !self.events.empty() {
0
} else if let Some(first) = self.timers.first() {
let timeout = first.timeout.get();
if timeout <= self.now {
0
} else {
(timeout - self.now).as_millis() as i32
}
} else {
-1
}
}
fn dispatch_events(&mut self) {
let mut list = event_list_new();
self.events.move_head(&mut list);
self.hi_pos = None;
while let Some(first) = list.first() {
unsafe { list.del(first) };
first.handle(first.events.get(), &mut self.base);
}
}
fn dispatch_fd_events(&mut self, list: &mut List<Event>) {
while let Some(first) = list.first() {
unsafe { list.del(first) };
first.handle(first.events.get(), &mut self.base);
}
}
fn dispatch_timers(&mut self) {
while let Some(first) = self.timers.first() {
if first.timeout.get() <= self.now {
self.timers.remove(first);
first.handle(&mut self.base);
} else {
break;
}
}
}
}
impl<'a, A: Allocator + 'a> SchedImpl<'a, A> {
fn run(&mut self) {
let mut list = List::<Event>::new(|e| unsafe { ptr::addr_of!((*e).node) });
while !self.stop.stopped {
self.now = time::now();
let timeout = self.wait_timeout();
let Ok(_) = self.poll.wait(timeout, |events, event| {
let event = unsafe { &*(event as *const Event) };
event.events.set(events);
unsafe { list.add_tail(event) };
}) else {
continue;
};
if timeout > 0 {
self.now = time::now();
}
self.dispatch_timers();
self.dispatch_fd_events(&mut list);
self.dispatch_events();
}
}
fn stop(&self) {
self.stop.waker.wake();
}
fn stopped(&self) -> bool {
self.stop.stopped
}
fn now(&self) -> Duration {
self.now
}
unsafe fn add_event(&mut self, e: &Event, events: u32, priority: i32) {
self.del_event(e);
if priority == 0 {
self.events.add_tail(e);
} else if let Some(hi_pos) = self.hi_pos {
self.events.add_after(e, hi_pos.as_ref());
self.hi_pos = Some(e.into());
} else {
self.events.add_head(e);
self.hi_pos = Some(e.into());
}
e.events.set(events);
}
unsafe fn add_event_list(&mut self, events: &mut List<Event>, priority: i32) {
if priority == 0 {
events.move_tail(&mut self.events);
} else if events.empty() {
return;
} else if let Some(hi_pos) = self.hi_pos {
self.hi_pos = Some(events.last().unwrap().into());
events.move_after(&mut self.events, hi_pos.as_ref());
} else {
self.hi_pos = Some(events.last().unwrap().into());
events.move_head(&mut self.events);
}
}
unsafe fn del_event(&mut self, e: &Event) {
if self.hi_pos == Some(e.into()) {
self.hi_pos = self.events.rev_iter_from(e).nth(1).map(|e| e.into());
}
self.events.del(e);
}
unsafe fn add_fd_event(&self, e: &Event, events: u32, fd: i32) -> Result<()> {
self.poll.add_event(fd, events, e as *const Event as u64)
}
unsafe fn mod_fd_event(&self, e: &Event, events: u32, fd: i32) -> Result<()> {
self.poll.mod_event(fd, events, e as *const Event as u64)
}
unsafe fn del_fd_event(&self, e: &Event, fd: i32) -> Result<()> {
self.poll.del_event(fd, e as *const Event as u64)
}
unsafe fn set_timer(&mut self, t: &Timer, micros: u32) {
self.timers.remove(t);
t.timeout
.set(self.now + Duration::from_micros(micros as u64));
self.timers.insert(t, false);
}
unsafe fn del_timer(&mut self, t: &Timer) {
self.timers.remove(t);
}
fn private_data(&self) -> *const () {
self.private_data
}
fn set_private_data(&mut self, data: *const ()) {
self.private_data = data;
}
}