Oregami
Repositories/oxedyne/fe2o3

oxedyne/fe2o3/fe2o3_steel/src/srv/fleet.rs

18.8 KiB, 1 run

created by r1870400018:61446, 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//! What the watcher saw, kept where the dashboard can read it.
2//!
3//! The watcher's beliefs about each peer are its own and stay private to
4//! [`crate::srv::watch`]. What it *saw* -- each probe's answer, its health body
5//! with the measured `probe_ms` folded in, and the belief that probe left behind
6//! -- goes into a bounded ring per peer behind the shared [`Fleet`] handle, which
7//! the admin state holds too. The view reads a projection of a judgement already
8//! made; it issues no request of its own, and nothing it does reaches back into
9//! the watcher.
10//!
11//! # A body that will not parse
12//!
13//! The design (2026-09-21) counted a `200` whose body will not parse as a failed
14//! probe. The watcher does not, and should not: a peer with no token serves an
15//! ordinary page at its URL, and calling that down would page someone about a
16//! healthy machine. So liveness keeps its rule, and the *page* marks a peer that
17//! has a token and answered without a readable body as stale -- its numbers are
18//! no longer current, which is the claim green would otherwise be making.
19//!
20//! [Written with AI entirely](https://need2know.ai/entirely-ai/code)\
21//! Anthropic Claude
22
23use crate::srv::{
24 cfg::{
25 WatchConfig,
26 WatchPeer,
27 },
28 health::{
29 F_SEALED_DBS,
30 HealthBody,
31 },
32 watch::{
33 is_cleared,
34 is_over,
35 },
36};
37
38use oxedyne_fe2o3_core::prelude::*;
39use oxedyne_fe2o3_data::ring::RingBuffer;
40
41use std::{
42 collections::BTreeMap,
43 sync::{
44 Arc,
45 RwLock,
46 atomic::{
47 AtomicBool,
48 Ordering,
49 },
50 },
51 time::{
52 SystemTime,
53 UNIX_EPOCH,
54 },
55};
56
57pub const RING_LEN: usize = 60; // an hour at the default 60 s round
58// A reading older than this many rounds is stale. Two, so one late or lost round
59// does not grey a row that the next round would have coloured.
60pub const FRESH_ROUNDS: u64 = 2;
61
62pub fn unix_secs() -> u64 {
63 SystemTime::now()
64 .duration_since(UNIX_EPOCH)
65 .map(|d| d.as_secs())
66 .unwrap_or(0)
67}
68
69/// The watcher's belief about a peer once a probe had been judged.
70#[derive(Clone, Copy, Debug, Eq, PartialEq)]
71pub enum PeerHealth {
72 Up,
73 Distressed,
74 Down,
75}
76
77/// One probe of one peer, as the watcher saw it.
78#[derive(Clone, Debug)]
79pub struct ProbeSample {
80 pub t_secs: u64, // unix seconds, when the probe was judged
81 pub ok: bool, // answered with a 2xx
82 pub probe_ms: u64, // measured by the watcher
83 pub body: Option<HealthBody>, // parsed, with `probe_ms` folded in
84 pub health: PeerHealth, // the belief this probe left behind
85}
86
87/// A watched peer as the page may show it. The token is reduced to whether there
88/// is one, so no path from here can put it on a page.
89#[derive(Clone, Debug, Eq, PartialEq)]
90pub struct FleetPeer {
91 pub name: String,
92 pub host: String,
93 pub url: String,
94 pub has_token: bool,
95 pub distress: BTreeMap<String, i64>,
96 pub clear: BTreeMap<String, i64>,
97}
98
99impl From<&WatchPeer> for FleetPeer {
100 fn from(p: &WatchPeer) -> Self {
101 Self {
102 name: p.name.clone(),
103 host: p.host.clone(),
104 url: p.url.clone(),
105 has_token: matches!(&p.token, Some(t) if !t.is_empty()),
106 distress: p.distress.clone(),
107 clear: p.clear.clone(),
108 }
109 }
110}
111
112/// The shared per-peer rings, and what the page needs to judge freshness.
113#[derive(Debug)]
114pub struct Fleet {
115 whoami: String, // this host, the page's own row
116 started_secs: u64,
117 interval_secs: u64, // the round, as the watcher runs it
118 fail_threshold: u32,
119 peers: Vec<FleetPeer>,
120 // Set by the watcher when its loop starts, so a watch list whose watcher never
121 // started -- no alerter, no outbound TLS -- says so rather than showing every
122 // peer as never heard from.
123 watching: AtomicBool,
124 // Set by the watcher after each round: no peer answered at all, so the round
125 // was read as its own link and kept out of the rings, and the page can say
126 // why its rows are ageing rather than leave that to be guessed.
127 link_down: AtomicBool,
128 rings: RwLock<Vec<RingBuffer<RING_LEN, ProbeSample>>>, // one per peer, in order
129}
130
131impl Fleet {
132 /// `cfg` is the watch list with its tokens already resolved, or `None` on a
133 /// host that watches nobody, which still has its own row to show.
134 pub fn new(whoami: String, cfg: Option<&WatchConfig>) -> Self {
135 let (peers, interval_secs, fail_threshold) = match cfg {
136 Some(c) => (
137 c.peers.iter().map(FleetPeer::from).collect::<Vec<_>>(),
138 // The watcher's own floor on its round.
139 c.interval_secs.max(5),
140 c.fail_threshold.max(1),
141 ),
142 None => {
143 let d = WatchConfig::default();
144 (Vec::new(), d.interval_secs, d.fail_threshold)
145 },
146 };
147 let rings = (0..peers.len()).map(|_| RingBuffer::default()).collect();
148 Self {
149 whoami,
150 started_secs: unix_secs(),
151 interval_secs,
152 fail_threshold,
153 peers,
154 watching: AtomicBool::new(false),
155 link_down: AtomicBool::new(false),
156 rings: RwLock::new(rings),
157 }
158 }
159
160 pub fn new_shared(whoami: String, cfg: Option<&WatchConfig>) -> Arc<Self> {
161 Arc::new(Self::new(whoami, cfg))
162 }
163
164 pub fn whoami(&self) -> &str { &self.whoami }
165 pub fn started_secs(&self) -> u64 { self.started_secs }
166 pub fn interval_secs(&self) -> u64 { self.interval_secs }
167 pub fn fail_threshold(&self) -> u32 { self.fail_threshold }
168 pub fn peers(&self) -> &[FleetPeer] { &self.peers }
169
170 pub fn set_watching(&self, on: bool) {
171 self.watching.store(on, Ordering::Release);
172 }
173
174 pub fn is_watching(&self) -> bool {
175 self.watching.load(Ordering::Acquire)
176 }
177
178 pub fn set_link_down(&self, down: bool) {
179 self.link_down.store(down, Ordering::Release);
180 }
181
182 /// Was the watcher's last round its own link going quiet?
183 pub fn is_link_down(&self) -> bool {
184 self.link_down.load(Ordering::Acquire)
185 }
186
187 /// Keep one probe of peer `i`, displacing the oldest once the ring is full.
188 pub fn record(&self, i: usize, sample: ProbeSample) -> Outcome<()> {
189 let mut rings = lock_write!(self.rings, "Recording a probe of peer {}.", i);
190 match rings.get_mut(i) {
191 Some(ring) => {
192 ring.set_and_adv(sample);
193 Ok(())
194 },
195 None => Err(err!(
196 "There is no ring for peer {}; the fleet was built for {} peer(s).",
197 i, self.peers.len();
198 Index, Missing)),
199 }
200 }
201
202 /// Every peer with its samples, oldest first, copied out so no lock is held
203 /// while a page is built from them.
204 pub fn snapshot(&self) -> Outcome<Vec<(FleetPeer, Vec<ProbeSample>)>> {
205 let rings = lock_read!(self.rings, "Reading the fleet's probe rings.");
206 let mut out = Vec::with_capacity(self.peers.len());
207 for (i, peer) in self.peers.iter().enumerate() {
208 let samples = match rings.get(i) {
209 Some(ring) => ring.iter_chrono().cloned().collect(),
210 None => Vec::new(),
211 };
212 out.push((peer.clone(), samples));
213 }
214 Ok(out)
215 }
216}
217
218// ┌───────────────────────────────────────────────────────────────────────────┐
219// │ WHAT THE PAGE MAY CLAIM │
220// └───────────────────────────────────────────────────────────────────────────┘
221
222/// The word at the head of a row.
223#[derive(Clone, Copy, Debug, Eq, PartialEq)]
224pub enum RowState {
225 Up,
226 Distressed,
227 Sealed, // answering, with databases the seal holds shut
228 Down, // the watcher's own call, the one that alarms
229 Stale, // no current reading, not yet down
230 Never, // no reading since this host's watcher started
231}
232
233impl RowState {
234 pub fn word(&self) -> &'static str {
235 match self {
236 Self::Up => "up",
237 Self::Distressed => "distressed",
238 Self::Sealed => "sealed",
239 Self::Down => "down",
240 Self::Stale => "stale",
241 Self::Never => "never",
242 }
243 }
244
245 /// May the row's numbers be shown as current and coloured? A stale row keeps
246 /// its last numbers, dimmed and grey; the others have none to show.
247 pub fn is_fresh(&self) -> bool {
248 matches!(self, Self::Up | Self::Distressed | Self::Sealed)
249 }
250}
251
252/// What the page may say about one peer, and why when that is not `up`.
253///
254/// Only the samples decide it: the latest belief the watcher recorded, the age of
255/// the last good reading against [`FRESH_ROUNDS`] rounds, and -- for a peer that
256/// has a token -- whether the latest answer carried a body that parsed.
257pub fn peer_state(
258 peer: &FleetPeer,
259 samples: &[ProbeSample],
260 now_secs: u64,
261 interval_secs: u64,
262)
263 -> (RowState, String)
264{
265 let latest = match samples.last() {
266 Some(s) => s,
267 None => return (RowState::Never, fmt!("not probed yet")),
268 };
269 if latest.health == PeerHealth::Down {
270 return (RowState::Down, fmt!("called down by this host's watcher"));
271 }
272 let window = interval_secs.saturating_mul(FRESH_ROUNDS);
273 if !peer.has_token {
274 // Liveness alone: a good reading is any 2xx.
275 return match samples.iter().rev().find(|s| s.ok) {
276 None => (RowState::Never, fmt!("no answer since this host's watcher started")),
277 Some(s) if now_secs.saturating_sub(s.t_secs) > window => (
278 RowState::Stale,
279 fmt!("no answer for {} s", now_secs.saturating_sub(s.t_secs))),
280 Some(_) => (RowState::Up, String::new()),
281 };
282 }
283 let last_read = samples.iter().rev().find(|s| s.body.is_some());
284 if latest.ok && latest.body.is_none() {
285 return match last_read {
286 None => (RowState::Never, fmt!("answers, but its health body does not parse")),
287 Some(_) => (RowState::Stale, fmt!("answered without a readable health body")),
288 };
289 }
290 let body = match last_read {
291 Some(s) if now_secs.saturating_sub(s.t_secs) > window => return (
292 RowState::Stale,
293 fmt!("no reading for {} s", now_secs.saturating_sub(s.t_secs))),
294 Some(s) => match &s.body {
295 Some(b) => b,
296 None => return (RowState::Never, String::new()),
297 },
298 None => return (RowState::Never, fmt!("no reading since this host's watcher started")),
299 };
300 if body.get(F_SEALED_DBS).unwrap_or(0) > 0 {
301 return (RowState::Sealed, fmt!("databases held shut awaiting an unseal"));
302 }
303 match latest.health {
304 PeerHealth::Distressed => (RowState::Distressed, String::new()),
305 _ => (RowState::Up, String::new()),
306 }
307}
308
309/// A cell's colour, from the same thresholds the alarm uses.
310#[derive(Clone, Copy, Debug, Eq, PartialEq)]
311pub enum Tone {
312 Plain, // no threshold, so no claim
313 Green,
314 Amber, // above the clear boundary, under distress: an alarm here would not yet clear
315 Red,
316}
317
318impl Tone {
319 pub fn word(&self) -> &'static str {
320 match self {
321 Self::Plain => "",
322 Self::Green => "green",
323 Self::Amber => "amber",
324 Self::Red => "red",
325 }
326 }
327}
328
329/// Red at or over `distress`, amber over the clear boundary, green at or under it.
330/// The clear boundary is `clear` where one is given and `distress` otherwise,
331/// exactly as the watcher reads it, so a field with no `clear` has no amber.
332pub fn tone(value: i64, distress: Option<i64>, clear: Option<i64>) -> Tone {
333 let d = match distress {
334 Some(d) => d,
335 None => return Tone::Plain,
336 };
337 if is_over(value, d) {
338 return Tone::Red;
339 }
340 if is_cleared(value, clear.unwrap_or(d)) {
341 Tone::Green
342 } else {
343 Tone::Amber
344 }
345}
346
347
348#[cfg(test)]
349mod tests {
350 use super::*;
351
352 fn peer(has_token: bool) -> FleetPeer {
353 FleetPeer {
354 name: fmt!("jarrah"),
355 host: fmt!("jarrah"),
356 url: fmt!("https://example.test/_steel/health"),
357 has_token,
358 distress: BTreeMap::new(),
359 clear: BTreeMap::new(),
360 }
361 }
362
363 fn body(fields: &[(&str, i64)]) -> HealthBody {
364 let mut b = HealthBody::new();
365 for (k, v) in fields {
366 b.set(k, *v);
367 }
368 b
369 }
370
371 fn sample(t_secs: u64, ok: bool, b: Option<HealthBody>, health: PeerHealth) -> ProbeSample {
372 ProbeSample { t_secs, ok, probe_ms: 80, body: b, health }
373 }
374
375 fn fleet_of(n: usize) -> Fleet {
376 let mut cfg = WatchConfig::default();
377 for i in 0..n {
378 cfg.peers.push(WatchPeer {
379 name: fmt!("peer{}", i),
380 host: fmt!("peer{}", i),
381 url: fmt!("https://peer{}.test/_steel/health", i),
382 plain_ok: false,
383 distress: BTreeMap::new(),
384 clear: BTreeMap::new(),
385 token: Some(fmt!("secret")),
386 repeat_secs: None,
387 });
388 }
389 Fleet::new(fmt!("karri"), Some(&cfg))
390 }
391
392 /// The ring keeps the last sixty probes and hands them back oldest first, however
393 /// many times it has wrapped.
394 #[test]
395 fn a_peer_ring_holds_the_last_sixty_probes_oldest_first() -> Outcome<()> {
396 let fleet = fleet_of(2);
397 for t in 1..=(RING_LEN as u64 + 10) {
398 res!(fleet.record(0, sample(t, true, None, PeerHealth::Up)));
399 }
400 res!(fleet.record(1, sample(500, true, None, PeerHealth::Up)));
401 let snap = res!(fleet.snapshot());
402 let times: Vec<u64> = snap[0].1.iter().map(|s| s.t_secs).collect();
403 assert_eq!(times.len(), RING_LEN, "the ring must be bounded at {}", RING_LEN);
404 assert_eq!(times.first(), Some(&11), "the ten oldest must have been displaced");
405 assert_eq!(times.last(), Some(&(RING_LEN as u64 + 10)));
406 assert!(times.windows(2).all(|w| w[0] < w[1]), "the samples must read oldest first");
407 assert_eq!(snap[1].1.len(), 1, "each peer has a ring of its own");
408 assert!(fleet.record(2, sample(1, true, None, PeerHealth::Up)).is_err(),
409 "a peer the fleet was not built for is refused, not silently dropped");
410 Ok(())
411 }
412
413 /// The token is never copied into what the page can reach.
414 #[test]
415 fn a_fleet_peer_keeps_whether_there_is_a_token_and_not_the_token() {
416 let fleet = fleet_of(1);
417 assert!(fleet.peers()[0].has_token);
418 assert!(!fmt!("{:?}", fleet.peers()).contains("secret"));
419 }
420
421 /// A peer that has a token and answers `200` with a body that does not parse is up to the
422 /// alarm and stale to the page: its last numbers are no longer current.
423 #[test]
424 fn an_unparseable_body_from_a_token_peer_reads_stale() {
425 let p = peer(true);
426 let good = sample(1_000, true, Some(body(&[("mem_pct", 40)])), PeerHealth::Up);
427 let garbled = sample(1_060, true, None, PeerHealth::Up);
428 let (state, note) = peer_state(&p, &[good.clone(), garbled.clone()], 1_070, 60);
429 assert_eq!(state, RowState::Stale, "the latest answer carried no readable body");
430 assert!(note.contains("readable"), "the note must say why, got '{}'", note);
431 assert!(!state.is_fresh());
432
433 // A peer that has never served a readable body has nothing stale to show.
434 let (state, _) = peer_state(&p, &[garbled.clone()], 1_070, 60);
435 assert_eq!(state, RowState::Never);
436
437 // The same answer from a peer with no token is ordinary liveness.
438 let (state, _) = peer_state(&peer(false), &[garbled], 1_070, 60);
439 assert_eq!(state, RowState::Up);
440
441 // And a readable body in the latest answer is fresh again.
442 let back = sample(1_120, true, Some(body(&[("mem_pct", 41)])), PeerHealth::Up);
443 let (state, _) = peer_state(&p, &[good, back], 1_130, 60);
444 assert_eq!(state, RowState::Up);
445 }
446
447 /// Older than two rounds without a good reading is stale, before the watcher has counted
448 /// enough failures to call the peer down; once it has, the row says down.
449 #[test]
450 fn freshness_follows_the_age_of_the_last_good_reading() {
451 let p = peer(true);
452 let good = sample(1_000, true, Some(body(&[("mem_pct", 40)])), PeerHealth::Up);
453 let miss = sample(1_120, false, None, PeerHealth::Up);
454 assert_eq!(peer_state(&p, &[good.clone()], 1_110, 60).0, RowState::Up);
455 assert_eq!(peer_state(&p, &[good.clone(), miss.clone()], 1_121, 60).0, RowState::Stale);
456 let dead = sample(1_180, false, None, PeerHealth::Down);
457 assert_eq!(peer_state(&p, &[good, miss, dead], 1_181, 60).0, RowState::Down);
458 assert_eq!(peer_state(&p, &[], 1_181, 60).0, RowState::Never);
459 }
460
461 /// Sealed is read from the databases the seal holds shut, so a box with no database -- which
462 /// runs sealed for ever by design -- reads up, not sealed.
463 #[test]
464 fn sealed_means_databases_held_shut() {
465 let p = peer(true);
466 let no_dbs = sample(1_000, true, Some(body(&[("sealed", 1), ("sealed_dbs", 0)])),
467 PeerHealth::Up);
468 assert_eq!(peer_state(&p, &[no_dbs], 1_001, 60).0, RowState::Up);
469 let shut = sample(1_000, true, Some(body(&[("sealed", 1), ("sealed_dbs", 2)])),
470 PeerHealth::Up);
471 assert_eq!(peer_state(&p, &[shut], 1_001, 60).0, RowState::Sealed);
472 let hot = sample(1_000, true, Some(body(&[("mem_pct", 95)])), PeerHealth::Distressed);
473 assert_eq!(peer_state(&p, &[hot], 1_001, 60).0, RowState::Distressed);
474 }
475
476 /// The colours are the alarm's own comparisons: red at the distress value, amber in the
477 /// dead-band, green at or under clear, and nothing where no threshold was set.
478 #[test]
479 fn tones_follow_the_alarm_thresholds() {
480 assert_eq!(tone(90, Some(90), Some(75)), Tone::Red, "at the threshold is over it");
481 assert_eq!(tone(89, Some(90), Some(75)), Tone::Amber);
482 assert_eq!(tone(76, Some(90), Some(75)), Tone::Amber);
483 assert_eq!(tone(75, Some(90), Some(75)), Tone::Green, "at clear is cleared");
484 assert_eq!(tone(89, Some(90), None), Tone::Green, "no clear map, no dead-band");
485 assert_eq!(tone(99, None, Some(75)), Tone::Plain, "a clear without a distress is no claim");
486 }
487}