use std::collections::HashMap;
use std::time::{Duration, Instant};
use zygo_core::pool::Status;
use zygo_core::supervisor::client::Client;
use zygo_core::supervisor::{Request, Response};
use crate::cli::Cli;
use crate::output::{self, Style};
struct Sample {
at: Instant,
counters: HashMap<String, (u64, u64)>,
}
pub fn run(cli: &Cli, interval: f64, once: bool) -> anyhow::Result<u8> {
anyhow::ensure!(interval > 0.0, "the interval must be greater than zero");
let interval = Duration::from_secs_f64(interval);
let paths = super::paths(cli);
let mut client = Client::connect(&paths)?;
let style = Style::stdout();
let animate = !once && !cli.json && std::io::IsTerminal::is_terminal(&std::io::stdout());
let mut previous: Option<Sample> = None;
loop {
let 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),
};
let now = Instant::now();
let rows: Vec<Row> = functions
.iter()
.map(|f| Row::of(f, previous.as_ref(), now))
.collect();
if cli.json {
output::json(&serde_json::json!({
"functions": rows,
"runtimes": pools,
}))?;
return Ok(0);
}
if animate {
print!("\x1b[H\x1b[J");
}
print(&rows, &style, previous.is_none() && pools.is_empty());
print_pools(&pools, &style);
previous = Some(Sample {
at: now,
counters: functions
.iter()
.map(|f| (f.name.clone(), (f.requests, f.failures)))
.collect(),
});
if once {
return Ok(0);
}
std::thread::sleep(interval);
}
}
#[derive(serde::Serialize)]
struct Row {
name: String,
state: String,
rss_kb: u64,
requests: u64,
failures: u64,
#[serde(skip_serializing_if = "Option::is_none")]
requests_per_s: Option<f64>,
#[serde(skip_serializing_if = "Option::is_none")]
failures_per_s: Option<f64>,
}
impl Row {
fn of(f: &Status, previous: Option<&Sample>, now: Instant) -> Row {
let rates = previous.and_then(|p| {
let (req, fail) = p.counters.get(&f.name)?;
let seconds = now.duration_since(p.at).as_secs_f64();
if seconds <= 0.0 {
return None;
}
Some((
f.requests.saturating_sub(*req) as f64 / seconds,
f.failures.saturating_sub(*fail) as f64 / seconds,
))
});
Row {
name: f.name.clone(),
state: f.state.as_str().to_string(),
rss_kb: f.rss_kb,
requests: f.requests,
failures: f.failures,
requests_per_s: rates.map(|(r, _)| r),
failures_per_s: rates.map(|(_, f)| f),
}
}
}
fn print_pools(pools: &[zygo_core::supervisor::RuntimeStatus], style: &Style) {
if pools.is_empty() {
return;
}
let name_w = pools.iter().map(|p| p.name.len()).max().unwrap_or(7).max(7);
println!();
println!(
"{:name_w$} {:>5} {:>6} {:>5} {:>7} {:>9} {:>8} {:>8}",
"RUNTIME", "WARM", "PAUSED", "ROOM", "IN/QUEUE", "RSS", "REQUESTS", "FAILURES"
);
for p in pools {
println!(
"{:name_w$} {:>5} {:>6} {:>5} {:>7} {:>9} {:>8} {:>8}",
p.name,
p.warm,
p.paused,
p.cold,
format!("{}/{}", p.in_flight, p.queued),
format!("{:.1} MB", p.rss_kb as f64 / 1024.0),
p.requests,
p.failures,
);
}
println!(
"{}",
style.dim(" ROOM is what is left between the zygotes that exist and max_warm")
);
}
fn print(rows: &[Row], style: &Style, first_frame: bool) {
if rows.is_empty() {
println!("{}", style.dim("no warm functions"));
return;
}
let name_w = rows.iter().map(|r| r.name.len()).max().unwrap_or(4).max(4);
println!(
"{:name_w$} {:6} {:>9} {:>8} {:>8} {:>8}",
"NAME", "STATE", "RSS", "REQ/S", "REQUESTS", "FAILURES"
);
for r in rows {
let rate = match r.requests_per_s {
Some(v) => format!("{v:.1}"),
None => "—".into(),
};
println!(
"{:name_w$} {:6} {:>9} {:>8} {:>8} {:>8}",
r.name,
r.state,
format!("{:.1} MB", r.rss_kb as f64 / 1024.0),
rate,
r.requests,
r.failures,
);
}
if first_frame {
println!();
println!(
"{}",
style.dim("a rate needs two samples; the next frame has one")
);
}
}
#[cfg(test)]
mod tests {
use super::*;
use zygo_core::sandbox::SandboxState;
fn status(name: &str, requests: u64, failures: u64) -> Status {
Status {
tenant: "default".into(),
name: name.into(),
image: String::new(),
state: SandboxState::Warm,
runtime: "python/3.12".into(),
rss_kb: 2048,
imports_ms: 1.0,
requests,
failures,
}
}
#[test]
fn the_first_frame_has_no_rate() {
let row = Row::of(&status("f", 10, 1), None, Instant::now());
assert_eq!(row.requests_per_s, None);
assert_eq!(row.requests, 10, "the counter is still shown");
}
#[test]
fn a_rate_is_the_difference_over_the_time_that_passed() {
let then = Instant::now();
let previous = Sample {
at: then,
counters: [("f".to_string(), (10, 1))].into_iter().collect(),
};
let now = then + Duration::from_secs(2);
let row = Row::of(&status("f", 30, 2), Some(&previous), now);
assert_eq!(row.requests_per_s, Some(10.0));
assert_eq!(row.failures_per_s, Some(0.5));
}
#[test]
fn a_function_that_was_not_there_before_gets_no_rate() {
let then = Instant::now();
let previous = Sample {
at: then,
counters: [("other".to_string(), (5, 0))].into_iter().collect(),
};
let row = Row::of(
&status("new", 900, 0),
Some(&previous),
then + Duration::from_secs(1),
);
assert_eq!(row.requests_per_s, None);
}
#[test]
fn counters_going_backwards_do_not_become_an_enormous_rate() {
let then = Instant::now();
let previous = Sample {
at: then,
counters: [("f".to_string(), (1000, 10))].into_iter().collect(),
};
let row = Row::of(
&status("f", 3, 0),
Some(&previous),
then + Duration::from_secs(1),
);
assert_eq!(row.requests_per_s, Some(0.0));
}
}