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, so sweep counts and
6//! allocations are fixed by the constants below rather than by scheduling.
7//! Changing the constants breaks comparability with previously uploaded reports.
8//!
9//! Run with:
10//!   cargo run -p hotpath-drain --example throughput --release
11//!   cargo run -p hotpath-drain --example throughput --release --features hotpath-meta
12//!   cargo run -p hotpath-drain --example throughput --release --features hotpath-meta,hotpath-alloc-meta
13
14use 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;
22/// Events each producer pushes per round; stays below the per-sweep chunk cap
23/// so one sweep always drains a whole round.
24const BATCH: u64 = 12_500;
25const EVENTS_PER_PRODUCER: u64 = ROUNDS * BATCH;
26const TOTAL_EVENTS: u64 = PRODUCERS as u64 * EVENTS_PER_PRODUCER;
27
28/// One event; the payload is a `Box` so the consumer's drop and the
29/// producer's allocation are both visible under allocation profiling.
30struct 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
42/// Producers and consumer meet here twice per round: once when every batch
43/// is pushed, once when the sweep is done.
44struct 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/// One sweep per round, then a final drain once the producers have exited
87/// (signalled through `producers_exited`) so no event is lost. Returns the
88/// total consumed.
89#[cfg_attr(feature = "hotpath-meta", hotpath_meta::measure)]
90fn consumer_loop(lockstep: Arc<Lockstep>, producers_exited: Arc<Barrier>) -> u64 {
91    // Sized for one round, the most a sweep can return since the vector is
92    // cleared after each one, so draining never reallocates.
93    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    // Producers have exited, so their thread-local `EventProducer`s are
142    // dropped and the queues are marked closed before the final drain.
143    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}