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