use std::collections::{HashMap, HashSet};
use futures_util::StreamExt;
use serde::Deserialize;
use crate::compose::types::ComposeFile;
use crate::error::{ComposeError, Result};
use crate::libpod::types::container::ContainerListEntry;
use crate::libpod::{parse_json_lines, urlencoded, API_PREFIX};
use super::Engine;
#[derive(Default)]
pub struct StatsOptions {
pub no_stream: bool,
pub all: bool,
pub json: bool,
pub no_trunc: bool,
}
impl StatsOptions {
pub fn new(no_stream: bool, all: bool, no_trunc: bool, json: bool) -> Self {
Self {
no_stream,
all,
no_trunc,
json,
}
}
}
const NAME_WIDTH: usize = 32;
fn containers_query(wanted: &HashSet<String>) -> String {
if wanted.is_empty() {
return String::new();
}
let mut names: Vec<&String> = wanted.iter().collect();
names.sort();
names
.iter()
.map(|n| format!("&containers={}", urlencoded(n)))
.collect::<String>()
}
fn stats_stream_broke_mid_sample(
sampled: &HashSet<String>,
still_running: &HashSet<String>,
) -> bool {
sampled.iter().any(|c| still_running.contains(c))
}
fn null_default<'de, D, T>(d: D) -> std::result::Result<T, D::Error>
where
D: serde::Deserializer<'de>,
T: Default + Deserialize<'de>,
{
Option::<T>::deserialize(d).map(|v| v.unwrap_or_default())
}
#[derive(Deserialize, Default)]
struct StatsReport {
#[serde(rename = "Stats", default)]
stats: Vec<ContainerStat>,
}
#[derive(Deserialize, Default, Clone)]
struct ContainerStat {
#[serde(rename = "Name", default)]
name: String,
#[serde(rename = "CPU", default)]
cpu: f64,
#[serde(rename = "MemUsage", default)]
mem_usage: u64,
#[serde(rename = "MemLimit", default)]
mem_limit: u64,
#[serde(rename = "MemPerc", default)]
mem_perc: f64,
#[serde(rename = "BlockInput", default)]
block_in: u64,
#[serde(rename = "BlockOutput", default)]
block_out: u64,
#[serde(rename = "PIDs", default)]
pids: u64,
#[serde(rename = "Network", default, deserialize_with = "null_default")]
network: HashMap<String, NetStat>,
}
#[derive(Deserialize, Default, Clone)]
struct NetStat {
#[serde(rename = "RxBytes", default)]
rx: u64,
#[serde(rename = "TxBytes", default)]
tx: u64,
}
fn format_bytes(bytes: u64) -> String {
const UNITS: [&str; 5] = ["B", "KiB", "MiB", "GiB", "TiB"];
let mut value = bytes as f64;
let mut unit = 0;
while value >= 1024.0 && unit < UNITS.len() - 1 {
value /= 1024.0;
unit += 1;
}
if unit == 0 {
format!("{bytes}B")
} else {
format!("{value:.1}{}", UNITS[unit])
}
}
fn net_totals(s: &ContainerStat) -> (u64, u64) {
s.network
.values()
.fold((0u64, 0u64), |(rx, tx), n| (rx + n.rx, tx + n.tx))
}
fn truncate_name(name: &str, no_trunc: bool) -> String {
if no_trunc || name.chars().count() <= NAME_WIDTH {
return name.to_string();
}
let head: String = name.chars().take(NAME_WIDTH - 1).collect();
format!("{head}…")
}
fn load_style(pct: f64) -> crate::ui::Style {
use crate::ui::AnsiColor;
let colour = if pct >= 90.0 {
AnsiColor::Red
} else if pct >= 70.0 {
AnsiColor::Yellow
} else {
AnsiColor::Green
};
crate::ui::Style::new().fg_color(Some(colour.into()))
}
fn format_row_with(s: &ContainerStat, no_trunc: bool, colour: bool) -> String {
use crate::ui::{identity_style, paint};
let (rx, tx) = net_totals(s);
let dim = crate::ui::Style::new().dimmed();
let name = format!("{:<NAME_WIDTH$}", truncate_name(&s.name, no_trunc));
let name = paint(identity_style(s.name.trim()), &name, colour);
let cpu = paint(load_style(s.cpu), &format!("{:>7.2}%", s.cpu), colour);
let mem_pct = paint(
load_style(s.mem_perc),
&format!("{:>6.2}%", s.mem_perc),
colour,
);
let mem = paint(
dim,
&format!(
"{:>10} / {:<10}",
format_bytes(s.mem_usage),
format_bytes(s.mem_limit)
),
colour,
);
let net = paint(
dim,
&format!("{:>9} / {:<9}", format_bytes(rx), format_bytes(tx)),
colour,
);
let block = paint(
dim,
&format!(
"{:>9} / {:<9}",
format_bytes(s.block_in),
format_bytes(s.block_out)
),
colour,
);
format!("{name} {cpu} {mem} {mem_pct} {net} {block} {:>5}", s.pids)
}
fn stat_json_row(s: &ContainerStat) -> serde_json::Value {
let (rx, tx) = net_totals(s);
serde_json::json!({
"Name": s.name,
"CPUPerc": s.cpu,
"MemUsage": s.mem_usage,
"MemLimit": s.mem_limit,
"MemPerc": s.mem_perc,
"NetInput": rx,
"NetOutput": tx,
"BlockInput": s.block_in,
"BlockOutput": s.block_out,
"PIDs": s.pids,
})
}
const HEADER: &str = "NAME CPU % MEM USAGE / LIMIT MEM % NET I/O BLOCK I/O PIDS";
impl Engine {
pub async fn stats(
&self,
file: &ComposeFile,
target_services: &[String],
no_stream: bool,
) -> Result<()> {
self.stats_with_options(
file,
target_services,
StatsOptions {
no_stream,
..StatsOptions::default()
},
)
.await
}
pub async fn stats_with_options(
&self,
file: &ComposeFile,
target_services: &[String],
opts: StatsOptions,
) -> Result<()> {
if let Some(unknown) = first_unknown_service(file, target_services) {
return Err(ComposeError::ServiceNotFound(unknown.into()));
}
let targets = self.target_containers(file, target_services).await?;
let running: HashSet<String> = targets
.iter()
.filter(|t| t.running)
.map(|t| t.name.clone())
.collect();
let stopped: Vec<String> = if opts.all {
targets
.iter()
.filter(|t| !t.running)
.map(|t| t.name.clone())
.collect()
} else {
Vec::new()
};
let containers = containers_query(&running);
if opts.no_stream || running.is_empty() {
let report = if running.is_empty() {
StatsReport::default()
} else {
self.client
.get_json(&format!(
"{API_PREFIX}/containers/stats?stream=false{containers}"
))
.await
.map_err(ComposeError::Podman)?
};
print_frame(&report, &running, &stopped, &opts);
return Ok(());
}
let resp = self
.client
.get_stream(&format!(
"{API_PREFIX}/containers/stats?stream=true{containers}"
))
.await
.map_err(ComposeError::Podman)?;
let mut frames = parse_json_lines::<StatsReport>(resp.into_body());
while let Some(frame) = frames.next().await {
match frame {
Ok(report) => print_frame(&report, &running, &stopped, &opts),
Err(e) => {
let still_running = match self.target_containers(file, target_services).await {
Ok(targets) => targets
.into_iter()
.filter(|c| c.running)
.map(|c| c.name)
.collect::<HashSet<String>>(),
Err(_) => {
tracing::warn!(
"stats: stream ended and the running set could not be \
re-checked [{}]: {e}",
e.stream_end_kind()
);
return Err(ComposeError::Podman(e));
}
};
if stats_stream_broke_mid_sample(&running, &still_running) {
tracing::warn!(
"stats: stream broke while a container was still running [{}]: {e}",
e.stream_end_kind()
);
return Err(ComposeError::Podman(e));
}
tracing::debug!(
"stats: stream ended as its containers stopped [{}]",
e.stream_end_kind()
);
break;
}
}
}
Ok(())
}
async fn target_containers(
&self,
file: &ComposeFile,
target_services: &[String],
) -> Result<Vec<TargetContainer>> {
let filters = serde_json::json!({ "label": [format!("podup.project={}", self.project)] });
let path = format!(
"{API_PREFIX}/containers/json?all=true&filters={}",
urlencoded(&filters.to_string()),
);
let entries = self
.client
.get_json::<Vec<ContainerListEntry>>(&path)
.await
.map_err(ComposeError::Podman)?;
let mut out = Vec::new();
for e in entries {
let service = e
.labels
.get("podup.service")
.map(String::as_str)
.unwrap_or("");
if !file.services.contains_key(service) {
continue;
}
if !target_services.is_empty() && !target_services.iter().any(|t| t == service) {
continue;
}
if let Some(raw) = e.names.first() {
out.push(TargetContainer {
name: raw.trim_start_matches('/').to_string(),
running: e.state == "running",
});
}
}
Ok(out)
}
}
struct TargetContainer {
name: String,
running: bool,
}
fn first_unknown_service<'a>(file: &ComposeFile, targets: &'a [String]) -> Option<&'a str> {
targets
.iter()
.map(String::as_str)
.find(|t| !file.services.contains_key(*t))
}
fn frame_rows(
report: &StatsReport,
running: &HashSet<String>,
stopped: &[String],
) -> Vec<ContainerStat> {
let mut rows: Vec<ContainerStat> = report
.stats
.iter()
.filter(|s| running.contains(&s.name))
.cloned()
.collect();
for name in stopped {
rows.push(ContainerStat {
name: name.clone(),
..ContainerStat::default()
});
}
rows.sort_by(|a, b| a.name.cmp(&b.name));
rows
}
fn print_frame(
report: &StatsReport,
running: &HashSet<String>,
stopped: &[String],
opts: &StatsOptions,
) {
let rows = frame_rows(report, running, stopped);
if opts.json {
let json: Vec<_> = rows.iter().map(stat_json_row).collect();
let text = if opts.no_stream {
serde_json::to_string_pretty(&json)
} else {
serde_json::to_string(&json)
};
println!("{}", text.unwrap_or_default());
return;
}
crate::ui::print_bold_header(HEADER);
let colour = crate::ui::stdout_colored();
for s in &rows {
println!("{}", format_row_with(s, opts.no_trunc, colour));
}
println!();
}
#[cfg(test)]
#[path = "stats_tests.rs"]
mod tests;