mod data;
mod entry;
mod ring;
use std::fmt;
use std::time::{Duration, Instant};
use self::data::DataCell;
use self::entry::Geometry;
use self::ring::{Deq, IndexRing};
use crate::error::{
PopError, PopTimeoutError, PushError, PushTimeoutError, TryPopError, TryPushError,
};
use crate::sync::{Backoff, WaitQueue};
use crate::traits::forward_bounded_queue;
pub struct ScqQueue<T> {
aq: IndexRing,
fq: IndexRing,
consumers: WaitQueue,
producers: WaitQueue,
data: Box<[DataCell<T>]>,
}
impl<T> ScqQueue<T> {
#[must_use]
pub fn new(capacity: usize) -> Self {
Self::with_start_position(capacity, None)
}
fn with_start_position(capacity: usize, start: Option<usize>) -> Self {
assert!(capacity > 0, "capacity must be non-zero");
assert!(capacity <= 1 << 32, "capacity too large");
let geo = Geometry::new(capacity.next_power_of_two().trailing_zeros());
let start = start.unwrap_or(geo.ring_len());
Self {
aq: IndexRing::new_empty(geo, start),
fq: IndexRing::new_full(geo, start),
consumers: WaitQueue::new(),
producers: WaitQueue::new(),
data: (0..geo.n()).map(|_| DataCell::new()).collect(),
}
}
#[must_use]
pub fn capacity(&self) -> usize {
self.data.len()
}
pub fn try_push(&self, item: T) -> Result<(), TryPushError<T>> {
if self.aq.is_closed() {
return Err(TryPushError::Closed(item));
}
let Deq::Item(index) = self.fq.dequeue(false) else {
return Err(if self.aq.is_closed() {
TryPushError::Closed(item)
} else {
TryPushError::Full(item)
});
};
unsafe { self.data[index].write(item) };
if self.aq.enqueue(index, true).is_ok() {
self.consumers.notify_one();
return Ok(());
}
let item = unsafe { self.data[index].read() };
let _ = self.fq.enqueue(index, false);
Err(TryPushError::Closed(item))
}
pub fn try_pop(&self) -> Result<T, TryPopError> {
let result = match self.aq.dequeue(false) {
Deq::EmptyByThreshold if self.aq.is_closed() => self.aq.dequeue(true),
other => other,
};
match result {
Deq::Item(index) => {
let item = unsafe { self.data[index].read() };
let _ = self.fq.enqueue(index, false);
self.producers.notify_one();
Ok(item)
}
Deq::EmptyAtTail { closed: true } => Err(TryPopError::Closed),
Deq::EmptyAtTail { closed: false } | Deq::EmptyByThreshold => Err(TryPopError::Empty),
}
}
fn push_until(
&self,
mut item: T,
deadline: Option<Instant>,
) -> Result<(), PushTimeoutError<T>> {
let mut backoff = Backoff::new();
loop {
match self.try_push(item) {
Ok(()) => return Ok(()),
Err(TryPushError::Closed(v)) => return Err(PushTimeoutError::Closed(v)),
Err(TryPushError::Full(v)) => item = v,
}
if deadline.is_some_and(|d| Instant::now() >= d) {
return Err(PushTimeoutError::Timeout(item));
}
if backoff.is_completed() {
self.producers
.wait_until(|| self.fq.ready() || self.aq.is_closed(), deadline);
backoff.reset();
} else {
backoff.snooze();
}
}
}
fn pop_until(&self, deadline: Option<Instant>) -> Result<T, PopTimeoutError> {
let mut backoff = Backoff::new();
loop {
match self.try_pop() {
Ok(v) => return Ok(v),
Err(TryPopError::Closed) => return Err(PopTimeoutError::Closed),
Err(TryPopError::Empty) => {}
}
if deadline.is_some_and(|d| Instant::now() >= d) {
return Err(PopTimeoutError::Timeout);
}
if backoff.is_completed() {
self.consumers.wait_until(|| self.aq.ready(), deadline);
backoff.reset();
} else {
backoff.snooze();
}
}
}
pub fn push(&self, item: T) -> Result<(), PushError<T>> {
match self.push_until(item, None) {
Ok(()) => Ok(()),
Err(PushTimeoutError::Closed(v)) => Err(PushError(v)),
Err(PushTimeoutError::Timeout(_)) => unreachable!("no deadline was set"),
}
}
pub fn pop(&self) -> Result<T, PopError> {
match self.pop_until(None) {
Ok(v) => Ok(v),
Err(PopTimeoutError::Closed) => Err(PopError),
Err(PopTimeoutError::Timeout) => unreachable!("no deadline was set"),
}
}
pub fn push_timeout(&self, item: T, timeout: Duration) -> Result<(), PushTimeoutError<T>> {
self.push_until(item, Instant::now().checked_add(timeout))
}
pub fn pop_timeout(&self, timeout: Duration) -> Result<T, PopTimeoutError> {
self.pop_until(Instant::now().checked_add(timeout))
}
pub fn close(&self) -> bool {
let newly_closed = self.aq.close();
if newly_closed {
self.consumers.notify_all();
self.producers.notify_all();
}
newly_closed
}
pub fn is_closed(&self) -> bool {
self.aq.is_closed()
}
pub fn len(&self) -> usize {
self.aq.len()
}
pub fn is_empty(&self) -> bool {
self.len() == 0
}
pub fn is_full(&self) -> bool {
self.len() == self.capacity()
}
}
impl<T> Drop for ScqQueue<T> {
fn drop(&mut self) {
if !std::mem::needs_drop::<T>() {
return;
}
let data = &self.data;
self.aq.for_each_index(|index| {
unsafe { data[index].drop_in_place() };
});
}
}
impl<T> fmt::Debug for ScqQueue<T> {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.debug_struct("ScqQueue")
.field("capacity", &self.capacity())
.field("len", &self.len())
.field("closed", &self.is_closed())
.field("aq", &self.aq)
.field("fq", &self.fq)
.finish_non_exhaustive()
}
}
forward_bounded_queue!(ScqQueue);
#[cfg(all(test, not(loom)))]
mod tests {
use super::*;
#[test]
fn empty_pops_then_push_pop_round_trips() {
for cap in [1, 2, 4, 16] {
let q = ScqQueue::new(cap);
for round in 0..200 {
for _ in 0..(3 * q.capacity() + 2) {
assert_eq!(q.try_pop(), Err(TryPopError::Empty));
}
q.try_push(round).unwrap();
assert_eq!(q.try_pop(), Ok(round));
}
}
}
#[test]
fn full_pushes_then_pop_frees_a_cell() {
for cap in [1, 2, 8] {
let q = ScqQueue::new(cap);
let n = q.capacity();
for i in 0..n {
q.try_push(i).unwrap();
}
for round in 0..100 {
for _ in 0..(3 * n + 2) {
assert!(q.try_push(usize::MAX).unwrap_err().is_full());
}
assert_eq!(q.try_pop(), Ok(round));
q.try_push(n + round).unwrap();
}
}
}
#[test]
fn huge_positions_do_not_overflow() {
let geo = Geometry::new(3);
let start = (1usize << 62) / geo.ring_len() * geo.ring_len();
let q = ScqQueue::with_start_position(8, Some(start));
for round in 0..50 {
for i in 0..8 {
q.try_push(round * 8 + i).unwrap();
}
assert!(q.try_push(0).unwrap_err().is_full());
for i in 0..8 {
assert_eq!(q.try_pop(), Ok(round * 8 + i));
}
assert_eq!(q.try_pop(), Err(TryPopError::Empty));
}
}
}