Skip to main content

throughput/
throughput.rs

1//! Sample use of hotpath-drain: several producer threads push one million
2//! events in total into a shared registry while a single consumer thread
3//! sweeps the queues, then drains what is left at shutdown.
4//!
5//! Producers and consumer run in lockstep rounds: each producer pushes
6//! `BATCH` events, every thread meets at a barrier, the consumer sweeps once,
7//! and a second barrier releases the producers for the next round. Sweep
8//! count, events per sweep and the batch vector's capacity are therefore
9//! fixed by the constants below rather than by thread scheduling, so two runs
10//! of this benchmark differ only in timing. Changing the constants invalidates
11//! comparability with previously uploaded reports.
12//!
13//! Run with:
14//!   cargo run -p hotpath-drain --example throughput --release
15//!   cargo run -p hotpath-drain --example throughput --release --features hotpath-meta
16//!   cargo run -p hotpath-drain --example throughput --release --features hotpath-meta,hotpath-alloc-meta
17
18use 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;
26/// Events each producer pushes per round; stays below the per-sweep chunk cap
27/// so one sweep always drains a whole round.
28const BATCH: u64 = 12_500;
29const EVENTS_PER_PRODUCER: u64 = ROUNDS * BATCH;
30const TOTAL_EVENTS: u64 = PRODUCERS as u64 * EVENTS_PER_PRODUCER;
31
32/// One event; the payload is a `Box` so the consumer's drop and the
33/// producer's allocation are both visible under allocation profiling.
34struct 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
46/// Producers and consumer meet here twice per round: once when every batch
47/// is pushed, once when the sweep is done.
48struct 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/// One sweep per round, then a final drain once the producers have exited
91/// (signalled through `producers_exited`) so no event is lost. Returns the
92/// total consumed.
93#[cfg_attr(feature = "hotpath-meta", hotpath_meta::measure)]
94fn consumer_loop(lockstep: Arc<Lockstep>, producers_exited: Arc<Barrier>) -> u64 {
95    // Sized for one round, the most a sweep can return since the vector is
96    // cleared after each one, so draining never reallocates.
97    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    // Producers have exited, so their thread-local `EventProducer`s are
146    // dropped and the queues are marked closed before the final drain.
147    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}