Skip to main content

isb_core/
metrics.rs

1//! Live numbers for dashboards: the host's CPU, memory, storage, disk I/O
2//! and network, and every instance's status, address, CPU, memory, disk and
3//! network, each with a short history for sparklines, plus an hour of the
4//! host's own numbers for the monitor.
5//!
6//! One `GET /1.0/instances?recursion=2` per sample. CPU is a rate, so it needs
7//! two samples: the first one reports none. Instance CPU is a percentage of
8//! one core, as `docker stats` shows it (four busy cores read 400%).
9
10use std::collections::{BTreeMap, VecDeque};
11use std::time::{Duration, Instant};
12
13use serde::Serialize;
14use serde_json::Value;
15
16use crate::client::Client;
17use crate::error::Result;
18
19mod sys;
20
21/// Samples of history kept per series.
22pub const HISTORY: usize = 40;
23
24/// Disk counters come from `/1.0/metrics`, a heavier call: every this many
25/// samples.
26pub const DISK_EVERY: u32 = 5;
27
28/// Host points kept: an hour at the controller's two-second sample.
29pub const HOST_POINTS: usize = 1800;
30
31#[derive(Debug, Clone, Serialize, Default, PartialEq)]
32pub struct HostSample {
33    pub hostname: String,
34    pub cpus: u32,
35    /// Busy share of all CPUs, 0-100.
36    pub cpu_pct: Option<f32>,
37    pub cpu_history: Vec<f32>,
38    pub mem_used: u64,
39    pub mem_total: u64,
40    /// incus storage pools, used and total bytes. Pools on one filesystem
41    /// (several `dir` pools, say) are counted once.
42    pub disk_used: u64,
43    pub disk_total: u64,
44    pub load1: f32,
45    pub load5: Option<f32>,
46    pub load15: Option<f32>,
47    /// Busy share of each CPU, 0-100 (Linux only).
48    pub cpu_cores: Vec<f32>,
49    pub uptime_secs: Option<u64>,
50    pub swap_used: Option<u64>,
51    pub swap_total: Option<u64>,
52    /// Each storage pool, with its own used and total bytes.
53    pub pools: Vec<PoolSample>,
54    /// Bytes per second over the host's whole disks (not partitions, device
55    /// mapper or zvols, which would count the same bytes twice).
56    pub disk_read_rate: Option<f64>,
57    pub disk_write_rate: Option<f64>,
58    /// Bytes per second over [`HostSample::interfaces`].
59    pub net_rx_rate: Option<f64>,
60    pub net_tx_rate: Option<f64>,
61    /// Network interfaces but loopback and instances' own ends (veth, tap).
62    pub interfaces: Vec<IfaceSample>,
63}
64
65#[derive(Debug, Clone, Serialize, Default, PartialEq)]
66pub struct PoolSample {
67    pub name: String,
68    pub driver: String,
69    pub used: u64,
70    pub total: u64,
71}
72
73#[derive(Debug, Clone, Serialize, Default, PartialEq)]
74pub struct IfaceSample {
75    pub name: String,
76    /// Global addresses, IPv4 and IPv6.
77    pub addresses: Vec<String>,
78    pub up: bool,
79    pub rx_bytes: u64,
80    pub tx_bytes: u64,
81    pub rx_rate: Option<f64>,
82    pub tx_rate: Option<f64>,
83}
84
85/// One host sample, kept [`HOST_POINTS`] deep for the monitor's charts.
86#[derive(Debug, Clone, Copy, Serialize, Default, PartialEq)]
87pub struct HostPoint {
88    /// Unix seconds.
89    pub t: u64,
90    pub cpu: Option<f32>,
91    pub mem_used: Option<u64>,
92    pub net_rx: Option<f64>,
93    pub net_tx: Option<f64>,
94}
95
96#[derive(Debug, Clone, Serialize, Default, PartialEq)]
97pub struct InstanceSample {
98    pub name: String,
99    pub status: String,
100    /// `container`, `virtual-machine`, or `oci` for an application container.
101    pub kind: String,
102    pub ip: Option<String>,
103    /// Percent of one core.
104    pub cpu_pct: Option<f32>,
105    pub cpu_history: Vec<f32>,
106    pub mem_bytes: Option<u64>,
107    /// The root disk's usage, where the storage driver reports it (ZFS,
108    /// Btrfs, LVM; not `dir`).
109    pub disk_bytes: Option<u64>,
110    /// `user.*` config keys without the prefix, isb's own included.
111    pub labels: BTreeMap<String, String>,
112    pub image: String,
113    pub created_at: String,
114    /// The incus project; with isb's orgs, `isb-<org>` (or `default`).
115    pub project: String,
116    /// Bytes received and sent on every interface but loopback, since the
117    /// instance started (counters: the history turns them into rates).
118    #[serde(skip_serializing_if = "Option::is_none")]
119    pub net_rx_bytes: Option<u64>,
120    #[serde(skip_serializing_if = "Option::is_none")]
121    pub net_tx_bytes: Option<u64>,
122    /// Bytes read from and written to disks (counters), from incus'
123    /// `/1.0/metrics`, taken every [`DISK_EVERY`]th sample only.
124    #[serde(skip_serializing_if = "Option::is_none")]
125    pub disk_read_bytes: Option<u64>,
126    #[serde(skip_serializing_if = "Option::is_none")]
127    pub disk_write_bytes: Option<u64>,
128    /// The counters above as bytes per second, from the previous sample.
129    #[serde(skip_serializing_if = "Option::is_none")]
130    pub net_rx_rate: Option<f64>,
131    #[serde(skip_serializing_if = "Option::is_none")]
132    pub net_tx_rate: Option<f64>,
133    #[serde(skip_serializing_if = "Option::is_none")]
134    pub disk_read_rate: Option<f64>,
135    #[serde(skip_serializing_if = "Option::is_none")]
136    pub disk_write_rate: Option<f64>,
137}
138
139impl InstanceSample {
140    pub fn running(&self) -> bool {
141        self.status.eq_ignore_ascii_case("running")
142    }
143    /// The stack this instance is a replica of.
144    pub fn stack(&self) -> Option<&str> {
145        self.labels.get("isb.stack").map(String::as_str)
146    }
147}
148
149/// Turns successive snapshots into rates and histories.
150#[derive(Debug, Default)]
151pub struct Sampler {
152    cpu: BTreeMap<String, (u64, Instant)>,
153    hist: BTreeMap<String, VecDeque<f32>>,
154    host_cpu: Option<(u64, u64)>,
155    host_hist: VecDeque<f32>,
156    /// Pool usage moves slowly and costs a request per pool: refreshed every
157    /// [`POOLS_EVERY`].
158    pools: Option<(Instant, Vec<PoolSample>)>,
159    /// Samples taken, for [`DISK_EVERY`].
160    n: u32,
161    /// Per-core (busy, total) ticks.
162    host_cores: Vec<(u64, u64)>,
163    /// Per interface (rx, tx) bytes, and when they were read.
164    host_net: Option<(Instant, IfaceCounters)>,
165    host_disk: Option<(Instant, (u64, u64))>,
166    /// Per instance key: the last counters, for rates.
167    io: BTreeMap<String, Io>,
168    points: VecDeque<HostPoint>,
169}
170
171/// Per interface: (rx, tx) bytes.
172type IfaceCounters = BTreeMap<String, (u64, u64)>;
173
174#[derive(Debug, Default)]
175struct Io {
176    net: Option<(Instant, u64, u64)>,
177    disk: Option<(Instant, u64, u64)>,
178    /// Disk counters are read every [`DISK_EVERY`]th sample: the rate holds
179    /// between reads.
180    disk_rate: (Option<f64>, Option<f64>),
181}
182
183/// Bytes per second between two counter readings; none across a reset.
184fn per_sec(
185    prev: Option<(Instant, u64, u64)>,
186    now: Instant,
187    a: u64,
188    b: u64,
189) -> (Option<f64>, Option<f64>) {
190    match prev {
191        Some((at, pa, pb)) if a >= pa && b >= pb => {
192            let dt = now.duration_since(at).as_secs_f64();
193            if dt > 0.0 {
194                (Some((a - pa) as f64 / dt), Some((b - pb) as f64 / dt))
195            } else {
196                (None, None)
197            }
198        }
199        _ => (None, None),
200    }
201}
202
203/// How often storage pool usage is re-read.
204const POOLS_EVERY: Duration = Duration::from_secs(30);
205
206fn push(h: &mut VecDeque<f32>, v: f32) {
207    if h.len() == HISTORY {
208        h.pop_front();
209    }
210    h.push_back(v);
211}
212
213impl Sampler {
214    pub fn new() -> Self {
215        Self::default()
216    }
217
218    /// Take one sample of the host and of every instance in the client's
219    /// project.
220    pub fn sample(&mut self, client: &Client) -> Result<(HostSample, Vec<InstanceSample>)> {
221        // Every project, so one sample covers every org.
222        let v = client.get("/1.0/instances?recursion=2&all-projects=true")?;
223        let now = Instant::now();
224        let mut out = Vec::new();
225        for i in v.as_array().into_iter().flatten() {
226            out.push(self.instance(i, now));
227        }
228        if self.n % DISK_EVERY == 0 {
229            // Disk I/O is optional: the rest of the sample stands without it.
230            if let Ok(text) = client.get_raw("/1.0/metrics") {
231                let io = parse_disk_metrics(&String::from_utf8_lossy(&text));
232                for i in out.iter_mut().filter(|i| i.running()) {
233                    let (r, w) = io
234                        .get(&(i.project.clone(), i.name.clone()))
235                        .copied()
236                        .unwrap_or((0, 0));
237                    i.disk_read_bytes = Some(r);
238                    i.disk_write_bytes = Some(w);
239                    let e = self
240                        .io
241                        .entry(format!("{}/{}", i.project, i.name))
242                        .or_default();
243                    e.disk_rate = per_sec(e.disk, now, r, w);
244                    e.disk = Some((now, r, w));
245                }
246            }
247        }
248        for i in out.iter_mut().filter(|i| i.running()) {
249            if let Some(e) = self.io.get(&format!("{}/{}", i.project, i.name)) {
250                (i.disk_read_rate, i.disk_write_rate) = e.disk_rate;
251            }
252        }
253        self.n = self.n.wrapping_add(1);
254        out.sort_by(|a, b| (&a.project, &a.name).cmp(&(&b.project, &b.name)));
255        let keys: std::collections::BTreeSet<String> = out
256            .iter()
257            .map(|i| format!("{}/{}", i.project, i.name))
258            .collect();
259        self.cpu.retain(|k, _| keys.contains(k));
260        self.hist.retain(|k, _| keys.contains(k));
261        self.io.retain(|k, _| keys.contains(k));
262        if self
263            .pools
264            .as_ref()
265            .is_none_or(|(at, _)| now.duration_since(*at) >= POOLS_EVERY)
266        {
267            // Storage is a nicety: a pool that cannot be read leaves the
268            // last numbers in place rather than failing the sample.
269            if let Ok(p) = pools(client) {
270                self.pools = Some((now, p));
271            }
272        }
273        Ok((self.host(now), out))
274    }
275
276    /// The host's last hour, oldest first.
277    pub fn host_points(&self) -> Vec<HostPoint> {
278        self.points.iter().copied().collect()
279    }
280
281    fn instance(&mut self, i: &Value, now: Instant) -> InstanceSample {
282        let project = i["project"].as_str().unwrap_or("default").to_string();
283        // Names are unique per project only: key by both.
284        let name = i["name"].as_str().unwrap_or_default().to_string();
285        let key = format!("{project}/{name}");
286        let config = i["config"].as_object();
287        let cfg = |k: &str| {
288            config
289                .and_then(|c| c.get(k))
290                .and_then(Value::as_str)
291                .unwrap_or_default()
292                .to_string()
293        };
294        let labels = config
295            .map(|c| {
296                c.iter()
297                    .filter_map(|(k, v)| {
298                        k.strip_prefix("user.")
299                            .filter(|k| *k != "isb.create-token")
300                            .map(|k| (k.to_string(), v.as_str().unwrap_or_default().to_string()))
301                    })
302                    .collect()
303            })
304            .unwrap_or_default();
305        let kind = if cfg("volatile.container.oci") == "true" {
306            "oci".to_string()
307        } else {
308            i["type"].as_str().unwrap_or_default().to_string()
309        };
310        let status = i["status"].as_str().unwrap_or_default().to_string();
311        let state = &i["state"];
312        let running = status.eq_ignore_ascii_case("running");
313        let usage = state["cpu"]["usage"].as_u64().filter(|_| running);
314        let cpu_pct = match (usage, self.cpu.get(&key)) {
315            (Some(u), Some((prev, at))) if u >= *prev => {
316                let wall = now.duration_since(*at).as_nanos() as f64;
317                (wall > 0.0).then(|| ((u - prev) as f64 / wall * 100.0) as f32)
318            }
319            _ => None,
320        };
321        match usage {
322            Some(u) => {
323                self.cpu.insert(key.clone(), (u, now));
324            }
325            None => {
326                self.cpu.remove(&key);
327            }
328        }
329        let h = self.hist.entry(key.clone()).or_default();
330        if running {
331            push(h, cpu_pct.unwrap_or(0.0));
332        } else {
333            h.clear();
334        }
335        let image = {
336            let d = cfg("image.description");
337            if d.is_empty() { cfg("image.id") } else { d }
338        };
339        let (rx, tx) = net_counters(state);
340        let (net_rx_rate, net_tx_rate) = match (rx, tx) {
341            (Some(r), Some(t)) if running => {
342                let e = self.io.entry(key.clone()).or_default();
343                let v = per_sec(e.net, now, r, t);
344                e.net = Some((now, r, t));
345                v
346            }
347            _ => {
348                self.io.remove(&key);
349                (None, None)
350            }
351        };
352        InstanceSample {
353            net_rx_bytes: rx.filter(|_| running),
354            net_tx_bytes: tx.filter(|_| running),
355            disk_read_bytes: None,
356            disk_write_bytes: None,
357            net_rx_rate,
358            net_tx_rate,
359            disk_read_rate: None,
360            disk_write_rate: None,
361            ip: running.then(|| first_ip(state)).flatten(),
362            cpu_pct,
363            cpu_history: h.iter().copied().collect(),
364            mem_bytes: state["memory"]["usage"].as_u64().filter(|_| running),
365            // -1 (or 0) when the driver cannot tell.
366            disk_bytes: state["disk"]["root"]["usage"]
367                .as_i64()
368                .filter(|u| *u > 0)
369                .map(|u| u as u64),
370            name,
371            status,
372            kind,
373            labels,
374            image,
375            created_at: i["created_at"].as_str().unwrap_or_default().to_string(),
376            project,
377        }
378    }
379
380    fn host(&mut self, now: Instant) -> HostSample {
381        let mut h = HostSample {
382            hostname: sys::hostname(),
383            cpus: std::thread::available_parallelism()
384                .map(|n| n.get() as u32)
385                .unwrap_or(1),
386            ..Default::default()
387        };
388        if let Some((busy, total)) = sys::cpu_ticks() {
389            if let Some((pb, pt)) = self.host_cpu {
390                if total > pt {
391                    let pct = (busy.saturating_sub(pb)) as f32 / (total - pt) as f32 * 100.0;
392                    h.cpu_pct = Some(pct);
393                    push(&mut self.host_hist, pct);
394                }
395            }
396            self.host_cpu = Some((busy, total));
397        }
398        h.cpu_history = self.host_hist.iter().copied().collect();
399        let cores = sys::cpu_cores();
400        if cores.len() == self.host_cores.len() {
401            h.cpu_cores = cores
402                .iter()
403                .zip(&self.host_cores)
404                .map(|((b, t), (pb, pt))| {
405                    if t > pt {
406                        b.saturating_sub(*pb) as f32 / (t - pt) as f32 * 100.0
407                    } else {
408                        0.0
409                    }
410                })
411                .collect();
412        }
413        self.host_cores = cores;
414        if let Some((total, avail)) = sys::memory() {
415            h.mem_total = total;
416            h.mem_used = total.saturating_sub(avail);
417        }
418        if let Some((total, free)) = sys::swap() {
419            h.swap_total = Some(total);
420            h.swap_used = Some(total.saturating_sub(free));
421        }
422        if let Some((_, pools)) = &self.pools {
423            (h.disk_used, h.disk_total) =
424                sum_pools(pools.iter().map(|p| (p.used, p.total)).collect());
425            h.pools = pools.clone();
426        }
427        if let Some([l1, l5, l15]) = sys::loadavg() {
428            h.load1 = l1;
429            h.load5 = Some(l5);
430            h.load15 = Some(l15);
431        }
432        h.uptime_secs = sys::uptime();
433        if let Some((r, w)) = sys::disk_io() {
434            let prev = self.host_disk.map(|(at, (a, b))| (at, a, b));
435            (h.disk_read_rate, h.disk_write_rate) = per_sec(prev, now, r, w);
436            self.host_disk = Some((now, (r, w)));
437        }
438        self.interfaces(&mut h, now);
439        push_point(
440            &mut self.points,
441            HostPoint {
442                t: unix_secs(),
443                cpu: h.cpu_pct,
444                mem_used: (h.mem_total > 0).then_some(h.mem_used),
445                net_rx: h.net_rx_rate,
446                net_tx: h.net_tx_rate,
447            },
448        );
449        h
450    }
451
452    fn interfaces(&mut self, h: &mut HostSample, now: Instant) {
453        let counters: BTreeMap<String, (u64, u64)> = sys::net_dev()
454            .into_iter()
455            .filter(|(name, _, _)| shown_iface(name))
456            .map(|(name, rx, tx)| (name, (rx, tx)))
457            .collect();
458        let addrs = sys::addresses();
459        let (mut rx_sum, mut tx_sum, mut rated) = (0.0, 0.0, false);
460        for (name, (rx, tx)) in &counters {
461            let prev = self
462                .host_net
463                .as_ref()
464                .and_then(|(at, m)| m.get(name).map(|(a, b)| (*at, *a, *b)));
465            let (rx_rate, tx_rate) = per_sec(prev, now, *rx, *tx);
466            if let (Some(r), Some(t)) = (rx_rate, tx_rate) {
467                rx_sum += r;
468                tx_sum += t;
469                rated = true;
470            }
471            h.interfaces.push(IfaceSample {
472                name: name.clone(),
473                addresses: addrs.get(name).cloned().unwrap_or_default(),
474                up: sys::iface_up(name),
475                rx_bytes: *rx,
476                tx_bytes: *tx,
477                rx_rate,
478                tx_rate,
479            });
480        }
481        if rated {
482            h.net_rx_rate = Some(rx_sum);
483            h.net_tx_rate = Some(tx_sum);
484        }
485        self.host_net = Some((now, counters));
486    }
487}
488
489fn push_point(points: &mut VecDeque<HostPoint>, p: HostPoint) {
490    if points.len() == HOST_POINTS {
491        points.pop_front();
492    }
493    points.push_back(p);
494}
495
496fn unix_secs() -> u64 {
497    std::time::SystemTime::now()
498        .duration_since(std::time::UNIX_EPOCH)
499        .map(|d| d.as_secs())
500        .unwrap_or(0)
501}
502
503/// Interfaces the host view lists: loopback and the host ends of instances'
504/// NICs (one per instance, which would drown the rest) are left out.
505fn shown_iface(name: &str) -> bool {
506    name != "lo"
507        && !name.starts_with("veth")
508        && !name.starts_with("tap")
509        && !name.starts_with("macvtap")
510}
511
512/// Points within the last `range` seconds of `now`, averaged into buckets so
513/// at most `max` come back. Returns the bucket step in seconds.
514pub fn downsample(points: &[HostPoint], range: u64, now: u64, max: usize) -> (u64, Vec<HostPoint>) {
515    let max = max.max(1) as u64;
516    let step = range.div_ceil(max).max(2);
517    let from = now.saturating_sub(range);
518    let mut out: Vec<HostPoint> = Vec::new();
519    let mut acc: Option<(u64, [f64; 4], [u32; 4])> = None;
520    let flush = |out: &mut Vec<HostPoint>, a: (u64, [f64; 4], [u32; 4])| {
521        let avg = |i: usize| (a.2[i] > 0).then(|| a.1[i] / f64::from(a.2[i]));
522        out.push(HostPoint {
523            t: a.0,
524            cpu: avg(0).map(|v| v as f32),
525            mem_used: avg(1).map(|v| v as u64),
526            net_rx: avg(2),
527            net_tx: avg(3),
528        });
529    };
530    for p in points.iter().filter(|p| p.t > from && p.t <= now) {
531        let bucket = p.t / step * step;
532        if let Some(a) = acc.filter(|a| a.0 != bucket) {
533            flush(&mut out, a);
534            acc = None;
535        }
536        let a = acc.get_or_insert((bucket, [0.0; 4], [0; 4]));
537        for (i, v) in [
538            p.cpu.map(f64::from),
539            p.mem_used.map(|m| m as f64),
540            p.net_rx,
541            p.net_tx,
542        ]
543        .into_iter()
544        .enumerate()
545        {
546            if let Some(v) = v {
547                a.1[i] += v;
548                a.2[i] += 1;
549            }
550        }
551    }
552    if let Some(a) = acc {
553        flush(&mut out, a);
554    }
555    (step, out)
556}
557
558/// Each storage pool's used and total bytes.
559fn pools(client: &Client) -> Result<Vec<PoolSample>> {
560    let list = client.get("/1.0/storage-pools?recursion=1")?;
561    let mut out = Vec::new();
562    for p in list.as_array().into_iter().flatten() {
563        let name = p["name"].as_str().unwrap_or_default();
564        let r = client.get(&format!("/1.0/storage-pools/{name}/resources"))?;
565        out.push(PoolSample {
566            name: name.to_string(),
567            driver: p["driver"].as_str().unwrap_or_default().to_string(),
568            used: r["space"]["used"].as_u64().unwrap_or(0),
569            total: r["space"]["total"].as_u64().unwrap_or(0),
570        });
571    }
572    Ok(out)
573}
574
575/// Pools that share a filesystem report the same total: count it once (with
576/// the larger used, as their reads are a moment apart).
577fn sum_pools(spaces: Vec<(u64, u64)>) -> (u64, u64) {
578    let mut by_total: BTreeMap<u64, u64> = BTreeMap::new();
579    for (used, total) in spaces.into_iter().filter(|(_, t)| *t > 0) {
580        let u = by_total.entry(total).or_default();
581        *u = (*u).max(used);
582    }
583    (by_total.values().sum(), by_total.keys().sum())
584}
585
586/// Bytes received and sent, summed over every interface but loopback.
587fn net_counters(state: &Value) -> (Option<u64>, Option<u64>) {
588    let Some(ifs) = state["network"].as_object() else {
589        return (None, None);
590    };
591    let (mut rx, mut tx, mut any) = (0u64, 0u64, false);
592    for (name, n) in ifs {
593        if name == "lo" {
594            continue;
595        }
596        let c = &n["counters"];
597        if let (Some(r), Some(t)) = (c["bytes_received"].as_u64(), c["bytes_sent"].as_u64()) {
598            rx = rx.saturating_add(r);
599            tx = tx.saturating_add(t);
600            any = true;
601        }
602    }
603    if any {
604        (Some(rx), Some(tx))
605    } else {
606        (None, None)
607    }
608}
609
610/// Per `(project, instance)`: disk bytes read and written, summed over its
611/// devices, from incus' OpenMetrics text.
612pub fn parse_disk_metrics(text: &str) -> BTreeMap<(String, String), (u64, u64)> {
613    let mut out: BTreeMap<(String, String), (u64, u64)> = BTreeMap::new();
614    for line in text.lines() {
615        let (read, rest) = if let Some(r) = line.strip_prefix("incus_disk_read_bytes_total{") {
616            (true, r)
617        } else if let Some(r) = line.strip_prefix("incus_disk_written_bytes_total{") {
618            (false, r)
619        } else {
620            continue;
621        };
622        let Some((labels, value)) = rest.rsplit_once('}') else {
623            continue;
624        };
625        let label = |k: &str| {
626            labels.split(',').find_map(|kv| {
627                let (a, b) = kv.split_once('=')?;
628                (a.trim() == k).then(|| b.trim().trim_matches('"').to_string())
629            })
630        };
631        let (Some(name), Some(project)) = (label("name"), label("project")) else {
632            continue;
633        };
634        let Ok(v) = value.trim().parse::<f64>() else {
635            continue;
636        };
637        let e = out.entry((project, name)).or_default();
638        let v = v.max(0.0) as u64;
639        if read {
640            e.0 = e.0.saturating_add(v);
641        } else {
642            e.1 = e.1.saturating_add(v);
643        }
644    }
645    out
646}
647
648/// First global address, IPv4 preferred, on any interface but loopback.
649/// The instance's address: `eth0`'s first (its NIC on the org's bridge),
650/// then any other interface's. A bridge made inside the instance, such as
651/// Docker's `docker0`, sorts before `eth0` but reaches nothing.
652fn first_ip(state: &Value) -> Option<String> {
653    let mut v6 = None;
654    let nets = state["network"].as_object()?;
655    let eth0 = nets.get_key_value("eth0");
656    for (ifname, n) in eth0
657        .into_iter()
658        .chain(nets.iter().filter(|(k, _)| *k != "eth0"))
659    {
660        if ifname == "lo" {
661            continue;
662        }
663        for a in n["addresses"].as_array().into_iter().flatten() {
664            if a["scope"] != "global" {
665                continue;
666            }
667            match a["family"].as_str() {
668                Some("inet") => return a["address"].as_str().map(String::from),
669                Some("inet6") if v6.is_none() => v6 = a["address"].as_str().map(String::from),
670                _ => {}
671            }
672        }
673    }
674    v6
675}
676
677#[cfg(test)]
678mod tests {
679    use super::*;
680
681    #[test]
682    fn the_address_is_eth0s_even_with_docker_inside() {
683        let addr = |a: &str| serde_json::json!({"addresses": [{"family": "inet", "scope": "global", "address": a}]});
684        let state = serde_json::json!({"network": {
685            "docker0": addr("172.17.0.1"),
686            "eth0": addr("10.81.189.151"),
687            "lo": addr("127.0.0.1"),
688        }});
689        assert_eq!(first_ip(&state).as_deref(), Some("10.81.189.151"));
690        let other = serde_json::json!({"network": {"enp5s0": addr("10.0.0.9")}});
691        assert_eq!(first_ip(&other).as_deref(), Some("10.0.0.9"));
692    }
693    use serde_json::json;
694
695    #[test]
696    fn shown_interfaces() {
697        assert!(shown_iface("eth0") && shown_iface("incusbr0") && shown_iface("tailscale0"));
698        assert!(!shown_iface("lo") && !shown_iface("veth1a2b") && !shown_iface("tap0"));
699    }
700
701    #[test]
702    fn host_points_downsample() {
703        let p = |t: u64, cpu: f32| HostPoint {
704            t,
705            cpu: Some(cpu),
706            mem_used: Some(100),
707            net_rx: None,
708            net_tx: Some(10.0),
709        };
710        let pts: Vec<HostPoint> = (0..60).map(|i| p(1000 + i * 2, i as f32)).collect();
711        // The last 60 s at most 10 points: 6 s buckets of three samples each.
712        let (step, out) = downsample(&pts, 60, 1118, 10);
713        assert_eq!(step, 6);
714        assert!(out.len() <= 11, "{}", out.len());
715        assert!(out.iter().all(|x| x.t > 1118 - 60 - step));
716        assert_eq!(out.last().unwrap().net_rx, None);
717        assert_eq!(out.last().unwrap().net_tx, Some(10.0));
718        // Short ranges keep the samples themselves.
719        let (step, out) = downsample(&pts, 10, 1118, 300);
720        assert_eq!(step, 2);
721        assert_eq!(out.len(), 5);
722        assert_eq!(out.last().unwrap().cpu, Some(59.0));
723    }
724
725    #[test]
726    fn instance_rates_from_counters() {
727        let mut s = Sampler::new();
728        let inst = |rx: u64| {
729            json!({"name": "a", "status": "Running", "type": "container",
730                   "state": {"network": {"eth0": {"counters": {"bytes_received": rx, "bytes_sent": 0}}}}})
731        };
732        let t0 = Instant::now();
733        assert_eq!(s.instance(&inst(1000), t0).net_rx_rate, None);
734        let b = s.instance(&inst(3000), t0 + Duration::from_secs(2));
735        assert_eq!((b.net_rx_rate, b.net_tx_rate), (Some(1000.0), Some(0.0)));
736        // A counter that went down (a restart) gives no rate.
737        let c = s.instance(&inst(10), t0 + Duration::from_secs(4));
738        assert_eq!(c.net_rx_rate, None);
739    }
740
741    #[test]
742    fn disk_metrics() {
743        let t = "# HELP x\n\
744                 incus_disk_read_bytes_total{device=\"vda\",name=\"web\",project=\"isb-acme\",type=\"container\"} 4096\n\
745                 incus_disk_read_bytes_total{device=\"vdb\",name=\"web\",project=\"isb-acme\",type=\"container\"} 1.5e+03\n\
746                 incus_disk_written_bytes_total{device=\"vda\",name=\"web\",project=\"isb-acme\",type=\"container\"} 10\n\
747                 incus_cpu_seconds_total{cpu=\"0\",mode=\"user\",name=\"web\",project=\"isb-acme\",type=\"container\"} 1\n";
748        let m = parse_disk_metrics(t);
749        assert_eq!(m[&("isb-acme".to_string(), "web".to_string())], (5596, 10));
750        assert_eq!(m.len(), 1);
751    }
752
753    #[test]
754    fn instance_rates_and_kinds() {
755        let mut s = Sampler::new();
756        let inst = |usage: u64| {
757            json!({
758                "name": "a", "status": "Running", "type": "container",
759                "config": {"volatile.container.oci": "true", "user.isb.stack": "app", "user.isb.create-token": "x"},
760                "state": {"cpu": {"usage": usage}, "memory": {"usage": 1024},
761                          "network": {"lo": {"addresses": [{"family": "inet", "address": "127.0.0.1", "scope": "local"}]},
762                                      "eth0": {"counters": {"bytes_received": 100, "bytes_sent": 50}, "addresses": [{"family": "inet6", "address": "fd42::1", "scope": "global"},
763                                                             {"family": "inet", "address": "10.0.0.2", "scope": "global"}]}}}
764            })
765        };
766        let t0 = Instant::now();
767        let a = s.instance(&inst(1_000_000_000), t0);
768        assert_eq!(a.cpu_pct, None);
769        assert_eq!(a.kind, "oci");
770        assert_eq!(a.ip.as_deref(), Some("10.0.0.2"));
771        assert_eq!(a.stack(), Some("app"));
772        assert_eq!((a.net_rx_bytes, a.net_tx_bytes), (Some(100), Some(50)));
773        assert!(!a.labels.contains_key("isb.create-token"));
774        let b = s.instance(&inst(1_500_000_000), t0 + std::time::Duration::from_secs(1));
775        assert!((b.cpu_pct.unwrap() - 50.0).abs() < 0.1, "{b:?}");
776        assert_eq!(b.cpu_history.len(), 2);
777        assert_eq!(b.disk_bytes, None, "no disk state: unknown");
778    }
779
780    #[test]
781    fn disk_usage() {
782        let mut s = Sampler::new();
783        let inst = |usage: i64| {
784            json!({"name": "a", "status": "Stopped", "type": "container",
785                   "state": {"disk": {"root": {"usage": usage, "total": 0}}}})
786        };
787        // Reported for stopped instances too; `dir` pools say -1.
788        assert_eq!(
789            s.instance(&inst(5 << 20), Instant::now()).disk_bytes,
790            Some(5 << 20)
791        );
792        assert_eq!(s.instance(&inst(-1), Instant::now()).disk_bytes, None);
793        // Two `dir` pools on one filesystem count once; a second disk adds.
794        assert_eq!(sum_pools(vec![(60, 100), (61, 100)]), (61, 100));
795        assert_eq!(sum_pools(vec![(60, 100), (5, 50), (0, 0)]), (65, 150));
796    }
797}