oxedyne/fe2o3/fe2o3_steel/src/srv/admin/host_sampler.rs
15.3 KiB, 121 runs
created by r1870400018:10767, which is this file's identity for as long as the history lasts, whatever it is later renamed to
download · who wrote it · its history
| 1 | //! Host resource sampler for the admin dashboard. |
| 2 | //! |
| 3 | //! Periodically takes a snapshot via `fe2o3_sys::Snapshot::sample` |
| 4 | //! and keeps a bounded ring of recent readings. The dashboard |
| 5 | //! reads the ring to draw host-resource charts (CPU, memory, |
| 6 | //! disk, network, load average). |
| 7 | //! |
| 8 | //! The sampler is parallel to [`super::traffic::TrafficRecorder`]: same |
| 9 | //! bounded-ring shape, same fixed-interval sampler task, same |
| 10 | //! `Arc`-shared ownership between the server and the dashboard. |
| 11 | //! Constructed once in the TUI startup path and carried through |
| 12 | //! [`AdminState`](super::state::AdminState). |
| 13 | //! |
| 14 | //! [Written with AI entirely](https://need2know.ai/entirely-ai/code)\ |
| 15 | //! Anthropic Claude |
| 16 | |
| 17 | use oxedyne_fe2o3_core::prelude::*; |
| 18 | use oxedyne_fe2o3_sys::{ |
| 19 | resident::Resident, |
| 20 | snapshot::Snapshot, |
| 21 | }; |
| 22 | |
| 23 | use std::{ |
| 24 | collections::VecDeque, |
| 25 | path::{ |
| 26 | Path, |
| 27 | PathBuf, |
| 28 | }, |
| 29 | sync::{ |
| 30 | Arc, |
| 31 | RwLock, |
| 32 | }, |
| 33 | time::{ |
| 34 | SystemTime, |
| 35 | UNIX_EPOCH, |
| 36 | }, |
| 37 | }; |
| 38 | |
| 39 | // At the default sample interval this is one hour of history, matching |
| 40 | // `TrafficRecorder::DEFAULT_HISTORY_CAPACITY`. |
| 41 | pub const DEFAULT_HISTORY_CAPACITY: usize = 720; |
| 42 | pub const DEFAULT_SAMPLE_INTERVAL_SECS: u64 = 5; |
| 43 | // A resident scan reads every process's status file, so it rides every sixth |
| 44 | // tick rather than every one: half a watcher's default round, so no probe reads |
| 45 | // a figure more than thirty seconds old. |
| 46 | pub const RESIDENT_SAMPLE_SECS: u64 = 30; |
| 47 | |
| 48 | /// Pairs a timestamp with the raw [`Snapshot`]. Rate-derived figures (CPU busy, |
| 49 | /// disk throughput) are computed by the consumer against the previous entry, |
| 50 | /// which keeps the sampler hot path free of arithmetic. |
| 51 | #[derive(Clone, Debug)] |
| 52 | pub struct HostSample { |
| 53 | pub when_secs: u64, // unix seconds |
| 54 | pub snapshot: Snapshot, |
| 55 | } |
| 56 | |
| 57 | /// A timestamp plus the four already-derived series values: the shape emitted |
| 58 | /// by `/admin/host.json` and the shape persisted to ozone so history survives a |
| 59 | /// restart. |
| 60 | /// |
| 61 | /// Derived because the useful figures need a pair of adjacent raw samples (CPU |
| 62 | /// busy fraction, disk B/s, net B/s). Persisting the reduced form keeps the |
| 63 | /// on-disk footprint small and sidesteps the need for ozone encoders over the |
| 64 | /// full `/proc`-derived struct tree. |
| 65 | #[derive(Clone, Copy, Debug)] |
| 66 | pub struct DerivedHostPoint { |
| 67 | pub t_secs: u64, // unix seconds of the later sample of the pair |
| 68 | pub cpu_pct: f64, // busy fraction over the preceding interval, per cent |
| 69 | pub mem_pct: f64, // used fraction of total RAM, per cent |
| 70 | pub disk_bps: f64, // aggregate disk throughput, bytes per second |
| 71 | pub net_bps: f64, // aggregate non-loopback rx + tx, bytes per second |
| 72 | } |
| 73 | |
| 74 | /// The four host figures the health body carries, already reduced to the |
| 75 | /// integers it emits. See [`HostSampler::health_metrics`]. |
| 76 | #[derive(Clone, Copy, Debug)] |
| 77 | pub struct HealthHostMetrics { |
| 78 | pub mem_pct: i64, |
| 79 | pub swap_pct: i64, |
| 80 | pub disk_iops: i64, |
| 81 | pub load1: i64, // 1-minute load average x100 |
| 82 | } |
| 83 | |
| 84 | /// Bounded ring of host snapshots, cheaply cloneable via `Arc` and shared |
| 85 | /// between the periodic sampler task spawned in [`crate::srv::server::Server::start`] and every |
| 86 | /// dashboard request handler. |
| 87 | #[derive(Debug)] |
| 88 | pub struct HostSampler { |
| 89 | history_capacity: usize, |
| 90 | history: RwLock<VecDeque<HostSample>>, // newest last |
| 91 | // Pre-restart points, loaded from ozone at start-up and rendered alongside |
| 92 | // the live derived history so the Overview sparkline strip does not reset to |
| 93 | // blank when Steel is restarted. |
| 94 | persisted: RwLock<Vec<DerivedHostPoint>>, |
| 95 | residents: Vec<String>, // `health_residents`, the processes to report |
| 96 | // The last resident scan and when it was taken, in unix seconds. Read by the |
| 97 | // health body, so a request formats figures rather than walking `/proc`. |
| 98 | resident_scan: RwLock<(u64, Vec<Resident>)>, |
| 99 | // The filesystem whose fullness the health body reports as `disk_pct`: the |
| 100 | // app root, where the databases and logs grow. `None` reports none. |
| 101 | disk_path: Option<PathBuf>, |
| 102 | disk_pct: RwLock<Option<i64>>, |
| 103 | } |
| 104 | |
| 105 | impl HostSampler { |
| 106 | pub fn new() -> Self { |
| 107 | Self { |
| 108 | history_capacity: DEFAULT_HISTORY_CAPACITY, |
| 109 | history: RwLock::new( |
| 110 | VecDeque::with_capacity(DEFAULT_HISTORY_CAPACITY), |
| 111 | ), |
| 112 | persisted: RwLock::new(Vec::new()), |
| 113 | residents: Vec::new(), |
| 114 | resident_scan: RwLock::new((0, Vec::new())), |
| 115 | disk_path: None, |
| 116 | disk_pct: RwLock::new(None), |
| 117 | } |
| 118 | } |
| 119 | |
| 120 | pub fn new_shared() -> Arc<Self> { |
| 121 | Arc::new(Self::new()) |
| 122 | } |
| 123 | |
| 124 | /// A sampler that also reads what the health body needs beyond the snapshot: |
| 125 | /// the named processes, every [`RESIDENT_SAMPLE_SECS`], and how full the |
| 126 | /// filesystem holding `disk_path` is, every tick. |
| 127 | pub fn new_shared_for_health( |
| 128 | residents: Vec<String>, |
| 129 | disk_path: Option<PathBuf>, |
| 130 | ) |
| 131 | -> Arc<Self> |
| 132 | { |
| 133 | let mut s = Self::new(); |
| 134 | s.residents = residents; |
| 135 | s.disk_path = disk_path; |
| 136 | Arc::new(s) |
| 137 | } |
| 138 | |
| 139 | pub fn history_capacity(&self) -> usize { |
| 140 | self.history_capacity |
| 141 | } |
| 142 | |
| 143 | /// Trims the oldest entry when the ring is already at capacity. |
| 144 | pub fn sample_now(&self) -> Outcome<()> { |
| 145 | let snap = res!(Snapshot::sample()); |
| 146 | let entry = HostSample { |
| 147 | when_secs: SystemTime::now() |
| 148 | .duration_since(UNIX_EPOCH) |
| 149 | .map(|d| d.as_secs()) |
| 150 | .unwrap_or(0), |
| 151 | snapshot: snap, |
| 152 | }; |
| 153 | let when_secs = entry.when_secs; |
| 154 | { |
| 155 | let mut hist = lock_write!(self.history); |
| 156 | if hist.len() == self.history_capacity { |
| 157 | hist.pop_front(); |
| 158 | } |
| 159 | hist.push_back(entry); |
| 160 | } |
| 161 | // Each read below is tried whatever became of the one before it, and the |
| 162 | // first failure is reported once both have had their turn. |
| 163 | let mut failed = None; |
| 164 | if let Some(path) = &self.disk_path { |
| 165 | // One `statvfs` call. A failure leaves the field absent rather than |
| 166 | // stale, so a watcher never reads an old figure as a current one. |
| 167 | let pct = disk_used_pct(path); |
| 168 | let mut slot = lock_write!(self.disk_pct); |
| 169 | match pct { |
| 170 | Ok(p) => *slot = Some(p), |
| 171 | Err(e) => { |
| 172 | *slot = None; |
| 173 | failed = Some(e); |
| 174 | }, |
| 175 | } |
| 176 | } |
| 177 | if !self.residents.is_empty() { |
| 178 | let due = { |
| 179 | let scan = lock_read!(self.resident_scan); |
| 180 | scan.0 == 0 || when_secs.saturating_sub(scan.0) >= RESIDENT_SAMPLE_SECS |
| 181 | }; |
| 182 | if due { |
| 183 | let found = Resident::sample(&self.residents); |
| 184 | let mut scan = lock_write!(self.resident_scan); |
| 185 | match found { |
| 186 | Ok(f) => *scan = (when_secs, f), |
| 187 | Err(e) => { |
| 188 | // The attempt is stamped, so a failing scan waits out the |
| 189 | // cadence like a good one and warns once a scan, not once a |
| 190 | // tick. |
| 191 | scan.0 = when_secs; |
| 192 | if failed.is_none() { |
| 193 | failed = Some(e); |
| 194 | } |
| 195 | }, |
| 196 | } |
| 197 | } |
| 198 | } |
| 199 | match failed { |
| 200 | Some(e) => Err(e), |
| 201 | None => Ok(()), |
| 202 | } |
| 203 | } |
| 204 | |
| 205 | /// The last resident scan; empty until the first has run, or when no |
| 206 | /// residents are configured. |
| 207 | pub fn residents(&self) -> Outcome<Vec<Resident>> { |
| 208 | let scan = lock_read!(self.resident_scan); |
| 209 | Ok(scan.1.clone()) |
| 210 | } |
| 211 | |
| 212 | /// How full the app root's filesystem was at the last tick, or `None` when |
| 213 | /// it has not been read or the read failed. |
| 214 | pub fn disk_pct(&self) -> Outcome<Option<i64>> { |
| 215 | let slot = lock_read!(self.disk_pct); |
| 216 | Ok(*slot) |
| 217 | } |
| 218 | |
| 219 | /// Chronological order, oldest first. |
| 220 | pub fn history_snapshot(&self) -> Outcome<Vec<HostSample>> { |
| 221 | let hist = lock_read!(self.history); |
| 222 | let mut out = Vec::with_capacity(hist.len()); |
| 223 | for s in hist.iter() { |
| 224 | out.push(s.clone()); |
| 225 | } |
| 226 | Ok(out) |
| 227 | } |
| 228 | |
| 229 | pub fn latest(&self) -> Outcome<Option<HostSample>> { |
| 230 | let hist = lock_read!(self.history); |
| 231 | Ok(hist.back().cloned()) |
| 232 | } |
| 233 | |
| 234 | /// The integer-ready host figures for the health body, or `None` when the |
| 235 | /// ring is empty. `disk_iops` needs two adjacent samples for its rate and is |
| 236 | /// zero until a second sample has landed; the level figures need only the |
| 237 | /// latest. |
| 238 | pub fn health_metrics(&self) -> Outcome<Option<HealthHostMetrics>> { |
| 239 | let hist = lock_read!(self.history); |
| 240 | let last = match hist.back() { |
| 241 | Some(s) => s, |
| 242 | None => return Ok(None), |
| 243 | }; |
| 244 | let mem_pct = (last.snapshot.mem.used_fraction() * 100.0).round() as i64; |
| 245 | let swap_pct = if last.snapshot.mem.swap_total == 0 { |
| 246 | 0 |
| 247 | } else { |
| 248 | (last.snapshot.mem.swap_used() as f64 |
| 249 | / last.snapshot.mem.swap_total as f64 * 100.0).round() as i64 |
| 250 | }; |
| 251 | // Load average times a hundred, so a fractional load survives the |
| 252 | // integer contract: 250 is a load of 2.50. |
| 253 | let load1 = (last.snapshot.load.one * 100.0).round() as i64; |
| 254 | // Completed reads + writes per second, summed over every device, from the |
| 255 | // two most recent samples. |
| 256 | let disk_iops = match hist.len() >= 2 { |
| 257 | true => { |
| 258 | let prev = &hist[hist.len() - 2]; |
| 259 | let elapsed = last.when_secs.saturating_sub(prev.when_secs); |
| 260 | if elapsed == 0 { |
| 261 | 0 |
| 262 | } else { |
| 263 | let mut ops: u64 = 0; |
| 264 | for dev in &last.snapshot.disk.devices { |
| 265 | if let Some(p) = prev.snapshot.disk.devices.iter() |
| 266 | .find(|d| d.name == dev.name) |
| 267 | { |
| 268 | ops = ops |
| 269 | .saturating_add(dev.reads.saturating_sub(p.reads)) |
| 270 | .saturating_add(dev.writes.saturating_sub(p.writes)); |
| 271 | } |
| 272 | } |
| 273 | (ops / elapsed) as i64 |
| 274 | } |
| 275 | }, |
| 276 | false => 0, |
| 277 | }; |
| 278 | Ok(Some(HealthHostMetrics { mem_pct, swap_pct, disk_iops, load1 })) |
| 279 | } |
| 280 | |
| 281 | /// Each entry carries the later-of-pair timestamp, because the rate-based |
| 282 | /// figures need two consecutive samples. |
| 283 | pub fn derived_history(&self) -> Outcome<Vec<DerivedHostPoint>> { |
| 284 | let hist = lock_read!(self.history); |
| 285 | if hist.len() < 2 { |
| 286 | return Ok(Vec::new()); |
| 287 | } |
| 288 | let mut out = Vec::with_capacity(hist.len() - 1); |
| 289 | let mut iter = hist.iter(); |
| 290 | let mut prev = match iter.next() { |
| 291 | Some(p) => p, |
| 292 | None => return Ok(out), |
| 293 | }; |
| 294 | for curr in iter { |
| 295 | let delta = curr.snapshot.delta(&prev.snapshot); |
| 296 | let disk_bps: f64 = delta.disk.iter() |
| 297 | .map(|d| d.read_bps + d.write_bps).sum(); |
| 298 | let net_bps: f64 = delta.net.iter() |
| 299 | .filter(|n| n.name != "lo") |
| 300 | .map(|n| n.rx_bps + n.tx_bps).sum(); |
| 301 | out.push(DerivedHostPoint { |
| 302 | t_secs: curr.when_secs, |
| 303 | cpu_pct: delta.cpu_busy * 100.0, |
| 304 | mem_pct: curr.snapshot.mem.used_fraction() * 100.0, |
| 305 | disk_bps, |
| 306 | net_bps, |
| 307 | }); |
| 308 | prev = curr; |
| 309 | } |
| 310 | Ok(out) |
| 311 | } |
| 312 | |
| 313 | pub fn seed_persisted(&self, points: Vec<DerivedHostPoint>) -> Outcome<()> { |
| 314 | let mut slot = lock_write!(self.persisted); |
| 315 | *slot = points; |
| 316 | Ok(()) |
| 317 | } |
| 318 | |
| 319 | /// Persisted plus live, capped at the ring's history capacity. The merge |
| 320 | /// drops persisted points at or after the oldest live derived timestamp, so |
| 321 | /// a sample still present in the live ring is not double-counted. |
| 322 | pub fn merged_derived_history(&self) -> Outcome<Vec<DerivedHostPoint>> { |
| 323 | let live = res!(self.derived_history()); |
| 324 | let persisted = { |
| 325 | let g = lock_read!(self.persisted); |
| 326 | g.clone() |
| 327 | }; |
| 328 | if live.is_empty() { |
| 329 | return Ok(persisted); |
| 330 | } |
| 331 | if persisted.is_empty() { |
| 332 | return Ok(live); |
| 333 | } |
| 334 | let cutoff = live.first().map(|p| p.t_secs).unwrap_or(0); |
| 335 | let mut out: Vec<DerivedHostPoint> = persisted.into_iter() |
| 336 | .filter(|p| p.t_secs < cutoff) |
| 337 | .collect(); |
| 338 | out.extend(live); |
| 339 | if out.len() > self.history_capacity { |
| 340 | let excess = out.len() - self.history_capacity; |
| 341 | out.drain(..excess); |
| 342 | } |
| 343 | Ok(out) |
| 344 | } |
| 345 | } |
| 346 | |
| 347 | impl Default for HostSampler { |
| 348 | fn default() -> Self { |
| 349 | Self::new() |
| 350 | } |
| 351 | } |
| 352 | |
| 353 | /// Per cent of a filesystem in use, as `df` reports it: blocks used over the |
| 354 | /// blocks an unprivileged writer could ever have, rounded up. The blocks reserved |
| 355 | /// for root count as neither, which is why a disk `df` shows at 100% still takes |
| 356 | /// root's writes -- and why an application that is not root finds it full. |
| 357 | pub fn used_pct(blocks: u64, free: u64, avail: u64) -> i64 { |
| 358 | let used = blocks.saturating_sub(free); |
| 359 | let usable = used.saturating_add(avail); |
| 360 | if usable == 0 { |
| 361 | return 0; |
| 362 | } |
| 363 | let pct = used.saturating_mul(100).saturating_add(usable - 1) / usable; |
| 364 | pct.min(100) as i64 |
| 365 | } |
| 366 | |
| 367 | #[cfg(unix)] |
| 368 | fn disk_used_pct(path: &Path) -> Outcome<i64> { |
| 369 | let st = match nix::sys::statvfs::statvfs(path) { |
| 370 | Ok(s) => s, |
| 371 | Err(e) => return Err(err!(e, |
| 372 | "Reading the filesystem usage of {:?} for the health body's disk_pct.", path; |
| 373 | IO, File, Read)), |
| 374 | }; |
| 375 | Ok(used_pct(st.blocks() as u64, st.blocks_free() as u64, st.blocks_available() as u64)) |
| 376 | } |
| 377 | |
| 378 | #[cfg(not(unix))] |
| 379 | fn disk_used_pct(path: &Path) -> Outcome<i64> { |
| 380 | Err(err!( |
| 381 | "Reading the filesystem usage of {:?} is implemented for Unix only.", path; |
| 382 | Unimplemented)) |
| 383 | } |
| 384 | |
| 385 | |
| 386 | #[cfg(test)] |
| 387 | mod tests { |
| 388 | use super::*; |
| 389 | |
| 390 | /// The figure `df` would print for the same counts, reserved blocks and all. |
| 391 | #[test] |
| 392 | fn disk_use_is_read_as_df_reads_it() { |
| 393 | // 100 blocks, 20 free of which 5 are reserved for root: 80 used of 95 usable. |
| 394 | assert_eq!(used_pct(100, 20, 15), 85, "80 of 95 is 84.2, which df rounds up"); |
| 395 | assert_eq!(used_pct(100, 5, 0), 100, "full to an unprivileged writer"); |
| 396 | assert_eq!(used_pct(100, 100, 100), 0); |
| 397 | assert_eq!(used_pct(0, 0, 0), 0, "an empty filesystem is not a division by zero"); |
| 398 | } |
| 399 | |
| 400 | /// The sampler reads the filesystem it was pointed at on its own tick, so the |
| 401 | /// health body only formats the figure. |
| 402 | #[cfg(unix)] |
| 403 | #[test] |
| 404 | fn a_tick_reads_the_disk_the_sampler_was_given() -> Outcome<()> { |
| 405 | let sampler = HostSampler::new_shared_for_health(Vec::new(), Some(PathBuf::from("/"))); |
| 406 | assert_eq!(res!(sampler.disk_pct()), None, "nothing read before the first tick"); |
| 407 | res!(sampler.sample_now()); |
| 408 | let pct = res!(res!(sampler.disk_pct()).ok_or_else(|| err!( |
| 409 | "the first tick did not read the disk"; Test, Missing))); |
| 410 | assert!((0..=100).contains(&pct), "got {}", pct); |
| 411 | Ok(()) |
| 412 | } |
| 413 | } |