1use std::collections::{BTreeMap, VecDeque};
10use std::time::{Duration, Instant};
11
12use serde::Serialize;
13use serde_json::Value;
14
15use crate::client::Client;
16use crate::error::Result;
17
18pub const HISTORY: usize = 40;
20
21pub const DISK_EVERY: u32 = 5;
24
25#[derive(Debug, Clone, Serialize, Default, PartialEq)]
26pub struct HostSample {
27 pub hostname: String,
28 pub cpus: u32,
29 pub cpu_pct: Option<f32>,
31 pub cpu_history: Vec<f32>,
32 pub mem_used: u64,
33 pub mem_total: u64,
34 pub disk_used: u64,
37 pub disk_total: u64,
38 pub load1: f32,
39}
40
41#[derive(Debug, Clone, Serialize, Default, PartialEq)]
42pub struct InstanceSample {
43 pub name: String,
44 pub status: String,
45 pub kind: String,
47 pub ip: Option<String>,
48 pub cpu_pct: Option<f32>,
50 pub cpu_history: Vec<f32>,
51 pub mem_bytes: Option<u64>,
52 pub disk_bytes: Option<u64>,
55 pub labels: BTreeMap<String, String>,
57 pub image: String,
58 pub created_at: String,
59 pub project: String,
61 #[serde(skip_serializing_if = "Option::is_none")]
64 pub net_rx_bytes: Option<u64>,
65 #[serde(skip_serializing_if = "Option::is_none")]
66 pub net_tx_bytes: Option<u64>,
67 #[serde(skip_serializing_if = "Option::is_none")]
70 pub disk_read_bytes: Option<u64>,
71 #[serde(skip_serializing_if = "Option::is_none")]
72 pub disk_write_bytes: Option<u64>,
73}
74
75impl InstanceSample {
76 pub fn running(&self) -> bool {
77 self.status.eq_ignore_ascii_case("running")
78 }
79 pub fn stack(&self) -> Option<&str> {
81 self.labels.get("isb.stack").map(String::as_str)
82 }
83}
84
85#[derive(Debug, Default)]
87pub struct Sampler {
88 cpu: BTreeMap<String, (u64, Instant)>,
89 hist: BTreeMap<String, VecDeque<f32>>,
90 host_cpu: Option<(u64, u64)>,
91 host_hist: VecDeque<f32>,
92 pools: Option<(Instant, (u64, u64))>,
95 n: u32,
97}
98
99const POOLS_EVERY: Duration = Duration::from_secs(30);
101
102fn push(h: &mut VecDeque<f32>, v: f32) {
103 if h.len() == HISTORY {
104 h.pop_front();
105 }
106 h.push_back(v);
107}
108
109impl Sampler {
110 pub fn new() -> Self {
111 Self::default()
112 }
113
114 pub fn sample(&mut self, client: &Client) -> Result<(HostSample, Vec<InstanceSample>)> {
117 let v = client.get("/1.0/instances?recursion=2&all-projects=true")?;
119 let now = Instant::now();
120 let mut out = Vec::new();
121 for i in v.as_array().into_iter().flatten() {
122 out.push(self.instance(i, now));
123 }
124 if self.n % DISK_EVERY == 0 {
125 if let Ok(text) = client.get_raw("/1.0/metrics") {
127 let io = parse_disk_metrics(&String::from_utf8_lossy(&text));
128 for i in out.iter_mut().filter(|i| i.running()) {
129 let (r, w) = io
130 .get(&(i.project.clone(), i.name.clone()))
131 .copied()
132 .unwrap_or((0, 0));
133 i.disk_read_bytes = Some(r);
134 i.disk_write_bytes = Some(w);
135 }
136 }
137 }
138 self.n = self.n.wrapping_add(1);
139 out.sort_by(|a, b| (&a.project, &a.name).cmp(&(&b.project, &b.name)));
140 let keys: Vec<String> = out
141 .iter()
142 .map(|i| format!("{}/{}", i.project, i.name))
143 .collect();
144 self.cpu.retain(|k, _| keys.contains(k));
145 self.hist.retain(|k, _| keys.contains(k));
146 if self
147 .pools
148 .is_none_or(|(at, _)| now.duration_since(at) >= POOLS_EVERY)
149 {
150 if let Ok(p) = pools(client) {
153 self.pools = Some((now, p));
154 }
155 }
156 Ok((self.host(), out))
157 }
158
159 fn instance(&mut self, i: &Value, now: Instant) -> InstanceSample {
160 let project = i["project"].as_str().unwrap_or("default").to_string();
161 let name = i["name"].as_str().unwrap_or_default().to_string();
163 let key = format!("{project}/{name}");
164 let config = i["config"].as_object();
165 let cfg = |k: &str| {
166 config
167 .and_then(|c| c.get(k))
168 .and_then(Value::as_str)
169 .unwrap_or_default()
170 .to_string()
171 };
172 let labels = config
173 .map(|c| {
174 c.iter()
175 .filter_map(|(k, v)| {
176 k.strip_prefix("user.")
177 .filter(|k| *k != "isb.create-token")
178 .map(|k| (k.to_string(), v.as_str().unwrap_or_default().to_string()))
179 })
180 .collect()
181 })
182 .unwrap_or_default();
183 let kind = if cfg("volatile.container.oci") == "true" {
184 "oci".to_string()
185 } else {
186 i["type"].as_str().unwrap_or_default().to_string()
187 };
188 let status = i["status"].as_str().unwrap_or_default().to_string();
189 let state = &i["state"];
190 let running = status.eq_ignore_ascii_case("running");
191 let usage = state["cpu"]["usage"].as_u64().filter(|_| running);
192 let cpu_pct = match (usage, self.cpu.get(&key)) {
193 (Some(u), Some((prev, at))) if u >= *prev => {
194 let wall = now.duration_since(*at).as_nanos() as f64;
195 (wall > 0.0).then(|| ((u - prev) as f64 / wall * 100.0) as f32)
196 }
197 _ => None,
198 };
199 match usage {
200 Some(u) => {
201 self.cpu.insert(key.clone(), (u, now));
202 }
203 None => {
204 self.cpu.remove(&key);
205 }
206 }
207 let h = self.hist.entry(key.clone()).or_default();
208 if running {
209 push(h, cpu_pct.unwrap_or(0.0));
210 } else {
211 h.clear();
212 }
213 let image = {
214 let d = cfg("image.description");
215 if d.is_empty() { cfg("image.id") } else { d }
216 };
217 let (rx, tx) = net_counters(state);
218 InstanceSample {
219 net_rx_bytes: rx.filter(|_| running),
220 net_tx_bytes: tx.filter(|_| running),
221 disk_read_bytes: None,
222 disk_write_bytes: None,
223 ip: running.then(|| first_ip(state)).flatten(),
224 cpu_pct,
225 cpu_history: h.iter().copied().collect(),
226 mem_bytes: state["memory"]["usage"].as_u64().filter(|_| running),
227 disk_bytes: state["disk"]["root"]["usage"]
229 .as_i64()
230 .filter(|u| *u > 0)
231 .map(|u| u as u64),
232 name,
233 status,
234 kind,
235 labels,
236 image,
237 created_at: i["created_at"].as_str().unwrap_or_default().to_string(),
238 project,
239 }
240 }
241
242 fn host(&mut self) -> HostSample {
243 let mut h = HostSample {
244 hostname: sys::hostname(),
245 cpus: std::thread::available_parallelism()
246 .map(|n| n.get() as u32)
247 .unwrap_or(1),
248 ..Default::default()
249 };
250 if let Some((busy, total)) = sys::cpu_ticks() {
251 if let Some((pb, pt)) = self.host_cpu {
252 if total > pt {
253 let pct = (busy.saturating_sub(pb)) as f32 / (total - pt) as f32 * 100.0;
254 h.cpu_pct = Some(pct);
255 push(&mut self.host_hist, pct);
256 }
257 }
258 self.host_cpu = Some((busy, total));
259 }
260 h.cpu_history = self.host_hist.iter().copied().collect();
261 if let Some((total, avail)) = sys::memory() {
262 h.mem_total = total;
263 h.mem_used = total.saturating_sub(avail);
264 }
265 if let Some((_, (used, total))) = self.pools {
266 h.disk_used = used;
267 h.disk_total = total;
268 }
269 h.load1 = sys::load1().unwrap_or(0.0);
270 h
271 }
272}
273
274fn pools(client: &Client) -> Result<(u64, u64)> {
276 let names = client.get("/1.0/storage-pools")?;
277 let mut spaces = Vec::new();
278 for url in names
279 .as_array()
280 .into_iter()
281 .flatten()
282 .filter_map(Value::as_str)
283 {
284 let name = url.rsplit('/').next().unwrap_or_default();
285 let r = client.get(&format!("/1.0/storage-pools/{name}/resources"))?;
286 spaces.push((
287 r["space"]["used"].as_u64().unwrap_or(0),
288 r["space"]["total"].as_u64().unwrap_or(0),
289 ));
290 }
291 Ok(sum_pools(spaces))
292}
293
294fn sum_pools(spaces: Vec<(u64, u64)>) -> (u64, u64) {
297 let mut by_total: BTreeMap<u64, u64> = BTreeMap::new();
298 for (used, total) in spaces.into_iter().filter(|(_, t)| *t > 0) {
299 let u = by_total.entry(total).or_default();
300 *u = (*u).max(used);
301 }
302 (by_total.values().sum(), by_total.keys().sum())
303}
304
305#[cfg(target_os = "linux")]
307mod sys {
308 pub fn hostname() -> String {
309 std::fs::read_to_string("/proc/sys/kernel/hostname")
310 .map(|s| s.trim().to_string())
311 .unwrap_or_default()
312 }
313
314 pub fn cpu_ticks() -> Option<(u64, u64)> {
316 super::parse_proc_stat(&std::fs::read_to_string("/proc/stat").ok()?)
317 }
318
319 pub fn memory() -> Option<(u64, u64)> {
321 Some(super::parse_meminfo(
322 &std::fs::read_to_string("/proc/meminfo").ok()?,
323 ))
324 }
325
326 pub fn load1() -> Option<f32> {
327 std::fs::read_to_string("/proc/loadavg")
328 .ok()?
329 .split_whitespace()
330 .next()?
331 .parse()
332 .ok()
333 }
334}
335
336#[cfg(target_os = "macos")]
338mod sys {
339 use std::mem::{MaybeUninit, size_of};
340
341 pub fn hostname() -> String {
342 let mut buf = [0u8; 256];
343 if unsafe { libc::gethostname(buf.as_mut_ptr().cast(), buf.len()) } != 0 {
345 return String::new();
346 }
347 let end = buf.iter().position(|&b| b == 0).unwrap_or(buf.len());
348 String::from_utf8_lossy(&buf[..end]).into_owned()
349 }
350
351 #[allow(deprecated)]
354 fn host() -> libc::mach_port_t {
355 static HOST: std::sync::OnceLock<libc::mach_port_t> = std::sync::OnceLock::new();
356 *HOST.get_or_init(|| unsafe { libc::mach_host_self() })
358 }
359
360 pub fn cpu_ticks() -> Option<(u64, u64)> {
362 let mut info = MaybeUninit::<libc::host_cpu_load_info>::zeroed();
363 let mut count = libc::HOST_CPU_LOAD_INFO_COUNT;
364 let kr = unsafe {
366 libc::host_statistics(
367 host(),
368 libc::HOST_CPU_LOAD_INFO,
369 info.as_mut_ptr().cast(),
370 &mut count,
371 )
372 };
373 if kr != libc::KERN_SUCCESS {
374 return None;
375 }
376 let t = unsafe { info.assume_init() }.cpu_ticks.map(u64::from);
378 let idle = t[libc::CPU_STATE_IDLE as usize];
379 let total: u64 = t.iter().sum();
380 Some((total - idle, total))
381 }
382
383 pub fn memory() -> Option<(u64, u64)> {
386 let mut total = 0u64;
387 let mut len = size_of::<u64>();
388 let rc = unsafe {
390 libc::sysctlbyname(
391 c"hw.memsize".as_ptr(),
392 (&raw mut total).cast(),
393 &mut len,
394 std::ptr::null_mut(),
395 0,
396 )
397 };
398 if rc != 0 {
399 return None;
400 }
401 let mut vm = MaybeUninit::<libc::vm_statistics64>::zeroed();
402 let mut count = libc::HOST_VM_INFO64_COUNT;
403 let kr = unsafe {
406 libc::host_statistics64(
407 host(),
408 libc::HOST_VM_INFO64,
409 vm.as_mut_ptr().cast(),
410 &mut count,
411 )
412 };
413 if kr != libc::KERN_SUCCESS {
414 return Some((total, 0));
415 }
416 let vm = unsafe { vm.assume_init() };
418 let page = unsafe { libc::sysconf(libc::_SC_PAGESIZE) }.max(0) as u64;
420 let avail = (u64::from(vm.free_count) + u64::from(vm.inactive_count)) * page;
421 Some((total, avail.min(total)))
422 }
423
424 pub fn load1() -> Option<f32> {
425 let mut l = [0f64; 1];
426 (unsafe { libc::getloadavg(l.as_mut_ptr(), 1) } == 1).then_some(l[0] as f32)
428 }
429}
430
431fn net_counters(state: &Value) -> (Option<u64>, Option<u64>) {
433 let Some(ifs) = state["network"].as_object() else {
434 return (None, None);
435 };
436 let (mut rx, mut tx, mut any) = (0u64, 0u64, false);
437 for (name, n) in ifs {
438 if name == "lo" {
439 continue;
440 }
441 let c = &n["counters"];
442 if let (Some(r), Some(t)) = (c["bytes_received"].as_u64(), c["bytes_sent"].as_u64()) {
443 rx = rx.saturating_add(r);
444 tx = tx.saturating_add(t);
445 any = true;
446 }
447 }
448 if any {
449 (Some(rx), Some(tx))
450 } else {
451 (None, None)
452 }
453}
454
455pub fn parse_disk_metrics(text: &str) -> BTreeMap<(String, String), (u64, u64)> {
458 let mut out: BTreeMap<(String, String), (u64, u64)> = BTreeMap::new();
459 for line in text.lines() {
460 let (read, rest) = if let Some(r) = line.strip_prefix("incus_disk_read_bytes_total{") {
461 (true, r)
462 } else if let Some(r) = line.strip_prefix("incus_disk_written_bytes_total{") {
463 (false, r)
464 } else {
465 continue;
466 };
467 let Some((labels, value)) = rest.rsplit_once('}') else {
468 continue;
469 };
470 let label = |k: &str| {
471 labels.split(',').find_map(|kv| {
472 let (a, b) = kv.split_once('=')?;
473 (a.trim() == k).then(|| b.trim().trim_matches('"').to_string())
474 })
475 };
476 let (Some(name), Some(project)) = (label("name"), label("project")) else {
477 continue;
478 };
479 let Ok(v) = value.trim().parse::<f64>() else {
480 continue;
481 };
482 let e = out.entry((project, name)).or_default();
483 let v = v.max(0.0) as u64;
484 if read {
485 e.0 = e.0.saturating_add(v);
486 } else {
487 e.1 = e.1.saturating_add(v);
488 }
489 }
490 out
491}
492
493fn first_ip(state: &Value) -> Option<String> {
498 let mut v6 = None;
499 let nets = state["network"].as_object()?;
500 let eth0 = nets.get_key_value("eth0");
501 for (ifname, n) in eth0
502 .into_iter()
503 .chain(nets.iter().filter(|(k, _)| *k != "eth0"))
504 {
505 if ifname == "lo" {
506 continue;
507 }
508 for a in n["addresses"].as_array().into_iter().flatten() {
509 if a["scope"] != "global" {
510 continue;
511 }
512 match a["family"].as_str() {
513 Some("inet") => return a["address"].as_str().map(String::from),
514 Some("inet6") if v6.is_none() => v6 = a["address"].as_str().map(String::from),
515 _ => {}
516 }
517 }
518 }
519 v6
520}
521
522#[cfg(target_os = "linux")]
524fn parse_proc_stat(s: &str) -> Option<(u64, u64)> {
525 let line = s.lines().find(|l| l.starts_with("cpu "))?;
526 let f: Vec<u64> = line
527 .split_whitespace()
528 .skip(1)
529 .filter_map(|x| x.parse().ok())
530 .collect();
531 if f.len() < 4 {
532 return None;
533 }
534 let total: u64 = f.iter().take(8).sum();
536 let idle = f[3] + f.get(4).copied().unwrap_or(0);
537 Some((total - idle, total))
538}
539
540#[cfg(target_os = "linux")]
542fn parse_meminfo(s: &str) -> (u64, u64) {
543 let get = |k: &str| {
544 s.lines()
545 .find(|l| l.starts_with(k))
546 .and_then(|l| l.split_whitespace().nth(1))
547 .and_then(|x| x.parse::<u64>().ok())
548 .unwrap_or(0)
549 * 1024
550 };
551 (get("MemTotal:"), get("MemAvailable:"))
552}
553
554#[cfg(test)]
555mod tests {
556 use super::*;
557
558 #[test]
559 fn the_address_is_eth0s_even_with_docker_inside() {
560 let addr = |a: &str| serde_json::json!({"addresses": [{"family": "inet", "scope": "global", "address": a}]});
561 let state = serde_json::json!({"network": {
562 "docker0": addr("172.17.0.1"),
563 "eth0": addr("10.81.189.151"),
564 "lo": addr("127.0.0.1"),
565 }});
566 assert_eq!(first_ip(&state).as_deref(), Some("10.81.189.151"));
567 let other = serde_json::json!({"network": {"enp5s0": addr("10.0.0.9")}});
568 assert_eq!(first_ip(&other).as_deref(), Some("10.0.0.9"));
569 }
570 use serde_json::json;
571
572 #[test]
573 fn host_counters() {
574 assert!(!sys::hostname().is_empty());
575 let (busy, total) = sys::cpu_ticks().unwrap();
576 assert!(total > 0 && busy <= total);
577 let (total, avail) = sys::memory().unwrap();
578 assert!(total > 0 && avail > 0 && avail <= total);
579 assert!(sys::load1().is_some());
580 }
581
582 #[cfg(target_os = "linux")]
583 #[test]
584 fn proc_parsers() {
585 assert_eq!(
586 parse_proc_stat("cpu 100 0 50 800 50 0 0 0 0 0\ncpu0 1 2 3 4\n"),
587 Some((150, 1000))
588 );
589 assert_eq!(
590 parse_meminfo("MemTotal: 1000 kB\nMemFree: 1 kB\nMemAvailable: 400 kB\n"),
591 (1024000, 409600)
592 );
593 }
594
595 #[test]
596 fn disk_metrics() {
597 let t = "# HELP x\n\
598 incus_disk_read_bytes_total{device=\"vda\",name=\"web\",project=\"isb-acme\",type=\"container\"} 4096\n\
599 incus_disk_read_bytes_total{device=\"vdb\",name=\"web\",project=\"isb-acme\",type=\"container\"} 1.5e+03\n\
600 incus_disk_written_bytes_total{device=\"vda\",name=\"web\",project=\"isb-acme\",type=\"container\"} 10\n\
601 incus_cpu_seconds_total{cpu=\"0\",mode=\"user\",name=\"web\",project=\"isb-acme\",type=\"container\"} 1\n";
602 let m = parse_disk_metrics(t);
603 assert_eq!(m[&("isb-acme".to_string(), "web".to_string())], (5596, 10));
604 assert_eq!(m.len(), 1);
605 }
606
607 #[test]
608 fn instance_rates_and_kinds() {
609 let mut s = Sampler::new();
610 let inst = |usage: u64| {
611 json!({
612 "name": "a", "status": "Running", "type": "container",
613 "config": {"volatile.container.oci": "true", "user.isb.stack": "app", "user.isb.create-token": "x"},
614 "state": {"cpu": {"usage": usage}, "memory": {"usage": 1024},
615 "network": {"lo": {"addresses": [{"family": "inet", "address": "127.0.0.1", "scope": "local"}]},
616 "eth0": {"counters": {"bytes_received": 100, "bytes_sent": 50}, "addresses": [{"family": "inet6", "address": "fd42::1", "scope": "global"},
617 {"family": "inet", "address": "10.0.0.2", "scope": "global"}]}}}
618 })
619 };
620 let t0 = Instant::now();
621 let a = s.instance(&inst(1_000_000_000), t0);
622 assert_eq!(a.cpu_pct, None);
623 assert_eq!(a.kind, "oci");
624 assert_eq!(a.ip.as_deref(), Some("10.0.0.2"));
625 assert_eq!(a.stack(), Some("app"));
626 assert_eq!((a.net_rx_bytes, a.net_tx_bytes), (Some(100), Some(50)));
627 assert!(!a.labels.contains_key("isb.create-token"));
628 let b = s.instance(&inst(1_500_000_000), t0 + std::time::Duration::from_secs(1));
629 assert!((b.cpu_pct.unwrap() - 50.0).abs() < 0.1, "{b:?}");
630 assert_eq!(b.cpu_history.len(), 2);
631 assert_eq!(b.disk_bytes, None, "no disk state: unknown");
632 }
633
634 #[test]
635 fn disk_usage() {
636 let mut s = Sampler::new();
637 let inst = |usage: i64| {
638 json!({"name": "a", "status": "Stopped", "type": "container",
639 "state": {"disk": {"root": {"usage": usage, "total": 0}}}})
640 };
641 assert_eq!(
643 s.instance(&inst(5 << 20), Instant::now()).disk_bytes,
644 Some(5 << 20)
645 );
646 assert_eq!(s.instance(&inst(-1), Instant::now()).disk_bytes, None);
647 assert_eq!(sum_pools(vec![(60, 100), (61, 100)]), (61, 100));
649 assert_eq!(sum_pools(vec![(60, 100), (5, 50), (0, 0)]), (65, 150));
650 }
651}