use std::thread;
use subms::{SubMsBenchParams, SubMsPerfHarness, SubMsRecipe, SubMsStageKind, SubMsTimer};
use crate::SpscRingBuffer;
pub struct SpscRingBufferRecipe;
impl SubMsRecipe for SpscRingBufferRecipe {
fn name(&self) -> &str {
"spsc-ring-buffer"
}
fn run(&self, h: &mut SubMsPerfHarness, params: &SubMsBenchParams) {
let entries = params.entries;
let warmup = params.warmup;
let (mut tx, mut rx) = SpscRingBuffer::with_capacity::<u64>(1024);
for i in 0..warmup as u64 {
while tx.try_push(i).is_err() {}
while rx.try_pop().is_none() {}
}
let consumer = thread::spawn(move || {
let mut dq = Vec::with_capacity(entries);
let mut next = 0u64;
while next < entries as u64 {
let t0 = SubMsTimer::tick();
if let Some(v) = rx.try_pop() {
let ns = t0.elapsed_ns();
debug_assert_eq!(v, next);
dq.push(ns);
next += 1;
}
}
dq
});
let enqueue_samples: Vec<u64> = {
let mut samples = Vec::with_capacity(entries);
let mut i = 0u64;
while i < entries as u64 {
let t0 = SubMsTimer::tick();
if tx.try_push(i).is_ok() {
samples.push(t0.elapsed_ns());
i += 1;
}
}
samples
};
let dequeue_samples = consumer.join().expect("consumer thread");
let s_enq = h
.stage("enqueue", enqueue_samples.len())
.with_kind(SubMsStageKind::HotPath);
for ns in &enqueue_samples {
s_enq.record(*ns);
}
let s_deq = h
.stage("dequeue", dequeue_samples.len())
.with_kind(SubMsStageKind::HotPath);
for ns in &dequeue_samples {
s_deq.record(*ns);
}
h.add_meta("capacity", "1024");
}
}