Skip to main content

subms_spsc_ring_buffer/
recipe.rs

1//! `SubMsRecipe` impl. Behind the `harness` feature.
2
3use std::thread;
4
5use subms::{SubMsBenchParams, SubMsPerfHarness, SubMsRecipe, SubMsStageKind, SubMsTimer};
6
7use crate::SpscRingBuffer;
8
9/// Stages: `enqueue`, `dequeue`. Each measured on the thread that owns the
10/// respective end; sample arrays merged back into the harness on the producer
11/// thread after the consumer joins.
12pub struct SpscRingBufferRecipe;
13
14impl SubMsRecipe for SpscRingBufferRecipe {
15    fn name(&self) -> &str {
16        "spsc-ring-buffer"
17    }
18
19    fn run(&self, h: &mut SubMsPerfHarness, params: &SubMsBenchParams) {
20        let entries = params.entries;
21        let warmup = params.warmup;
22        let (mut tx, mut rx) = SpscRingBuffer::with_capacity::<u64>(1024);
23
24        // Warm-up: a small round-trip so branch predictors / page faults are not
25        // in the timed samples.
26        for i in 0..warmup as u64 {
27            while tx.try_push(i).is_err() {}
28            while rx.try_pop().is_none() {}
29        }
30
31        let consumer = thread::spawn(move || {
32            let mut dq = Vec::with_capacity(entries);
33            let mut next = 0u64;
34            while next < entries as u64 {
35                let t0 = SubMsTimer::tick();
36                if let Some(v) = rx.try_pop() {
37                    let ns = t0.elapsed_ns();
38                    debug_assert_eq!(v, next);
39                    dq.push(ns);
40                    next += 1;
41                }
42            }
43            dq
44        });
45
46        let enqueue_samples: Vec<u64> = {
47            let mut samples = Vec::with_capacity(entries);
48            let mut i = 0u64;
49            while i < entries as u64 {
50                let t0 = SubMsTimer::tick();
51                if tx.try_push(i).is_ok() {
52                    samples.push(t0.elapsed_ns());
53                    i += 1;
54                }
55            }
56            samples
57        };
58        let dequeue_samples = consumer.join().expect("consumer thread");
59
60        let s_enq = h
61            .stage("enqueue", enqueue_samples.len())
62            .with_kind(SubMsStageKind::HotPath);
63        for ns in &enqueue_samples {
64            s_enq.record(*ns);
65        }
66        let s_deq = h
67            .stage("dequeue", dequeue_samples.len())
68            .with_kind(SubMsStageKind::HotPath);
69        for ns in &dequeue_samples {
70            s_deq.record(*ns);
71        }
72        h.add_meta("capacity", "1024");
73    }
74}