use std::io::{self, Write};
use subms::{SubMsPerfHarness, summarize, summary_to_json};
const ENTRIES: usize = 50_000;
const SEED: u64 = 0;
const BULK_BATCH: usize = 32;
fn main() -> io::Result<()> {
let mut h = SubMsPerfHarness::new("spsc-ring-buffer-features", "rust");
h.input("entries", &ENTRIES.to_string());
h.input("seed", &SEED.to_string());
base(&mut h);
#[cfg(feature = "bulk")]
bulk(&mut h);
#[cfg(feature = "wait-strategies")]
wait_strategies(&mut h);
#[cfg(feature = "mpsc-fan-in")]
mpsc_fan_in(&mut h);
#[cfg(feature = "mpmc-disruptor")]
mpmc_disruptor(&mut h);
#[cfg(feature = "metrics")]
metrics(&mut h);
let summary = summarize(&h);
let mut stdout = io::stdout();
summary_to_json(&summary, &mut stdout)?;
writeln!(stdout)?;
Ok(())
}
fn base(h: &mut SubMsPerfHarness) {
use subms_spsc_ring_buffer::SpscRingBuffer;
let (mut tx, mut rx) = SpscRingBuffer::with_capacity::<u64>(ENTRIES);
let stage = h.stage("base_enqueue", ENTRIES);
for i in 0..ENTRIES as u64 {
stage.time(|| {
let _ = tx.try_push(i);
});
}
let stage = h.stage("base_dequeue", ENTRIES);
for _ in 0..ENTRIES {
stage.time(|| {
let _ = rx.try_pop();
});
}
}
#[cfg(feature = "bulk")]
fn bulk(h: &mut SubMsPerfHarness) {
use subms_spsc_ring_buffer::SpscRingBuffer;
let (mut tx, mut rx) = SpscRingBuffer::with_capacity::<u64>(ENTRIES);
let batch: Vec<u64> = (0..BULK_BATCH as u64).collect();
let calls = ENTRIES / BULK_BATCH;
let stage = h.stage("enqueue_bulk", calls);
for _ in 0..calls {
stage.time(|| {
let _ = tx.try_enqueue_bulk(&batch);
});
}
let mut out = [0u64; BULK_BATCH];
let stage = h.stage("dequeue_bulk", calls);
for _ in 0..calls {
stage.time(|| {
let _ = rx.try_dequeue_bulk(&mut out);
});
}
}
#[cfg(feature = "wait-strategies")]
fn wait_strategies(h: &mut SubMsPerfHarness) {
use subms_spsc_ring_buffer::{
BlockingSpscConsumer, BlockingSpscProducer, BusySpin, SpscRingBuffer,
};
let (tx, rx) = SpscRingBuffer::with_capacity::<u64>(1024);
let mut p = BlockingSpscProducer::new(tx, BusySpin);
let mut c = BlockingSpscConsumer::new(rx, BusySpin);
let stage = h.stage("wait_enqueue", ENTRIES);
for i in 0..ENTRIES as u64 {
stage.time(|| p.push(i));
c.pop();
}
let stage = h.stage("wait_dequeue", ENTRIES);
for i in 0..ENTRIES as u64 {
p.push(i);
stage.time(|| {
let _ = c.pop();
});
}
}
#[cfg(feature = "mpsc-fan-in")]
fn mpsc_fan_in(h: &mut SubMsPerfHarness) {
use subms_spsc_ring_buffer::MpscFanIn;
const PRODUCERS: usize = 4;
let per_ring = ENTRIES / PRODUCERS;
let (mut producers, mut consumer) =
MpscFanIn::with_capacity::<u64>(PRODUCERS, per_ring.next_power_of_two());
let stage = h.stage("fanin_enqueue", per_ring * PRODUCERS);
for i in 0..per_ring {
for p in producers.iter_mut() {
let v = i as u64;
stage.time(|| {
let _ = p.try_push(v);
});
}
}
let stage = h.stage("fanin_dequeue", per_ring * PRODUCERS);
for _ in 0..per_ring * PRODUCERS {
stage.time(|| {
let _ = consumer.try_pop();
});
}
}
#[cfg(feature = "mpmc-disruptor")]
fn mpmc_disruptor(h: &mut SubMsPerfHarness) {
use subms_spsc_ring_buffer::MpmcDisruptor;
let (producer, mut consumers) = MpmcDisruptor::with_consumers::<u64>(1024, 1);
let mut consumer = consumers.remove(0);
h.stage("publish", ENTRIES);
h.stage("consume", ENTRIES);
for i in 0..ENTRIES as u64 {
h.stage_mut("publish").unwrap().time(|| {
let _ = producer.try_publish(i);
});
h.stage_mut("consume").unwrap().time(|| {
let _ = consumer.try_consume();
});
}
}
#[cfg(feature = "metrics")]
fn metrics(h: &mut SubMsPerfHarness) {
use subms_spsc_ring_buffer::{InstrumentedSpsc, SpscRingBuffer};
let (tx, rx) = SpscRingBuffer::with_capacity::<u64>(ENTRIES);
let (mut tx, mut rx, _m) = InstrumentedSpsc::wrap(tx, rx);
let stage = h.stage("metrics_enqueue", ENTRIES);
for i in 0..ENTRIES as u64 {
stage.time(|| {
let _ = tx.try_push(i);
});
}
let stage = h.stage("metrics_dequeue", ENTRIES);
for _ in 0..ENTRIES {
stage.time(|| {
let _ = rx.try_pop();
});
}
}