1use std::sync::{Arc, Barrier};
19use std::thread;
20use std::time::Instant;
21
22use hotpath_drain::{EventProducer, EventQueueRegistry};
23
24const PRODUCERS: usize = 4;
25const ROUNDS: u64 = 20;
26const BATCH: u64 = 12_500;
29const EVENTS_PER_PRODUCER: u64 = ROUNDS * BATCH;
30const TOTAL_EVENTS: u64 = PRODUCERS as u64 * EVENTS_PER_PRODUCER;
31
32struct Event {
35 producer: usize,
36 seq: u64,
37 payload: Box<[u8; 16]>,
38}
39
40static REGISTRY: EventQueueRegistry<Event> = EventQueueRegistry::new();
41
42thread_local! {
43 static PRODUCER: EventProducer<Event> = REGISTRY.register();
44}
45
46struct Lockstep {
49 batch_pushed: Barrier,
50 sweep_done: Barrier,
51}
52
53impl Lockstep {
54 fn new() -> Self {
55 Self {
56 batch_pushed: Barrier::new(PRODUCERS + 1),
57 sweep_done: Barrier::new(PRODUCERS + 1),
58 }
59 }
60}
61
62#[cfg_attr(feature = "hotpath-meta", hotpath_meta::measure)]
63fn produce(producer: usize, lockstep: &Lockstep) {
64 for round in 0..ROUNDS {
65 for i in 0..BATCH {
66 let event = Event {
67 producer,
68 seq: round * BATCH + i,
69 payload: Box::new([i as u8; 16]),
70 };
71 let _ = PRODUCER.try_with(|p| p.push(event));
72 }
73 lockstep.batch_pushed.wait();
74 lockstep.sweep_done.wait();
75 }
76}
77
78#[cfg_attr(feature = "hotpath-meta", hotpath_meta::measure)]
79fn consume(batch: &[Event]) -> u64 {
80 let mut checksum = 0u64;
81 for event in batch {
82 checksum = checksum
83 .wrapping_add(event.producer as u64)
84 .wrapping_add(event.seq)
85 .wrapping_add(event.payload[0] as u64);
86 }
87 checksum
88}
89
90#[cfg_attr(feature = "hotpath-meta", hotpath_meta::measure)]
94fn consumer_loop(lockstep: Arc<Lockstep>, producers_exited: Arc<Barrier>) -> u64 {
95 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}
115
116#[cfg_attr(feature = "hotpath-meta", hotpath_meta::main(percentiles = [95, 99, 99.9]))]
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 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}