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