Skip to main content

SpscRing

Struct SpscRing 

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

A bounded SPSC ring buffer owned in place, split into borrowing halves.

use moirai_core::channel::{Consumer, Producer, SpscRing};

let mut ring = SpscRing::<u64>::new(64);
let (tx, rx) = ring.split();

std::thread::scope(|scope| {
    scope.spawn(move || {
        for value in 0..1000 {
            if tx.send(value).is_err() {
                break;
            }
        }
    });

    let mut sum = 0;
    for _ in 0..1000 {
        match rx.recv() {
            Ok(value) => sum += value,
            Err(_) => break,
        }
    }
    assert_eq!(sum, (0..1000).sum::<u64>());
});

Implementations§

Source§

impl<T> SpscRing<T>
where T: Send,

Source

pub fn new(capacity: usize) -> SpscRing<T>

Create a ring with at least capacity slots, rounded up to a power of two so the index-to-slot mapping is a mask rather than a division.

Source

pub fn split(&mut self) -> (SpscProducer<'_, T>, SpscConsumer<'_, T>)

Borrow the ring as a producer and a consumer.

Both halves borrow self, so the ring cannot be split again, moved, or dropped until they are.

A ring may be split again once its previous halves are gone. Values left queued by an earlier round stay queued — the counters are not reset — which is what makes a ring reusable across phases without reallocating.

Two details make re-splitting sound, and both are why the caches are seeded from the live counters rather than from zero:

  • closed is cleared. A half sets it on drop so its peer stops blocking; leaving it set would make the next round’s first operation fail immediately.
  • The consumer’s cache must satisfy tail <= cached_head <= head. A zeroed cached_head against a non-zero tail breaks the left side, and has_value would then report a value present in an empty ring and read a slot that was never written. Seeding from head restores it; seeding the producer’s cache from tail is exact for the same reason.

Taking &mut self is what makes this safe to do without atomics: no half exists, so nothing else can be touching either counter.

Source

pub fn len(&self) -> usize

Number of values currently queued.

Source

pub fn is_empty(&self) -> bool

Whether the ring holds no values.

Source

pub fn capacity(&self) -> usize

Total slots, after the power-of-two rounding applied at construction.

Auto Trait Implementations§

§

impl<T> !Freeze for SpscRing<T>

§

impl<T> !RefUnwindSafe for SpscRing<T>

§

impl<T> Send for SpscRing<T>
where SpscChannel<T>: Send,

§

impl<T> Sync for SpscRing<T>
where SpscChannel<T>: Sync,

§

impl<T> Unpin for SpscRing<T>
where SpscChannel<T>: Unpin,

§

impl<T> UnsafeUnpin for SpscRing<T>
where SpscChannel<T>: UnsafeUnpin,

§

impl<T> UnwindSafe for SpscRing<T>
where SpscChannel<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.