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

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.