use alloc::boxed::Box;
use core::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
use crate::config::Config;
use crate::control::{AtomicControl, WakerSlot};
use crate::pubqueue::{PQContents, Priority, PubQueue, Stealer};
use crate::task::Task;
use crate::util::CachePadded;
pub struct Queue<T: Task> {
pub(crate) config: Config,
pub(crate) started: Box<CachePadded<AtomicUsize>>,
pub(crate) pending: Box<CachePadded<AtomicUsize>>,
pub(crate) shutdown: Box<CachePadded<AtomicBool>>,
pub(crate) priority: Box<[CachePadded<Priority<T>>]>,
pub(crate) stealer: Box<[CachePadded<Stealer>]>,
pub(crate) pq_contents: Box<[Box<PQContents<T>>]>,
pub(crate) control: Box<[CachePadded<AtomicControl>]>,
pub(crate) waker_slot: Box<[CachePadded<WakerSlot>]>,
}
impl<T: Task> Queue<T> {
pub fn new(config: Config) -> Self {
let num_workers = config.num_workers.get();
let started = Box::new(CachePadded::new(AtomicUsize::new(0)));
let pending = Box::new(CachePadded::new(AtomicUsize::new(0)));
let shutdown = Box::new(CachePadded::new(AtomicBool::new(false)));
let priority = (0..num_workers)
.map(|_| CachePadded::new(Priority::default()))
.collect();
let stealer = (0..num_workers)
.map(|id| CachePadded::new(unsafe { Stealer::new(id, &config) }))
.collect();
let pq_contents = (0..num_workers)
.map(|_| PQContents::new_boxed(&config))
.collect();
let control = (0..num_workers)
.map(|_| CachePadded::new(AtomicControl::new()))
.collect();
let waker_slot = (0..num_workers)
.map(|_| CachePadded::new(WakerSlot::new()))
.collect();
Self {
config,
started,
pending,
shutdown,
priority,
stealer,
pq_contents,
control,
waker_slot,
}
}
#[inline]
pub const fn config(&self) -> &Config {
&self.config
}
pub fn shutdown(&self) {
self.shutdown.store(true, Ordering::Relaxed);
for (control, waker_slot) in self.control.iter().zip(&self.waker_slot) {
unsafe { control.shutdown(waker_slot) };
}
}
pub fn has_shut_down(&self) -> bool {
self.shutdown.load(Ordering::Relaxed)
}
pub(crate) fn all_started(&self) -> bool {
self.started.load(Ordering::Relaxed) == self.config.num_workers.get()
}
pub(crate) unsafe fn steal(
&self,
min_priority: Option<T::Priority>,
replacement: usize,
) -> Option<PubQueue> {
debug_assert!(replacement < self.config.num_workers.get());
let num_workers = self.config.num_workers.get();
let index = (0..num_workers).find(|&index| unsafe {
self.priority[index].steal_if_above(min_priority).is_some()
})?;
Some(unsafe { self.stealer[index].steal(replacement, &self.config) })
}
pub(crate) fn wake(&self) {
for (control, waker_slot) in self.control.iter().zip(&self.waker_slot) {
if unsafe { control.try_wake(waker_slot) } {
break;
}
}
}
}