hirun 0.1.16

A concurrent framework for asynchronous programming based on event-driven, non-blocking I/O mechanism
Documentation
use super::{RawTask, TaskRef, Worker};
use crate::event::Scheduler;
use core::ptr::NonNull;
use core::task::{Context, Waker};
use hioff::container_of_mut;

/// 为Context<'_>提供增强接口
/// 注意:只对本crate的Waker生成的Context<'_>有效.
pub trait TaskContext {
    /// 获取当前调度周期内底层的Scheduler,
    /// 注意: 因为Task的生命周期可能在不同线程内调度,如果注册资源,应该结合set_affinity使用.
    fn sched(&mut self) -> &mut Scheduler;
    /// 设置当前task只在当前worker中调度,如果使用到了sched(),应该总是调用这个接口
    fn set_affinity(&mut self);
    /// 当前Future是否被提前终止,用于全局资源的清理,比如注册到sched()上的资源,应该确保在abort的时候取消注册
    /// 可能是JoinHandle::abort或者task内部调用Context::exit, Context::set_aborted所导致.
    fn aborted(&mut self) -> bool;
    /// 修改abort标志,一般应用于Future的封装层,可能在封装层提前abort内部的Future.
    /// 比如Or/Deadline等就利用了这个接口
    /// 注意:封装层的设置并不影响task的真实运行状态, 如果要退出当前task,请调用exit()接口.
    fn set_aborted(&mut self, aborted: bool);
    /// 在task运行时主动退出,退出之后Future::poll会被再调度一次,再调度时aborted()返回true
    fn exit(&mut self);
}

#[allow(dead_code)]
pub(crate) trait RawTaskContext: TaskContext {
    fn task(&mut self) -> &mut RawTask;
    fn task_ref(&mut self) -> TaskRef;
    // 后续可考虑放入TaskContext中.
    fn freeze_local(&mut self);
    fn unfreeze_local(&mut self);
    // 设置的值可以在同一个task内多个future之间共享.
    fn set_private(&mut self, data: *const ());
    fn get_private(&mut self) -> *const ();
}

pub(crate) struct TaskWaker<'a> {
    waker: Waker,
    task: NonNull<RawTask>,
    sched: &'a mut Scheduler,
    aborted: bool,
}

impl<'a> TaskWaker<'a> {
    pub(crate) fn new(task: &mut RawTask, sched: &'a mut Scheduler) -> Self {
        let worker = unsafe { Worker::from_sched(sched) };
        worker.set_current_task(Some(NonNull::from(&*task)));
        Self {
            waker: task.task_waker(),
            task: NonNull::from(task),
            sched,
            aborted: false,
        }
    }

    pub(crate) fn from_ctx(ctx: &mut Context<'a>) -> &'a mut Self {
        unsafe { container_of_mut!(ctx.waker(), Self, waker) }
    }

    pub(crate) fn waker(&self) -> &Waker {
        &self.waker
    }
}

impl Drop for TaskWaker<'_> {
    fn drop(&mut self) {
        let worker = unsafe { Worker::from_sched(self.sched) };
        worker.set_current_task(None);
    }
}

impl RawTaskContext for Context<'_> {
    fn task(&mut self) -> &mut RawTask {
        let this = TaskWaker::from_ctx(self);
        unsafe { this.task.as_mut() }
    }

    fn task_ref(&mut self) -> TaskRef {
        self.task().task_ref()
    }

    fn unfreeze_local(&mut self) {
        self.task().status.unfreeze_local();
    }

    fn freeze_local(&mut self) {
        let worker = unsafe { Worker::from_sched(self.sched()) };
        self.task().status.freeze_local(worker.worker_id());
    }

    fn set_private(&mut self, data: *const ()) {
        self.task().private = data;
    }

    fn get_private(&mut self) -> *const () {
        self.task().private
    }
}

impl TaskContext for Context<'_> {
    fn exit(&mut self) {
        let this = TaskWaker::from_ctx(self);
        let task = unsafe { this.task.as_mut() };
        task.exit();
    }

    fn set_aborted(&mut self, aborted: bool) {
        let this = TaskWaker::from_ctx(self);
        this.aborted = aborted;
    }

    fn aborted(&mut self) -> bool {
        let this = TaskWaker::from_ctx(self);
        this.aborted
    }

    fn sched(&mut self) -> &mut Scheduler {
        let this = TaskWaker::from_ctx(self);
        this.sched
    }

    fn set_affinity(&mut self) {
        if self.task().status.get_local().is_none() {
            let worker = unsafe { Worker::from_sched(self.sched()) };
            self.task().status.set_local(worker.worker_id());
        }
    }
}