use std::collections::BTreeMap;
use std::hint::black_box;
use std::io::{self, Write};
use std::path::PathBuf;
use subms::{
SubMsFeatureManifest, SubMsP99Source, SubMsPerfHarness, SubMsStage, classify_feature, summarize,
};
use subms_mpsc_queue::MpscQueue;
const SIZES: [usize; 3] = [4_096, 32_768, 262_144];
const CANON: usize = SIZES[SIZES.len() - 1];
const ITEMS_PER_SAMPLE: usize = 1_024;
const SAMPLES: usize = 256;
const SWEEP_REPEATS: usize = 5;
const BATCH: usize = 256;
const OPS: usize = 50_000;
const WARM_NANOS: u64 = 300_000_000;
const WARM_MAX_SAMPLES: usize = 5_000;
const BURN_NANOS: u64 = 1_000_000_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 {
let mut i = 0usize;
let start = std::time::Instant::now();
for _ in 0..WARM_MAX_SAMPLES {
for _ in 0..reps {
op(i);
i += 1;
}
if start.elapsed().as_nanos() as u64 >= WARM_NANOS {
break;
}
}
let mut h = SubMsPerfHarness::new("mpsc-queue-feature", "rust");
let st = h.stage("op", SAMPLES);
for _ in 0..SAMPLES {
st.time(|| {
for _ in 0..reps {
op(i);
i += 1;
}
});
}
stat(&h, true)
}
fn single(mut body: impl FnMut(&mut SubMsStage, usize)) -> u64 {
let mut warm = SubMsPerfHarness::new("mpsc-queue-feature", "rust");
{
let st = warm.stage("op", OPS);
for i in 0..OPS {
body(&mut *st, i);
}
}
let mut h = SubMsPerfHarness::new("mpsc-queue-feature", "rust");
{
let st = h.stage("op", OPS);
for i in 0..OPS {
body(&mut *st, i);
}
}
stat(&h, false)
}
fn sweep(label: &str, mut at: impl FnMut(usize) -> u64) -> Vec<(usize, u64)> {
let mut runs = vec![Vec::with_capacity(SWEEP_REPEATS); SIZES.len()];
for _ in 0..SWEEP_REPEATS {
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:?} repeats {runs:?}");
rows
}
fn filled(n: usize) -> MpscQueue<u64> {
let q = MpscQueue::new();
for i in 0..n {
q.push(i as u64);
}
q
}
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());
#[cfg(feature = "affinity")]
if let Some(core) = std::env::var("SUBMS_PIN").ok().and_then(|v| v.parse().ok()) {
let _ = subms_mpsc_queue::set_affinity(&[core]);
}
{
let mut q = filled(CANON);
let start = std::time::Instant::now();
while (start.elapsed().as_nanos() as u64) < BURN_NANOS {
for i in 0..ITEMS_PER_SAMPLE {
q.push(i as u64);
black_box(q.try_pop());
}
}
}
let base_sweep = sweep("base/push+pop", |n| {
let mut q = filled(n);
batched(ITEMS_PER_SAMPLE, |i| {
q.push(i as u64);
black_box(q.try_pop());
})
});
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 = "bounded")]
{
use subms_mpsc_queue::BoundedMpscQueue;
fn ring(n: usize) -> BoundedMpscQueue<u64> {
let q = BoundedMpscQueue::new(n);
for i in 0..n / 2 {
let _ = q.try_enqueue(i as u64);
}
q
}
let sw = sweep("bounded/enqueue+dequeue", |n| {
let mut q = ring(n);
batched(ITEMS_PER_SAMPLE, |i| {
let _ = q.try_enqueue(i as u64);
black_box(q.try_dequeue());
})
});
let (cat, reason) = classify_feature(&sw, Some(base_p50), None);
let mut q = ring(CANON);
let mut p99 = BTreeMap::new();
p99.insert(
"enqueue".to_string(),
single(|st, i| {
st.time(|| {
let _ = q.try_enqueue(i as u64);
});
black_box(q.try_dequeue());
}),
);
p99.insert(
"dequeue".to_string(),
single(|st, i| {
let _ = q.try_enqueue(i as u64);
st.time(|| black_box(q.try_dequeue()));
}),
);
let full: BoundedMpscQueue<u64> = BoundedMpscQueue::new(CANON);
while full.try_enqueue(0).is_ok() {}
p99.insert(
"enqueue_full".to_string(),
single(|st, i| {
st.time(|| {
let _ = full.try_enqueue(i as u64);
});
}),
);
manifest.set_feature("bounded", cat, &p99, &reason);
}
#[cfg(feature = "mpmc")]
{
use subms_mpsc_queue::MpmcQueue;
fn ring(n: usize) -> MpmcQueue<u64> {
let q = MpmcQueue::new(n);
for i in 0..n / 2 {
let _ = q.try_enqueue(i as u64);
}
q
}
let sw = sweep("mpmc/enqueue+dequeue", |n| {
let q = ring(n);
batched(ITEMS_PER_SAMPLE, |i| {
let _ = q.try_enqueue(i as u64);
black_box(q.try_dequeue());
})
});
let (cat, reason) = classify_feature(&sw, Some(base_p50), None);
let q = ring(CANON);
let mut p99 = BTreeMap::new();
p99.insert(
"enqueue".to_string(),
single(|st, i| {
st.time(|| {
let _ = q.try_enqueue(i as u64);
});
black_box(q.try_dequeue());
}),
);
p99.insert(
"dequeue".to_string(),
single(|st, i| {
let _ = q.try_enqueue(i as u64);
st.time(|| black_box(q.try_dequeue()));
}),
);
manifest.set_feature("mpmc", cat, &p99, &reason);
}
#[cfg(feature = "batch")]
{
use subms_mpsc_queue::BatchMpscQueue;
fn filled_batch(n: usize) -> BatchMpscQueue<u64> {
let q = BatchMpscQueue::new();
for i in 0..n {
q.push(i as u64);
}
q
}
let sw = sweep("batch/push+dequeue_batch", |n| {
let mut q = filled_batch(n);
let mut buf: Vec<Option<u64>> = (0..BATCH).map(|_| None).collect();
batched(ITEMS_PER_SAMPLE / BATCH, |i| {
for j in 0..BATCH {
q.push((i + j) as u64);
}
black_box(q.try_dequeue_batch(&mut buf));
})
});
let (cat, reason) = classify_feature(&sw, Some(base_p50), None);
let mut q = filled_batch(CANON);
let mut buf: Vec<Option<u64>> = (0..BATCH).map(|_| None).collect();
let mut p99 = BTreeMap::new();
p99.insert(
"dequeue_batch".to_string(),
single(|st, i| {
st.time(|| black_box(q.try_dequeue_batch(&mut buf)));
for j in 0..BATCH {
q.push((i + j) as u64);
}
}),
);
p99.insert(
"enqueue".to_string(),
single(|st, i| {
st.time(|| q.push(i as u64));
let _ = q.try_dequeue_batch(&mut buf[..1]);
}),
);
p99.insert(
"enqueue_batch".to_string(),
single(|st, i| {
let base = i as u64;
st.time(|| black_box(q.push_batch(base..base + BATCH as u64)));
let _ = q.try_dequeue_batch(&mut buf);
}),
);
manifest.set_feature("batch", cat, &p99, &reason);
}
#[cfg(feature = "metrics")]
{
use subms_mpsc_queue::MetricsMpscQueue;
fn filled_metrics(n: usize) -> MetricsMpscQueue<u64> {
let q = MetricsMpscQueue::new();
for i in 0..n {
q.push(i as u64);
}
q
}
let sw = sweep("metrics/push+pop", |n| {
let mut q = filled_metrics(n);
batched(ITEMS_PER_SAMPLE, |i| {
q.push(i as u64);
black_box(q.try_pop());
})
});
let (cat, reason) = classify_feature(&sw, Some(base_p50), None);
let mut q = filled_metrics(CANON);
let mut p99 = BTreeMap::new();
p99.insert(
"enqueue".to_string(),
single(|st, i| {
st.time(|| q.push(i as u64));
black_box(q.try_pop());
}),
);
p99.insert(
"dequeue".to_string(),
single(|st, i| {
q.push(i as u64);
st.time(|| black_box(q.try_pop()));
}),
);
p99.insert(
"snapshot".to_string(),
single(|st, _| {
st.time(|| black_box(q.snapshot()));
}),
);
manifest.set_feature("metrics", cat, &p99, &reason);
}
#[cfg(feature = "affinity")]
{
use subms_mpsc_queue::set_affinity;
let sw = sweep("affinity/set_affinity", |_| {
batched(ITEMS_PER_SAMPLE, |_| {
let _ = set_affinity(&[0]);
})
});
let (cat, reason) = classify_feature(
&sw,
Some(base_p50),
Some(subms::SubMsFeatureCategory::Auxiliary),
);
let mut p99 = BTreeMap::new();
p99.insert(
"set_affinity".to_string(),
single(|st, _| {
st.time(|| {
let _ = set_affinity(&[0]);
});
}),
);
manifest.set_feature("affinity", cat, &p99, &reason);
let cores: Vec<usize> = (0..std::thread::available_parallelism()
.map_or(1, std::num::NonZeroUsize::get))
.collect();
let _ = set_affinity(&cores);
}
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(())
}