Skip to main content

EventQueueRegistry

Struct EventQueueRegistry 

Source
pub struct EventQueueRegistry<M> { /* private fields */ }
Expand description

All live per-thread queues for one event type, plus the gate that tells producers whether a worker is consuming.

Implementations§

Source§

impl<M: Send> EventQueueRegistry<M>

Source

pub const fn new() -> Self

Examples found in repository?
examples/throughput.rs (line 36)
36static REGISTRY: EventQueueRegistry<Event> = EventQueueRegistry::new();
Source

pub fn is_active(&self) -> bool

Whether a worker is consuming. Producers check this before pushing so events cannot pile up unbounded when no one will ever drain them.

Source

pub fn set_active(&self, active: bool)

Examples found in repository?
examples/throughput.rs (line 114)
113fn main() {
114    REGISTRY.set_active(true);
115    let start = Instant::now();
116    let lockstep = Arc::new(Lockstep::new());
117    let producers_exited = Arc::new(Barrier::new(2));
118
119    let consumer = thread::Builder::new()
120        .name("consumer".into())
121        .spawn({
122            let lockstep = Arc::clone(&lockstep);
123            let producers_exited = Arc::clone(&producers_exited);
124            move || consumer_loop(lockstep, producers_exited)
125        })
126        .expect("spawn consumer");
127
128    let producers: Vec<_> = (0..PRODUCERS)
129        .map(|id| {
130            let lockstep = Arc::clone(&lockstep);
131            thread::Builder::new()
132                .name(format!("producer-{id}"))
133                .spawn(move || produce(id, &lockstep))
134                .expect("spawn producer")
135        })
136        .collect();
137    for handle in producers {
138        handle.join().expect("producer panicked");
139    }
140
141    // Producers have exited, so their thread-local `EventProducer`s are
142    // dropped and the queues are marked closed before the final drain.
143    REGISTRY.set_active(false);
144    producers_exited.wait();
145    let consumed = consumer.join().expect("consumer panicked");
146    let elapsed = start.elapsed();
147
148    assert_eq!(consumed, TOTAL_EVENTS, "events lost in transit");
149    println!(
150        "{PRODUCERS} producers sent {TOTAL_EVENTS} events, consumer received {consumed} in {elapsed:.2?} ({:.1} M events/s)",
151        TOTAL_EVENTS as f64 / elapsed.as_secs_f64() / 1e6
152    );
153}
Source

pub fn register(&self) -> EventProducer<M>

Creates and registers a queue for the calling thread.

Source

pub fn sweep(&self, out: &mut Vec<M>)

Drains every registered queue into out and releases queues whose producer thread has exited. Holding the registry lock for the whole sweep is what makes this the single consumer. Each queue is capped at MAX_CHUNKS_PER_SWEEP; leftovers are picked up on the next tick.

Examples found in repository?
examples/throughput.rs (line 98)
90fn consumer_loop(lockstep: Arc<Lockstep>, producers_exited: Arc<Barrier>) -> u64 {
91    // Sized for one round, the most a sweep can return since the vector is
92    // cleared after each one, so draining never reallocates.
93    let mut batch = Vec::with_capacity((PRODUCERS as u64 * BATCH) as usize);
94    let mut consumed = 0u64;
95    let mut checksum = 0u64;
96    for _ in 0..ROUNDS {
97        lockstep.batch_pushed.wait();
98        REGISTRY.sweep(&mut batch);
99        consumed += batch.len() as u64;
100        checksum = checksum.wrapping_add(consume(&batch));
101        batch.clear();
102        lockstep.sweep_done.wait();
103    }
104    producers_exited.wait();
105    REGISTRY.drain_all(&mut batch);
106    consumed += batch.len() as u64;
107    checksum = checksum.wrapping_add(consume(&batch));
108    std::hint::black_box(checksum);
109    consumed
110}
Source

pub fn drain_all(&self, out: &mut Vec<M>)

Uncapped variant for the worker’s final sweep at shutdown: there is no next tick to pick up leftovers, so every queue is drained to its tail. Terminates because producers are deactivated (set_active(false)) before shutdown is signalled, so queues can no longer grow.

Examples found in repository?
examples/throughput.rs (line 105)
90fn consumer_loop(lockstep: Arc<Lockstep>, producers_exited: Arc<Barrier>) -> u64 {
91    // Sized for one round, the most a sweep can return since the vector is
92    // cleared after each one, so draining never reallocates.
93    let mut batch = Vec::with_capacity((PRODUCERS as u64 * BATCH) as usize);
94    let mut consumed = 0u64;
95    let mut checksum = 0u64;
96    for _ in 0..ROUNDS {
97        lockstep.batch_pushed.wait();
98        REGISTRY.sweep(&mut batch);
99        consumed += batch.len() as u64;
100        checksum = checksum.wrapping_add(consume(&batch));
101        batch.clear();
102        lockstep.sweep_done.wait();
103    }
104    producers_exited.wait();
105    REGISTRY.drain_all(&mut batch);
106    consumed += batch.len() as u64;
107    checksum = checksum.wrapping_add(consume(&batch));
108    std::hint::black_box(checksum);
109    consumed
110}

Trait Implementations§

Source§

impl<M: Send> Default for EventQueueRegistry<M>

Source§

fn default() -> Self

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

Auto Trait Implementations§

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.