use std::collections::BTreeMap;
use std::fmt::Write as _;
use std::io::{self, Write};
use crate::{
SubMsBenchDiff, SubMsBenchParams, SubMsBenchSummary, SubMsBenchSweep, SubMsMetricDiff,
SubMsPerfHarness, SubMsRecipe, SubMsStageDiff, SubMsStageSummary, stats,
};
pub fn summarize(h: &SubMsPerfHarness) -> SubMsBenchSummary {
let summary = summarize_internal(
h,
true,
0,
h.sample_cap(),
);
if let Some(obs) = h.observer() {
obs.on_summarize(&summary);
}
summary
}
pub fn summarize_lean(h: &SubMsPerfHarness) -> SubMsBenchSummary {
summarize_internal(
h,
false,
0,
h.sample_cap(),
)
}
pub fn summarize_skipping(h: &SubMsPerfHarness, skip_warmup: usize) -> SubMsBenchSummary {
summarize_internal(
h,
true,
skip_warmup,
h.sample_cap(),
)
}
pub fn summarize_windowed(h: &SubMsPerfHarness, window: usize) -> Vec<SubMsBenchSummary> {
let window = window.max(1);
let max_len = h
.stages()
.iter()
.map(|s| s.samples().len())
.max()
.unwrap_or(0);
if max_len == 0 {
return Vec::new();
}
let n_windows = max_len.div_ceil(window);
let mut out = Vec::with_capacity(n_windows);
for w in 0..n_windows {
let start = w * window;
let stages = h
.stages()
.iter()
.map(|s| {
let samples = s.samples();
let end = (start + window).min(samples.len());
let slice = if start < samples.len() {
&samples[start..end]
} else {
&[][..]
};
summarize_stage(
s.name(),
slice,
false,
500,
)
})
.collect();
out.push(SubMsBenchSummary {
workload: h.workload().to_string(),
lang: h.lang().to_string(),
timestamp: h.timestamp(),
cpu_core: None,
cpu_affinity: None,
inputs: {
let mut m = clone_map(h.inputs());
m.insert("__window_index".to_string(), w.to_string());
m.insert("__window_size".to_string(), window.to_string());
m
},
meta: clone_map(h.meta()),
stages,
});
}
out
}
fn summarize_internal(
h: &SubMsPerfHarness,
include_samples: bool,
skip_warmup: usize,
sample_cap: usize,
) -> SubMsBenchSummary {
let stages = h
.stages()
.iter()
.map(|s| {
let trimmed = if skip_warmup > 0 && s.samples().len() > skip_warmup {
&s.samples()[skip_warmup..]
} else {
s.samples()
};
summarize_stage(s.name(), trimmed, include_samples, sample_cap)
})
.collect();
let (cpu_core, cpu_affinity) = cpu_placement();
SubMsBenchSummary {
workload: h.workload().to_string(),
lang: h.lang().to_string(),
timestamp: h.timestamp(),
cpu_core,
cpu_affinity,
inputs: clone_map(h.inputs()),
meta: clone_map(h.meta()),
stages,
}
}
fn cpu_placement() -> (Option<u32>, Option<String>) {
let core = std::fs::read_to_string("/proc/self/stat")
.ok()
.and_then(|s| {
let start = s.rfind(')').map(|i| i + 1)?;
s[start..]
.split_whitespace()
.nth(36)
.and_then(|v| v.parse::<u32>().ok())
});
let affinity = std::fs::read_to_string("/proc/self/status")
.ok()
.and_then(|s| {
s.lines()
.find_map(|l| l.strip_prefix("Cpus_allowed_list:"))
.map(|v| v.trim().to_string())
});
(core, affinity)
}
fn summarize_stage(
name: &str,
chronological: &[u64],
include_samples: bool,
sample_cap: usize,
) -> SubMsStageSummary {
let mut sorted = chronological.to_vec();
sorted.sort_unstable();
let samples_ns = if include_samples {
Some(downsample(chronological, sample_cap))
} else {
None
};
SubMsStageSummary {
name: name.to_string(),
count: sorted.len(),
p50_ns: stats::percentile(&sorted, 0.50),
p99_ns: stats::percentile(&sorted, 0.99),
p999_ns: stats::percentile(&sorted, 0.999),
max_ns: sorted.last().copied().unwrap_or(0),
mean_ns: stats::mean(chronological),
stddev_ns: stats::stddev(chronological),
cdf_buckets_ns: stats::cdf_buckets(chronological),
jitter_score: stats::jitter_score(chronological),
samples_ns,
}
}
pub(crate) fn downsample(chronological: &[u64], cap: usize) -> Vec<u64> {
let n = chronological.len();
if n == 0 {
return Vec::new();
}
let step = (n / cap.max(1)).max(1);
chronological.iter().copied().step_by(step).collect()
}
fn clone_map(src: &BTreeMap<String, String>) -> BTreeMap<String, String> {
src.iter().map(|(k, v)| (k.clone(), v.clone())).collect()
}
pub fn print_summary<W: Write>(s: &SubMsBenchSummary, out: &mut W) -> io::Result<()> {
writeln!(
out,
" {:<9} {:>9} {:>9} {:>9} {:>9} {:>9}",
"stage", "p50", "p99", "p99.9", "max", "mean"
)?;
for stage in &s.stages {
writeln!(
out,
" {:<9} {:>9} {:>9} {:>9} {:>9} {:>9}",
stage.name,
format_ns(stage.p50_ns),
format_ns(stage.p99_ns),
format_ns(stage.p999_ns),
format_ns(stage.max_ns),
format_ns(stage.mean_ns),
)?;
}
Ok(())
}
pub fn format_ns(ns: u64) -> String {
if ns < 1_000 {
format!("{}ns", ns)
} else if ns < 1_000_000 {
format!("{:.1}us", ns as f64 / 1_000.0)
} else {
format!("{:.2}ms", ns as f64 / 1_000_000.0)
}
}
#[derive(Debug, Clone, Copy)]
pub struct SubMsBenchAssertion {
pub stage: &'static str,
pub p99_ns_max: u64,
}
pub trait SubMsAssertionTarget {
fn lookup_p99_ns(&self, stage: &str) -> Option<u64>;
}
impl SubMsAssertionTarget for SubMsBenchSummary {
fn lookup_p99_ns(&self, stage: &str) -> Option<u64> {
self.stage(stage).map(|s| s.p99_ns)
}
}
impl SubMsAssertionTarget for SubMsPerfHarness {
fn lookup_p99_ns(&self, stage: &str) -> Option<u64> {
let st = self.stage_by_name(stage)?;
let mut sorted = st.samples().to_vec();
sorted.sort_unstable();
Some(stats::percentile(&sorted, 0.99))
}
}
pub fn assert_p99_under<T: SubMsAssertionTarget + ?Sized>(
target: &T,
assertions: &[SubMsBenchAssertion],
) -> Result<(), String> {
for a in assertions {
let p99 = target
.lookup_p99_ns(a.stage)
.ok_or_else(|| format!("stage '{}' not found", a.stage))?;
if p99 > a.p99_ns_max {
return Err(format!(
"stage '{}' p99 = {} ns exceeded limit {} ns",
a.stage, p99, a.p99_ns_max
));
}
}
Ok(())
}
pub fn run_bench<R: SubMsRecipe + ?Sized>(
recipe: &R,
params: &SubMsBenchParams,
) -> SubMsPerfHarness {
crate::recipe::benchmark(recipe, params)
}
pub fn contended_warmup<F>(threads: usize, iterations_per_thread: usize, work: F)
where
F: Fn(usize, usize) + Send + Sync + 'static + Copy,
{
let mut handles = Vec::with_capacity(threads);
for tid in 0..threads {
handles.push(std::thread::spawn(move || {
for i in 0..iterations_per_thread {
work(tid, i);
}
}));
}
for h in handles {
h.join().expect("contended_warmup thread");
}
}
pub fn summary_to_json<W: Write>(s: &SubMsBenchSummary, out: &mut W) -> io::Result<()> {
let mut buf = String::with_capacity(64 * 1024);
append_summary_json(&mut buf, s);
out.write_all(buf.as_bytes())?;
out.write_all(b"\n")?;
Ok(())
}
pub(crate) fn append_summary_json(out: &mut String, s: &SubMsBenchSummary) {
out.push('{');
json_kv_str(out, "workload", &s.workload);
out.push(',');
json_kv_str(out, "lang", &s.lang);
out.push(',');
json_kv_str(out, "timestamp", &s.timestamp);
out.push(',');
out.push_str("\"inputs\":");
json_map(out, &s.inputs);
out.push(',');
out.push_str("\"meta\":");
json_map(out, &s.meta);
out.push(',');
out.push_str("\"cpu\":");
cpu_json(out, s.cpu_core, s.cpu_affinity.as_deref());
out.push(',');
out.push_str("\"stages\":{");
for (i, stage) in s.stages.iter().enumerate() {
if i > 0 {
out.push(',');
}
json_str(out, &stage.name);
out.push(':');
stage_json(out, stage);
}
out.push_str("}}");
}
fn cpu_json(out: &mut String, core: Option<u32>, affinity: Option<&str>) {
if core.is_none() && affinity.is_none() {
out.push_str("null");
return;
}
out.push('{');
match core {
Some(c) => {
let _ = write!(out, "\"core\":{c}");
}
None => out.push_str("\"core\":null"),
}
out.push_str(",\"affinity\":");
match affinity {
Some(a) => json_str(out, a),
None => out.push_str("null"),
}
out.push('}');
}
fn stage_json(out: &mut String, stage: &SubMsStageSummary) {
out.push('{');
let _ = write!(out, "\"count\":{},", stage.count);
let _ = write!(out, "\"p50_ns\":{},", stage.p50_ns);
let _ = write!(out, "\"p99_ns\":{},", stage.p99_ns);
let _ = write!(out, "\"p999_ns\":{},", stage.p999_ns);
let _ = write!(out, "\"max_ns\":{},", stage.max_ns);
let _ = write!(out, "\"mean_ns\":{},", stage.mean_ns);
let _ = write!(out, "\"stddev_ns\":{},", stage.stddev_ns);
let _ = write!(out, "\"jitter_score\":{:.4},", stage.jitter_score);
out.push_str("\"cdf_buckets_ns\":[");
for (i, c) in stage.cdf_buckets_ns.iter().enumerate() {
if i > 0 {
out.push(',');
}
let _ = write!(out, "{}", c);
}
out.push_str("],");
out.push_str("\"samples_ns\":[");
if let Some(samples) = &stage.samples_ns {
for (i, x) in samples.iter().enumerate() {
if i > 0 {
out.push(',');
}
let _ = write!(out, "{}", x);
}
}
out.push_str("]}");
}
fn json_str(out: &mut String, s: &str) {
out.push('"');
for c in s.chars() {
match c {
'"' => out.push_str("\\\""),
'\\' => out.push_str("\\\\"),
'\n' => out.push_str("\\n"),
'\r' => out.push_str("\\r"),
'\t' => out.push_str("\\t"),
c if (c as u32) < 0x20 => {
let _ = write!(out, "\\u{:04x}", c as u32);
}
c => out.push(c),
}
}
out.push('"');
}
fn json_kv_str(out: &mut String, k: &str, v: &str) {
json_str(out, k);
out.push(':');
json_str(out, v);
}
fn json_map(out: &mut String, m: &BTreeMap<String, String>) {
out.push('{');
for (i, (k, v)) in m.iter().enumerate() {
if i > 0 {
out.push(',');
}
json_kv_str(out, k, v);
}
out.push('}');
}
pub fn run_sweep<R: SubMsRecipe + ?Sized>(
recipe: &R,
params_list: &[SubMsBenchParams],
varied_input_key: Option<&str>,
) -> SubMsBenchSweep {
let runs = params_list
.iter()
.map(|p| summarize(&run_bench(recipe, p)))
.collect();
SubMsBenchSweep {
workload: recipe.name().to_string(),
lang: "rust".to_string(),
varied_input_key: varied_input_key.map(|s| s.to_string()),
runs,
}
}
pub fn summarize_sweep(
summaries: Vec<SubMsBenchSummary>,
varied_input_key: Option<&str>,
) -> SubMsBenchSweep {
assert!(
!summaries.is_empty(),
"summarize_sweep requires at least one run"
);
SubMsBenchSweep {
workload: summaries[0].workload.clone(),
lang: summaries[0].lang.clone(),
varied_input_key: varied_input_key.map(|s| s.to_string()),
runs: summaries,
}
}
pub fn print_sweep<W: Write>(sweep: &SubMsBenchSweep, out: &mut W) -> io::Result<()> {
if sweep.runs.is_empty() {
writeln!(out, "(empty sweep)")?;
return Ok(());
}
let first = &sweep.runs[0];
let header_label = sweep.varied_input_key.as_deref().unwrap_or("run");
for stage in &first.stages {
writeln!(out, "stage: {}", stage.name)?;
writeln!(
out,
" {:<15} {:>9} {:>9} {:>9} {:>9} {:>9} {:>9}",
header_label, "count", "p50", "p99", "p99.9", "max", "mean"
)?;
for (i, run) in sweep.runs.iter().enumerate() {
let label = match &sweep.varied_input_key {
Some(k) => run
.inputs
.get(k)
.cloned()
.unwrap_or_else(|| "?".to_string()),
None => format!("run {}", i + 1),
};
match run.stage(&stage.name) {
None => writeln!(out, " {:<15} (stage missing)", label)?,
Some(s) => writeln!(
out,
" {:<15} {:>9} {:>9} {:>9} {:>9} {:>9} {:>9}",
label,
s.count,
format_ns(s.p50_ns),
format_ns(s.p99_ns),
format_ns(s.p999_ns),
format_ns(s.max_ns),
format_ns(s.mean_ns)
)?,
}
}
writeln!(out)?;
}
Ok(())
}
pub fn sweep_to_json<W: Write>(sweep: &SubMsBenchSweep, out: &mut W) -> io::Result<()> {
let mut buf = String::with_capacity(64 * 1024);
buf.push('[');
for (i, run) in sweep.runs.iter().enumerate() {
if i > 0 {
buf.push(',');
}
append_summary_json(&mut buf, run);
}
buf.push(']');
out.write_all(buf.as_bytes())?;
out.write_all(b"\n")?;
Ok(())
}
pub const DEFAULT_REGRESSION_THRESHOLD_PCT: f64 = 10.0;
pub fn diff_summary(baseline: &SubMsBenchSummary, candidate: &SubMsBenchSummary) -> SubMsBenchDiff {
diff_summary_with(baseline, candidate, DEFAULT_REGRESSION_THRESHOLD_PCT)
}
pub fn diff_summary_with(
baseline: &SubMsBenchSummary,
candidate: &SubMsBenchSummary,
regression_threshold_pct: f64,
) -> SubMsBenchDiff {
let baseline_names: Vec<&str> = baseline.stages.iter().map(|s| s.name.as_str()).collect();
let candidate_names: std::collections::BTreeSet<&str> =
candidate.stages.iter().map(|s| s.name.as_str()).collect();
let mut stage_diffs = Vec::new();
for cand in &candidate.stages {
if let Some(base) = baseline.stage(&cand.name) {
stage_diffs.push(diff_stage(base, cand));
}
}
let candidate_name_set: std::collections::BTreeSet<&str> = candidate_names.clone();
let baseline_name_set: std::collections::BTreeSet<&str> =
baseline_names.iter().copied().collect();
let baseline_only: Vec<String> = baseline_names
.iter()
.filter(|n| !candidate_name_set.contains(*n))
.map(|s| s.to_string())
.collect();
let candidate_only: Vec<String> = candidate
.stages
.iter()
.map(|s| s.name.clone())
.filter(|n| !baseline_name_set.contains(n.as_str()))
.collect();
SubMsBenchDiff {
baseline_workload: baseline.workload.clone(),
candidate_workload: candidate.workload.clone(),
lang: candidate.lang.clone(),
stages: stage_diffs,
baseline_only_stages: baseline_only,
candidate_only_stages: candidate_only,
regression_threshold_pct,
}
}
fn diff_stage(baseline: &SubMsStageSummary, candidate: &SubMsStageSummary) -> SubMsStageDiff {
let metrics = vec![
metric_diff("p50", baseline.p50_ns, candidate.p50_ns),
metric_diff("p99", baseline.p99_ns, candidate.p99_ns),
metric_diff("p99.9", baseline.p999_ns, candidate.p999_ns),
metric_diff("max", baseline.max_ns, candidate.max_ns),
metric_diff("mean", baseline.mean_ns, candidate.mean_ns),
];
let worst = metrics
.iter()
.filter(|m| m.delta_pct.is_finite())
.map(|m| m.delta_pct)
.fold(0.0_f64, f64::max);
SubMsStageDiff {
stage: baseline.name.clone(),
metrics,
worst_regression_pct: worst,
}
}
fn metric_diff(name: &str, baseline: u64, candidate: u64) -> SubMsMetricDiff {
let delta_ns = candidate as i64 - baseline as i64;
let delta_pct = if baseline == 0 {
if candidate == 0 { 0.0 } else { f64::INFINITY }
} else {
(100.0 * delta_ns as f64) / baseline as f64
};
SubMsMetricDiff {
metric: name.to_string(),
baseline_ns: baseline,
candidate_ns: candidate,
delta_ns,
delta_pct,
}
}
pub fn print_diff<W: Write>(diff: &SubMsBenchDiff, out: &mut W) -> io::Result<()> {
writeln!(
out,
"diff: {} vs {} ({}) threshold=+{:.1}%",
diff.baseline_workload, diff.candidate_workload, diff.lang, diff.regression_threshold_pct
)?;
writeln!(
out,
" {:<12} {:<7} {:>9} {:>9} {:>9} {:>9} verdict",
"stage", "metric", "baseline", "candidate", "delta", "%delta"
)?;
for stage in &diff.stages {
for m in &stage.metrics {
let pct_str = if m.delta_pct.is_finite() {
format!("{:+.1}%", m.delta_pct)
} else {
"+inf%".to_string()
};
let verdict = if m.delta_pct.is_finite() && m.delta_pct > diff.regression_threshold_pct
{
"REGRESSED"
} else {
"ok"
};
let abs = m.delta_ns.unsigned_abs();
let delta_str = if m.delta_ns >= 0 {
format!("+{}", format_ns(abs))
} else {
format!("-{}", format_ns(abs))
};
writeln!(
out,
" {:<12} {:<7} {:>9} {:>9} {:>9} {:>9} {}",
stage.stage,
m.metric,
format_ns(m.baseline_ns),
format_ns(m.candidate_ns),
delta_str,
pct_str,
verdict,
)?;
}
}
if !diff.baseline_only_stages.is_empty() {
writeln!(
out,
" stages only in baseline: {}",
diff.baseline_only_stages.join(", ")
)?;
}
if !diff.candidate_only_stages.is_empty() {
writeln!(
out,
" stages only in candidate: {}",
diff.candidate_only_stages.join(", ")
)?;
}
Ok(())
}
pub fn diff_to_json<W: Write>(diff: &SubMsBenchDiff, out: &mut W) -> io::Result<()> {
let mut buf = String::with_capacity(8 * 1024);
append_diff_json(&mut buf, diff);
out.write_all(buf.as_bytes())?;
out.write_all(b"\n")?;
Ok(())
}
fn append_diff_json(out: &mut String, diff: &SubMsBenchDiff) {
out.push('{');
json_kv_str(out, "baseline_workload", &diff.baseline_workload);
out.push(',');
json_kv_str(out, "candidate_workload", &diff.candidate_workload);
out.push(',');
json_kv_str(out, "lang", &diff.lang);
out.push(',');
let _ = write!(
out,
"\"regression_threshold_pct\":{},",
diff.regression_threshold_pct
);
let _ = write!(out, "\"has_regression\":{},", diff.has_regression());
out.push_str("\"stages\":[");
for (i, s) in diff.stages.iter().enumerate() {
if i > 0 {
out.push(',');
}
out.push('{');
json_kv_str(out, "stage", &s.stage);
out.push(',');
let _ = write!(
out,
"\"worst_regression_pct\":{},",
json_number(s.worst_regression_pct)
);
out.push_str("\"metrics\":[");
for (j, m) in s.metrics.iter().enumerate() {
if j > 0 {
out.push(',');
}
out.push('{');
json_kv_str(out, "metric", &m.metric);
out.push(',');
let _ = write!(out, "\"baseline_ns\":{},", m.baseline_ns);
let _ = write!(out, "\"candidate_ns\":{},", m.candidate_ns);
let _ = write!(out, "\"delta_ns\":{},", m.delta_ns);
let _ = write!(out, "\"delta_pct\":{}", json_number(m.delta_pct));
out.push('}');
}
out.push_str("]}");
}
out.push_str("],");
out.push_str("\"baseline_only_stages\":[");
for (i, n) in diff.baseline_only_stages.iter().enumerate() {
if i > 0 {
out.push(',');
}
json_str(out, n);
}
out.push_str("],");
out.push_str("\"candidate_only_stages\":[");
for (i, n) in diff.candidate_only_stages.iter().enumerate() {
if i > 0 {
out.push(',');
}
json_str(out, n);
}
out.push_str("]}");
}
fn json_number(d: f64) -> String {
if d.is_finite() {
d.to_string()
} else {
"null".to_string()
}
}
#[cfg(test)]
#[path = "bench_tests.rs"]
mod tests;