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, stored inline: it owns no heap memory, so the only allocations
29/// left under profiling are the queue's own chunk storage.
30struct Event {
31    producer: usize,
32    seq: u64,
33    payload: [u8; 8],
34}
35
36// The policy's `push` budget is sized for 24 B slots.
37const _: () = assert!(size_of::<Event>() == 24);
38
39static REGISTRY: EventQueueRegistry<Event> = EventQueueRegistry::new();
40
41thread_local! {
42    static PRODUCER: EventProducer<Event> = REGISTRY.register();
43}
44
45/// Producers and consumer meet here twice per round: once when every batch
46/// is pushed, once when the sweep is done.
47struct 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/// One sweep per round, then a final drain once the producers have exited
90/// (signalled through `producers_exited`) so no event is lost. Returns the
91/// total consumed.
92#[cfg_attr(feature = "hotpath-meta", hotpath_meta::measure)]
93fn consumer_loop(lockstep: Arc<Lockstep>, producers_exited: Arc<Barrier>) -> u64 {
94    // Sized for one round, the most a sweep can return since the vector is
95    // cleared after each one, so draining never reallocates.
96    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    // Producers have exited, so their thread-local `EventProducer`s are
145    // dropped and the queues are marked closed before the final drain.
146    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}