hirun 0.1.22

A concurrent framework for asynchronous programming based on event-driven, non-blocking I/O mechanism
Documentation
use super::{TaskQueue, TaskRef};
use core::sync::atomic::{AtomicUsize, Ordering};

pub(crate) struct WorkerSender {
    pos: u16,
    id: u16,
}

impl WorkerSender {
    pub(crate) fn new(id: u16) -> Self {
        Self { pos: 0, id }
    }

    pub(crate) fn send(&mut self, queue: &[TaskQueue], task: TaskRef) -> Option<TaskRef> {
        let pos = self.pos as usize;
        self.pos = if pos < queue.len() - 1 {
            (pos + 1) as u16
        } else {
            0
        };
        if pos != self.id as usize {
            queue[pos].push(task);
            None
        } else {
            Some(task)
        }
    }
}

pub(crate) struct GroupSender {
    pos: AtomicUsize,
}

impl GroupSender {
    pub(crate) fn new() -> Self {
        Self {
            pos: AtomicUsize::new(0),
        }
    }

    pub(crate) fn send(&self, queue: &[TaskQueue], task: TaskRef) {
        let pos = self.pos.fetch_add(1, Ordering::Relaxed) % queue.len();
        queue[pos].push(task);
    }
}