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: Box<[u8; 16]>,
34}
35
36static REGISTRY: EventQueueRegistry<Event> = EventQueueRegistry::new();
37
38thread_local! {
39 static PRODUCER: EventProducer<Event> = REGISTRY.register();
40}
41
42struct Lockstep {
45 batch_pushed: Barrier,
46 sweep_done: Barrier,
47}
48
49impl Lockstep {
50 fn new() -> Self {
51 Self {
52 batch_pushed: Barrier::new(PRODUCERS + 1),
53 sweep_done: Barrier::new(PRODUCERS + 1),
54 }
55 }
56}
57
58#[cfg_attr(feature = "hotpath-meta", hotpath_meta::measure)]
59fn produce(producer: usize, lockstep: &Lockstep) {
60 for round in 0..ROUNDS {
61 for i in 0..BATCH {
62 let event = Event {
63 producer,
64 seq: round * BATCH + i,
65 payload: Box::new([i as u8; 16]),
66 };
67 let _ = PRODUCER.try_with(|p| p.push(event));
68 }
69 lockstep.batch_pushed.wait();
70 lockstep.sweep_done.wait();
71 }
72}
73
74#[cfg_attr(feature = "hotpath-meta", hotpath_meta::measure)]
75fn consume(batch: &[Event]) -> u64 {
76 let mut checksum = 0u64;
77 for event in batch {
78 checksum = checksum
79 .wrapping_add(event.producer as u64)
80 .wrapping_add(event.seq)
81 .wrapping_add(event.payload[0] as u64);
82 }
83 checksum
84}
85
86#[cfg_attr(feature = "hotpath-meta", hotpath_meta::measure)]
90fn consumer_loop(lockstep: Arc<Lockstep>, producers_exited: Arc<Barrier>) -> u64 {
91 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}
111
112#[cfg_attr(feature = "hotpath-meta", hotpath_meta::main(percentiles = [95, 99, 99.9]))]
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 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}