Skip to main content

LockFreeQueue

Struct LockFreeQueue 

Source
pub struct LockFreeQueue<T> { /* private fields */ }
Expand description

A bounded, genuinely lock-free multi-producer multi-consumer queue.

This is an array-based MPMC queue using per-slot sequence numbers (the Vyukov algorithm). Producers and consumers operate through independent atomic head/tail cursors and never acquire a mutex or spinlock. The sequence-number protocol eliminates the ABA problem without tagged pointers or epoch-based reclamation: slots are reused in place, so no node allocation or deallocation occurs during enqueue/dequeue.

§Capacity

The queue is bounded. LockFreeQueue::new creates a queue with DEFAULT_QUEUE_CAPACITY usable slots. LockFreeQueue::with_capacity accepts any request of one slot or more: the ring is sized to the next power of two at least two, while capacity keeps reporting the request, so a non-power-of-two capacity bounds the queue exactly and only the ring’s unused tail is wasted. When the queue is full, enqueue retries with exponential backoff (preserving the unblocked-sender contract of the previous API), while try_enqueue returns Err(item) for callers that prefer explicit backpressure. Full means capacity items queued: a slot a dequeue has emptied but not yet reopened is waited for, never reported as fullness.

§Memory safety

Each slot’s MaybeUninit<T> is written by the producer — the only writer, between sequence == pos and sequence == pos + 1 — and moved out by the consumer, which is the only reader, between sequence == pos + 1 and sequence == pos + ring_len. The sequence-number protocol therefore guarantees that only one thread ever touches a slot’s payload.

Implementations§

Source§

impl<T> LockFreeQueue<T>

Source

pub fn new() -> Self

Create a new queue with the default capacity.

Source

pub fn with_capacity(capacity: usize) -> Self

Create a new queue holding up to capacity items.

The ring is the next power of two at least two, which the sequence protocol needs to tell a slot’s empty generation from its full one; a one-slot request therefore gets a two-slot ring and still bounds the queue at one item.

Source

pub fn try_enqueue(&self, item: T) -> Result<(), T>

Try to enqueue an item without waiting for space.

Returns Ok(()) if the item was enqueued, or Err(item) if the queue holds capacity items. It takes no lock. When the queue has room but the item’s slot still belongs to a dequeue that has moved its item out and not yet reopened the slot, it waits for that reopening instead of reporting a full queue; the wait is one store unless the dequeuing thread was preempted.

Source

pub fn enqueue(&self, item: T)

Enqueue an item, retrying with exponential backoff if the queue is full.

This preserves the unblocked-sender contract of the previous API: the call always eventually succeeds (assuming consumers make progress). The backoff path uses core::hint::spin_loop and, on std targets, std::thread::yield_now after heavy contention, but never acquires a global lock, so multiple producers can enqueue concurrently.

Source

pub fn try_dequeue(&self) -> Option<T>

Try to dequeue an item from the front of the queue. Returns None if the queue is empty.

This is the lock-free fast path: no spinlock, no mutex.

Source

pub fn is_empty(&self) -> bool

Check if the queue is empty.

This is a best-effort check: the queue may have items added or removed between this call and the next operation. It is safe to call concurrently with enqueue/dequeue.

Source

pub fn is_full(&self) -> bool

Check if the queue holds capacity items.

Best-effort in the same sense as is_empty.

Source

pub const fn capacity(&self) -> usize

Number of items the queue accepts before try_enqueue reports full.

Source

pub fn len(&self) -> usize

Number of items currently queued.

Best-effort in the same sense as is_empty: it is the cursor difference, so a push that has reserved its position but not yet published its slot is already counted, and the value can change under the caller immediately after the read.

Trait Implementations§

Source§

impl<T> Default for LockFreeQueue<T>

Source§

fn default() -> Self

Returns the “default value” for a type. Read more
Source§

impl<T> Drop for LockFreeQueue<T>

Source§

fn drop(&mut self)

Executes the destructor for this type. Read more
Source§

fn pin_drop(self: Pin<&mut Self>)

🔬This is a nightly-only experimental API. (pin_ergonomics)
Execute the destructor for this type, but different to Drop::drop, it requires self to be pinned. Read more
Source§

impl<T: Send> Send for LockFreeQueue<T>

Source§

impl<T: Send> Sync for LockFreeQueue<T>

Auto Trait Implementations§

§

impl<T> !Freeze for LockFreeQueue<T>

§

impl<T> !RefUnwindSafe for LockFreeQueue<T>

§

impl<T> Unpin for LockFreeQueue<T>
where Box<[Slot<T>]>: Unpin,

§

impl<T> UnsafeUnpin for LockFreeQueue<T>
where Box<[Slot<T>]>: UnsafeUnpin,

§

impl<T> UnwindSafe for LockFreeQueue<T>
where Box<[Slot<T>]>: UnwindSafe,

Blanket Implementations§

Source§

impl<T> Any for T
where T: 'static + ?Sized,

Source§

fn type_id(&self) -> TypeId

Gets the TypeId of self. Read more
Source§

impl<T> Borrow<T> for T
where T: ?Sized,

Source§

fn borrow(&self) -> &T

Immutably borrows from an owned value. Read more
Source§

impl<T> BorrowMut<T> for T
where T: ?Sized,

Source§

fn borrow_mut(&mut self) -> &mut T

Mutably borrows from an owned value. Read more
Source§

impl<T> From<T> for T

Source§

fn from(t: T) -> T

Returns the argument unchanged.

Source§

impl<T, U> Into<U> for T
where U: From<T>,

Source§

fn into(self) -> U

Calls U::from(self).

That is, this conversion is whatever the implementation of From<T> for U chooses to do.

Source§

impl<T, U> TryFrom<U> for T
where U: Into<T>,

Source§

type Error = !

The type returned in the event of a conversion error.
Source§

fn try_from(value: U) -> Result<T, !>

Performs the conversion.
Source§

impl<T, U> TryInto<U> for T
where U: TryFrom<T>,

Source§

type Error = <U as TryFrom<T>>::Error

The type returned in the event of a conversion error.
Source§

fn try_into(self) -> Result<U, <U as TryFrom<T>>::Error>

Performs the conversion.