hirun 0.1.13

A concurrent framework for asynchronous programming based on event-driven, non-blocking I/O mechanism
Documentation
use super::{Attr, THREADS_MAX, THREADS_MASK};
use core::sync::atomic::{
    fence, AtomicU8, AtomicU16, AtomicU32,
    Ordering::{self, Acquire, Relaxed, Release},
};

#[repr(C)]
pub(crate) struct AtomicStatus {
    pcnt: u64,
    wcnt: AtomicU32,
    rcnt: AtomicU16,
    stat: AtomicU8,
    wexp: u8,
    group: u8,
    sext: u8,
    worker: u16,
}

pub(crate) const STAT_RUN: u8 = 0;
pub(crate) const STAT_JOIN: u8 = 1;
pub(crate) const STAT_RETURN: u8 = 2;
pub(crate) const STAT_FINISH: u8 = 3;
pub(crate) const STAT_ABORT: u8 = 4;
pub(crate) const STAT_EXIT: u8 = 5;

const PRI_SHIFT: u8 = 0;
// const YIELD_SHIFT: u8 = 1;

impl AtomicStatus {
    pub(crate) fn new(attr: &Attr) -> Self {
        Self {
            stat: AtomicU8::new(STAT_RUN),
            rcnt: AtomicU16::new(1),
            wcnt: AtomicU32::new(1),
            wexp: 1,
            group: attr.group_id,
            sext: sext_new(0, PRI_SHIFT, attr.priority != 0),
            worker: 0,
            pcnt: 0,
        }
    }

    // 应该仅在Future::poll中调用,
    // 目的是此task不可能被外部激活了,
    // 后续应该是当前Worker中的事件或者定时器处理中调用sched_waked重新激活
    pub(crate) fn freeze_wake(&self) {
        self.wcnt.fetch_add(0x10000, Relaxed);
    }

    fn frozen_wake(&self, cnt: u32) -> bool {
        (cnt & 0xFFFF000) > 0
    }

    // wait——all接口使用,目的尽可能减少task空转
    // 这是在Future::poll时才被调用,之后框架一定会调用test_dec_wake. 
    // test_dec_wake和test_inc_wake会有并发,利用release/acquire,可以消除wexp在两者间的数据竞争
    pub(crate) fn set_wake_expect(&mut self, wexp: u8) {
        assert!(wexp > 0);
        self.wexp = wexp;
    }

    pub(crate) fn test_dec_wake(&self, cnt: u32) -> bool {
        let cnt = self.wcnt.fetch_sub(cnt, Release) - cnt;
        // self.wexp只在当前线程被修改,一定可以读取到最新值
        !self.frozen_wake(cnt) && cnt >= self.wexp as u32
    }
    pub(crate) fn test_inc_wake(&self) -> bool {
        let cnt = self.wcnt.fetch_add(1, Acquire);
        cnt == self.wexp as u32 - 1
    }

    // 在调用Future::poll之前先获取当前wake_count,作为poll之后test_dec_wake的输入参数,表示本次消耗的次数
    pub(crate) fn wake_count(&self) -> u32 {
        self.wcnt.load(Relaxed)
    }

    // 设置退出标志,调度结束的时候快速判断是否退出,相比从status获取状态性能更优
    pub(crate) fn set_exited(&mut self) {
        self.pcnt |= 0x01 << 63;
    }

    pub(crate) fn exited(&self) -> bool {
        (self.pcnt & (0x01 << 63)) != 0
    }

    // 设置abort标志,确保只会abort一次
    pub(crate) fn set_aborted(&mut self) {
        self.pcnt |= 0x01 << 62;
    }

    pub(crate) fn aborted(&self) -> bool {
        (self.pcnt & (0x01 << 62)) != 0
    }

    // 实际只用到62位,不必考虑溢出
    pub(crate) fn poll_inc(&mut self) {
        self.pcnt += 1;
    }

    pub(crate) fn poll_cnt(&self) -> u64 {
        self.pcnt
    }

    pub(crate) fn inc_ref(&self) {
        let cnt = self.rcnt.fetch_add(1, Relaxed);
        assert!(cnt < 254);
        assert!(cnt > 0);
    }

    pub(crate) fn test_dec_ref(&self) -> bool {
        if self.rcnt.fetch_sub(1, Release) == 1 {
            fence(Acquire);
            return true;
        }
        false
    }

    pub(crate) fn status(&self, order: Ordering) -> u8 {
        self.stat.load(order)
    }

    pub(crate) fn set_status(&self, new: u8, order: Ordering) {
        self.stat.store(new, order);
    }

    pub(crate) fn cmp_xchg_status(
        &self,
        current: u8,
        new: u8,
        set_order: Ordering,
    ) -> Result<u8, u8> {
        self.stat.compare_exchange(current, new, set_order, Relaxed)
    }

    pub(crate) fn priority(&self) -> u8 {
        if sext_get(self.sext, PRI_SHIFT) {
            1
        } else {
            0
        }
    }
    pub(crate) fn group(&self) -> u8 {
        self.group
    }

    // worker: THREAD_MAX == 1024(0x0400), THREAD_MASK == 0x03FF
    // 0x8000: 表示spawn_local调度,绑定当前worker, future不支持Send场景
    // 0x4000: 表示临时绑定当前worker,适用于资源和scheduler挂钩场景,future是支持Send的
    pub(crate) fn get_local(&self) -> Option<u16> {
        if self.worker >= THREADS_MAX as u16 {
            return Some(self.worker & THREADS_MASK as u16);
        }
        None
    }
    pub(crate) fn set_local(&mut self, id: u16) {
        self.worker = id | 0x8000;
    }

    // Or操作可能一次调度有多个Future调用freeze_local/unfreeze_local,因此需要用计数机制
    pub(crate) fn freeze_local(&mut self, id: u16) {
        self.worker += 0x0400;
        self.worker |= id;
    }
    pub(crate) fn unfreeze_local(&mut self) {
        self.worker -= 0x0400;
        if self.worker < THREADS_MAX as u16 {
            self.worker = 0;
        }
    }
}

fn sext_new(old: u8, shift: u8, val: bool) -> u8 {
    if val {
        old | (0x01 << shift)
    } else {
        old & !(0x01 << shift)
    }
}

fn sext_get(sext: u8, shift: u8) -> bool {
    (sext & (0x01 << shift)) != 0
}