use crate::sync::Ordering::{AcqRel, Acquire, Relaxed, Release, SeqCst};
use crate::sync::{
spin_loop, yield_now, AtomicPtr, AtomicU8, AtomicUsize, CachePadded, Mutex, UnsafeCell,
};
use std::alloc::{alloc, dealloc, handle_alloc_error, Layout};
use std::mem::MaybeUninit;
use std::ptr;
#[cfg(any(feature = "loom", fcrs_small_blocks))]
pub(crate) const LAP_SHIFT: usize = 2;
#[cfg(not(any(feature = "loom", fcrs_small_blocks)))]
pub(crate) const LAP_SHIFT: usize = 6;
pub(crate) const LAP: usize = 1 << LAP_SHIFT;
pub(crate) const BLOCK_CAP: usize = LAP - 1;
const OFFSET_MASK: usize = LAP - 1;
const TAG_WRITTEN: usize = 1;
#[inline(always)]
fn written_tag(lap: usize) -> usize {
(lap << 2) | TAG_WRITTEN
}
#[inline(always)]
fn items_between(head: usize, tail: usize) -> usize {
if head >= tail {
0
} else {
(tail - head) - ((tail >> LAP_SHIFT) - (head >> LAP_SHIFT))
}
}
pub(crate) struct Backoff(u32);
const SPIN_LIMIT: u32 = 6;
const TIGHT_POLLS: u32 = 64;
const SNOOZE_LIMIT: u32 = TIGHT_POLLS + 8;
impl Backoff {
#[inline(always)]
pub(crate) fn new() -> Self {
Backoff(0)
}
#[cold]
#[inline(never)]
pub(crate) fn snooze_bounded(&mut self) -> bool {
if self.0 >= SNOOZE_LIMIT {
return false;
}
self.snooze();
true
}
#[inline(always)]
fn spin(&mut self) {
for _ in 0..(1u32 << self.0.min(SPIN_LIMIT)) {
spin_loop();
}
if self.0 <= SPIN_LIMIT {
self.0 += 1;
}
}
#[inline(always)]
pub(crate) fn snooze(&mut self) {
if self.0 < TIGHT_POLLS {
spin_loop();
self.0 += 1;
} else if self.0 < SNOOZE_LIMIT {
for _ in 0..(1u32 << (self.0 - TIGHT_POLLS)) {
spin_loop();
}
self.0 += 1;
} else {
yield_now();
}
}
}
#[inline(always)]
fn publish(state: &AtomicUsize, tag: usize) {
#[cfg(all(
target_arch = "aarch64",
fcrs_publish_asm,
not(any(feature = "loom", miri))
))]
{
unsafe { core::arch::asm!("dmb ishst", options(nostack, preserves_flags)) };
state.store(tag, Relaxed);
}
#[cfg(all(fcrs_publish_fence, not(any(feature = "loom", miri))))]
{
std::sync::atomic::fence(Release);
state.store(tag, Relaxed);
}
#[cfg(all(fcrs_publish_rmw, not(any(feature = "loom", miri))))]
{
state.swap(tag, Release);
}
#[cfg(not(any(
all(
target_arch = "aarch64",
fcrs_publish_asm,
not(any(feature = "loom", miri))
),
all(fcrs_publish_fence, not(any(feature = "loom", miri))),
all(fcrs_publish_rmw, not(any(feature = "loom", miri)))
)))]
{
state.store(tag, Release);
}
}
#[inline(always)]
fn mark_read(mark: &AtomicU8) {
#[cfg(all(fcrs_mark_rmw, not(any(feature = "loom", miri))))]
{
mark.swap(1, Release);
}
#[cfg(not(all(fcrs_mark_rmw, not(any(feature = "loom", miri)))))]
{
mark.store(1, Release);
}
}
#[repr(C)]
struct Slot<T> {
value: UnsafeCell<MaybeUninit<T>>,
state: AtomicUsize,
}
const POOLED: usize = usize::MAX;
#[repr(C)]
struct ProducerHeader<T> {
start: AtomicUsize,
next: AtomicPtr<Block<T>>,
prev: AtomicPtr<Block<T>>,
}
#[repr(C)]
struct ConsumerHeader {
read_marks: [AtomicU8; BLOCK_CAP],
}
#[repr(C)]
struct Block<T> {
phdr: CachePadded<ProducerHeader<T>>,
chdr: CachePadded<ConsumerHeader>,
slots: [Slot<T>; BLOCK_CAP],
}
impl<T> Block<T> {
const LAYOUT: Layout = Layout::new::<Block<T>>();
fn allocate() -> *mut Block<T> {
unsafe {
let block = alloc(Self::LAYOUT) as *mut Block<T>;
if block.is_null() {
handle_alloc_error(Self::LAYOUT);
}
ptr::addr_of_mut!((*block).phdr.0.start).write(AtomicUsize::new(POOLED));
ptr::addr_of_mut!((*block).phdr.0.next).write(AtomicPtr::new(ptr::null_mut()));
ptr::addr_of_mut!((*block).phdr.0.prev).write(AtomicPtr::new(ptr::null_mut()));
let marks = ptr::addr_of_mut!((*block).chdr.0.read_marks) as *mut AtomicU8;
for i in 0..BLOCK_CAP {
marks.add(i).write(AtomicU8::new(0));
}
let slots = ptr::addr_of_mut!((*block).slots) as *mut Slot<T>;
for i in 0..BLOCK_CAP {
slots.add(i).write(Slot {
value: UnsafeCell::new(MaybeUninit::uninit()),
state: AtomicUsize::new(0),
});
}
block
}
}
unsafe fn release(block: *mut Block<T>) {
ptr::drop_in_place(block);
dealloc(block as *mut u8, Self::LAYOUT);
}
#[inline]
fn readers_done(&self) -> bool {
self.chdr.0.read_marks[..BLOCK_CAP - 1]
.iter()
.all(|mark| mark.load(Acquire) != 0)
}
unsafe fn reset(&self) {
for mark in &self.chdr.0.read_marks {
mark.store(0, Relaxed);
}
self.phdr.0.start.store(POOLED, Relaxed);
}
}
struct BlockPool<T> {
ready: Vec<*mut Block<T>>,
retired: Vec<*mut Block<T>>,
}
struct Position<T> {
index: AtomicUsize,
block: AtomicPtr<Block<T>>,
cached: AtomicUsize,
contended: AtomicUsize,
}
const CONTENTION_WINDOW: usize = 64;
const UNBOUNDED: usize = usize::MAX;
pub(crate) struct Queue<T> {
head: CachePadded<Position<T>>,
tail: CachePadded<Position<T>>,
spare: CachePadded<AtomicPtr<Block<T>>>,
pool: Mutex<BlockPool<T>>,
capacity: usize,
}
unsafe impl<T: Send> Send for Queue<T> {}
unsafe impl<T: Send> Sync for Queue<T> {}
impl<T> Queue<T> {
pub(crate) fn new(capacity: Option<usize>) -> Self {
let first = Block::<T>::allocate();
unsafe { (*first).phdr.0.start.store(0, Relaxed) };
Queue {
head: CachePadded(Position {
index: AtomicUsize::new(0),
block: AtomicPtr::new(first),
cached: AtomicUsize::new(0),
contended: AtomicUsize::new(0),
}),
tail: CachePadded(Position {
index: AtomicUsize::new(0),
block: AtomicPtr::new(first),
cached: AtomicUsize::new(0),
contended: AtomicUsize::new(0),
}),
spare: CachePadded(AtomicPtr::new(ptr::null_mut())),
pool: Mutex::new(BlockPool {
ready: Vec::new(),
retired: Vec::new(),
}),
capacity: capacity.unwrap_or(UNBOUNDED),
}
}
#[inline(always)]
pub(crate) fn capacity(&self) -> Option<usize> {
if self.capacity == UNBOUNDED {
None
} else {
Some(self.capacity)
}
}
#[inline]
pub(crate) fn receiver_parking(&self) -> usize {
self.tail.index.fetch_add(0, AcqRel)
}
#[inline]
pub(crate) fn claims_below(&self, tail_snapshot: usize) -> bool {
items_between(self.head.index.load(Acquire), tail_snapshot) != 0
}
#[inline]
pub(crate) fn sender_parking(&self) {
self.head.index.fetch_add(0, AcqRel);
}
#[inline]
pub(crate) fn is_empty_seqcst(&self) -> bool {
let head = self.head.index.load(SeqCst);
let tail = self.tail.index.load(SeqCst);
items_between(head, tail) == 0
}
#[inline]
pub(crate) fn len(&self) -> usize {
let head = self.head.index.load(Relaxed);
let tail = self.tail.index.load(Relaxed);
items_between(head, tail).min(self.capacity)
}
#[inline]
fn lock_pool(&self) -> crate::sync::MutexGuard<'_, BlockPool<T>> {
self.pool.lock().unwrap_or_else(|e| e.into_inner())
}
#[inline(always)]
fn take_block(&self) -> *mut Block<T> {
let block = self.spare.swap(ptr::null_mut(), Acquire);
if !block.is_null() {
return block;
}
self.take_block_slow()
}
#[cold]
fn take_block_slow(&self) -> *mut Block<T> {
let retired = {
let mut pool = self.lock_pool();
if let Some(block) = pool.ready.pop() {
return block;
}
let ready = pool.retired.iter().position(|&block| {
unsafe { (*block).readers_done() }
});
ready.map(|i| pool.retired.swap_remove(i))
};
if let Some(block) = retired {
unsafe { (*block).reset() };
return block;
}
Block::allocate()
}
#[inline(always)]
unsafe fn recycle(&self, block: *mut Block<T>) {
if self
.spare
.compare_exchange(ptr::null_mut(), block, Release, Relaxed)
.is_err()
{
self.recycle_slow(block);
}
}
#[cold]
fn recycle_slow(&self, block: *mut Block<T>) {
self.lock_pool().ready.push(block);
}
#[inline(always)]
pub(crate) fn push(&self, value: T, wake_flag: &AtomicUsize) -> Result<bool, T> {
let cap = self.capacity;
let tail = if cap == UNBOUNDED {
if self.tail.contended.load(Relaxed) == 0 {
self.claim_fetch_add()
} else {
self.claim_cas()
}
} else {
match self.claim_bounded(cap) {
Some(tail) => tail,
None => return Err(value),
}
};
let wake = wake_flag.load(Relaxed) != 0;
let offset = tail & OFFSET_MASK;
let block = self.find_block(tail);
unsafe {
if offset + 1 == BLOCK_CAP {
self.install_next(block, tail);
}
let slot = (*block).slots.get_unchecked(offset);
slot.value.with_mut(|p| p.write(MaybeUninit::new(value)));
publish(&slot.state, written_tag(tail >> LAP_SHIFT));
}
Ok(wake)
}
#[inline(always)]
fn claim_fetch_add(&self) -> usize {
loop {
let expected = self.tail.index.load(Relaxed);
let tail = self.tail.index.fetch_add(1, AcqRel);
if tail != expected {
self.tail.contended.store(CONTENTION_WINDOW, Relaxed);
}
if tail & OFFSET_MASK != BLOCK_CAP {
return tail;
}
}
}
#[inline(never)]
fn claim_cas(&self) -> usize {
let mut backoff = Backoff::new();
let mut tail = self.tail.index.load(Acquire);
let mut failed = false;
loop {
match self
.tail
.index
.compare_exchange_weak(tail, tail + 1, AcqRel, Acquire)
{
Ok(_) => {
if tail & OFFSET_MASK == BLOCK_CAP {
tail += 1;
continue;
}
let level = self.tail.contended.load(Relaxed);
if failed {
self.tail.contended.store(CONTENTION_WINDOW, Relaxed);
} else if level > 0 {
self.tail.contended.store(level - 1, Relaxed);
}
return tail;
}
Err(actual) => {
failed = true;
tail = actual;
backoff.spin();
}
}
}
}
#[inline(never)]
fn claim_bounded(&self, cap: usize) -> Option<usize> {
let mut backoff = Backoff::new();
let mut tail = self.tail.index.load(Acquire);
loop {
if tail & OFFSET_MASK == BLOCK_CAP {
match self
.tail
.index
.compare_exchange_weak(tail, tail + 1, AcqRel, Acquire)
{
Ok(_) => tail += 1,
Err(actual) => tail = actual,
}
continue;
}
let mut head = self.tail.cached.load(Relaxed);
if items_between(head, tail) >= cap {
head = self.head.index.load(Acquire);
if items_between(head, tail) >= cap {
return None;
}
self.tail.cached.store(head, Relaxed);
}
match self
.tail
.index
.compare_exchange_weak(tail, tail + 1, AcqRel, Acquire)
{
Ok(_) => return Some(tail),
Err(actual) => {
tail = actual;
backoff.spin();
}
}
}
}
#[inline(always)]
fn find_block(&self, pos: usize) -> *mut Block<T> {
let start = pos & !OFFSET_MASK;
let block = self.tail.block.load(Acquire);
if unsafe { (*block).phdr.0.start.load(Acquire) } == start {
return block;
}
self.find_block_slow(pos)
}
#[cold]
#[inline(never)]
fn find_block_slow(&self, pos: usize) -> *mut Block<T> {
let start = pos & !OFFSET_MASK;
let mut backoff = Backoff::new();
let mut block = self.tail.block.load(Acquire);
#[cfg(fcrs_debug)]
let mut iters: u64 = 0;
loop {
#[cfg(fcrs_debug)]
{
iters += 1;
if iters == 20_000_000 {
unsafe {
let tb = self.tail.block.load(Relaxed);
let mut chain = String::new();
let mut b = tb;
for _ in 0..6 {
if b.is_null() {
break;
}
chain += &format!(
"[{:p} start={} next={:p}] <-prev- ",
b,
(*b).phdr.0.start.load(Relaxed) as isize,
(*b).phdr.0.next.load(Relaxed)
);
b = (*b).phdr.0.prev.load(Relaxed);
}
eprintln!(
"FIND_BLOCK STUCK pos={} (lap {} off {}) target_start={} cur_block={:p} cur_start={} tail.index={} tail.block={:p} chain-from-tail: {}",
pos, pos >> LAP_SHIFT, pos & OFFSET_MASK, start, block,
(*block).phdr.0.start.load(Relaxed) as isize,
self.tail.index.load(Relaxed), tb, chain
);
}
}
}
let s = unsafe { (*block).phdr.0.start.load(Acquire) };
if s == start {
return block;
}
if s > start && s != POOLED {
block = unsafe { (*block).phdr.0.prev.load(Acquire) };
debug_assert!(!block.is_null(), "prev of a block ahead of the target");
continue;
}
backoff.snooze();
block = self.tail.block.load(Acquire);
}
}
#[cold]
unsafe fn install_next(&self, block: *mut Block<T>, tail: usize) {
let next = self.take_block();
(*next).phdr.0.prev.store(block, Relaxed);
(*next).phdr.0.next.store(ptr::null_mut(), Relaxed);
(*next)
.phdr
.0
.start
.store((tail & !OFFSET_MASK) + LAP, Release);
(*block).phdr.0.next.store(next, Release);
self.tail.block.store(next, Release);
}
#[inline(always)]
pub(crate) fn pop(&self, wake_flag: &AtomicUsize) -> Option<(T, bool)> {
let mut backoff = Backoff::new();
let mut head = self.head.index.load(Acquire);
let mut block = self.head.block.load(Acquire);
loop {
let offset = head & OFFSET_MASK;
if offset == BLOCK_CAP {
if !backoff.snooze_bounded() {
return None;
}
head = self.head.index.load(Acquire);
block = self.head.block.load(Acquire);
continue;
}
let lap = head >> LAP_SHIFT;
let expected = written_tag(lap);
let slot = unsafe { (*block).slots.get_unchecked(offset) };
if slot.state.load(Acquire) != expected {
let now = self.head.index.load(Acquire);
if now == head {
return None;
}
head = now;
block = self.head.block.load(Acquire);
continue;
}
match self
.head
.index
.compare_exchange_weak(head, head + 1, AcqRel, Acquire)
{
Ok(_) => {
let wake = self.capacity != UNBOUNDED && wake_flag.load(Relaxed) != 0;
unsafe {
let value = slot.value.with(|p| p.read().assume_init());
if offset + 1 == BLOCK_CAP {
self.finish_block(block, head);
} else {
mark_read((*block).chdr.0.read_marks.get_unchecked(offset));
}
return Some((value, wake));
}
}
Err(actual) => {
head = actual;
block = self.head.block.load(Acquire);
backoff.spin();
}
}
}
}
}
impl<T> Queue<T> {
#[inline(always)]
pub(crate) unsafe fn pop_single(&self, wake_flag: &AtomicUsize) -> Option<(T, bool)> {
let head = self.head.index.load(Relaxed);
let block = self.head.block.load(Relaxed);
let offset = head & OFFSET_MASK;
debug_assert!(offset < BLOCK_CAP);
let slot = (*block).slots.get_unchecked(offset);
if slot.state.load(Acquire) != written_tag(head >> LAP_SHIFT) {
return None;
}
let value = slot.value.with(|p| p.read().assume_init());
let last = offset + 1 == BLOCK_CAP;
if last {
let next = (*block).phdr.0.next.load(Acquire);
debug_assert!(!next.is_null());
self.head.block.store(next, Release);
}
let advance = if last { 2 } else { 1 };
let wake = if self.capacity == UNBOUNDED {
self.head.index.store(head + advance, Release);
false
} else {
self.head.index.fetch_add(advance, AcqRel);
wake_flag.load(Relaxed) != 0
};
if last {
self.recycle_single(block);
}
Some((value, wake))
}
#[cold]
#[inline(never)]
unsafe fn recycle_single(&self, block: *mut Block<T>) {
(*block).phdr.0.start.store(POOLED, Relaxed);
self.recycle(block);
}
#[inline]
pub(crate) unsafe fn pop_single_batch(
&self,
buffer: &mut Vec<T>,
limit: usize,
wake_flag: &AtomicUsize,
) -> (usize, usize) {
if limit == 0 {
return (0, 0);
}
let head = self.head.index.load(Relaxed);
let block = self.head.block.load(Relaxed);
let offset = head & OFFSET_MASK;
let tag = written_tag(head >> LAP_SHIFT);
let max = limit.min(BLOCK_CAP - offset);
if (*block).slots.get_unchecked(offset).state.load(Acquire) != tag {
return (0, 0);
}
buffer.reserve(max);
let mut count = 0;
while count < max {
let slot = (*block).slots.get_unchecked(offset + count);
if slot.state.load(Acquire) != tag {
break;
}
buffer.push(slot.value.with(|p| p.read().assume_init()));
count += 1;
}
let last = offset + count == BLOCK_CAP;
if last {
let next = (*block).phdr.0.next.load(Acquire);
self.head.block.store(next, Release);
}
let advance = count + usize::from(last);
let wake = if self.capacity == UNBOUNDED {
self.head.index.store(head + advance, Release);
0
} else {
self.head.index.fetch_add(advance, AcqRel);
wake_flag.load(Relaxed).min(count)
};
if last {
self.recycle_single(block);
}
(count, wake)
}
#[cold]
#[inline(never)]
unsafe fn finish_block(&self, block: *mut Block<T>, head: usize) {
let next = (*block).phdr.0.next.load(Acquire);
debug_assert!(
!next.is_null(),
"last slot was published before its next block"
);
self.head.block.store(next, Release);
self.head.index.swap(head + 2, AcqRel);
if !(*block).readers_done() {
self.lock_pool().retired.push(block);
return;
}
(*block).reset();
self.recycle(block);
}
#[doc(hidden)]
pub(crate) fn debug_dump(&self) -> String {
unsafe {
let h = self.head.index.load(Relaxed);
let t = self.tail.index.load(Relaxed);
let hb = self.head.block.load(Relaxed);
let tb = self.tail.block.load(Relaxed);
let hs = (*hb).phdr.0.start.load(Relaxed);
let ts = (*tb).phdr.0.start.load(Relaxed);
let hoff = h & OFFSET_MASK;
let hstate = if hoff < BLOCK_CAP {
(*hb).slots.get_unchecked(hoff).state.load(Relaxed)
} else {
0
};
let mut chain = String::new();
let mut b = tb;
for _ in 0..6 {
if b.is_null() {
break;
}
chain += &format!(
"[{:p} start={} next={:p}] <-prev- ",
b,
(*b).phdr.0.start.load(Relaxed) as isize,
(*b).phdr.0.next.load(Relaxed)
);
b = (*b).phdr.0.prev.load(Relaxed);
}
format!(
"head={} (lap {} off {}) head.block={:p} start={} head_slot_state={} expected={} head.contended={} | tail={} tail.block={:p} start={} tail.contended={} | spare={:p} pool={} | chain: {}",
h, h >> LAP_SHIFT, hoff, hb, hs as isize, hstate, written_tag(h >> LAP_SHIFT),
self.head.contended.load(Relaxed),
t, tb, ts as isize, self.tail.contended.load(Relaxed),
self.spare.load(Relaxed), {
let pool = self.lock_pool();
pool.ready.len() + pool.retired.len()
}, chain
)
}
}
}
impl<T> Drop for Queue<T> {
fn drop(&mut self) {
unsafe {
let mut head = self.head.index.load(Relaxed);
let tail = self.tail.index.load(Relaxed);
let mut block = self.head.block.load(Relaxed);
while head < tail && !block.is_null() {
let offset = head & OFFSET_MASK;
if offset == BLOCK_CAP {
block = (*block).phdr.0.next.load(Relaxed);
head += 1;
continue;
}
let slot = (*block).slots.get_unchecked(offset);
if slot.state.load(Relaxed) == written_tag(head >> LAP_SHIFT) {
slot.value
.with_mut(|p| ptr::drop_in_place((*p).as_mut_ptr()));
}
head += 1;
}
let mut block = self.head.block.load(Relaxed);
while !block.is_null() {
let next = (*block).phdr.0.next.load(Relaxed);
Block::release(block);
block = next;
}
let spare = self.spare.load(Relaxed);
if !spare.is_null() {
Block::release(spare);
}
let pool = &mut *self.lock_pool();
for block in pool.ready.drain(..).chain(pool.retired.drain(..)) {
Block::release(block);
}
}
}
}
#[cfg(all(test, not(feature = "loom")))]
mod tests {
use super::*;
#[test]
fn items_between_skips_sentinels() {
assert_eq!(items_between(0, 0), 0);
assert_eq!(items_between(0, BLOCK_CAP), BLOCK_CAP);
assert_eq!(items_between(0, LAP), BLOCK_CAP);
assert_eq!(items_between(0, LAP + 1), BLOCK_CAP + 1);
assert_eq!(items_between(BLOCK_CAP, LAP), 0);
assert_eq!(items_between(5, 3), 0);
assert_eq!(items_between(LAP - 2, LAP + 2), 3);
}
#[test]
fn push_pop_across_many_blocks() {
let q = Queue::new(None);
let flag = AtomicUsize::new(0);
for i in 0..10 * LAP {
q.push(i, &flag).unwrap();
}
assert_eq!(q.len(), 10 * LAP);
for i in 0..10 * LAP {
assert_eq!(q.pop(&flag).map(|v| v.0), Some(i));
}
assert_eq!(q.pop(&flag).map(|v| v.0), None);
assert_eq!(q.len(), 0);
}
#[test]
fn interleaved_push_pop_recycles_blocks() {
let q = Queue::new(None);
let flag = AtomicUsize::new(0);
for round in 0..1000usize {
q.push(round, &flag).unwrap();
q.push(round + 1, &flag).unwrap();
assert_eq!(q.pop(&flag).map(|v| v.0), Some(round));
assert_eq!(q.pop(&flag).map(|v| v.0), Some(round + 1));
}
assert_eq!(q.pop(&flag).map(|v| v.0), None);
assert!(q.lock_pool().ready.len() <= 1);
}
#[test]
fn backlog_grows_and_shrinks_through_pool() {
let q = Queue::new(None);
let flag = AtomicUsize::new(0);
for round in 0..3 {
for i in 0..20 * LAP {
q.push(round * 100_000 + i, &flag).unwrap();
}
for i in 0..20 * LAP {
assert_eq!(q.pop(&flag).map(|v| v.0), Some(round * 100_000 + i));
}
assert_eq!(q.pop(&flag).map(|v| v.0), None);
}
assert!(q.lock_pool().ready.len() >= 15);
}
#[test]
fn bounded_is_exact() {
let q = Queue::new(Some(3));
let flag = AtomicUsize::new(0);
q.push(1, &flag).unwrap();
q.push(2, &flag).unwrap();
q.push(3, &flag).unwrap();
assert_eq!(q.push(4, &flag), Err(4));
assert_eq!(q.pop(&flag).map(|v| v.0), Some(1));
q.push(4, &flag).unwrap();
assert_eq!(q.push(5, &flag), Err(5));
for expected in 2..=4 {
assert_eq!(q.pop(&flag).map(|v| v.0), Some(expected));
}
assert_eq!(q.pop(&flag).map(|v| v.0), None);
}
#[test]
fn bounded_across_block_boundary() {
let cap = LAP + 3;
let q = Queue::new(Some(cap));
let flag = AtomicUsize::new(0);
for i in 0..cap {
q.push(i, &flag).unwrap();
}
assert_eq!(q.push(usize::MAX, &flag), Err(usize::MAX));
for i in 0..cap {
assert_eq!(q.pop(&flag).map(|v| v.0), Some(i));
q.push(cap + i, &flag).unwrap();
assert_eq!(q.push(usize::MAX, &flag), Err(usize::MAX));
}
}
#[test]
fn parking_snapshot_reports_claims() {
let q = Queue::new(None);
let flag = AtomicUsize::new(0);
assert!(!q.claims_below(q.receiver_parking()));
q.push(7usize, &flag).unwrap();
let tail = q.receiver_parking();
assert!(q.claims_below(tail));
assert_eq!(q.pop(&flag).map(|v| v.0), Some(7));
assert!(!q.claims_below(tail));
}
#[test]
fn retired_block_stays_live_until_its_reader_finishes() {
for reclaim in [false, true] {
let q = Queue::new(None);
let flag = AtomicUsize::new(0);
for i in 0..LAP {
q.push(Box::new(i), &flag).unwrap();
}
let block = q.head.block.load(Acquire);
let slot = unsafe { &(*block).slots[0] };
assert_eq!(slot.state.load(Acquire), written_tag(0));
q.head
.index
.compare_exchange(0, 1, AcqRel, Acquire)
.unwrap();
for expected in 1..LAP {
assert_eq!(*q.pop(&flag).unwrap().0, expected);
}
for value in LAP..8 * LAP {
q.push(Box::new(value), &flag).unwrap();
assert_eq!(*q.pop(&flag).unwrap().0, value);
}
assert_eq!(q.lock_pool().retired.as_slice(), &[block]);
let held = unsafe { slot.value.with(|p| p.read().assume_init()) };
assert_eq!(*held, 0);
unsafe { mark_read(&(*block).chdr.0.read_marks[0]) };
drop(held);
if reclaim {
assert!(q.lock_pool().ready.is_empty());
let reused = q.take_block_slow();
assert_eq!(reused, block);
assert!(q.lock_pool().retired.is_empty());
unsafe { q.recycle(reused) };
}
drop(q);
}
}
#[test]
fn stalled_head_transition_returns_to_the_caller() {
let q = Queue::new(None);
let flag = AtomicUsize::new(0);
for i in 0..LAP {
q.push(i, &flag).unwrap();
}
for expected in 0..BLOCK_CAP - 1 {
assert_eq!(q.pop(&flag).unwrap().0, expected);
}
let head = BLOCK_CAP - 1;
let block = q.head.block.load(Acquire);
let slot = unsafe { &(*block).slots[head] };
assert_eq!(slot.state.load(Acquire), written_tag(0));
q.head
.index
.compare_exchange(head, head + 1, AcqRel, Acquire)
.unwrap();
let held = unsafe { slot.value.with(|p| p.read().assume_init()) };
assert_eq!(held, head);
assert!(q.pop(&flag).is_none());
assert!(q.claims_below(q.receiver_parking()));
unsafe { q.finish_block(block, head) };
assert_eq!(q.pop(&flag).unwrap().0, BLOCK_CAP);
}
#[test]
fn drop_releases_unread_values() {
use std::sync::atomic::{AtomicUsize, Ordering};
static DROPS: AtomicUsize = AtomicUsize::new(0);
struct D;
impl Drop for D {
fn drop(&mut self) {
DROPS.fetch_add(1, Ordering::Relaxed);
}
}
{
let q = Queue::new(None);
let flag = AtomicUsize::new(0);
for _ in 0..(3 * LAP + 7) {
assert!(q.push(D, &flag).is_ok());
}
for _ in 0..(LAP + 2) {
drop(q.pop(&flag));
}
assert_eq!(DROPS.load(Ordering::Relaxed), LAP + 2);
}
assert_eq!(DROPS.load(Ordering::Relaxed), 3 * LAP + 7);
}
#[test]
fn threads_bounded_mpmc_all_values_arrive_once() {
use std::sync::Arc;
use std::thread;
const PRODUCERS: usize = 4;
const CONSUMERS: usize = 2;
const PER_PRODUCER: usize = 10_000;
const CAP: usize = 5;
let q = Arc::new(Queue::new(Some(CAP)));
static FLAG: AtomicUsize = AtomicUsize::new(0);
let mut handles = Vec::new();
for p in 0..PRODUCERS {
let q = q.clone();
handles.push(thread::spawn(move || {
for i in 0..PER_PRODUCER {
let mut v = p * PER_PRODUCER + i;
loop {
match q.push(v, &FLAG) {
Ok(_) => break,
Err(back) => {
v = back;
std::hint::spin_loop();
}
}
}
}
}));
}
let received = Arc::new(std::sync::atomic::AtomicUsize::new(0));
let mut consumers = Vec::new();
for _ in 0..CONSUMERS {
let q = q.clone();
let received = received.clone();
consumers.push(thread::spawn(move || {
let mut seen = Vec::new();
loop {
if let Some((v, _)) = q.pop(&FLAG) {
seen.push(v);
if received.fetch_add(1, std::sync::atomic::Ordering::Relaxed) + 1
== PRODUCERS * PER_PRODUCER
{
break;
}
} else if received.load(std::sync::atomic::Ordering::Relaxed)
== PRODUCERS * PER_PRODUCER
{
break;
} else {
std::hint::spin_loop();
}
}
seen
}));
}
for h in handles {
h.join().unwrap();
}
let mut all: Vec<usize> = consumers
.into_iter()
.flat_map(|h| h.join().unwrap())
.collect();
all.sort_unstable();
assert_eq!(all.len(), PRODUCERS * PER_PRODUCER);
for (i, v) in all.iter().enumerate() {
assert_eq!(*v, i);
}
assert_eq!(q.pop(&FLAG).map(|v| v.0), None);
}
#[test]
fn threads_mpmc_all_values_arrive_once() {
use std::sync::Arc;
use std::thread;
const PRODUCERS: usize = 4;
const CONSUMERS: usize = 4;
const PER_PRODUCER: usize = 20_000;
let q = Arc::new(Queue::new(None));
static FLAG: AtomicUsize = AtomicUsize::new(0);
let mut handles = Vec::new();
for p in 0..PRODUCERS {
let q = q.clone();
handles.push(thread::spawn(move || {
for i in 0..PER_PRODUCER {
q.push(p * PER_PRODUCER + i, &FLAG).unwrap();
}
}));
}
let received = Arc::new(std::sync::atomic::AtomicUsize::new(0));
let mut consumers = Vec::new();
for _ in 0..CONSUMERS {
let q = q.clone();
let received = received.clone();
consumers.push(thread::spawn(move || {
let mut seen = Vec::new();
loop {
if let Some((v, _)) = q.pop(&FLAG) {
seen.push(v);
if received.fetch_add(1, std::sync::atomic::Ordering::Relaxed) + 1
== PRODUCERS * PER_PRODUCER
{
break;
}
} else if received.load(std::sync::atomic::Ordering::Relaxed)
== PRODUCERS * PER_PRODUCER
{
break;
} else {
std::hint::spin_loop();
}
}
seen
}));
}
for h in handles {
h.join().unwrap();
}
let mut all: Vec<usize> = consumers
.into_iter()
.flat_map(|h| h.join().unwrap())
.collect();
all.sort_unstable();
assert_eq!(all.len(), PRODUCERS * PER_PRODUCER);
for (i, v) in all.iter().enumerate() {
assert_eq!(*v, i);
}
assert_eq!(q.pop(&FLAG).map(|v| v.0), None);
}
}