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>
impl<M: Send> EventQueueRegistry<M>
Sourcepub fn is_active(&self) -> bool
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.
Sourcepub fn set_active(&self, active: bool)
pub fn set_active(&self, active: bool)
Examples found in repository?
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}Sourcepub fn register(&self) -> EventProducer<M>
pub fn register(&self) -> EventProducer<M>
Creates and registers a queue for the calling thread.
Sourcepub fn sweep(&self, out: &mut Vec<M>)
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?
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}Sourcepub fn drain_all(&self, out: &mut Vec<M>)
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?
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}