Oregami
Repositories/oxedyne/fe2o3

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
17use oxedyne_fe2o3_core::prelude::*;
18use oxedyne_fe2o3_sys::{
19 resident::Resident,
20 snapshot::Snapshot,
21};
22
23use 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`.
41pub const DEFAULT_HISTORY_CAPACITY: usize = 720;
42pub 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.
46pub 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)]
52pub 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)]
66pub 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)]
77pub 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)]
88pub 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
105impl 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
347impl 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.
357pub 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)]
368fn 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))]
379fn 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)]
387mod 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}