use serde::Serialize;
use zygo_core::pool::{LogEntry, LogKind, Status};
use zygo_core::sandbox::SandboxState;
use zygo_core::supervisor::client::Client;
use zygo_core::supervisor::{Request, Response};
use crate::cli::Cli;
use crate::output::{self, Style};
const ENOUGH_FOR_P99: usize = 100;
const WHOLE_LOG: u32 = u32::MAX;
#[derive(Debug, Serialize)]
struct FunctionStats {
name: String,
state: SandboxState,
runtime: String,
requests: u64,
failures: u64,
samples: usize,
#[serde(skip_serializing_if = "Option::is_none")]
p50_ms: Option<f64>,
#[serde(skip_serializing_if = "Option::is_none")]
p99_ms: Option<f64>,
#[serde(skip_serializing_if = "Option::is_none")]
max_ms: Option<f64>,
timeouts: usize,
other_kills: usize,
}
pub fn run(cli: &Cli, name: Option<&str>) -> anyhow::Result<u8> {
let paths = super::paths(cli);
let mut client = Client::connect(&paths)?;
let mut functions = match client.send(&Request::List)? {
Response::Functions { functions } => functions,
other => return super::supervisor::report_failure(cli, &other),
};
let pools = match client.send(&Request::Runtimes)? {
Response::Runtimes { runtimes } => runtimes,
other => return super::supervisor::report_failure(cli, &other),
};
functions.extend(pools.iter().map(|p| Status {
name: p.name.clone(),
tenant: p.tenant.clone(),
image: p.image.clone(),
state: if p.warm > 0 {
SandboxState::Warm
} else if p.paused > 0 {
SandboxState::Paused
} else {
SandboxState::Cold
},
runtime: format!("{} ×{}", p.runtime, p.warm + p.paused),
rss_kb: p.rss_kb,
imports_ms: 0.0,
requests: p.requests,
failures: p.failures,
}));
let wanted: Vec<Status> = match name {
Some(n) => functions.into_iter().filter(|f| f.name == n).collect(),
None => functions,
};
if let Some(n) = name
&& wanted.is_empty()
{
anyhow::bail!("no function or runtime named `{n}`; `zygo ps` lists the warm ones");
}
let mut stats = Vec::with_capacity(wanted.len());
for status in wanted {
let entries = match client.send(&Request::Logs {
name: status.name.clone(),
after: 0,
limit: WHOLE_LOG,
failed: false,
tenant: None,
})? {
Response::Logs { entries, .. } => entries,
_ => Vec::new(),
};
stats.push(summarise(status, &entries));
}
if cli.json {
output::json(&stats)?;
return Ok(0);
}
print_table(&stats);
Ok(0)
}
fn summarise(status: Status, entries: &[LogEntry]) -> FunctionStats {
let mut wall: Vec<f64> = Vec::new();
let mut timeouts = 0;
let mut other_kills = 0;
for entry in entries {
let LogKind::Request {
exit_code,
wall_ms,
timed_out,
..
} = &entry.kind
else {
continue;
};
wall.push(*wall_ms);
if *timed_out {
timeouts += 1;
} else if *exit_code == 137 {
other_kills += 1;
}
}
wall.sort_by(f64::total_cmp);
let samples = wall.len();
FunctionStats {
name: status.name,
state: status.state,
runtime: status.runtime,
requests: status.requests,
failures: status.failures,
samples,
p50_ms: percentile(&wall, 0.50),
p99_ms: (samples >= ENOUGH_FOR_P99)
.then(|| percentile(&wall, 0.99))
.flatten(),
max_ms: wall.last().copied(),
timeouts,
other_kills,
}
}
fn percentile(sorted: &[f64], q: f64) -> Option<f64> {
if sorted.is_empty() {
return None;
}
let rank = (q * sorted.len() as f64).ceil().max(1.0) as usize;
sorted.get(rank - 1).copied()
}
fn print_table(stats: &[FunctionStats]) {
let style = Style::stdout();
if stats.is_empty() {
println!("{}", style.dim("no warm functions"));
return;
}
let ms = |v: Option<f64>| {
v.map(|v| format!("{v:.1} ms"))
.unwrap_or_else(|| "—".into())
};
let name_w = stats.iter().map(|s| s.name.len()).max().unwrap_or(4).max(4);
println!(
"{:name_w$} {:6} {:>8} {:>8} {:>7} {:>9} {:>9} {:>8} {:>6}",
"NAME", "STATE", "REQUESTS", "FAILURES", "SAMPLES", "p50", "p99", "MAX", "KILLED"
);
let mut any_short = false;
for s in stats {
if s.samples > 0 && s.p99_ms.is_none() {
any_short = true;
}
let killed = s.timeouts + s.other_kills;
println!(
"{:name_w$} {:6} {:>8} {:>8} {:>7} {:>9} {:>9} {:>8} {:>6}",
s.name,
s.state.as_str(),
s.requests,
s.failures,
s.samples,
ms(s.p50_ms),
s.p99_ms
.map(|v| format!("{v:.1} ms"))
.unwrap_or_else(|| "—".into()),
ms(s.max_ms),
killed,
);
}
let timeouts: usize = stats.iter().map(|s| s.timeouts).sum();
let others: usize = stats.iter().map(|s| s.other_kills).sum();
println!();
println!(
"{}",
style.dim(
"requests and failures are counted since the function was warmed; \
the latencies are over the entries still in its log"
)
);
if any_short {
println!(
"{}",
style.dim(&format!(
"p99 is shown from {ENOUGH_FOR_P99} samples up; below that it would be the max"
))
);
}
if timeouts + others > 0 {
println!(
"{}",
style.dim(&format!(
"of the killed: {timeouts} overran a deadline, {others} were killed by \
something else — usually the memory limit"
))
);
}
}
#[cfg(test)]
mod tests {
use super::*;
fn request(wall_ms: f64, exit_code: i32, timed_out: bool) -> LogEntry {
LogEntry {
seq: 1,
at_ms: 0,
kind: LogKind::Request {
id: "x".into(),
exit_code,
wall_ms,
timed_out,
error: None,
stderr: String::new(),
},
text: String::new(),
}
}
fn status() -> Status {
Status {
tenant: "default".into(),
name: "f".into(),
image: String::new(),
state: SandboxState::Warm,
runtime: "python/3.12".into(),
rss_kb: 1000,
imports_ms: 1.0,
requests: 1_000_000,
failures: 3,
}
}
#[test]
fn the_percentile_is_one_of_the_observed_values() {
let sorted = [1.0, 2.0, 3.0, 4.0];
assert_eq!(percentile(&sorted, 0.5), Some(2.0));
assert_eq!(percentile(&sorted, 0.99), Some(4.0));
assert_eq!(percentile(&[], 0.5), None);
assert_eq!(percentile(&[7.0], 0.5), Some(7.0));
}
#[test]
fn p99_is_withheld_until_there_are_enough_samples() {
let few: Vec<LogEntry> = (0..10).map(|i| request(i as f64, 0, false)).collect();
let short = summarise(status(), &few);
assert_eq!(short.samples, 10);
assert!(short.p50_ms.is_some());
assert_eq!(short.p99_ms, None, "ten samples have no p99");
assert!(short.max_ms.is_some(), "the slowest is still reported");
let many: Vec<LogEntry> = (0..ENOUGH_FOR_P99)
.map(|i| request(i as f64, 0, false))
.collect();
assert!(summarise(status(), &many).p99_ms.is_some());
}
#[test]
fn a_deadline_kill_is_counted_apart_from_every_other_kill() {
let entries = vec![
request(1.0, 0, false),
request(2.0, 137, true),
request(3.0, 137, false),
request(4.0, 1, false),
];
let s = summarise(status(), &entries);
assert_eq!(s.timeouts, 1);
assert_eq!(s.other_kills, 1);
assert_eq!(s.samples, 4, "every request is a sample, killed or not");
}
#[test]
fn the_counters_are_not_the_sample_count() {
let s = summarise(status(), &[request(1.0, 0, false)]);
assert_eq!(s.requests, 1_000_000, "the counter is since the warm-up");
assert_eq!(s.samples, 1, "the log holds one");
}
}