hirun 0.1.16

A concurrent framework for asynchronous programming based on event-driven, non-blocking I/O mechanism
Documentation
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();
        // 注意内存布局,base必须在最前面,且是#[repr(C)],避免编译器修改
        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
    }

    /// # Safety
    /// 调用之后e不能转移所有权
    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);
    }

    /// # Safety
    /// 调用之后所有event不能转移所有权
    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);
        }
    }

    /// # Safety
    /// e应该调用了add_event
    unsafe fn del_event(&mut self, e: &Event) {
        if self.hi_pos == Some(e.into()) {
            // #14. rev_iter_from(e).next返回的仍然是e本身,应该调用nth(1)取得其前一个
            self.hi_pos = self.events.rev_iter_from(e).nth(1).map(|e| e.into());
        }
        self.events.del(e);
    }

    /// # Safety
    /// 调用之后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)
    }

    /// # Safety
    /// e应该调用了add_fd_event
    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)
    }

    /// # Safety
    /// e应该调用了add_fd_event
    unsafe fn del_fd_event(&self, e: &Event, fd: i32) -> Result<()> {
        self.poll.del_event(fd, e as *const Event as u64)
    }

    /// # Safety
    /// 调用之后t不能转移所有权
    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);
    }

    /// # Safety
    /// t应该调用set_timer了
    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;
    }
}