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?
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}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?
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}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?
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}