use std::time::{Duration, Instant};
use anyhow::Context;
use zygo_core::pool::{CpuAccounting, Pool, PoolConfig};
use zygo_core::spec::{Layer, ResolveOptions, Spec};
use crate::cli::{BenchCommand, Cli};
use crate::output::{self, Style};
const WARM_P50_BUDGET_US: f64 = 2_000.0;
const WARM_P99_BUDGET_US: f64 = 10_000.0;
const EXEC_P50_BUDGET_US: f64 = 3_000.0;
const POOL_P99_BUDGET_US: f64 = 5_000.0;
pub fn run(cli: &Cli, command: &BenchCommand) -> anyhow::Result<u8> {
let measured = match command {
BenchCommand::Warm {
n,
no_cgroup,
rate,
cpu,
pool,
scripts,
cmd,
} => warm(
cli,
*n,
!no_cgroup,
*rate,
*cpu,
(!cmd.is_empty()).then_some(cmd.as_slice()),
pool.then_some(*scripts),
)?,
BenchCommand::Cold { n, image, command } => cold(cli, *n, image, command.as_deref())?,
BenchCommand::Load {
seconds,
concurrency,
cpu,
} => load(cli, *seconds, *concurrency, *cpu)?,
BenchCommand::All { quick } => return all(cli, *quick),
};
if cli.json {
output::json(&measured.1)?;
}
Ok(measured.0)
}
const NOT_A_VERDICT: u8 = 2;
fn warm(
cli: &Cli,
n: u32,
per_request_cgroup: bool,
rate: Option<f64>,
cpu: Option<f64>,
cmd: Option<&[String]>,
pool_scripts: Option<u32>,
) -> anyhow::Result<(u8, serde_json::Value)> {
let paths = super::paths(cli);
let pool = Pool::new(PoolConfig {
per_request_cgroup,
..PoolConfig::new(paths.clone())
})?;
let dir = tempfile::tempdir()?;
let handler = dir.path().join("handler.py");
std::fs::write(&handler, "def handler(event):\n return None\n")?;
let scripts: Vec<zygo_core::protocol::Script> = (0..pool_scripts.unwrap_or(0))
.map(|i| {
zygo_core::protocol::Script::inline(format!(
"CONSTANT = {i}\n\n\ndef handler(event):\n return None\n"
))
})
.collect();
let spec = Spec::default();
let resolved = spec.resolve(
None,
&Layer {
entry: (cmd.is_none() && pool_scripts.is_none()).then(|| handler.clone()),
cmd: cmd.map(<[String]>::to_vec),
image: (cmd.is_some() || pool_scripts.is_some())
.then(|| "python:3.12-slim".to_string()),
runtime: pool_scripts.map(|_| {
zygo_core::spec::Runtime::Builtin(zygo_core::spec::BuiltinRuntime::Python)
}),
cpu: cpu.map(zygo_core::spec::Cpu),
..Default::default()
},
&ResolveOptions {
one_shot: true,
pool: pool_scripts.is_some(),
..Default::default()
},
)?;
let style = Style::stdout();
eprintln!(
"{} {} {}{}",
style.dim("image"),
resolved.image,
style.dim(if per_request_cgroup {
"per-request cgroup"
} else {
"no per-request cgroup"
}),
match cmd {
Some(c) => format!(" {}", style.dim(&format!("warm-exec: {}", c.join(" ")))),
None => String::new(),
}
);
let warmup_started = Instant::now();
let function = pool
.serve(&resolved)
.with_context(|| format!("could not warm `{}`", resolved.image))?;
let warmup = warmup_started.elapsed();
let status = function.status();
eprintln!(
"{} {} in {:.0} ms (imports {:.1} ms, rss {} MB)",
style.dim("warm"),
status.runtime,
warmup.as_secs_f64() * 1000.0,
status.imports_ms,
status.rss_kb / 1024
);
let settle = (n / 20).clamp(50, 500);
if scripts.is_empty() {
for _ in 0..settle {
function.call(serde_json::Value::Null)?;
}
} else {
for script in &scripts {
function.call_script_with_timeout(
serde_json::Value::Null,
Some(script.clone()),
Duration::from_secs(30),
)?;
}
}
let mut samples = Vec::with_capacity(n as usize);
let mut phases: Vec<zygo_core::pool::CallTiming> = Vec::with_capacity(n as usize);
let mut handler_us: Vec<f64> = Vec::with_capacity(n as usize);
let cpu_before = function.cpu_accounting();
let interval = request_interval(rate);
let started = Instant::now();
for i in 0..n {
if let Some(interval) = interval {
let due = interval.mul_f64(f64::from(i));
if let Some(wait) = due.checked_sub(started.elapsed()) {
std::thread::sleep(wait);
}
}
let t0 = Instant::now();
let (outcome, timing) = match scripts.is_empty() {
true => function.call_timed(serde_json::Value::Null, Duration::from_secs(30))?,
false => function.call_script_timed(
serde_json::Value::Null,
Some(scripts[i as usize % scripts.len()].clone()),
Duration::from_secs(30),
)?,
};
samples.push(t0.elapsed().as_secs_f64() * 1e6);
phases.push(timing);
handler_us.push(outcome.metrics.wall_ms * 1000.0);
anyhow::ensure!(
outcome.succeeded(),
"request {i} failed: {}",
outcome.error.unwrap_or_else(|| "no error given".into())
);
if n >= 1000 && i > 0 && i % (n / 10) == 0 {
eprint!("\r {}%", i * 100 / n);
}
}
let elapsed = started.elapsed();
let quota = match (cpu_before, function.cpu_accounting()) {
(Some(before), Some(after)) => Some(after.since(&before)),
_ => None,
};
if n >= 1000 {
eprintln!("\r ");
}
let floor = measure_fork_floor(500);
let mut report = Report::of(&samples, elapsed, n, quota);
if cmd.is_some() {
report.p50_budget = EXEC_P50_BUDGET_US;
report.label = "warm-exec request overhead";
}
if !scripts.is_empty() {
report.p50_budget = POOL_P99_BUDGET_US;
report.p99_budget = POOL_P99_BUDGET_US;
report.label = "pooled script request overhead";
}
let mut json = report.to_json();
if !scripts.is_empty() {
json["scripts"] = scripts.len().into();
}
json["warm_ms"] = (warmup.as_secs_f64() * 1000.0).into();
json["imports_ms"] = status.imports_ms.into();
json["rss_kb"] = status.rss_kb.into();
json["runtime"] = status.runtime.clone().into();
if !cli.json {
report.print(&style);
print_phases(&phases);
if cmd.is_none() {
print_handler_share(&phases, &handler_us);
}
print_floor(&style, &floor, &report);
print_cgroup_note(&phases, &report, &style);
}
let _ = function.shutdown();
Ok((u8::from(!report.within_budget()), json))
}
mod published {
pub const WARM_P50_MS: f64 = 1.44;
pub const WARM_P99_MS: f64 = 10.53;
pub const WARM_RATE: f64 = 250.0;
pub const LOAD_PER_SECOND: f64 = 1108.0;
pub const COLD_P50_MS: f64 = 12.3;
pub const EXEC_P50_MS: f64 = 1.40;
pub const POOL_P50_MS: f64 = 1.91;
pub const POOL_P99_MS: f64 = 11.36;
pub const SERVE_MS: f64 = 34.0;
}
const PUBLISHED_TOLERANCE: f64 = 2.0;
fn all(cli: &Cli, quick: bool) -> anyhow::Result<u8> {
let style = Style::stdout();
let host = Host::describe(&super::paths(cli));
if !cli.json {
host.print(&style);
println!();
}
let before = Thermal::sample();
let (warm_n, cold_n, load_seconds) = if quick {
(500, 10, 3)
} else {
(10_000, 50, 10)
};
let mut results = serde_json::Map::new();
let mut failed = 0u8;
let mut step = |name: &str, outcome: anyhow::Result<(u8, serde_json::Value)>| {
match outcome {
Ok((code, json)) => {
failed |= code;
results.insert(name.to_string(), json);
}
Err(e) => {
eprintln!(" {} {name}: {e:#}", Style::stdout().red("could not run"));
failed |= 1;
results.insert(
name.to_string(),
serde_json::json!({ "error": format!("{e:#}") }),
);
}
}
};
if !cli.json {
println!("{}", style.bold("1/5 the warm path"));
}
step(
"warm",
warm(
cli,
warm_n,
true,
Some(published::WARM_RATE),
None,
None,
None,
),
);
if !cli.json {
println!();
println!("{}", style.bold("2/5 warm-exec"));
}
let exec_cmd: Vec<String> = ["sh", "-c", "cat"].iter().map(|s| s.to_string()).collect();
step(
"warm_exec",
warm(
cli,
warm_n,
true,
Some(published::WARM_RATE),
None,
Some(&exec_cmd),
None,
),
);
if !cli.json {
println!();
println!(
"{}",
style.bold("3/5 a runtime pool, a different script each request")
);
}
step(
"pool",
warm(
cli,
warm_n,
true,
Some(published::WARM_RATE),
None,
None,
Some(if quick { 50 } else { 1_000 }),
),
);
if !cli.json {
println!();
println!("{}", style.bold("4/5 a cold start"));
}
step("cold", cold(cli, cold_n, "python:3.12-slim", None));
let load_cores = host.cores.clamp(1, 4) as f64;
if !cli.json {
println!();
println!(
"{} {}",
style.bold("5/5 sustained throughput"),
style.dim(&format!(
"with the tenant's CPU quota raised to {load_cores:.0} cores, so this measures the runtime and not the quota"
))
);
}
step("load", load(cli, load_seconds, 4, Some(load_cores)));
let after = Thermal::sample();
let disturbance = Thermal::compare(&before, &after, &host);
let comparison = compare_with_published(&results);
if cli.json {
output::json(&serde_json::json!({
"host": host.to_json(),
"results": results,
"published": comparison.iter().map(Claim::to_json).collect::<Vec<_>>(),
"disturbance": disturbance,
"verdict": if !disturbance.is_empty() {
"not a verdict: the host was throttled or busy"
} else if failed != 0 {
"a budget was missed"
} else {
"within budget"
},
}))?;
} else {
println!();
println!("{}", style.bold("against the published numbers"));
println!(" {:<34} {:>12} {:>12}", "", "published", "here");
for claim in &comparison {
let mark = if claim.close() {
style.green("ok")
} else {
style.yellow("differs")
};
println!(
" {:<34} {:>12} {:>12} {mark}",
claim.what,
format!("{:.2} {}", claim.published, claim.unit),
match claim.measured {
Some(v) => format!("{v:.2} {}", claim.unit),
None => "—".to_string(),
}
);
}
println!();
println!(
"{}",
style.dim(
" * the throughput run lifts the tenant's CPU quota; the latency runs keep it.\n\
\n \
\"differs\" is not a failure. These were measured on the machines in\n \
docs/book/25-performance.md; a different host produces different numbers, which\n \
is why the one above is printed. What a budget says is in each section."
)
);
println!();
if !disturbance.is_empty() {
println!("{} these numbers are not a verdict:", style.red("✗"));
for reason in &disturbance {
println!(" {}", style.yellow(reason));
}
println!(
"{}",
style.dim(
" A run on a host that was throttled or busy measures the host. \n \
Repeat it on an idle machine before drawing a conclusion."
)
);
} else if failed != 0 {
println!(
"{} a budget was missed; see the sections above",
style.red("✗")
);
} else {
println!(
"{} every budget met, on an undisturbed host",
style.green("✓")
);
}
}
if !disturbance.is_empty() {
return Ok(NOT_A_VERDICT);
}
Ok(failed)
}
struct Claim {
what: &'static str,
unit: &'static str,
published: f64,
measured: Option<f64>,
higher_is_better: bool,
}
impl Claim {
fn close(&self) -> bool {
let Some(measured) = self.measured else {
return false;
};
let (a, b) = if self.higher_is_better {
(self.published, measured)
} else {
(measured, self.published)
};
a <= b * PUBLISHED_TOLERANCE
}
fn to_json(&self) -> serde_json::Value {
serde_json::json!({
"what": self.what,
"unit": self.unit,
"published": self.published,
"measured": self.measured,
"close": self.close(),
})
}
}
fn compare_with_published(results: &serde_json::Map<String, serde_json::Value>) -> Vec<Claim> {
let number =
|section: &str, key: &str| -> Option<f64> { results.get(section)?.get(key)?.as_f64() };
let micros = |section: &str, key: &str| number(section, key).map(|v| v / 1000.0);
vec![
Claim {
what: "warm request, median",
unit: "ms",
published: published::WARM_P50_MS,
measured: micros("warm", "p50_us"),
higher_is_better: false,
},
Claim {
what: "warm request, 99th percentile",
unit: "ms",
published: published::WARM_P99_MS,
measured: micros("warm", "p99_us"),
higher_is_better: false,
},
Claim {
what: "warm-exec request, median",
unit: "ms",
published: published::EXEC_P50_MS,
measured: micros("warm_exec", "p50_us"),
higher_is_better: false,
},
Claim {
what: "pooled script request, median",
unit: "ms",
published: published::POOL_P50_MS,
measured: micros("pool", "p50_us"),
higher_is_better: false,
},
Claim {
what: "pooled script request, 99th percentile",
unit: "ms",
published: published::POOL_P99_MS,
measured: micros("pool", "p99_us"),
higher_is_better: false,
},
Claim {
what: "cold `zygo run`, median",
unit: "ms",
published: published::COLD_P50_MS,
measured: number("cold", "p50_ms"),
higher_is_better: false,
},
Claim {
what: "throughput at concurrency 4 *",
unit: "req/s",
published: published::LOAD_PER_SECOND,
measured: number("load", "requests_per_second"),
higher_is_better: true,
},
Claim {
what: "`zygo serve`, once",
unit: "ms",
published: published::SERVE_MS,
measured: number("warm", "warm_ms"),
higher_is_better: false,
},
]
}
struct Host {
kernel: String,
arch: &'static str,
cpu_model: String,
cores: usize,
memory_kb: Option<u64>,
governor: Option<String>,
max_mhz: Option<f64>,
virtualised: Option<String>,
data_root: String,
version: &'static str,
}
impl Host {
fn describe(paths: &zygo_core::paths::Paths) -> Host {
Host {
kernel: read_first_line("/proc/sys/kernel/osrelease")
.unwrap_or_else(|| std::env::consts::OS.to_string()),
arch: std::env::consts::ARCH,
cpu_model: cpu_model().unwrap_or_else(|| "unknown".into()),
cores: std::thread::available_parallelism()
.map(std::num::NonZeroUsize::get)
.unwrap_or(0),
memory_kb: meminfo("MemTotal"),
governor: read_first_line("/sys/devices/system/cpu/cpu0/cpufreq/scaling_governor"),
max_mhz: read_first_line("/sys/devices/system/cpu/cpu0/cpufreq/cpuinfo_max_freq")
.and_then(|v| v.parse::<f64>().ok())
.map(|khz| khz / 1000.0),
virtualised: virtualisation(),
data_root: paths.data().display().to_string(),
version: env!("CARGO_PKG_VERSION"),
}
}
fn print(&self, style: &Style) {
println!("{}", style.bold("the machine these numbers are about"));
let line = |k: &str, v: String| println!(" {:<12} {v}", style.dim(k));
line("zygo", format!("{} ({})", self.version, self.arch));
line("kernel", self.kernel.clone());
line(
"cpu",
format!(
"{} × {}{}",
self.cores,
self.cpu_model,
match self.max_mhz {
Some(mhz) => format!(", up to {mhz:.0} MHz"),
None => String::new(),
}
),
);
if let Some(kb) = self.memory_kb {
line("memory", format!("{:.1} GiB", kb as f64 / 1024.0 / 1024.0));
}
if let Some(g) = &self.governor {
line("governor", g.clone());
}
if let Some(v) = &self.virtualised {
line("virtual", v.clone());
}
line("data root", self.data_root.clone());
}
fn to_json(&self) -> serde_json::Value {
serde_json::json!({
"zygo": self.version,
"arch": self.arch,
"kernel": self.kernel,
"cpu_model": self.cpu_model,
"cores": self.cores,
"memory_kb": self.memory_kb,
"governor": self.governor,
"max_mhz": self.max_mhz,
"virtualised": self.virtualised,
"data_root": self.data_root,
})
}
}
struct Thermal {
core_throttles: u64,
pi_flags: Option<u32>,
loadavg1: Option<f64>,
}
impl Thermal {
fn sample() -> Thermal {
Thermal {
core_throttles: core_throttle_count(),
pi_flags: pi_throttled_flags(),
loadavg1: read_first_line("/proc/loadavg")
.and_then(|l| l.split_whitespace().next()?.parse().ok()),
}
}
fn compare(before: &Thermal, after: &Thermal, host: &Host) -> Vec<String> {
let mut reasons = Vec::new();
if after.core_throttles > before.core_throttles {
reasons.push(format!(
"the CPU throttled {} time(s) during the run (/sys/devices/system/cpu/*/thermal_throttle)",
after.core_throttles - before.core_throttles
));
}
for (sample, when) in [(before, "before"), (after, "after")] {
if let Some(flags) = sample.pi_flags
&& flags & 0xF != 0
{
reasons.push(format!(
"the firmware reports under-voltage or capped frequency {when} the run \
(vcgencmd get_throttled = {flags:#x})"
));
}
}
if let Some(load) = before.loadavg1
&& host.cores > 0
&& load > host.cores as f64 / 2.0
{
reasons.push(format!(
"the machine was already busy when the run started (load {load:.2} \
on {} cores)",
host.cores
));
}
reasons
}
}
fn core_throttle_count() -> u64 {
let Ok(entries) = std::fs::read_dir("/sys/devices/system/cpu") else {
return 0;
};
let mut total = 0;
for entry in entries.flatten() {
for name in ["core_throttle_count", "package_throttle_count"] {
let path = entry.path().join("thermal_throttle").join(name);
if let Ok(text) = std::fs::read_to_string(&path)
&& let Ok(n) = text.trim().parse::<u64>()
{
total += n;
}
}
}
total
}
fn pi_throttled_flags() -> Option<u32> {
let output = std::process::Command::new("vcgencmd")
.arg("get_throttled")
.output()
.ok()?;
let text = String::from_utf8_lossy(&output.stdout);
let value = text.trim().strip_prefix("throttled=")?;
let digits = value.strip_prefix("0x").unwrap_or(value);
u32::from_str_radix(digits, 16).ok()
}
fn read_first_line(path: &str) -> Option<String> {
let text = std::fs::read_to_string(path).ok()?;
Some(text.lines().next()?.trim().to_string())
}
fn meminfo(key: &str) -> Option<u64> {
let text = std::fs::read_to_string("/proc/meminfo").ok()?;
text.lines()
.find(|l| l.starts_with(key))?
.split_whitespace()
.nth(1)?
.parse()
.ok()
}
fn cpu_model() -> Option<String> {
if let Ok(text) = std::fs::read_to_string("/proc/cpuinfo") {
for key in ["model name", "Model", "Hardware", "cpu model", "Processor"] {
if let Some(line) = text
.lines()
.find(|l| l.trim_start().starts_with(key) && l.contains(':'))
&& let Some((_, value)) = line.split_once(':')
&& !value.trim().is_empty()
{
return Some(value.trim().to_string());
}
}
}
for path in [
"/sys/firmware/devicetree/base/model",
"/proc/device-tree/model",
] {
if let Ok(bytes) = std::fs::read(path) {
let text = String::from_utf8_lossy(&bytes)
.trim_end_matches('\0')
.trim()
.to_string();
if !text.is_empty() {
return Some(text);
}
}
}
None
}
fn virtualisation() -> Option<String> {
if let Ok(output) = std::process::Command::new("systemd-detect-virt").output()
&& output.status.success()
{
let what = String::from_utf8_lossy(&output.stdout).trim().to_string();
if !what.is_empty() && what != "none" {
return Some(what);
}
return None;
}
read_first_line("/sys/class/dmi/id/product_name")
}
const COLD_BUDGET_MS: f64 = 50.0;
const LOAD_TARGET_PER_SECOND: f64 = 600.0;
fn cold(
cli: &Cli,
n: u32,
image: &str,
command: Option<&[String]>,
) -> anyhow::Result<(u8, serde_json::Value)> {
use zygo_core::image::{Reference, Store};
use zygo_core::sandbox::SandboxConfig;
let paths = super::paths(cli);
paths.ensure()?;
let store = Store::new(paths.clone());
let reference: Reference = image.parse()?;
let entry = store.get(&reference).with_context(|| {
format!("`{image}` is not in the store\n → pull it first: zygo pull {image}")
})?;
let argv: Vec<String> = match command {
Some(cmd) => cmd.to_vec(),
None => ["python3", "-c", "pass"]
.iter()
.map(|s| s.to_string())
.collect(),
};
let spec = Spec::default();
let resolved = spec.resolve(
None,
&Layer {
image: Some(image.to_string()),
cmd: Some(argv.clone()),
..Default::default()
},
&ResolveOptions {
one_shot: true,
..Default::default()
},
)?;
let overlay = zygo_core::doctor::cached(&paths)
.checks
.iter()
.any(|c| c.name == "overlayfs (userns)" && c.status == zygo_core::doctor::Status::Ok);
let mount_points = zygo_core::sandbox::mount::required_mount_points(&resolved.mounts);
let view = store.rootfs_view(&entry.layers, overlay, &mount_points)?;
let backend = zygo_core::backend::for_isolation(resolved.isolation, &paths)?;
let style = Style::stdout();
eprintln!(
"{} {image} {} {}",
style.dim("cold"),
argv.join(" "),
style.dim(if overlay { "overlayfs" } else { "flattened" })
);
let mut samples = Vec::with_capacity(n as usize);
for i in 0..n {
let newroot = paths
.tmp()
.join(format!("bench-cold-{}-{i}", std::process::id()));
std::fs::create_dir_all(&newroot)?;
let config = SandboxConfig::from_resolved(&resolved, &view, &newroot, argv.clone(), &[]);
let t0 = Instant::now();
let mut sandbox = backend
.start(&config)
.with_context(|| format!("iteration {i} could not start"))?;
let code = sandbox.wait()?;
samples.push(t0.elapsed().as_secs_f64() * 1000.0);
anyhow::ensure!(code == 0, "iteration {i} exited {code}");
let _ = std::fs::remove_dir_all(&newroot);
}
samples.sort_by(|a, b| a.partial_cmp(b).expect("no NaN"));
let q = |p: f64| percentile(&samples, p);
let p50 = q(50.0);
let json = serde_json::json!({
"image": image,
"command": argv,
"runs": n,
"p50_ms": p50,
"p90_ms": q(90.0),
"p99_ms": q(99.0),
"max_ms": samples.last().copied().unwrap_or(0.0),
"budget_ms": COLD_BUDGET_MS,
"pass": p50 < COLD_BUDGET_MS,
"rootfs": if overlay { "overlayfs" } else { "flattened" },
});
if !cli.json {
println!("cold start over {n} runs, image already in the store");
println!(
" p50 {:>7.1} ms p90 {:>7.1} p99 {:>7.1} max {:>7.1}",
p50,
q(90.0),
q(99.0),
samples.last().copied().unwrap_or(0.0)
);
println!();
let ok = p50 < COLD_BUDGET_MS;
println!(
" p50 < {COLD_BUDGET_MS:.0} ms {}",
if ok {
style.green("PASS")
} else {
style.red("FAIL")
}
);
if !overlay {
println!(
"{}",
style.dim(
" note: this host has no unprivileged overlayfs, so the rootfs is\n \
flattened — the same bind mount every run, which is the cheap case"
)
);
}
}
let within_budget = p50 < COLD_BUDGET_MS;
Ok((u8::from(!within_budget), json))
}
fn load(
cli: &Cli,
seconds: u32,
concurrency: u32,
cpu: Option<f64>,
) -> anyhow::Result<(u8, serde_json::Value)> {
use std::sync::Arc;
use std::sync::atomic::{AtomicBool, Ordering};
anyhow::ensure!(concurrency >= 1, "concurrency must be at least 1");
let paths = super::paths(cli);
let pool = Pool::new(PoolConfig::new(paths.clone()))?;
let dir = tempfile::tempdir()?;
let handler = dir.path().join("handler.py");
std::fs::write(&handler, "def handler(event):\n return None\n")?;
let spec = Spec::default();
let resolved = spec.resolve(
None,
&Layer {
entry: Some(handler.clone()),
cpu: cpu.map(zygo_core::spec::Cpu),
concurrency: Some(concurrency),
..Default::default()
},
&ResolveOptions {
one_shot: true,
..Default::default()
},
)?;
let style = Style::stdout();
let function = Arc::new(
pool.serve(&resolved)
.with_context(|| format!("could not warm `{}`", resolved.image))?,
);
eprintln!(
"{} {} clients for {seconds}s",
style.dim("load"),
concurrency
);
for _ in 0..50 {
function.call(serde_json::Value::Null)?;
}
let stop = Arc::new(AtomicBool::new(false));
let cpu_before = function.cpu_accounting();
let started = Instant::now();
let workers: Vec<_> = (0..concurrency)
.map(|_| {
let (function, stop) = (Arc::clone(&function), Arc::clone(&stop));
std::thread::spawn(move || {
let mut samples = Vec::new();
let mut lock_us = Vec::new();
let mut failures = 0u64;
let mut served = 0u64;
while !stop.load(Ordering::Relaxed) {
served += 1;
let t0 = Instant::now();
match function.call_timed(serde_json::Value::Null, Duration::from_secs(30)) {
Ok((outcome, timing)) => {
samples.push(t0.elapsed().as_secs_f64() * 1e6);
lock_us.push(timing.lock.as_secs_f64() * 1e6);
if !outcome.succeeded() {
failures += 1;
}
}
Err(_) => failures += 1,
}
}
(samples, lock_us, failures, served)
})
})
.collect();
std::thread::sleep(Duration::from_secs(u64::from(seconds)));
stop.store(true, Ordering::Relaxed);
let mut samples = Vec::new();
let mut lock_us = Vec::new();
let mut failures = 0u64;
let mut per_worker = Vec::new();
for w in workers {
let (s, l, f, served) = w.join().map_err(|_| anyhow::anyhow!("a worker panicked"))?;
samples.extend(s);
lock_us.extend(l);
failures += f;
per_worker.push(served);
}
per_worker.sort_unstable();
let elapsed = started.elapsed();
let quota = match (cpu_before, function.cpu_accounting()) {
(Some(before), Some(after)) => Some(after.since(&before)),
_ => None,
};
let count = samples.len() as u32;
let report = Report::of(&samples, elapsed, count, quota);
lock_us.sort_by(|a, b| a.partial_cmp(b).expect("no NaN"));
let met = report.per_second >= LOAD_TARGET_PER_SECOND;
let mut json = report.to_json();
json["concurrency"] = concurrency.into();
json["failures"] = failures.into();
json["target_per_second"] = LOAD_TARGET_PER_SECOND.into();
json["meets_target"] = met.into();
json["lock_p50_us"] = percentile(&lock_us, 50.0).into();
json["lock_max_us"] = lock_us.last().copied().unwrap_or(0.0).into();
json["requests_per_client"] = per_worker.clone().into();
if !cli.json {
println!(
"sustained load: {concurrency} clients, {:.1}s, empty handler",
elapsed.as_secs_f64()
);
println!(
" {:.0} requests/s {count} requests {failures} failed",
report.per_second
);
println!(
" p50 {:>7.0} µs p90 {:>7.0} p99 {:>7.0} max {:>8.0}",
report.p50, report.p90, report.p99, report.max
);
println!();
println!(
" >= {LOAD_TARGET_PER_SECOND:.0} requests/s {}",
if met {
style.green("PASS")
} else {
style.red("FAIL")
}
);
print_contention(&style, &lock_us, &report, concurrency);
print_fairness(&style, &per_worker);
report.print_quota(&style);
}
let _ = function.shutdown();
Ok((u8::from(!met), json))
}
fn print_fairness(style: &Style, per_worker: &[u64]) {
if per_worker.len() < 2 {
return;
}
let (min, max) = (per_worker[0], per_worker[per_worker.len() - 1]);
let total: u64 = per_worker.iter().sum();
let fair = total / per_worker.len() as u64;
println!();
println!(" requests per client: min {min}, max {max}, even would be {fair}");
if min * 4 < max {
println!(
"{}",
style.yellow(
" The work is not being shared. A warm function serialises on one
connection, and the lock guarding it is not fair, so a client that
releases and immediately re-acquires can starve the others for seconds
at a time. Until the agent can hold several forks at once, concurrency
above 1 buys nothing and costs fairness."
)
);
}
}
fn print_contention(style: &Style, lock_us: &[f64], report: &Report, concurrency: u32) {
if lock_us.is_empty() {
return;
}
let lock_p50 = percentile(lock_us, 50.0);
let share = lock_p50 / report.p50.max(1.0) * 100.0;
println!();
println!(
" waiting for the connection: p50 {lock_p50:>6.0} µs p99 {:>7.0} max {:>9.0} ({share:.0}% of p50)",
percentile(lock_us, 99.0),
lock_us.last().copied().unwrap_or(0.0)
);
let lock_max = lock_us.last().copied().unwrap_or(0.0);
if concurrency > 1 && lock_max > report.p50 * 20.0 {
println!(
"{}",
style.yellow(
" One client waited far longer for the connection than a request takes.\n \
A warm function serves one request at a time — the agent handles one\n \
`EXEC` to completion before reading the next — and the lock guarding\n \
that connection is not fair, so waiting is unbounded and badly skewed.\n \
Until the agent can hold several forks at once, concurrency above 1\n \
buys no throughput and costs a long latency tail."
)
);
} else if concurrency > 1 && share > 30.0 {
println!(
"{}",
style.yellow(
" Most of each request is spent queueing behind another one: the agent\n \
handles one `EXEC` to completion before reading the next, so\n \
`concurrency` bounds what the supervisor admits, not what the agent\n \
can overlap."
)
);
}
}
fn request_interval(rate: Option<f64>) -> Option<Duration> {
let rate = rate?;
if !rate.is_finite() || rate <= 0.0 {
return None;
}
Some(Duration::from_secs_f64(1.0 / rate))
}
struct Report {
p50: f64,
p90: f64,
p99: f64,
p999: f64,
max: f64,
mean: f64,
per_second: f64,
count: u32,
elapsed: Duration,
quota: Option<CpuAccounting>,
p50_budget: f64,
p99_budget: f64,
label: &'static str,
}
impl Report {
fn of(samples: &[f64], elapsed: Duration, count: u32, quota: Option<CpuAccounting>) -> Report {
let mut sorted = samples.to_vec();
sorted.sort_by(|a, b| a.partial_cmp(b).expect("no NaN in a duration"));
Report {
p50: percentile(&sorted, 50.0),
p90: percentile(&sorted, 90.0),
p99: percentile(&sorted, 99.0),
p999: percentile(&sorted, 99.9),
max: sorted.last().copied().unwrap_or(0.0),
mean: sorted.iter().sum::<f64>() / sorted.len().max(1) as f64,
per_second: count as f64 / elapsed.as_secs_f64(),
count,
elapsed,
quota,
p50_budget: WARM_P50_BUDGET_US,
p99_budget: WARM_P99_BUDGET_US,
label: "warm request overhead",
}
}
fn saturated(&self) -> bool {
self.quota.is_some_and(|q| q.saturated())
}
fn p99_is_meaningful(&self) -> bool {
!self.saturated()
}
fn within_budget(&self) -> bool {
self.p50 < self.p50_budget && (!self.p99_is_meaningful() || self.p99 < self.p99_budget)
}
fn headroom_percent(&self) -> f64 {
(self.p50_budget - self.p50) / self.p50_budget * 100.0
}
fn print(&self, style: &Style) {
println!("{} over {} requests, empty handler", self.label, self.count);
println!(
" p50 {:>7.0} µs p90 {:>7.0} p99 {:>7.0} p99.9 {:>8.0} max {:>8.0} mean {:>7.0}",
self.p50, self.p90, self.p99, self.p999, self.max, self.mean
);
println!(" {:.0} requests/s, sequential", self.per_second);
println!();
let verdict = |ok: bool, text: &str| {
if ok {
style.green(text)
} else {
style.red(text)
}
};
println!(
" p50 < {:.0} µs {}",
self.p50_budget,
verdict(
self.p50 < self.p50_budget,
if self.p50 < self.p50_budget {
"PASS"
} else {
"FAIL"
}
)
);
if self.p99_is_meaningful() {
println!(
" p99 < {:.0} µs {}",
self.p99_budget,
verdict(
self.p99 < self.p99_budget,
if self.p99 < self.p99_budget {
"PASS"
} else {
"FAIL"
}
)
);
} else {
println!(
" p99 < {:.0} µs {}",
self.p99_budget,
style.yellow("NOT MEASURED — the tenant was at its CPU quota")
);
}
let headroom = self.headroom_percent();
let note = format!(" headroom at p50: {headroom:.0}%");
println!(
"{}",
if headroom < 15.0 {
style.yellow(¬e)
} else {
note
}
);
self.print_quota(style);
}
fn print_quota(&self, style: &Style) {
let Some(q) = self.quota else {
return;
};
println!();
let demand = q.demand_cores(self.elapsed);
let per_request = if self.count > 0 {
q.usage_us as f64 / f64::from(self.count) / 1000.0
} else {
0.0
};
match q.quota_cores {
Some(cores) => println!(
" cpu {:.2} cores used of {:.2} quota ({:.0}%), {:.2} ms per request",
demand,
cores,
demand / cores * 100.0,
per_request
),
None => println!(
" cpu {demand:.2} cores used, no quota, {per_request:.2} ms per request"
),
}
if q.periods > 0 {
println!(
" throttled in {} of {} periods, {:.0} ms stopped in total",
q.throttled_periods,
q.periods,
q.throttled_us as f64 / 1000.0
);
}
if !self.saturated() {
return;
}
let bound = q.period.as_secs_f64() * 1e6 / 2.0;
println!();
println!(
"{}",
style.yellow(
" This run saturated the tenant's own CPU quota, so the tail above is\n \
the quota being enforced, not Zygo's cost. A throttled request waits out\n \
the rest of the enforcement period:"
)
);
println!(
" period {:.0} ms → an expected {:.0} µs of added latency per throttled request",
q.period.as_secs_f64() * 1000.0,
bound
);
let capacity = q.quota_cores.unwrap_or(1.0) * 1000.0 / per_request.max(0.01);
println!(
" this tenant sustains about {capacity:.0} requests/s; \
measure latency below that with --rate, or raise --cpu"
);
}
fn to_json(&self) -> serde_json::Value {
serde_json::json!({
"requests": self.count,
"p50_us": self.p50,
"p90_us": self.p90,
"p99_us": self.p99,
"p999_us": self.p999,
"max_us": self.max,
"mean_us": self.mean,
"requests_per_second": self.per_second,
"budget": { "p50_us": self.p50_budget, "p99_us": self.p99_budget },
"headroom_p50_percent": self.headroom_percent(),
"pass": self.within_budget(),
"p99_measured": self.p99_is_meaningful(),
"cpu": self.quota.map(|q| serde_json::json!({
"quota_cores": q.quota_cores,
"period_us": q.period.as_micros() as u64,
"used_cores": q.demand_cores(self.elapsed),
"ms_per_request": q.usage_us as f64 / f64::from(self.count.max(1)) / 1000.0,
"periods": q.periods,
"throttled_periods": q.throttled_periods,
"throttled_us": q.throttled_us,
"saturated": q.saturated(),
})),
})
}
}
fn print_handler_share(phases: &[zygo_core::pool::CallTiming], handler_us: &[f64]) {
if handler_us.len() != phases.len() {
return;
}
let mut handler = handler_us.to_vec();
let mut plumbing: Vec<f64> = phases
.iter()
.zip(handler_us)
.map(|(t, h)| (t.run.as_secs_f64() * 1e6 - h).max(0.0))
.collect();
handler.sort_by(|a, b| a.partial_cmp(b).expect("no NaN"));
plumbing.sort_by(|a, b| a.partial_cmp(b).expect("no NaN"));
println!(
"{:<24} p50 {:>7.0} µs p99 {:>7.0} max {:>8.0}",
" of which handler",
percentile(&handler, 50.0),
percentile(&handler, 99.0),
handler.last().copied().unwrap_or(0.0)
);
println!(
"{:<24} p50 {:>7.0} µs p99 {:>7.0} max {:>8.0}",
" of which plumbing",
percentile(&plumbing, 50.0),
percentile(&plumbing, 99.0),
plumbing.last().copied().unwrap_or(0.0)
);
}
#[cfg(unix)]
fn measure_fork_floor(iterations: usize) -> Vec<f64> {
let mut samples = Vec::with_capacity(iterations);
for _ in 0..iterations {
let t0 = Instant::now();
unsafe {
match libc::fork() {
0 => libc::_exit(0),
-1 => return samples,
pid => {
let mut status = 0;
libc::waitpid(pid, &mut status, 0);
}
}
}
samples.push(t0.elapsed().as_secs_f64() * 1e6);
}
samples
}
#[cfg(not(unix))]
fn measure_fork_floor(_iterations: usize) -> Vec<f64> {
Vec::new()
}
fn print_floor(style: &Style, floor: &[f64], report: &Report) {
if floor.is_empty() {
return;
}
let mut sorted = floor.to_vec();
sorted.sort_by(|a, b| a.partial_cmp(b).expect("no NaN"));
let floor_p50 = percentile(&sorted, 50.0);
let floor_p99 = percentile(&sorted, 99.0);
println!();
println!(
" {} bare fork+wait on this host: p50 {:.0} µs, p99 {:.0} µs",
style.dim("floor"),
floor_p50,
floor_p99
);
if floor_p99 > WARM_P99_BUDGET_US * 0.2 {
println!(
"{}",
style.yellow(&format!(
" the host's own fork p99 is {:.0}% of the p99 budget — this measurement is \n dominated by the machine, not by the warm path. Re-run on an idle host.",
floor_p99 / WARM_P99_BUDGET_US * 100.0
))
);
}
let _ = report;
}
fn cgroup_share_of_p99(phases: &[zygo_core::pool::CallTiming]) -> Option<f64> {
if phases.is_empty() {
return None;
}
let p99_of = |get: fn(&zygo_core::pool::CallTiming) -> Duration| {
let mut v: Vec<f64> = phases.iter().map(|t| get(t).as_secs_f64() * 1e6).collect();
v.sort_by(|a, b| a.partial_cmp(b).expect("no NaN"));
percentile(&v, 99.0)
};
let total = p99_of(|t| t.lock + t.fork + t.admit + t.run + t.release);
if total <= 0.0 {
return None;
}
Some((p99_of(|t| t.admit) + p99_of(|t| t.release)) / total)
}
fn print_cgroup_note(phases: &[zygo_core::pool::CallTiming], report: &Report, style: &Style) {
let Some(share) = cgroup_share_of_p99(phases) else {
return;
};
if share < 0.5 || !report.p99_is_meaningful() || report.p99 < WARM_P99_BUDGET_US {
return;
}
println!();
println!(
"{}",
style.yellow(&format!(
" {:.0}% of the p99 above is `admit` and `release` — creating this request's\n \
own cgroup and removing it, not running the handler. That is the cost of\n \
per-request containment, not a property of the warm path.",
share * 100.0
))
);
println!(
" measure without it: zygo bench warm --no-cgroup{}",
match report.count {
0 => String::new(),
n => format!(" --n {n}"),
}
);
let mountinfo = std::fs::read_to_string("/proc/self/mountinfo").unwrap_or_default();
if zygo_core::doctor::cgroup2_favors_moves(&mountinfo) == Some(false) {
println!(
"{}",
style.yellow(
" cgroup2 here has no favordynmods, so on Linux 6.0+ moving a request into\n \
its cgroup can wait several ms for the kernel. `zygo doctor` says how to\n \
change that, and what it costs."
)
);
}
}
fn print_phases(phases: &[zygo_core::pool::CallTiming]) {
let micros = |d: Duration| d.as_secs_f64() * 1e6;
type Column = (&'static str, fn(&zygo_core::pool::CallTiming) -> Duration);
let columns: [Column; 5] = [
(" lock (contention)", |t| t.lock),
(" fork (EXEC→FORKED)", |t| t.fork),
(" admit (cgroup)", |t| t.admit),
(" run (GO→DONE)", |t| t.run),
(" release (cgroup rm)", |t| t.release),
];
println!();
for (label, get) in columns {
let mut values: Vec<f64> = phases.iter().map(|t| micros(get(t))).collect();
values.sort_by(|a, b| a.partial_cmp(b).expect("no NaN"));
println!(
"{label:<24} p50 {:>7.0} µs p99 {:>7.0} max {:>8.0}",
percentile(&values, 50.0),
percentile(&values, 99.0),
values.last().copied().unwrap_or(0.0)
);
}
}
fn percentile(sorted: &[f64], p: f64) -> f64 {
if sorted.is_empty() {
return 0.0;
}
let rank = (p / 100.0) * (sorted.len() - 1) as f64;
let lo = rank.floor() as usize;
let hi = rank.ceil() as usize;
if lo == hi {
sorted[lo]
} else {
sorted[lo] + (sorted[hi] - sorted[lo]) * (rank - lo as f64)
}
}
#[cfg(test)]
mod tests {
use super::*;
fn report(samples: &[f64]) -> Report {
Report::of(samples, Duration::from_secs(1), samples.len() as u32, None)
}
fn report_with_quota(samples: &[f64], periods: u64, throttled: u64) -> Report {
Report::of(
samples,
Duration::from_secs(1),
samples.len() as u32,
Some(CpuAccounting {
usage_us: 1_000_000,
periods,
throttled_periods: throttled,
throttled_us: throttled * 50_000,
quota_cores: Some(1.0),
period: Duration::from_millis(100),
}),
)
}
#[test]
fn a_rate_becomes_the_gap_between_requests() {
let gap = request_interval(Some(300.0)).expect("paced");
assert!(
(gap.as_secs_f64() - 1.0 / 300.0).abs() < 1e-9,
"300 requests/s is 3.3 ms apart, not 300 s; got {gap:?}"
);
assert_eq!(request_interval(Some(1.0)), Some(Duration::from_secs(1)));
assert_eq!(
request_interval(Some(1000.0)),
Some(Duration::from_millis(1))
);
}
#[test]
fn a_rate_that_cannot_be_paced_runs_unpaced() {
assert_eq!(request_interval(None), None, "no --rate at all");
assert_eq!(request_interval(Some(0.0)), None, "0/s would never finish");
assert_eq!(request_interval(Some(-5.0)), None);
assert_eq!(request_interval(Some(f64::NAN)), None);
assert_eq!(
request_interval(Some(f64::INFINITY)),
None,
"infinitely fast is just unpaced"
);
}
#[test]
fn percentiles_interpolate() {
let sorted: Vec<f64> = (1..=100).map(|n| n as f64).collect();
assert_eq!(percentile(&sorted, 0.0), 1.0);
assert_eq!(percentile(&sorted, 100.0), 100.0);
assert!((percentile(&sorted, 50.0) - 50.5).abs() < 0.01);
assert_eq!(percentile(&[], 50.0), 0.0);
}
#[test]
fn the_budget_is_the_design_documents() {
assert_eq!(WARM_P50_BUDGET_US, 2_000.0);
assert_eq!(WARM_P99_BUDGET_US, 10_000.0);
}
#[test]
fn a_run_inside_the_budget_passes_and_one_outside_does_not() {
let fast: Vec<f64> = (0..1000).map(|i| 1000.0 + (i % 100) as f64).collect();
assert!(report(&fast).within_budget());
let slow: Vec<f64> = (0..1000).map(|_| 2500.0).collect();
assert!(!report(&slow).within_budget());
let mut tailed: Vec<f64> = (0..1000).map(|_| 1000.0).collect();
tailed.extend((0..50).map(|_| 50_000.0));
assert!(!report(&tailed).within_budget(), "p99 was ignored");
}
#[test]
fn headroom_is_reported_as_a_percentage_of_the_budget() {
let r = report(&[1_880.0; 100]);
assert!(
(r.headroom_percent() - 6.0).abs() < 0.5,
"{}",
r.headroom_percent()
);
let r = report(&[1_000.0; 100]);
assert!((r.headroom_percent() - 50.0).abs() < 0.5);
}
#[test]
fn the_json_report_carries_the_budget_it_was_judged_against() {
let v = report(&[1_500.0; 100]).to_json();
assert_eq!(v["budget"]["p50_us"], 2_000.0);
assert_eq!(v["pass"], true);
assert!(v["headroom_p50_percent"].as_f64().unwrap() > 0.0);
}
#[test]
fn exit_status_follows_the_verdict() {
assert_eq!(u8::from(!report(&[1_000.0; 100]).within_budget()), 0);
assert_eq!(u8::from(!report(&[9_000.0; 100]).within_budget()), 1);
}
fn a_throttled_run() -> Report {
let mut samples: Vec<f64> = (0..980).map(|_| 1_400.0).collect();
samples.extend((0..20).map(|_| 48_000.0));
report_with_quota(&samples, 28, 27)
}
#[test]
fn a_tail_made_by_the_cpu_quota_is_not_reported_as_a_failure() {
let r = a_throttled_run();
assert!(r.saturated());
assert!(
r.p99 > WARM_P99_BUDGET_US,
"the fixture should be over budget at p99"
);
assert!(
!r.p99_is_meaningful(),
"a saturated run cannot judge the p99 budget"
);
assert!(
r.within_budget(),
"the quota doing its job is not a failure of this code"
);
assert_eq!(u8::from(!r.within_budget()), 0, "and CI should not go red");
}
#[test]
fn a_saturated_run_still_fails_on_p50() {
let r = report_with_quota(&[2_500.0; 1000], 28, 27);
assert!(r.saturated());
assert!(!r.within_budget());
}
#[test]
fn an_unsaturated_run_is_judged_on_both_budgets() {
let mut samples: Vec<f64> = (0..980).map(|_| 1_400.0).collect();
samples.extend((0..20).map(|_| 48_000.0));
let r = report_with_quota(&samples, 200, 1);
assert!(!r.saturated(), "1 period in 200 is noise, not saturation");
assert!(r.p99_is_meaningful());
assert!(!r.within_budget(), "a real tail must still fail");
}
#[test]
fn a_host_without_a_readable_cgroup_judges_both_budgets() {
let mut samples: Vec<f64> = (0..980).map(|_| 1_400.0).collect();
samples.extend((0..20).map(|_| 48_000.0));
let r = report(&samples);
assert!(!r.saturated());
assert!(!r.within_budget());
}
#[test]
fn the_json_report_says_whether_the_p99_was_measured() {
let v = a_throttled_run().to_json();
assert_eq!(v["p99_measured"], false);
assert_eq!(v["cpu"]["saturated"], true);
assert_eq!(v["cpu"]["throttled_periods"], 27);
assert_eq!(v["cpu"]["quota_cores"], 1.0);
assert_eq!(v["cpu"]["period_us"], 100_000);
assert!(v["p99_us"].as_f64().unwrap() > WARM_P99_BUDGET_US);
let v = report(&[1_500.0; 100]).to_json();
assert_eq!(v["p99_measured"], true);
assert!(v["cpu"].is_null(), "no cgroup, no cpu section");
}
}