use std::collections::BTreeMap;
use std::hint::black_box;
use std::io::{self, Write};
use std::path::PathBuf;
use subms::{SubMsFeatureManifest, SubMsP99Source, SubMsPerfHarness, classify_feature, summarize};
use subms_spsc_ring_buffer::{Consumer, Producer, SpscRingBuffer};
const SIZES: [usize; 3] = [1_024, 16_384, 262_144];
const CANON: usize = SIZES[SIZES.len() - 1];
const ITEMS_PER_SAMPLE: usize = 1_024;
const SAMPLES: usize = 256;
const ROUNDS: usize = 7;
const BULK_BATCH: usize = 32;
const OPS: usize = 50_000;
const WARM_NANOS: u64 = 300_000_000;
const WARM_MAX_SAMPLES: usize = 1_000;
fn stat(h: &SubMsPerfHarness, median: bool) -> u64 {
summarize(h)
.stages
.iter()
.find(|s| s.name == "op")
.map_or(0, |s| if median { s.p50_ns } else { s.p99_ns })
}
fn batched(reps: usize, mut op: impl FnMut(usize) -> u64) -> u64 {
let mut i = 0usize;
let start = std::time::Instant::now();
for _ in 0..WARM_MAX_SAMPLES {
let mut acc = 0u64;
for _ in 0..reps {
acc = acc.wrapping_add(op(i));
i += 1;
}
black_box(acc);
if start.elapsed().as_nanos() as u64 >= WARM_NANOS {
break;
}
}
let mut h = SubMsPerfHarness::new("spsc-feature", "rust");
let st = h.stage("op", SAMPLES);
for _ in 0..SAMPLES {
st.time(|| {
let mut acc = 0u64;
for _ in 0..reps {
acc = acc.wrapping_add(op(i));
i += 1;
}
black_box(acc);
});
}
stat(&h, true)
}
fn single(mut timed: impl FnMut(usize) -> u64, mut untimed: impl FnMut(usize) -> u64) -> u64 {
for i in 0..OPS {
black_box(timed(i));
black_box(untimed(i));
}
let mut h = SubMsPerfHarness::new("spsc-feature", "rust");
let st = h.stage("op", OPS);
for i in 0..OPS {
st.time(|| black_box(timed(i)));
black_box(untimed(i));
}
stat(&h, false)
}
fn sweep(label: &str, mut at: impl FnMut(usize) -> u64) -> Vec<(usize, u64)> {
let mut runs: Vec<Vec<u64>> = vec![Vec::new(); SIZES.len()];
for _ in 0..ROUNDS {
for (k, &n) in SIZES.iter().enumerate() {
runs[k].push(at(n));
}
}
let rows: Vec<(usize, u64)> = SIZES
.iter()
.enumerate()
.map(|(k, &n)| (n, runs[k].iter().copied().min().unwrap_or(0)))
.collect();
eprintln!("sweep {label}: {rows:?} raw {runs:?}");
rows
}
fn best(mut f: impl FnMut() -> u64) -> u64 {
(0..ROUNDS).map(|_| f()).min().unwrap_or(0)
}
fn pair(cap: usize) -> (Producer<u64>, Consumer<u64>) {
let (mut tx, rx) = SpscRingBuffer::with_capacity::<u64>(cap);
for i in 0..cap / 2 {
let _ = tx.try_push(i as u64);
}
(tx, rx)
}
fn main() -> io::Result<()> {
let path = PathBuf::from(env!("CARGO_MANIFEST_DIR"))
.join("..")
.join(".subms")
.join("features")
.join("rust.json");
let existing = std::fs::read_to_string(&path).unwrap_or_default();
let mut manifest = SubMsFeatureManifest::load_str("rust", &existing);
let (source, instance) = SubMsP99Source::from_env();
manifest.set_p99_source(source, instance.as_deref());
let base_sweep = sweep("base/push+pop", |cap| {
let (mut tx, mut rx) = pair(cap);
batched(ITEMS_PER_SAMPLE, |i| {
let _ = tx.try_push(i as u64);
rx.try_pop().unwrap_or(0)
})
});
let base_p50 = base_sweep
.iter()
.find(|(n, _)| *n == CANON)
.map_or(0, |(_, v)| *v);
eprintln!("base push+pop p50 per {ITEMS_PER_SAMPLE}-item sample: {base_p50}ns");
#[cfg(feature = "bulk")]
{
let sw = sweep("bulk/enqueue+dequeue", |cap| {
let (mut tx, mut rx) = pair(cap);
let batch = [0u64; BULK_BATCH];
let mut out = [0u64; BULK_BATCH];
batched(ITEMS_PER_SAMPLE / BULK_BATCH, |_| {
let n = tx.try_enqueue_bulk(&batch);
(n + rx.try_dequeue_bulk(&mut out)) as u64
})
});
let (cat, reason) = classify_feature(&sw, Some(base_p50), None);
let batch = [0u64; BULK_BATCH];
let mut out = [0u64; BULK_BATCH];
let (mut tx, mut rx) = pair(CANON);
let enq = single(
|_| tx.try_enqueue_bulk(&batch) as u64,
|_| rx.try_dequeue_bulk(&mut out) as u64,
);
let (mut tx, mut rx) = pair(CANON);
let deq = single(
|_| rx.try_dequeue_bulk(&mut out) as u64,
|_| tx.try_enqueue_bulk(&batch) as u64,
);
let mut p99 = BTreeMap::new();
p99.insert("enqueue_bulk".to_string(), enq);
p99.insert("dequeue_bulk".to_string(), deq);
manifest.set_feature("bulk", cat, &p99, &reason);
}
#[cfg(feature = "wait-strategies")]
{
use subms_spsc_ring_buffer::{
BlockingSpscConsumer, BlockingSpscProducer, BusySpin, ParkStrategy,
};
let sw = sweep("wait-strategies/push+pop", |cap| {
let (tx, rx) = pair(cap);
let mut p = BlockingSpscProducer::new(tx, BusySpin);
let mut c = BlockingSpscConsumer::new(rx, BusySpin);
batched(ITEMS_PER_SAMPLE, |i| {
p.push(i as u64);
c.pop()
})
});
let (cat, reason) = classify_feature(&sw, Some(base_p50), None);
let park_batched = best(|| {
let (ps, cs) = ParkStrategy::pair();
let (tx, rx) = pair(CANON);
let mut p = BlockingSpscProducer::new(tx, ps);
let mut c = BlockingSpscConsumer::new(rx, cs);
batched(ITEMS_PER_SAMPLE, |i| {
p.push(i as u64);
c.pop()
})
});
eprintln!(
"wait-strategies park fast path at {CANON}: {park_batched}ns per \
{ITEMS_PER_SAMPLE}-item sample (spin {}ns) - signal() takes a lock per op",
sw.iter().find(|(n, _)| *n == CANON).map_or(0, |(_, v)| *v)
);
let (tx, rx) = pair(CANON);
let mut p = BlockingSpscProducer::new(tx, BusySpin);
let mut c = BlockingSpscConsumer::new(rx, BusySpin);
let push_spin = single(
|i| {
p.push(i as u64);
0
},
|_| c.pop(),
);
let (tx, rx) = pair(CANON);
let mut p = BlockingSpscProducer::new(tx, BusySpin);
let mut c = BlockingSpscConsumer::new(rx, BusySpin);
let pop_spin = single(
|_| c.pop(),
|i| {
p.push(i as u64);
0
},
);
let (ps, cs) = ParkStrategy::pair();
let (tx, rx) = pair(CANON);
let mut p = BlockingSpscProducer::new(tx, ps);
let mut c = BlockingSpscConsumer::new(rx, cs);
let push_park = single(
|i| {
p.push(i as u64);
0
},
|_| c.pop(),
);
let (ps, cs) = ParkStrategy::pair();
let (tx, rx) = pair(CANON);
let mut p = BlockingSpscProducer::new(tx, ps);
let mut c = BlockingSpscConsumer::new(rx, cs);
let pop_park = single(
|_| c.pop(),
|i| {
p.push(i as u64);
0
},
);
let mut p99 = BTreeMap::new();
p99.insert("push_spin".to_string(), push_spin);
p99.insert("pop_spin".to_string(), pop_spin);
p99.insert("push_park".to_string(), push_park);
p99.insert("pop_park".to_string(), pop_park);
manifest.set_feature("wait-strategies", cat, &p99, &reason);
}
#[cfg(feature = "mpsc-fan-in")]
{
use subms_spsc_ring_buffer::{MpscFanIn, MpscFanInConsumer, MpscFanInProducer};
const PRODUCERS: usize = 4;
fn fanin(cap: usize) -> (Vec<MpscFanInProducer<u64>>, MpscFanInConsumer<u64>) {
let (mut ps, c) = MpscFanIn::with_capacity::<u64>(PRODUCERS, cap);
for p in ps.iter_mut() {
for i in 0..cap / 2 {
let _ = p.try_push(i as u64);
}
}
(ps, c)
}
let sw = sweep("mpsc-fan-in/push+pop", |cap| {
let (mut ps, mut c) = fanin(cap);
batched(ITEMS_PER_SAMPLE, |i| {
let _ = ps[i % PRODUCERS].try_push(i as u64);
c.try_pop().unwrap_or(0)
})
});
let (cat, reason) = classify_feature(&sw, Some(base_p50), None);
let (mut ps, mut c) = fanin(CANON);
let enq = single(
|i| {
let _ = ps[i % PRODUCERS].try_push(i as u64);
0
},
|_| c.try_pop().unwrap_or(0),
);
let (mut ps, mut c) = fanin(CANON);
let deq = single(
|_| c.try_pop().unwrap_or(0),
|i| {
let _ = ps[i % PRODUCERS].try_push(i as u64);
0
},
);
let mut p99 = BTreeMap::new();
p99.insert("fanin_enqueue".to_string(), enq);
p99.insert("fanin_dequeue".to_string(), deq);
manifest.set_feature("mpsc-fan-in", cat, &p99, &reason);
}
#[cfg(feature = "mpmc-disruptor")]
{
use subms_spsc_ring_buffer::{DisruptorConsumer, DisruptorProducer, MpmcDisruptor};
const CONSUMERS: usize = 1;
fn disruptor(cap: usize) -> (DisruptorProducer<u64>, DisruptorConsumer<u64>) {
let (p, mut cs) = MpmcDisruptor::with_consumers::<u64>(cap, CONSUMERS);
let c = cs.remove(0);
for i in 0..cap / 2 {
let _ = p.try_publish(i as u64);
}
(p, c)
}
let sw = sweep("mpmc-disruptor/publish+consume", |cap| {
let (p, mut c) = disruptor(cap);
batched(ITEMS_PER_SAMPLE, |i| {
let _ = p.try_publish(i as u64);
c.try_consume().unwrap_or(0)
})
});
let (cat, reason) = classify_feature(&sw, Some(base_p50), None);
let (p, mut c) = disruptor(CANON);
let published = single(
|i| {
let _ = p.try_publish(i as u64);
0
},
|_| c.try_consume().unwrap_or(0),
);
let (p, mut c) = disruptor(CANON);
let consumed = single(
|_| c.try_consume().unwrap_or(0),
|i| {
let _ = p.try_publish(i as u64);
0
},
);
let mut p99 = BTreeMap::new();
p99.insert("publish".to_string(), published);
p99.insert("consume".to_string(), consumed);
manifest.set_feature("mpmc-disruptor", cat, &p99, &reason);
}
#[cfg(feature = "metrics")]
{
use subms_spsc_ring_buffer::InstrumentedSpsc;
let sw = sweep("metrics/push+pop", |cap| {
let (tx, rx) = pair(cap);
let (mut tx, mut rx, _m) = InstrumentedSpsc::wrap(tx, rx);
batched(ITEMS_PER_SAMPLE, |i| {
let _ = tx.try_push(i as u64);
rx.try_pop().unwrap_or(0)
})
});
let (cat, reason) = classify_feature(&sw, Some(base_p50), None);
let (tx, rx) = pair(CANON);
let (mut tx, mut rx, _m) = InstrumentedSpsc::wrap(tx, rx);
let enq = single(
|i| {
let _ = tx.try_push(i as u64);
0
},
|_| rx.try_pop().unwrap_or(0),
);
let (tx, rx) = pair(CANON);
let (mut tx, mut rx, _m) = InstrumentedSpsc::wrap(tx, rx);
let deq = single(
|_| rx.try_pop().unwrap_or(0),
|i| {
let _ = tx.try_push(i as u64);
0
},
);
let mut p99 = BTreeMap::new();
p99.insert("metrics_enqueue".to_string(), enq);
p99.insert("metrics_dequeue".to_string(), deq);
manifest.set_feature("metrics", cat, &p99, &reason);
}
std::fs::create_dir_all(path.parent().unwrap())?;
std::fs::write(&path, manifest.to_json())?;
io::stdout().write_all(manifest.to_json().as_bytes())?;
Ok(())
}