use std::collections::BTreeMap;
use std::io::{self, Write};
use std::path::PathBuf;
use std::sync::atomic::{AtomicU64, Ordering};
use std::time::Instant;
use subms::{SubMsFeatureManifest, SubMsP99Source, SubMsPerfHarness, classify_feature, summarize};
#[cfg(feature = "keyed")]
use subms_rate_limiter::Acquire;
use subms_rate_limiter::{Clock, RateLimiter};
const SIZES: [usize; 3] = [1024, 8192, 65536];
const CANON_N: usize = SIZES[SIZES.len() - 1];
const OPS: usize = CANON_N;
const DIST_SIZES: [usize; 3] = [512, 2048, 8192];
const DIST_CANON: usize = DIST_SIZES[DIST_SIZES.len() - 1];
const DIST_OPS: usize = 4_096;
const DIST_HOT: usize = 256;
const DIST_LIMIT: u64 = 9;
const DIST_WINDOW_NS: u64 = 3_600_000_000_000;
const STRIDE_NS: u64 = 1_000;
const TB_CAP: u64 = 4;
const TB_RATE: f64 = 500_000.0;
const MET_RATE: f64 = 166_667.0;
const PRE_DRAIN: usize = 16;
const HIER_PARENT_CAP: u64 = 64;
const HIER_PARENT_RATE: f64 = 2_000_000.0;
const BASE_RATE: f64 = 1_000_000.0;
const BASE_BURST: u64 = 4;
const WARM_NANOS: u64 = 300_000_000;
const WARM_MAX_REPS: usize = 200_000;
const BATCH: usize = 16;
struct SteppingClock {
origin: Instant,
steps: AtomicU64,
sink: AtomicU64,
}
impl SteppingClock {
fn new() -> Self {
Self {
origin: Instant::now(),
steps: AtomicU64::new(0),
sink: AtomicU64::new(0),
}
}
}
impl Clock for SteppingClock {
fn now_ns(&self) -> u64 {
self.sink
.store(self.origin.elapsed().as_nanos() as u64, Ordering::Relaxed);
self.steps.fetch_add(STRIDE_NS, Ordering::Relaxed) + STRIDE_NS
}
}
struct Measured {
p50: u64,
p99: u64,
accept: f64,
}
fn stat(h: &SubMsPerfHarness, name: &str, median: bool) -> u64 {
summarize(h)
.stages
.iter()
.find(|s| s.name == name)
.map_or(0, |s| if median { s.p50_ns } else { s.p99_ns })
}
fn measure<T>(
mut setup: impl FnMut() -> T,
mut op: impl FnMut(&T, usize) -> bool,
ops: usize,
batch: usize,
) -> Measured {
let warm = setup();
let start = Instant::now();
for rep in 0..WARM_MAX_REPS {
if start.elapsed().as_nanos() as u64 >= WARM_NANOS {
break;
}
op(&warm, rep % ops);
}
drop(warm);
let target = setup();
let mut h = SubMsPerfHarness::new("rate-limiter-feature", "rust");
let mut granted = 0usize;
{
let st = h.stage("op", ops);
for i in 0..ops {
if st.time(|| op(&target, i)) {
granted += 1;
}
}
}
let p99 = stat(&h, "op", false);
let p50 = if batch > 1 {
let samples = ops / batch;
let st = h.stage("batched", samples);
for s in 0..samples {
st.time(|| {
for k in 0..batch {
op(&target, s * batch + k);
}
});
}
stat(&h, "batched", true) / batch as u64
} else {
stat(&h, "op", true)
};
Measured {
p50,
p99,
accept: granted as f64 / ops as f64,
}
}
fn sweep(
label: &str,
sizes: &[usize],
mut at: impl FnMut(usize) -> Measured,
) -> (Vec<(usize, u64)>, Measured) {
let mut rows = Vec::with_capacity(sizes.len());
let mut line = format!("sweep {label}:");
let mut last = None;
for &n in sizes {
let m = at(n);
line.push_str(&format!(
" (n={n} p50={}ns p99={}ns accept={:.0}%)",
m.p50,
m.p99,
m.accept * 100.0
));
rows.push((n, m.p50));
last = Some(m);
}
eprintln!("{line}");
(rows, last.expect("non-empty sizes"))
}
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 = measure(
|| {
(0..CANON_N)
.map(|_| RateLimiter::new(BASE_RATE, BASE_BURST))
.collect::<Vec<_>>()
},
|v, i| v[i % v.len()].try_acquire(),
OPS,
BATCH,
);
let base_p50 = base.p50;
eprintln!(
"base try_acquire over {CANON_N} limiters: p50={base_p50}ns p99={}ns accept={:.0}%",
base.p99,
base.accept * 100.0
);
let hot = measure(
|| RateLimiter::new(BASE_RATE, (2 * OPS) as u64),
|r, _| r.try_acquire(),
OPS,
BATCH,
);
eprintln!(
"base try_acquire on 1 hot limiter: p50={}ns p99={}ns (context only)",
hot.p50, hot.p99
);
#[cfg(feature = "token-bucket")]
{
use subms_rate_limiter::TokenBucket;
fn fleet(n: usize) -> Vec<TokenBucket> {
let v: Vec<TokenBucket> = (0..n)
.map(|_| TokenBucket::with_clock(TB_CAP, TB_RATE, Box::new(SteppingClock::new())))
.collect();
for (j, b) in v.iter().enumerate() {
for _ in 0..PRE_DRAIN + (j & 1) {
b.try_acquire(1);
}
}
v
}
let (rows, canon) = sweep("token-bucket/try_acquire", &SIZES, |n| {
measure(
|| fleet(n),
|v, i| v[i % v.len()].try_acquire(1),
OPS,
BATCH,
)
});
let (cat, reason) = classify_feature(&rows, Some(base_p50), None);
let avail = measure(
|| fleet(CANON_N),
|v, i| {
let _ = v[i % v.len()].available();
true
},
OPS,
BATCH,
);
let mut p99 = BTreeMap::new();
p99.insert("try_acquire".to_string(), canon.p99);
p99.insert("available".to_string(), avail.p99);
manifest.set_feature("token-bucket", cat, &p99, &reason);
}
#[cfg(feature = "hierarchical")]
{
use subms_rate_limiter::HierarchicalLimiter;
fn hier(n: usize) -> HierarchicalLimiter {
let h = HierarchicalLimiter::with_clock_fn(
HIER_PARENT_CAP,
HIER_PARENT_RATE,
n,
TB_CAP,
TB_RATE,
|| Box::new(SteppingClock::new()),
);
for c in 0..h.num_children() {
for _ in 0..PRE_DRAIN + (c & 1) {
h.try_acquire(c, 1);
}
}
h
}
let (rows, canon) = sweep("hierarchical/try_acquire", &SIZES, |n| {
measure(
|| hier(n),
|h, i| h.try_acquire(i % h.num_children(), 1),
OPS,
BATCH,
)
});
let (cat, reason) = classify_feature(&rows, Some(base_p50), None);
let mut p99 = BTreeMap::new();
p99.insert("try_acquire".to_string(), canon.p99);
manifest.set_feature("hierarchical", cat, &p99, &reason);
}
#[cfg(feature = "distributed-backend")]
{
use subms_rate_limiter::{DistributedLimiter, InMemoryBackend};
let keys: Vec<String> = (0..DIST_CANON).map(|i| format!("key-{i:06}")).collect();
let prefilled = |n: usize| {
let d = DistributedLimiter::with_clock(
Box::new(InMemoryBackend::new()),
DIST_LIMIT,
DIST_WINDOW_NS,
Box::new(SteppingClock::new()),
);
for k in &keys[..n] {
d.try_acquire(k);
}
d
};
let (rows, canon) = sweep("distributed-backend/try_acquire", &DIST_SIZES, |n| {
measure(
|| prefilled(n),
|d, i| d.try_acquire(&keys[i % DIST_HOT]),
DIST_OPS,
1,
)
});
let (cat, reason) = classify_feature(&rows, Some(base_p50), None);
let mut p99 = BTreeMap::new();
p99.insert("try_acquire".to_string(), canon.p99);
manifest.set_feature("distributed-backend", cat, &p99, &reason);
}
#[cfg(feature = "metrics")]
{
use subms_rate_limiter::MeteredTokenBucket;
fn fleet(n: usize) -> Vec<MeteredTokenBucket> {
let v: Vec<MeteredTokenBucket> = (0..n)
.map(|_| {
MeteredTokenBucket::with_clock(TB_CAP, MET_RATE, Box::new(SteppingClock::new()))
})
.collect();
for (j, b) in v.iter().enumerate() {
for _ in 0..PRE_DRAIN + (j & 1) {
b.try_acquire(1);
}
}
v
}
let (rows, canon) = sweep("metrics/try_acquire", &SIZES, |n| {
measure(
|| fleet(n),
|v, i| v[i % v.len()].try_acquire(1),
OPS,
BATCH,
)
});
let (cat, reason) = classify_feature(&rows, Some(base_p50), None);
let snap = measure(
|| fleet(CANON_N),
|v, i| {
let _ = v[i % v.len()].snapshot();
true
},
OPS,
BATCH,
);
let mut p99 = BTreeMap::new();
p99.insert("try_acquire".to_string(), canon.p99);
p99.insert("snapshot".to_string(), snap.p99);
manifest.set_feature("metrics", cat, &p99, &reason);
}
#[cfg(feature = "keyed")]
{
use subms_rate_limiter::KeyedRateLimiter;
let keys: Vec<String> = (0..CANON_N).map(|i| format!("key-{i:06}")).collect();
let burst = (OPS / SIZES[0] + 2) as u64;
let (rows, canon) = sweep("keyed/try_acquire", &SIZES, |n| {
measure(
|| KeyedRateLimiter::new(BASE_RATE, burst),
|k, i| matches!(k.try_acquire_at(0, &keys[i % n], 1), Acquire::Ok),
OPS,
BATCH,
)
});
let (cat, reason) = classify_feature(&rows, Some(base_p50), None);
let mut p99 = BTreeMap::new();
p99.insert("try_acquire".to_string(), canon.p99);
manifest.set_feature("keyed", 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(())
}