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 40)
40static 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 118)
117fn main() {
118    REGISTRY.set_active(true);
119    let start = Instant::now();
120    let lockstep = Arc::new(Lockstep::new());
121    let producers_exited = Arc::new(Barrier::new(2));
122
123    let consumer = thread::Builder::new()
124        .name("consumer".into())
125        .spawn({
126            let lockstep = Arc::clone(&lockstep);
127            let producers_exited = Arc::clone(&producers_exited);
128            move || consumer_loop(lockstep, producers_exited)
129        })
130        .expect("spawn consumer");
131
132    let producers: Vec<_> = (0..PRODUCERS)
133        .map(|id| {
134            let lockstep = Arc::clone(&lockstep);
135            thread::Builder::new()
136                .name(format!("producer-{id}"))
137                .spawn(move || produce(id, &lockstep))
138                .expect("spawn producer")
139        })
140        .collect();
141    for handle in producers {
142        handle.join().expect("producer panicked");
143    }
144
145    // Producers have exited, so their thread-local `EventProducer`s are
146    // dropped and the queues are marked closed before the final drain.
147    REGISTRY.set_active(false);
148    producers_exited.wait();
149    let consumed = consumer.join().expect("consumer panicked");
150    let elapsed = start.elapsed();
151
152    assert_eq!(consumed, TOTAL_EVENTS, "events lost in transit");
153    println!(
154        "{PRODUCERS} producers sent {TOTAL_EVENTS} events, consumer received {consumed} in {elapsed:.2?} ({:.1} M events/s)",
155        TOTAL_EVENTS as f64 / elapsed.as_secs_f64() / 1e6
156    );
157}
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 102)
94fn consumer_loop(lockstep: Arc<Lockstep>, producers_exited: Arc<Barrier>) -> u64 {
95    // Sized for one round, the most a sweep can return since the vector is
96    // cleared after each one, so draining never reallocates.
97    let mut batch = Vec::with_capacity((PRODUCERS as u64 * BATCH) as usize);
98    let mut consumed = 0u64;
99    let mut checksum = 0u64;
100    for _ in 0..ROUNDS {
101        lockstep.batch_pushed.wait();
102        REGISTRY.sweep(&mut batch);
103        consumed += batch.len() as u64;
104        checksum = checksum.wrapping_add(consume(&batch));
105        batch.clear();
106        lockstep.sweep_done.wait();
107    }
108    producers_exited.wait();
109    REGISTRY.drain_all(&mut batch);
110    consumed += batch.len() as u64;
111    checksum = checksum.wrapping_add(consume(&batch));
112    std::hint::black_box(checksum);
113    consumed
114}
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 109)
94fn consumer_loop(lockstep: Arc<Lockstep>, producers_exited: Arc<Barrier>) -> u64 {
95    // Sized for one round, the most a sweep can return since the vector is
96    // cleared after each one, so draining never reallocates.
97    let mut batch = Vec::with_capacity((PRODUCERS as u64 * BATCH) as usize);
98    let mut consumed = 0u64;
99    let mut checksum = 0u64;
100    for _ in 0..ROUNDS {
101        lockstep.batch_pushed.wait();
102        REGISTRY.sweep(&mut batch);
103        consumed += batch.len() as u64;
104        checksum = checksum.wrapping_add(consume(&batch));
105        batch.clear();
106        lockstep.sweep_done.wait();
107    }
108    producers_exited.wait();
109    REGISTRY.drain_all(&mut batch);
110    consumed += batch.len() as u64;
111    checksum = checksum.wrapping_add(consume(&batch));
112    std::hint::black_box(checksum);
113    consumed
114}

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.