Oregami
Repositories/oxedyne/fe2o3

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

73.0 KiB, 7 runs

created by r1870400018:21449, 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//! Watching the other machines, because a host cannot report its own death.
2//!
3//! # The problem this exists for
4//!
5//! [`crate::srv::alert`] says it plainly in its own header: Steel is reporting on itself, and a
6//! Steel that is wedged, unreachable or dead sends nothing. **Silence is indistinguishable from
7//! health.** On 2026-08-10 a payments gateway exited on its own and was down for about fifty
8//! minutes while the websites in front of it served perfectly; nothing said a word, because the
9//! only thing positioned to notice was the machine that had died.
10//!
11//! The fix is to invert the question. Do not detect a failure -- require a success, on a
12//! schedule, from somewhere else, and treat its absence as the alarm. That is a thing only
13//! another host can do.
14//!
15//! # The shape: a mesh, not a monitor
16//!
17//! Every node watches every other node it is told about, and any node can raise the alarm. There
18//! is no monitoring server, because a monitoring server is one more single point that fails
19//! silently. A three-node estate where each watches the other two survives losing any one of
20//! them, and survives losing any two as far as the third is concerned.
21//!
22//! **Adding or removing a machine is one line of configuration**, and nothing else changes: no
23//! code, no central registry, no re-deployment of the others beyond their own peer list. That is
24//! the whole reason the peer list is data rather than a compiled set.
25//!
26//! # Duplicate alarms are a feature
27//!
28//! When two nodes both notice that a third has gone, the operator gets told twice. That is left
29//! alone deliberately. Suppressing it would need the watchers to agree with each other, which
30//! means a protocol between them, which means a thing that can itself fail and take the alarm
31//! with it. Two messages saying the same true thing cost a few cents and a glance; one message
32//! that was suppressed by a consensus that broke costs an outage.
33//!
34//! # What is watched
35//!
36//! A URL that answers `200` when the machine is well. Nothing more clever, because anything more
37//! clever is a thing to keep in step with the machine it watches. What that URL means is the
38//! watched machine's business -- a gateway's `/api/health` already reports whether its store
39//! opened, which is a far better answer than whether a port accepts a connection.
40//!
41//! # Why a peer may say `http`, and why it must say so
42//!
43//! A probe is normally `https`, because it crosses the public internet and an unauthenticated
44//! answer can be forged by anybody on the path -- and the forgery that matters is a `200`, which
45//! hides the outage rather than inventing one. That stays the default and the absence of the key
46//! keeps it.
47//!
48//! The exception this admits is a machine that serves nothing else. birch holds the off-host copy
49//! of the forge and listens on 22 alone; giving it a certificate means giving it a public name, a
50//! port open to a certificate authority's validators, and a renewal that can quietly stop -- three
51//! moving parts on the one machine whose entire value is being boring. Its freshness endpoint is a
52//! high port admitted by its firewall from the watcher's address alone. So a peer may carry
53//! `"plain_ok": (true)`, and it is per peer rather than a switch on the watcher, because the
54//! judgement is about one wire and not about watching in general. On 2026-08-23 that copy failed
55//! every hour for twenty-one hours with nobody able to see it, which is the cost of the stricter
56//! answer.
57//!
58//! # What it deliberately does not do
59//!
60//! It does not restart anything. A watcher that repairs is a watcher that can flap a service in
61//! a loop at three in the morning and hide the fault it was built to reveal; and a decision to
62//! restart a payments process belongs to a person who has read why it stopped.
63//!
64//! # When nobody answers, the silence is the watcher's own
65//!
66//! A watcher on a home connection loses its link for an hour, or wakes from sleep before its
67//! Wi-Fi does. Judged peer by peer, that is every peer `DOWN` at once -- messages that cannot be
68//! sent -- and then every peer `is back` for outages that never happened. So a round in which
69//! peers on two or more hosts, none of them already `DOWN`, were asked and not one answered at
70//! all is read as this watcher's link and judged not at all (see [`is_own_link_down`]). An answer of any kind, an error page included,
71//! proves the link, so a round of refusals is still news about the peers.
72//!
73//! [Written with AI entirely](https://need2know.ai/entirely-ai/code)\
74//! Anthropic Claude
75
76use crate::srv::{
77 alert::{
78 AlertEvent,
79 Alerter,
80 },
81 cfg::{
82 WatchConfig,
83 WatchPeer,
84 },
85 fleet::{
86 Fleet,
87 PeerHealth,
88 ProbeSample,
89 unix_secs,
90 },
91 health::{
92 F_PROBE_MS,
93 HealthBody,
94 },
95};
96
97use oxedyne_fe2o3_core::prelude::*;
98
99use std::collections::{
100 BTreeMap,
101 BTreeSet,
102};
103use oxedyne_fe2o3_net::http::{
104 client::{
105 http_request_limited,
106 https_request_limited,
107 },
108 header::{
109 HttpHeadline,
110 HttpMethod,
111 },
112 loc::Url,
113 msg::ReadLimits,
114};
115
116use std::{
117 sync::Arc,
118 time::{
119 Duration,
120 Instant,
121 },
122};
123
124use tokio_rustls::rustls::ClientConfig;
125
126
127/// What this node currently believes about one peer.
128///
129/// Liveness (`Up` / `Down`) and distress (`Distressed`) are the two things a peer
130/// can be wrong about, and they are distinct: `Down` is a peer that stops
131/// answering, `Distressed` is one that answers but reports itself unwell. A peer
132/// cannot be both, because distress is read from a body a dead peer never sends.
133#[derive(Clone, Copy, Debug, Eq, PartialEq)]
134enum Health {
135 // The count on `Up` is consecutive probe failures since the last success, not yet enough to
136 // call the peer down.
137 Up { failures: u32 },
138 // Answering, but a health-body class has stayed over its distress threshold. `over` is the
139 // consecutive over-threshold readings that tipped it in; `failures` is the consecutive probe
140 // failures counted since, so a peer that stops answering while distressed still reaches `Down`
141 // by the same rule an `Up` peer does.
142 Distressed { over: u32, failures: u32 },
143 Down,
144}
145
146impl From<Health> for PeerHealth {
147 fn from(h: Health) -> Self {
148 match h {
149 Health::Up { .. } => Self::Up,
150 Health::Distressed { .. } => Self::Distressed,
151 Health::Down => Self::Down,
152 }
153 }
154}
155
156// The most of an answer a probe reads. A real health body is under a few KiB even with every
157// stamp and resident, and each one read is kept in the dashboard's ring for an hour, so a peer
158// that sends more is refused rather than held (D-06 audit D1).
159pub const PROBE_BODY_MAX: usize = 64 * 1024;
160pub const PROBE_HEADER_MAX: usize = 16 * 1024;
161
162/// What one probe came back with.
163#[derive(Clone, Debug, Eq, PartialEq)]
164enum Probe {
165 Up(Option<HealthBody>), // a `2xx`, with its body when one was served and parsed
166 Refused(u16), // an answer that is not a `2xx`: reachable, and not well
167 Silent, // no answer at all: no connection, no handshake, or the timeout
168}
169
170impl Probe {
171 fn is_up(&self) -> bool {
172 matches!(self, Self::Up(_))
173 }
174
175 /// Did anything answer? A refusal did: it proves the path to the peer works.
176 fn answered(&self) -> bool {
177 !matches!(self, Self::Silent)
178 }
179}
180
181/// Is a round in which nobody answered this watcher's own link, rather than news about peers?
182///
183/// Two or more hosts that were answering asked and not one answer back is the watcher's link,
184/// not a coincidence of outages. Hosts, not entries, since two entries on one machine fall silent
185/// together when it dies. And only peers not already `Down`, since a peer told dead is silent
186/// anyway: counting it would let one outage hide the next. A single live host's silence is judged
187/// as ever: with one the two cannot be told apart, and an outage must not be the thing explained
188/// away.
189fn is_own_link_down(peers: &[WatchPeer], state: &[PeerState], round: &[Probe]) -> bool {
190 if round.iter().any(Probe::answered) {
191 return false;
192 }
193 let live: BTreeSet<&str> = peers.iter()
194 .zip(state)
195 .filter(|(_, st)| st.health != Health::Down)
196 .map(|(p, _)| p.host.as_str())
197 .collect();
198 live.len() >= 2
199}
200
201/// The peers whose distress this host would notice and tell nobody.
202///
203/// Distress is told by mail alone (`Severity::Notice`), so a peer with thresholds, on a host
204/// whose alerter has no mail recipient, is a watch that looks like cover and is not.
205fn distress_unheard(peers: &[WatchPeer], mails: bool) -> Vec<String> {
206 if mails {
207 return Vec::new();
208 }
209 peers.iter()
210 .filter(|p| !p.distress.is_empty())
211 .map(|p| p.name.clone())
212 .collect()
213}
214
215/// One peer's running state.
216struct PeerState {
217 health: Health,
218 failing_at: Option<Instant>, // first seen failing, so a recovery can say how long
219 told_at: Option<Instant>, // last told down, so a lasting outage reminds not streams
220 // Consecutive over-threshold readings while still `Up`, so distress needs the same run of
221 // agreeing evidence that a death does before it wakes anyone.
222 distress_rising: u32,
223 distress_at: Option<Instant>, // when the distress run began, for the recovery line
224 distress_told_at: Option<Instant>, // last told distressed, for the repeat cadence
225}
226
227impl Default for PeerState {
228 fn default() -> Self {
229 Self {
230 health: Health::Up { failures: 0 },
231 failing_at: None,
232 told_at: None,
233 distress_rising: 0,
234 distress_at: None,
235 distress_told_at: None,
236 }
237 }
238}
239
240/// Is a reading over its distress threshold?
241///
242/// The alarm and the Fleet page's red both ask this one function, and whether a
243/// reading has cleared both ask [`is_cleared`], so the page cannot draw a colour
244/// the alarm would not have agreed with.
245pub(crate) fn is_over(value: i64, distress: i64) -> bool {
246 value >= distress
247}
248
249/// Is a reading at or under its clear boundary?
250pub(crate) fn is_cleared(value: i64, boundary: i64) -> bool {
251 value <= boundary
252}
253
254/// The distress classes a body is currently over, as `(field, value)` pairs.
255fn classes_breaching(distress: &BTreeMap<String, i64>, body: &HealthBody) -> Vec<(String, i64)> {
256 let mut out = Vec::new();
257 for (field, threshold) in distress {
258 if let Some(v) = body.get(field) {
259 if is_over(v, *threshold) {
260 out.push((field.clone(), v));
261 }
262 }
263 }
264 out
265}
266
267/// Is every distress class at or below its clear boundary?
268///
269/// The clear boundary is a field's `clear` value where one is given, and its
270/// `distress` value otherwise -- so a peer with no `clear` map has no dead-band
271/// and leaves distress the moment it drops back under the entry threshold.
272fn is_below_clear(peer: &WatchPeer, body: &HealthBody) -> bool {
273 for (field, dthresh) in &peer.distress {
274 let boundary = peer.clear.get(field).copied().unwrap_or(*dthresh);
275 if let Some(v) = body.get(field) {
276 if !is_cleared(v, boundary) {
277 return false;
278 }
279 }
280 }
281 true
282}
283
284/// The firing classes as one line for an alarm, e.g. `mem_pct 94, swap_pct 71`.
285fn classes_text(classes: &[(String, i64)]) -> String {
286 classes.iter()
287 .map(|(f, v)| fmt!("{} {}", f, v))
288 .collect::<Vec<_>>()
289 .join(", ")
290}
291
292/// The peer watcher.
293///
294/// Owns its configuration, its TLS client and its beliefs, and writes what it saw into the
295/// shared [`Fleet`] for the dashboard to draw. Constructed once at start-up and driven by
296/// [`Self::run`], which never returns.
297pub struct Watcher {
298 cfg: Arc<WatchConfig>,
299 alerter: Arc<Alerter>,
300 tls: Arc<ClientConfig>,
301 // This node's own name, so an alert says who noticed as well as what happened. Two nodes
302 // watching a third send two messages, and without this they are indistinguishable.
303 whoami: String,
304 // One per peer, in `cfg.peers` order. Kept by position rather than by name, because
305 // nothing makes a name unique, and two entries sharing one would share one set of beliefs.
306 state: Vec<PeerState>,
307 link_down: bool, // the last round was this watcher's own link, so the episode is said once
308 fleet: Arc<Fleet>,
309}
310
311/// Refuse a peer whose URL this cannot honestly probe.
312///
313/// Separated from [`Watcher::new`] so the rule can be tested on its own: it decides whether an
314/// operator is watched over an authenticated wire, and a rule that costly should not need an
315/// alerter and a TLS client standing up before it can be exercised.
316fn vet(p: &WatchPeer) -> Outcome<()> {
317 let url = res!(Url::parse(&p.url));
318 if !url.scheme.is_tls() && !p.plain_ok {
319 return Err(err!(
320 "The watch entry for '{}' names {}, which is not https. A health probe crosses the \
321 public internet and its answer decides whether an operator is woken, so it is \
322 authenticated or it is not worth making. A machine that cannot hold a certificate \
323 may say so with \"plain_ok\": (true) on its own entry, which is a judgement about \
324 that one wire.", p.name, p.url;
325 Configuration, Invalid, Input));
326 }
327 if p.plain_ok && url.scheme.is_tls() {
328 warn!("The watch entry for '{}' says plain_ok and names an https URL. The key does \
329 nothing here; remove it, so it does not read as a weakness that was accepted.",
330 p.name);
331 }
332 Ok(())
333}
334
335impl Watcher {
336 /// Build a watcher over a peer list.
337 ///
338 /// A peer whose URL cannot be parsed, or one that is not `https` and has not said
339 /// `plain_ok`, is refused at start-up rather than at the first poll: a watcher that silently
340 /// watches nothing is the failure this module exists to prevent, and start-up is when
341 /// somebody is looking.
342 pub fn new(
343 cfg: Arc<WatchConfig>,
344 alerter: Arc<Alerter>,
345 tls: Arc<ClientConfig>,
346 whoami: String,
347 fleet: Arc<Fleet>,
348 )
349 -> Outcome<Self>
350 {
351 let mut state = Vec::with_capacity(cfg.peers.len());
352 for p in &cfg.peers {
353 res!(vet(p));
354 state.push(PeerState::default());
355 }
356 let unheard = distress_unheard(&cfg.peers, alerter.sends_mail());
357 if !unheard.is_empty() {
358 warn!("Watch: {} carry distress thresholds, but this host's alerts have no mail \
359 recipient, and distress is told by mail alone. Their distress would be noticed \
360 and told to nobody: give 'alerts' a 'from' and a 'to', with a 'submission' relay \
361 where this host cannot send direct.", unheard.join(", "));
362 }
363 Ok(Self { cfg, alerter, tls, whoami, state, link_down: false, fleet })
364 }
365
366 /// Poll every peer for ever.
367 ///
368 /// Never returns. Errors from a single probe are the point of the exercise and are handled;
369 /// there is no failure here worth ending the loop for, because ending the loop is exactly
370 /// the silence this is meant to break.
371 pub async fn run(mut self) {
372 let every = Duration::from_secs(self.cfg.interval_secs.max(5));
373 // A peer probed over plain http is named as such here, so that the concession shows in
374 // the log an operator actually reads rather than only in a configuration file nobody
375 // opens between incidents.
376 info!("Watching {} peer(s) every {}s: {}.",
377 self.cfg.peers.len(),
378 every.as_secs(),
379 self.cfg.peers.iter()
380 .map(|p| if p.plain_ok { fmt!("{} (plain http)", p.name) } else { p.name.clone() })
381 .collect::<Vec<_>>().join(", "));
382
383 // The first probe waits one interval. A node that has just restarted is a node whose
384 // peers may still be restarting too -- a shared power event, a rolling deploy -- and an
385 // alarm raised in the first second of life is usually about the estate coming up, not
386 // about anything being wrong.
387 let beat = Duration::from_secs(self.cfg.heartbeat_secs);
388 let started = Instant::now();
389 let mut last_beat = started;
390 self.fleet.set_watching(true);
391 loop {
392 tokio::time::sleep(every).await;
393 // Every peer is asked before any is judged, so a round in which nobody answered can
394 // be told from one in which some did.
395 let mut round = Vec::with_capacity(self.cfg.peers.len());
396 for peer in self.cfg.peers.iter() {
397 let (probe, dur) = self.probe(peer).await;
398 debug!("Watch: {} up={} answered={} in {}ms.",
399 peer.name, probe.is_up(), probe.answered(), dur.as_millis());
400 round.push((probe, dur));
401 }
402 let ok_count = round.iter().filter(|(p, _)| p.is_up()).count();
403 let (events, samples) = judge_round(
404 &self.cfg, &self.whoami, &mut self.state, &mut self.link_down,
405 round, Instant::now(), unix_secs());
406 for event in events {
407 self.alerter.raise(event);
408 }
409 // The dashboard's copy is written after the alerts are raised, and a failure to write
410 // it is only logged, so nothing the page does can delay or silence an alarm.
411 self.fleet.set_link_down(self.link_down);
412 for (i, sample) in samples {
413 if let Err(e) = self.fleet.record(i, sample) {
414 let name = self.cfg.peers.get(i).map(|p| p.name.as_str()).unwrap_or("?");
415 warn!("Watch: the dashboard's copy of {}'s probe was not kept: {}", name, e);
416 }
417 }
418 // Proof of life, on the same loop that does the watching -- so a
419 // heartbeat arriving is evidence the watcher is running and not
420 // merely that a timer somewhere else still fires.
421 if self.cfg.heartbeat_secs > 0 && Instant::now().duration_since(last_beat) >= beat {
422 last_beat = Instant::now();
423 self.alerter.raise(AlertEvent::Heartbeat {
424 uptime_secs: Instant::now().duration_since(started).as_secs(),
425 peers_ok: ok_count,
426 peers_total: self.cfg.peers.len(),
427 });
428 }
429 }
430 }
431
432 async fn probe(&self, peer: &WatchPeer) -> (Probe, std::time::Duration) {
433 let timeout = Duration::from_secs(self.cfg.timeout_secs.max(2));
434 probe(peer, timeout, self.tls.clone()).await
435 }
436}
437
438/// Ask one peer whether it is well, returning what came back and how long the probe took.
439///
440/// Any answer that is not a `2xx` is a failure, including a `503`: a Steel that is up and
441/// sealed is answering, and it is still not serving the databases behind it. It is a
442/// [`Probe::Refused`] rather than [`Probe::Silent`] all the same, because an answer proves
443/// the path to the peer. The body is present only when the peer carries a `token` (so the
444/// gate opens) and the answer parsed; its absence never fails the liveness check, so a peer
445/// that answers `200` without a body is still up. The duration is measured here, not
446/// self-reported by the peer.
447async fn probe(
448 peer: &WatchPeer,
449 timeout: Duration,
450 tls: Arc<ClientConfig>,
451)
452 -> (Probe, std::time::Duration)
453{
454 let url = &peer.url;
455 let started = Instant::now();
456 let loc = match Url::parse(url) {
457 Ok(l) => l,
458 // Refused at construction, so this cannot happen -- and if it ever does, a peer
459 // that cannot be addressed is a peer that is not answering.
460 Err(e) => {
461 warn!("The watch URL {} stopped parsing: {}", url, e);
462 return (Probe::Silent, started.elapsed());
463 },
464 };
465 let host = loc.host.clone();
466 let port = loc.port;
467 let path = loc.target.clone();
468
469 // The token, when the operator gave one, opens the peer's health gate. Without it the
470 // peer answers a 404 for the health path, so the probe reads liveness only.
471 let mut headers: Vec<(&str, &str)> = vec![
472 ("Connection", "close"),
473 ("User-Agent", "steel-watch"),
474 ];
475 if let Some(tok) = &peer.token {
476 headers.push(("x-steel-health-token", tok.as_str()));
477 }
478 let limits = ReadLimits {
479 max_header_bytes: Some(PROBE_HEADER_MAX),
480 max_body_bytes: Some(PROBE_BODY_MAX),
481 header_read_timeout: None,
482 };
483 // The scheme decides, and `new` has already refused a plain URL that nobody opted in
484 // to -- so by the time a probe runs, `http` here means the operator wrote it down.
485 let reply = if loc.scheme.is_tls() {
486 let call = https_request_limited(
487 &host, port, HttpMethod::GET, &path, &headers, &[], tls, Some(&limits),
488 );
489 tokio::time::timeout(timeout, call).await
490 } else {
491 let call = http_request_limited(
492 &host, port, HttpMethod::GET, &path, &headers, &[], Some(&limits),
493 );
494 tokio::time::timeout(timeout, call).await
495 };
496 let elapsed = started.elapsed();
497 match reply {
498 Ok(Ok(reply)) => {
499 let code = match &reply.header.headline {
500 HttpHeadline::Response { status } => *status as u16,
501 // A response with a request headline is not an answer this
502 // can read, and an unreadable answer is not a healthy peer.
503 _ => 0,
504 };
505 if (200..300).contains(&code) {
506 let body = match String::from_utf8(reply.body.clone()) {
507 Ok(s) => match HealthBody::parse(&s) {
508 Ok(b) => Some(b),
509 // A 200 that does not parse as a health body is a peer that is up but
510 // served something else (no token, a plain page): still alive.
511 Err(_) => None,
512 },
513 Err(_) => None,
514 };
515 (Probe::Up(body), elapsed)
516 } else {
517 debug!("Watch: {} answered {}.", url, code);
518 (Probe::Refused(code), elapsed)
519 }
520 },
521 // An answer, so the path works, but not one this will hold: a failure like any refusal,
522 // and said at warn, since the outage it leads to would otherwise read as silence.
523 Ok(Err(e)) if e.tags().contains(&ErrTag::TooBig) => {
524 warn!("Watch: {} answered with more than a probe reads ({} bytes of body, {} of \
525 header), so the answer is refused: {}", url, PROBE_BODY_MAX, PROBE_HEADER_MAX, e);
526 (Probe::Refused(0), elapsed)
527 },
528 Ok(Err(e)) => {
529 debug!("Watch: {} did not answer: {}", url, e);
530 (Probe::Silent, elapsed)
531 },
532 Err(_) => {
533 debug!("Watch: {} did not answer within {}s.", url, timeout.as_secs());
534 (Probe::Silent, elapsed)
535 },
536 }
537}
538
539/// Fold one round of probes into what this node believes, returning the alerts owed and the
540/// samples the dashboard keeps.
541///
542/// `round` holds one probe, with the time it took, per entry of `cfg.peers`, in order; `state`
543/// holds one set of beliefs per entry in the same order. A round that is this watcher's own link
544/// (see [`is_own_link_down`]) is not judged at all: no failure is counted against any peer, and
545/// none is forgotten, so a link that drops for an hour yields neither `DOWN` messages that cannot
546/// be sent nor, afterwards, `is back` messages about outages that never happened. Nor does such
547/// a round reach the dashboard: each peer's ring keeps the last readings actually heard, which
548/// the page greys as they age, rather than an hour of silence displacing them. `link_down`
549/// carries the verdict from round to round, so each episode is logged once.
550///
551/// Separated from the polling so the rule can be tested without a network, on the same path the
552/// loop takes; each peer of a judged round goes through [`assess`].
553fn judge_round(
554 cfg: &WatchConfig,
555 whoami: &str,
556 state: &mut [PeerState],
557 link_down: &mut bool,
558 round: Vec<(Probe, Duration)>,
559 now: Instant,
560 t_secs: u64,
561)
562 -> (Vec<AlertEvent>, Vec<(usize, ProbeSample)>)
563{
564 let (probes, took): (Vec<Probe>, Vec<Duration>) = round.into_iter().unzip();
565 if is_own_link_down(&cfg.peers, state, &probes) {
566 if !*link_down {
567 *link_down = true;
568 warn!("Watch: none of the {} peers answered at all this round, across two or more \
569 hosts still thought up, so it is this watcher's own link that is down, not \
570 every peer at once. Nothing is judged until one answers.", probes.len());
571 }
572 return (Vec::new(), Vec::new());
573 }
574 if *link_down {
575 *link_down = false;
576 info!("Watch: peers are answering again, so judging resumes.");
577 }
578 let threshold = cfg.fail_threshold.max(1);
579 let repeat = Duration::from_secs(cfg.repeat_secs.max(60));
580 let mut events = Vec::new();
581 let mut samples = Vec::with_capacity(probes.len());
582 for (i, ((peer, probe), took)) in cfg.peers.iter().zip(probes).zip(took).enumerate() {
583 let st = match state.get_mut(i) {
584 Some(st) => st,
585 // Built one per peer at construction, so unreachable; a peer with no beliefs is
586 // one this cannot judge, and saying so beats judging it against another's.
587 None => {
588 warn!("Watch: no state for peer {} ('{}'); skipped.", i, peer.name);
589 continue;
590 },
591 };
592 let ok = probe.is_up();
593 let body = match probe {
594 Probe::Up(b) => b,
595 _ => None,
596 };
597 let (ev, sample) = assess(
598 threshold, whoami, peer, st, ok, body, took, repeat, now, t_secs);
599 events.extend(ev.into_iter().map(|e| (i, e)));
600 samples.push((i, sample));
601 }
602 (by_host(cfg, whoami, state, events, now), samples)
603}
604
605/// Fold a round's liveness events into one per host (D-06 audit A1).
606///
607/// A machine behind several entries -- its Steel and a forge copy, say -- dies as one, so it is
608/// told as one: a single `DOWN` or `is back` naming the entries, rather than a text for each.
609/// When any entry of a host is told `DOWN`, every entry of that host already down is named in
610/// the same message and its reminder restarted with it, so the host reminds on one cadence, the
611/// shortest of its entries', not once per entry. A host told about one entry alone reads
612/// exactly as before, and every other event passes through untouched.
613fn by_host(
614 cfg: &WatchConfig,
615 whoami: &str,
616 state: &mut [PeerState],
617 events: Vec<(usize, AlertEvent)>,
618 now: Instant,
619)
620 -> Vec<AlertEvent>
621{
622 let host_of = |i: usize| cfg.peers.get(i).map(|p| p.host.clone()).unwrap_or_default();
623 let mut out = Vec::with_capacity(events.len());
624 let mut told_down: BTreeSet<String> = BTreeSet::new(); // hosts given their one message
625 let mut told_back: BTreeSet<String> = BTreeSet::new();
626 for (i, ev) in &events {
627 let host = host_of(*i);
628 match ev {
629 AlertEvent::PeerDown { .. } => {
630 if !told_down.insert(host.clone()) {
631 continue;
632 }
633 let failures = events.iter()
634 .filter(|(j, _)| host_of(*j) == host)
635 .filter_map(|(_, e)| match e {
636 AlertEvent::PeerDown { failures, .. } => Some(*failures),
637 _ => None,
638 })
639 .max()
640 .unwrap_or(0);
641 // Every entry on the host that is down now, whether or not its own reminder
642 // fell due this round.
643 let mut names = Vec::new();
644 let mut urls = Vec::new();
645 let mut down_secs = 0;
646 for (j, p) in cfg.peers.iter().enumerate() {
647 if p.host != host {
648 continue;
649 }
650 if let Some(st) = state.get_mut(j) {
651 if st.health == Health::Down {
652 st.told_at = Some(now);
653 let secs = st.failing_at
654 .map(|t| now.duration_since(t).as_secs()).unwrap_or(0);
655 down_secs = down_secs.max(secs);
656 names.push(p.name.clone());
657 urls.push(p.url.clone());
658 }
659 }
660 }
661 if names.len() < 2 {
662 out.push(ev.clone());
663 } else {
664 out.push(AlertEvent::PeerDown {
665 peer: fmt!("{} ({})", host, names.join(", ")),
666 url: urls.join(", "),
667 failures,
668 down_secs,
669 noticed_by: whoami.to_string(),
670 });
671 }
672 },
673 AlertEvent::PeerRecovered { .. } => {
674 if !told_back.insert(host.clone()) {
675 continue;
676 }
677 let mut names = Vec::new();
678 let mut urls = Vec::new();
679 let mut away = 0;
680 for (j, e) in &events {
681 if let AlertEvent::PeerRecovered { peer, url, away_secs, .. } = e {
682 if host_of(*j) == host {
683 names.push(peer.clone());
684 urls.push(url.clone());
685 away = away.max(*away_secs);
686 }
687 }
688 }
689 if names.len() < 2 {
690 out.push(ev.clone());
691 } else {
692 out.push(AlertEvent::PeerRecovered {
693 peer: fmt!("{} ({})", host, names.join(", ")),
694 url: urls.join(", "),
695 away_secs: away,
696 noticed_by: whoami.to_string(),
697 });
698 }
699 },
700 _ => out.push(ev.clone()),
701 }
702 }
703 out
704}
705
706/// One probe result in; the alerts it raises and the sample the dashboard keeps out.
707///
708/// The probe time is measured here rather than reported by the peer, and it is written into the
709/// body before the body is judged, so a `probe_ms` threshold in a peer's `distress` map is an
710/// ordinary class with no special case. A peer that serves no body has nothing to carry it and
711/// no distress to judge; the sample still records the time.
712fn assess(
713 threshold: u32,
714 whoami: &str,
715 peer: &WatchPeer,
716 st: &mut PeerState,
717 ok: bool,
718 body: Option<HealthBody>,
719 took: Duration,
720 repeat: Duration,
721 now: Instant,
722 t_secs: u64,
723)
724 -> (Vec<AlertEvent>, ProbeSample)
725{
726 let probe_ms = took.as_millis().min(i64::MAX as u128) as u64;
727 let body = body.map(|mut b| {
728 b.set(F_PROBE_MS, probe_ms as i64);
729 b
730 });
731 let events = decide(threshold, whoami, peer, st, ok, body.clone(), repeat, now);
732 let sample = ProbeSample {
733 t_secs,
734 ok,
735 probe_ms,
736 body,
737 health: st.health.into(),
738 };
739 (events, sample)
740}
741
742/// The peer-state transition, separated from the polling and the alerter.
743///
744/// This is the whole of the interesting behaviour -- when an operator is told, and about what --
745/// and it is a pure function of the current state and one probe result. It returns the alerts to
746/// raise rather than raising them, so a test can drive it directly with a `PeerState` and no
747/// network, no SMTP client, and no alerter standing up behind it.
748fn decide(
749 threshold: u32,
750 whoami: &str,
751 peer: &WatchPeer,
752 st: &mut PeerState,
753 ok: bool,
754 body: Option<HealthBody>,
755 repeat: Duration,
756 now: Instant,
757)
758 -> Vec<AlertEvent>
759{
760 // The peer's own cadence, where it names one, for both kinds of reminder; floored as the
761 // watcher's is, so no entry can remind every round.
762 let repeat = match peer.repeat_secs {
763 Some(s) => Duration::from_secs(s.max(60)),
764 None => repeat,
765 };
766 let mut out = Vec::new();
767 let name = &peer.name;
768 let url = &peer.url;
769 // Whether this peer even asks to be watched for distress: a token to open the gate and at
770 // least one threshold to test. A peer without both is plain up/down, exactly as before.
771 let watches_distress = peer.token.is_some() && !peer.distress.is_empty();
772
773 // A peer that asked for distress watching but gave us nothing to test is a misconfiguration
774 // worth saying out loud every poll, not silently reading as well: the operator believes the
775 // class is watched. Only warned on a live answer, since a dead peer's missing body is the
776 // outage itself, not a config fault.
777 if ok && watches_distress {
778 match &body {
779 None => warn!("Watch: {} answered but served no parseable health body, so its \
780 distress thresholds cannot be evaluated -- check the health token and path.",
781 name),
782 Some(b) => {
783 let missing: Vec<String> = peer.distress.keys()
784 .filter(|f| b.get(f).is_none())
785 .cloned()
786 .collect();
787 if !missing.is_empty() {
788 warn!("Watch: {}'s health body is missing distress field(s) {:?}, so \
789 those classes cannot be evaluated.", name, missing);
790 }
791 },
792 }
793 }
794
795 if ok {
796 // Any successful probe clears the liveness-failure marker; distress is a separate run.
797 if st.health == Health::Down {
798 let away = st.failing_at.map(|t| now.duration_since(t).as_secs()).unwrap_or(0);
799 *st = PeerState::default();
800 out.push(AlertEvent::PeerRecovered {
801 peer: name.clone(),
802 url: url.clone(),
803 away_secs: away,
804 noticed_by: whoami.to_string(),
805 });
806 return out;
807 }
808 st.failing_at = None;
809
810 // The classes currently over threshold, and whether we have a reading at all.
811 let breaching = match &body {
812 Some(b) if watches_distress => classes_breaching(&peer.distress, b),
813 _ => Vec::new(),
814 };
815 let have_reading = watches_distress && body.is_some();
816
817 match st.health {
818 Health::Up { .. } => {
819 if !have_reading {
820 // Healthy liveness, nothing to assess.
821 st.health = Health::Up { failures: 0 };
822 st.distress_rising = 0;
823 return out;
824 }
825 if breaching.is_empty() {
826 // Under threshold: the rising run is broken.
827 st.distress_rising = 0;
828 st.health = Health::Up { failures: 0 };
829 return out;
830 }
831 // Over threshold: distress needs the same run of agreeing evidence a death does.
832 st.distress_rising += 1;
833 if st.distress_rising >= threshold {
834 st.health = Health::Distressed { over: st.distress_rising, failures: 0 };
835 st.distress_at = Some(now);
836 st.distress_told_at = Some(now);
837 out.push(AlertEvent::PeerDistress {
838 peer: name.clone(),
839 url: url.clone(),
840 classes: classes_text(&breaching),
841 since_secs: 0,
842 noticed_by: whoami.to_string(),
843 });
844 } else {
845 st.health = Health::Up { failures: 0 };
846 }
847 },
848 Health::Distressed { over, .. } => {
849 // A live probe resets the distress-phase failure counter. Whether distress clears
850 // is a question only a reading under the clear boundary can answer; with no reading
851 // we hold distress rather than guess it away.
852 let cleared = match &body {
853 Some(b) if watches_distress => is_below_clear(peer, b),
854 _ => false,
855 };
856 if cleared {
857 let were = st.distress_at
858 .map(|t| now.duration_since(t).as_secs()).unwrap_or(0);
859 out.push(AlertEvent::PeerDistressCleared {
860 peer: name.clone(),
861 url: url.clone(),
862 were_secs: were,
863 noticed_by: whoami.to_string(),
864 });
865 *st = PeerState::default();
866 } else {
867 st.health = Health::Distressed { over, failures: 0 };
868 // Remind on the repeat interval, as a lasting outage does.
869 let due = st.distress_told_at
870 .map(|t| now.duration_since(t) >= repeat).unwrap_or(true);
871 if due && !breaching.is_empty() {
872 st.distress_told_at = Some(now);
873 let since = st.distress_at
874 .map(|t| now.duration_since(t).as_secs()).unwrap_or(0);
875 out.push(AlertEvent::PeerDistress {
876 peer: name.clone(),
877 url: url.clone(),
878 classes: classes_text(&breaching),
879 since_secs: since,
880 noticed_by: whoami.to_string(),
881 });
882 }
883 }
884 },
885 // Handled by the recovery return above; defensive and inert.
886 Health::Down => {},
887 }
888 return out;
889 }
890
891 // The probe failed.
892 match st.health {
893 Health::Up { failures } => {
894 let failures = failures + 1;
895 if st.failing_at.is_none() {
896 st.failing_at = Some(now);
897 }
898 if failures >= threshold {
899 st.health = Health::Down;
900 st.told_at = Some(now);
901 let down_secs = st.failing_at
902 .map(|t| now.duration_since(t).as_secs()).unwrap_or(0);
903 out.push(AlertEvent::PeerDown {
904 peer: name.clone(),
905 url: url.clone(),
906 failures,
907 down_secs,
908 noticed_by: whoami.to_string(),
909 });
910 } else {
911 st.health = Health::Up { failures };
912 }
913 },
914 Health::Distressed { over, failures } => {
915 // A distressed peer that stops answering reaches Down by the same rule an Up one does:
916 // the distress episode gives way to the outage it foreshadowed.
917 let failures = failures + 1;
918 if st.failing_at.is_none() {
919 st.failing_at = Some(now);
920 }
921 if failures >= threshold {
922 st.health = Health::Down;
923 st.told_at = Some(now);
924 let down_secs = st.failing_at
925 .map(|t| now.duration_since(t).as_secs()).unwrap_or(0);
926 out.push(AlertEvent::PeerDown {
927 peer: name.clone(),
928 url: url.clone(),
929 failures,
930 down_secs,
931 noticed_by: whoami.to_string(),
932 });
933 } else {
934 st.health = Health::Distressed { over, failures };
935 }
936 },
937 Health::Down => {
938 // Still down. Remind, but only on the repeat interval -- an alarm that fires every
939 // poll is an alarm that gets silenced, and the SMS leg of this costs money per message.
940 let due = st.told_at.map(|t| now.duration_since(t) >= repeat).unwrap_or(true);
941 if due {
942 st.told_at = Some(now);
943 let down_secs = st.failing_at
944 .map(|t| now.duration_since(t).as_secs()).unwrap_or(0);
945 out.push(AlertEvent::PeerDown {
946 peer: name.clone(),
947 url: url.clone(),
948 failures: threshold,
949 down_secs,
950 noticed_by: whoami.to_string(),
951 });
952 }
953 },
954 }
955 out
956}
957
958
959#[cfg(test)]
960mod tests {
961 use super::*;
962
963 /// A plain up/down peer: no token and no thresholds, so distress is never assessed.
964 fn plain_peer(name: &str, url: &str) -> WatchPeer {
965 WatchPeer {
966 name: name.to_string(),
967 host: name.to_string(),
968 url: url.to_string(),
969 plain_ok: false,
970 distress: BTreeMap::new(),
971 clear: BTreeMap::new(),
972 token: None,
973 repeat_secs: None,
974 }
975 }
976
977 /// A peer that watches distress: a token opens the gate, and one or more thresholds are tested.
978 fn distress_peer(
979 name: &str,
980 distress: &[(&str, i64)],
981 clear: &[(&str, i64)],
982 )
983 -> WatchPeer
984 {
985 let mut p = plain_peer(name, "https://example.test/_steel/health");
986 p.token = Some(fmt!("a-shared-secret"));
987 p.distress = distress.iter().map(|(k, v)| (k.to_string(), *v)).collect();
988 p.clear = clear.iter().map(|(k, v)| (k.to_string(), *v)).collect();
989 p
990 }
991
992 fn body_of(fields: &[(&str, i64)]) -> HealthBody {
993 let mut b = HealthBody::new();
994 for (k, v) in fields {
995 b.set(k, *v);
996 }
997 b
998 }
999
1000 /// The state machine, without a network.
1001 ///
1002 /// `decide` is driven directly and no alerter stands up, so what is under test is the rule
1003 /// about when an operator is told -- which is the part that costs money when it is wrong in
1004 /// one direction and costs an outage when it is wrong in the other.
1005 fn machine(threshold: u32) -> (WatchConfig, PeerState) {
1006 let cfg = WatchConfig {
1007 enabled: true,
1008 peers: vec![plain_peer("jarrah", "https://example.test/api/health")],
1009 interval_secs: 60,
1010 fail_threshold: threshold,
1011 timeout_secs: 10,
1012 repeat_secs: 900,
1013 heartbeat_secs: 2_592_000,
1014 };
1015 (cfg, PeerState::default())
1016 }
1017
1018 /// A single failure is not an outage. This is the property that keeps a flaky minute from
1019 /// waking somebody, and it is the one most likely to be tuned away by accident.
1020 #[test]
1021 fn one_failure_below_the_threshold_is_not_yet_an_outage() {
1022 let (cfg, mut st) = machine(3);
1023 assert_eq!(st.health, Health::Up { failures: 0 });
1024 // Two failures against a threshold of three: still up, and counting.
1025 for expected in 1..=2u32 {
1026 let Health::Up { failures } = st.health else { panic!("went down too early") };
1027 st.health = Health::Up { failures: failures + 1 };
1028 assert_eq!(st.health, Health::Up { failures: expected });
1029 }
1030 assert!(cfg.fail_threshold == 3);
1031 }
1032
1033 /// A success clears the count, so isolated failures never accumulate into a false outage.
1034 #[test]
1035 fn a_success_forgets_the_run_rather_than_carrying_it() {
1036 let (_cfg, mut st) = machine(3);
1037 st.health = Health::Up { failures: 2 };
1038 st.failing_at = Some(Instant::now());
1039 st = PeerState::default();
1040 assert_eq!(st.health, Health::Up { failures: 0 });
1041 assert!(st.failing_at.is_none(), "a recovered peer still remembered when it failed");
1042 }
1043
1044 #[test]
1045 fn a_peer_list_is_data_so_the_estate_can_change_without_a_rebuild() {
1046 let (mut cfg, _) = machine(2);
1047 cfg.peers.push(plain_peer("conifer", "https://ontheism.org/health"));
1048 assert_eq!(cfg.peers.len(), 2,
1049 "adding a machine must be a configuration change and nothing else");
1050 }
1051
1052 /// The default must be the strict one, because a config that says nothing about a wire is a
1053 /// config whose author never thought about it.
1054 #[test]
1055 fn a_plain_url_is_refused_unless_the_operator_wrote_it_down() {
1056 let refused = vet(&plain_peer("birch", "http://65.21.145.109:9109/forge/fresh"));
1057 assert!(refused.is_err(), "a plain http peer was accepted without plain_ok");
1058 let msg = fmt!("{}", refused.err().unwrap());
1059 assert!(msg.contains("plain_ok"),
1060 "the refusal must name the key that would allow it, or the operator has to read \
1061 the source to find out: got '{}'", msg);
1062
1063 let mut allowed = plain_peer("birch", "http://65.21.145.109:9109/forge/fresh");
1064 allowed.plain_ok = true;
1065 assert!(vet(&allowed).is_ok(),
1066 "a plain http peer was refused although plain_ok was set");
1067 }
1068
1069 /// A malformed URL is refused whatever the peer says about its wire, so `plain_ok` cannot be
1070 /// read as a general relaxation.
1071 #[test]
1072 fn plain_ok_does_not_excuse_a_url_that_cannot_be_parsed() {
1073 let mut p = plain_peer("nowhere", "not a url at all");
1074 p.plain_ok = true;
1075 assert!(vet(&p).is_err(), "plain_ok let an unparseable URL through");
1076 }
1077
1078 // ── Distress state machine (A6) ──────────────────────────────────────────
1079
1080 fn kinds(events: &[AlertEvent]) -> Vec<&'static str> {
1081 events.iter().map(|e| match e {
1082 AlertEvent::PeerDown { .. } => "down",
1083 AlertEvent::PeerRecovered { .. } => "recovered",
1084 AlertEvent::PeerDistress { .. } => "distress",
1085 AlertEvent::PeerDistressCleared { .. } => "cleared",
1086 _ => "other",
1087 }).collect()
1088 }
1089
1090 /// A peer stays up until a run of over-threshold readings as long as the fail threshold, then
1091 /// enters distress once -- not on the first reading, and not repeatedly.
1092 #[test]
1093 fn up_to_distressed_needs_the_same_run_as_a_death() {
1094 let peer = distress_peer("jarrah", &[("mem_pct", 90)], &[("mem_pct", 80)]);
1095 let mut st = PeerState::default();
1096 let now = Instant::now();
1097 let hot = body_of(&[("mem_pct", 94)]);
1098
1099 // Two over-threshold readings against a threshold of three: still up, no alert.
1100 for _ in 0..2 {
1101 let ev = decide(3, "argonaut", &peer, &mut st, true, Some(hot.clone()),
1102 Duration::from_secs(900), now);
1103 assert!(ev.is_empty(), "distress fired before the run was long enough");
1104 assert!(matches!(st.health, Health::Up { .. }));
1105 }
1106 // The third tips it into distress, exactly once.
1107 let ev = decide(3, "argonaut", &peer, &mut st, true, Some(hot.clone()),
1108 Duration::from_secs(900), now);
1109 assert_eq!(kinds(&ev), vec!["distress"]);
1110 assert!(matches!(st.health, Health::Distressed { .. }));
1111 }
1112
1113 /// Hysteresis: once distressed, a reading must fall under the *clear* boundary, not merely back
1114 /// under the entry threshold, before the peer is called well again.
1115 #[test]
1116 fn distress_clears_only_under_the_clear_boundary() {
1117 let peer = distress_peer("jarrah", &[("mem_pct", 90)], &[("mem_pct", 80)]);
1118 let mut st = PeerState::default();
1119 let now = Instant::now();
1120 // Enter distress at threshold 1 for brevity.
1121 let _ = decide(1, "argonaut", &peer, &mut st, true, Some(body_of(&[("mem_pct", 95)])),
1122 Duration::from_secs(900), now);
1123 assert!(matches!(st.health, Health::Distressed { .. }));
1124
1125 // 85 is under the 90 entry threshold but still over the 80 clear boundary: holds distress.
1126 let ev = decide(1, "argonaut", &peer, &mut st, true, Some(body_of(&[("mem_pct", 85)])),
1127 Duration::from_secs(900), now);
1128 assert!(matches!(st.health, Health::Distressed { .. }),
1129 "distress cleared before crossing the clear boundary -- no hysteresis");
1130 assert!(!kinds(&ev).contains(&"cleared"));
1131
1132 // 78 is under the clear boundary: now it clears.
1133 let ev = decide(1, "argonaut", &peer, &mut st, true, Some(body_of(&[("mem_pct", 78)])),
1134 Duration::from_secs(900), now);
1135 assert_eq!(kinds(&ev), vec!["cleared"]);
1136 assert!(matches!(st.health, Health::Up { .. }));
1137 }
1138
1139 /// A distressed peer that stops answering reaches Down by the failure counter, the same rule an
1140 /// Up peer follows.
1141 #[test]
1142 fn distressed_to_down_via_the_failure_counter() {
1143 let peer = distress_peer("jarrah", &[("mem_pct", 90)], &[("mem_pct", 80)]);
1144 let mut st = PeerState::default();
1145 let now = Instant::now();
1146 let _ = decide(2, "argonaut", &peer, &mut st, true, Some(body_of(&[("mem_pct", 95)])),
1147 Duration::from_secs(900), now);
1148 let _ = decide(2, "argonaut", &peer, &mut st, true, Some(body_of(&[("mem_pct", 95)])),
1149 Duration::from_secs(900), now);
1150 assert!(matches!(st.health, Health::Distressed { .. }));
1151
1152 // One failure: still distressed, counting.
1153 let ev = decide(2, "argonaut", &peer, &mut st, false, None,
1154 Duration::from_secs(900), now);
1155 assert!(ev.is_empty());
1156 assert!(matches!(st.health, Health::Distressed { failures: 1, .. }));
1157 // Second failure reaches the threshold: Down.
1158 let ev = decide(2, "argonaut", &peer, &mut st, false, None,
1159 Duration::from_secs(900), now);
1160 assert_eq!(kinds(&ev), vec!["down"]);
1161 assert_eq!(st.health, Health::Down);
1162 }
1163
1164 /// The probe time the watcher measured reaches `decide` as an ordinary field: a slow answer
1165 /// over the peer's `probe_ms` threshold is distress, named with the time measured here, and
1166 /// a figure the peer put in its own body under that name is not believed.
1167 #[test]
1168 fn probe_ms_is_folded_in_before_the_body_is_judged() {
1169 let peer = distress_peer("jarrah", &[("probe_ms", 3000)], &[("probe_ms", 1500)]);
1170 let mut st = PeerState::default();
1171 let now = Instant::now();
1172 let repeat = Duration::from_secs(900);
1173 // The peer claims a 1 ms probe; the watcher measured 3.5 s.
1174 let claimed = body_of(&[("mem_pct", 40), ("probe_ms", 1)]);
1175
1176 let (ev, sample) = assess(1, "karri", &peer, &mut st, true, Some(claimed),
1177 Duration::from_millis(3_500), repeat, now, 1_000);
1178 assert_eq!(kinds(&ev), vec!["distress"],
1179 "a probe over its threshold must reach decide and raise distress");
1180 match &ev[0] {
1181 AlertEvent::PeerDistress { classes, .. } => assert_eq!(classes, "probe_ms 3500"),
1182 _ => panic!("expected a distress event"),
1183 }
1184 assert_eq!(sample.probe_ms, 3_500);
1185 assert_eq!(sample.body.as_ref().and_then(|b| b.get("probe_ms")), Some(3_500),
1186 "the body the dashboard keeps must carry the measured time, not the claimed one");
1187 assert_eq!(sample.health, PeerHealth::Distressed);
1188 assert_eq!(sample.t_secs, 1_000);
1189
1190 // A fast answer under the clear boundary lifts it.
1191 let (ev, sample) = assess(1, "karri", &peer, &mut st, true,
1192 Some(body_of(&[("mem_pct", 40)])), Duration::from_millis(200), repeat, now, 1_060);
1193 assert_eq!(kinds(&ev), vec!["cleared"]);
1194 assert_eq!(sample.health, PeerHealth::Up);
1195
1196 // No body, nothing to judge: the time is still kept.
1197 let (ev, sample) = assess(1, "karri", &peer, &mut st, true, None,
1198 Duration::from_millis(4_000), repeat, now, 1_120);
1199 assert!(ev.is_empty());
1200 assert!(sample.body.is_none());
1201 assert_eq!(sample.probe_ms, 4_000);
1202 }
1203
1204 /// A peer with no thresholds is plain up/down: a hot body never makes it distressed.
1205 #[test]
1206 fn a_peer_without_thresholds_never_goes_distressed() {
1207 let peer = plain_peer("jarrah", "https://example.test/api/health");
1208 let mut st = PeerState::default();
1209 let now = Instant::now();
1210 for _ in 0..5 {
1211 let ev = decide(1, "argonaut", &peer, &mut st, true, Some(body_of(&[("mem_pct", 99)])),
1212 Duration::from_secs(900), now);
1213 assert!(ev.is_empty());
1214 assert!(matches!(st.health, Health::Up { .. }));
1215 }
1216 }
1217
1218 // ── Per-peer reminder cadence (F2) ──────────────────────────────────────
1219
1220 /// A peer's own `repeat_secs` sets its reminders, down or distressed, in place of the
1221 /// watcher's; a peer without one keeps the watcher's.
1222 #[test]
1223 fn a_peers_own_repeat_overrides_the_watchers() {
1224 let global = Duration::from_secs(900);
1225 let t0 = Instant::now();
1226 let at = |secs: u64| t0 + Duration::from_secs(secs);
1227
1228 let mut slow = plain_peer("jarrah forge copy", "https://example.test/_steel/health");
1229 slow.repeat_secs = Some(21_600);
1230 let quick = plain_peer("jarrah", "https://example.test/api/health");
1231 let (mut st_slow, mut st_quick) = (PeerState::default(), PeerState::default());
1232 assert_eq!(kinds(&decide(1, "conifer", &slow, &mut st_slow, false, None, global, t0)),
1233 vec!["down"]);
1234 assert_eq!(kinds(&decide(1, "conifer", &quick, &mut st_quick, false, None, global, t0)),
1235 vec!["down"]);
1236
1237 // At the watcher's fifteen minutes the plain peer reminds and the slow one does not.
1238 assert_eq!(kinds(&decide(1, "conifer", &quick, &mut st_quick, false, None, global,
1239 at(900))), vec!["down"]);
1240 assert!(decide(1, "conifer", &slow, &mut st_slow, false, None, global, at(900)).is_empty(),
1241 "the peer's own six hours must beat the watcher's fifteen minutes");
1242 assert!(decide(1, "conifer", &slow, &mut st_slow, false, None, global, at(21_599))
1243 .is_empty());
1244 assert_eq!(kinds(&decide(1, "conifer", &slow, &mut st_slow, false, None, global,
1245 at(21_600))), vec!["down"], "and it must still remind once its own interval is up");
1246
1247 // Distress reminders follow the same cadence.
1248 let mut hot = distress_peer("jarrah forge copy", &[("forge_state_age_s", 10_800)], &[]);
1249 hot.repeat_secs = Some(21_600);
1250 let mut st = PeerState::default();
1251 let stale = body_of(&[("forge_state_age_s", 12_000)]);
1252 assert_eq!(kinds(&decide(1, "conifer", &hot, &mut st, true, Some(stale.clone()), global,
1253 t0)), vec!["distress"]);
1254 assert!(decide(1, "conifer", &hot, &mut st, true, Some(stale.clone()), global, at(900))
1255 .is_empty(), "a distress reminder at the watcher's cadence, not the peer's");
1256 assert_eq!(kinds(&decide(1, "conifer", &hot, &mut st, true, Some(stale), global,
1257 at(21_600))), vec!["distress"]);
1258 }
1259
1260 /// A cadence under a minute is floored, as the watcher's own is, so no entry can remind
1261 /// every round.
1262 #[test]
1263 fn a_peers_repeat_is_floored_at_a_minute() {
1264 let t0 = Instant::now();
1265 let mut p = plain_peer("jarrah", "https://example.test/api/health");
1266 p.repeat_secs = Some(0);
1267 let mut st = PeerState::default();
1268 let _ = decide(1, "conifer", &p, &mut st, false, None, Duration::from_secs(900), t0);
1269 assert!(decide(1, "conifer", &p, &mut st, false, None, Duration::from_secs(900),
1270 t0 + Duration::from_secs(59)).is_empty());
1271 assert_eq!(kinds(&decide(1, "conifer", &p, &mut st, false, None, Duration::from_secs(900),
1272 t0 + Duration::from_secs(60))), vec!["down"]);
1273 }
1274
1275 // ── The watcher's own link (F5) ─────────────────────────────────────────
1276
1277 fn watch_of(names: &[&str], threshold: u32) -> WatchConfig {
1278 let (mut cfg, _) = machine(threshold);
1279 cfg.peers = names.iter()
1280 .map(|n| plain_peer(n, &fmt!("https://{}.example.test/health", n)))
1281 .collect();
1282 cfg
1283 }
1284
1285 fn states_for(cfg: &WatchConfig) -> Vec<PeerState> {
1286 cfg.peers.iter().map(|_| PeerState::default()).collect()
1287 }
1288
1289 /// A round as the loop hands it over, each probe with an ordinary measured time.
1290 fn round_of(probes: Vec<Probe>) -> Vec<(Probe, Duration)> {
1291 probes.into_iter().map(|p| (p, Duration::from_millis(80))).collect()
1292 }
1293
1294 /// Rounds in which no peer answered at all judge nothing: no `DOWN` however long it lasts,
1295 /// no `is back` when the link returns, and no failure counted against any peer.
1296 #[test]
1297 fn a_round_in_which_nobody_answered_judges_nothing() {
1298 let cfg = watch_of(&["karri", "jarrah", "birch"], 3);
1299 let mut state = states_for(&cfg);
1300 let mut link_down = false;
1301 let t0 = Instant::now();
1302 for i in 0..30u64 {
1303 let (ev, samples) = judge_round(&cfg, "argonaut", &mut state, &mut link_down,
1304 round_of(vec![Probe::Silent, Probe::Silent, Probe::Silent]),
1305 t0 + Duration::from_secs(60 * i), 1_000 + 60 * i);
1306 assert!(ev.is_empty(), "round {} raised {:?} while the watcher's own link was down",
1307 i, kinds(&ev));
1308 assert!(samples.is_empty(),
1309 "a round that was the watcher's own link displaced readings in the dashboard's rings");
1310 assert!(link_down);
1311 }
1312 assert!(state.iter().all(|st| st.health == Health::Up { failures: 0 }),
1313 "a round that was the watcher's own link counted a failure against a peer");
1314
1315 let (ev, samples) = judge_round(&cfg, "argonaut", &mut state, &mut link_down,
1316 round_of(vec![Probe::Up(None), Probe::Up(None), Probe::Up(None)]),
1317 t0 + Duration::from_secs(1_860), 2_860);
1318 assert!(ev.is_empty(), "the link coming back is not a recovery of anything: {:?}",
1319 kinds(&ev));
1320 assert_eq!(samples.iter().map(|(i, _)| *i).collect::<Vec<_>>(), vec![0, 1, 2],
1321 "once judging resumes, every peer's probe is kept, each under its own position");
1322 assert!(!link_down);
1323 }
1324
1325 /// One peer failing while the others answer still reaches `DOWN` at the usual threshold, and
1326 /// a round of the watcher's own silence in between neither advances nor forgets its count.
1327 #[test]
1328 fn one_silent_peer_still_alarms() {
1329 let cfg = watch_of(&["karri", "jarrah", "birch"], 3);
1330 let mut state = states_for(&cfg);
1331 let mut link_down = false;
1332 let t0 = Instant::now();
1333 let jarrah_silent = || round_of(vec![Probe::Up(None), Probe::Silent, Probe::Up(None)]);
1334
1335 for i in 0..2u64 {
1336 let (ev, _) = judge_round(&cfg, "conifer", &mut state, &mut link_down, jarrah_silent(),
1337 t0 + Duration::from_secs(60 * i), 1_000 + 60 * i);
1338 assert!(ev.is_empty());
1339 }
1340 let (ev, _) = judge_round(&cfg, "conifer", &mut state, &mut link_down,
1341 round_of(vec![Probe::Silent, Probe::Silent, Probe::Silent]),
1342 t0 + Duration::from_secs(120), 1_120);
1343 assert!(ev.is_empty() && link_down);
1344
1345 let (ev, samples) = judge_round(&cfg, "conifer", &mut state, &mut link_down,
1346 jarrah_silent(), t0 + Duration::from_secs(180), 1_180);
1347 assert_eq!(kinds(&ev), vec!["down"], "the third judged failure must call jarrah down");
1348 match &ev[0] {
1349 AlertEvent::PeerDown { peer, .. } => assert_eq!(peer, "jarrah"),
1350 other => panic!("expected jarrah down, got {:?}", other),
1351 }
1352 // The dashboard keeps the belief the alarm reached, beside the peers that answered.
1353 let (i, jarrah) = &samples[1];
1354 assert_eq!(*i, 1);
1355 assert!(!jarrah.ok);
1356 assert_eq!(jarrah.health, PeerHealth::Down);
1357 assert_eq!(jarrah.t_secs, 1_180);
1358 assert_eq!(samples[0].1.health, PeerHealth::Up);
1359 assert_eq!(samples[0].1.probe_ms, 80, "the measured time must reach the sample");
1360 }
1361
1362 /// A refusal is an answer, so a round in which every peer refused is news about the peers --
1363 /// every Steel crash-looping behind its proxy, say -- and is judged.
1364 #[test]
1365 fn a_round_of_refusals_is_still_judged() {
1366 let cfg = watch_of(&["karri", "jarrah"], 2);
1367 let mut state = states_for(&cfg);
1368 let mut link_down = false;
1369 let t0 = Instant::now();
1370 let (_, samples) = judge_round(&cfg, "conifer", &mut state, &mut link_down,
1371 round_of(vec![Probe::Refused(502), Probe::Refused(503)]), t0, 1_000);
1372 assert_eq!(samples.len(), 2, "a round of refusals is news, so it reaches the dashboard");
1373 let (ev, _) = judge_round(&cfg, "conifer", &mut state, &mut link_down,
1374 round_of(vec![Probe::Refused(502), Probe::Silent]), t0 + Duration::from_secs(60),
1375 1_060);
1376 assert_eq!(kinds(&ev), vec!["down", "down"]);
1377 assert!(!link_down);
1378 }
1379
1380 /// With a single peer the watcher cannot tell its own link from the peer's death, and the
1381 /// outage is not the thing to explain away.
1382 #[test]
1383 fn a_single_silent_peer_is_judged() {
1384 let cfg = watch_of(&["jarrah"], 3);
1385 let mut state = states_for(&cfg);
1386 let mut link_down = false;
1387 let t0 = Instant::now();
1388 let mut raised = Vec::new();
1389 for i in 0..3u64 {
1390 let (ev, _) = judge_round(&cfg, "karri", &mut state, &mut link_down,
1391 round_of(vec![Probe::Silent]), t0 + Duration::from_secs(60 * i), 1_000 + 60 * i);
1392 raised.extend(ev);
1393 }
1394 assert_eq!(kinds(&raised), vec!["down"]);
1395 assert!(!link_down);
1396 }
1397
1398 /// A peer already `DOWN` is silent anyway, so it is no evidence about the watcher's link:
1399 /// with jarrah down, birch falling silent is birch's outage, and is told (D-06 audit W1).
1400 #[test]
1401 fn a_peer_already_down_does_not_hide_the_next_outage() {
1402 let cfg = watch_of(&["jarrah", "birch"], 2);
1403 let mut state = states_for(&cfg);
1404 let mut link_down = false;
1405 let t0 = Instant::now();
1406 let mut raised = Vec::new();
1407 for i in 0..2u64 {
1408 let (ev, _) = judge_round(&cfg, "karri", &mut state, &mut link_down,
1409 round_of(vec![Probe::Silent, Probe::Up(None)]),
1410 t0 + Duration::from_secs(60 * i), 1_000 + 60 * i);
1411 raised.extend(ev);
1412 }
1413 assert_eq!(kinds(&raised), vec!["down"]);
1414 assert_eq!(state[0].health, Health::Down);
1415
1416 let mut raised = Vec::new();
1417 for i in 2..4u64 {
1418 let (ev, _) = judge_round(&cfg, "karri", &mut state, &mut link_down,
1419 round_of(vec![Probe::Silent, Probe::Silent]),
1420 t0 + Duration::from_secs(60 * i), 1_000 + 60 * i);
1421 assert!(!link_down, "round {}: birch's silence was read as karri's own link", i);
1422 raised.extend(ev);
1423 }
1424 match raised.as_slice() {
1425 [AlertEvent::PeerDown { peer, .. }] => assert_eq!(peer, "birch"),
1426 other => panic!("expected birch down alone, got {:?}", kinds(other)),
1427 }
1428 assert_eq!(state[1].health, Health::Down);
1429 }
1430
1431 /// Two entries on one machine are one host: when it dies both fall silent together, and
1432 /// that is its death, not the watcher's link (D-06 audit W1), told once (A1).
1433 #[test]
1434 fn two_entries_on_one_host_are_judged_as_one_host() {
1435 let mut cfg = watch_of(&["jarrah", "forge"], 2);
1436 cfg.peers[1].host = fmt!("jarrah");
1437 let mut state = states_for(&cfg);
1438 let mut link_down = false;
1439 let t0 = Instant::now();
1440 let mut raised = Vec::new();
1441 for i in 0..2u64 {
1442 let (ev, _) = judge_round(&cfg, "karri", &mut state, &mut link_down,
1443 round_of(vec![Probe::Silent, Probe::Silent]),
1444 t0 + Duration::from_secs(60 * i), 1_000 + 60 * i);
1445 raised.extend(ev);
1446 }
1447 assert_eq!(kinds(&raised), vec!["down"]);
1448 assert!(state.iter().all(|st| st.health == Health::Down));
1449 assert!(!link_down);
1450 }
1451
1452 /// One host's death is one text, and so is each reminder and the recovery, however many
1453 /// entries sit on it and whatever their own cadences (D-06 audit A1). An entry on the same
1454 /// host that fails alone is still told alone, under its own name.
1455 #[test]
1456 fn a_hosts_death_is_told_once_however_many_entries_it_has() {
1457 let mut cfg = watch_of(&["jarrah", "forge", "birch"], 2);
1458 cfg.peers[1].host = fmt!("jarrah");
1459 cfg.peers[1].repeat_secs = Some(3_600);
1460 let mut state = states_for(&cfg);
1461 let mut link_down = false;
1462 let t0 = Instant::now();
1463 let jarrah_dead = || round_of(vec![Probe::Silent, Probe::Silent, Probe::Up(None)]);
1464
1465 // Two hours of jarrah dead, one round a minute.
1466 let mut texts = 0;
1467 for i in 0..120u64 {
1468 let (ev, _) = judge_round(&cfg, "karri", &mut state, &mut link_down, jarrah_dead(),
1469 t0 + Duration::from_secs(60 * i), 1_000 + 60 * i);
1470 assert!(ev.len() <= 1, "round {} sent {} texts for one host", i, ev.len());
1471 for e in &ev {
1472 match e {
1473 AlertEvent::PeerDown { peer, url, .. } => {
1474 assert_eq!(peer, "jarrah (jarrah, forge)");
1475 assert!(url.contains("jarrah.example") && url.contains("forge.example"));
1476 texts += 1;
1477 },
1478 other => panic!("round {}: expected jarrah down, got {:?}", i, kinds(&[other.clone()])),
1479 }
1480 }
1481 }
1482 // Told at the second round, then reminded every 15 minutes: jarrah's cadence, the
1483 // shortest, restarts forge's each time rather than adding its own.
1484 assert_eq!(texts, 1 + (119 - 1) / 15);
1485
1486 let (ev, _) = judge_round(&cfg, "karri", &mut state, &mut link_down,
1487 round_of(vec![Probe::Up(None), Probe::Up(None), Probe::Up(None)]),
1488 t0 + Duration::from_secs(7_200), 8_200);
1489 match ev.as_slice() {
1490 [AlertEvent::PeerRecovered { peer, .. }] => assert_eq!(peer, "jarrah (jarrah, forge)"),
1491 other => panic!("expected one recovery for the host, got {:?}", kinds(other)),
1492 }
1493
1494 // The forge copy alone: the host is up, so the entry is told under its own name.
1495 let mut raised = Vec::new();
1496 for i in 0..2u64 {
1497 let (ev, _) = judge_round(&cfg, "karri", &mut state, &mut link_down,
1498 round_of(vec![Probe::Up(None), Probe::Silent, Probe::Up(None)]),
1499 t0 + Duration::from_secs(7_260 + 60 * i), 8_260 + 60 * i);
1500 raised.extend(ev);
1501 }
1502 match raised.as_slice() {
1503 [AlertEvent::PeerDown { peer, .. }] => assert_eq!(peer, "forge"),
1504 other => panic!("expected forge down alone, got {:?}", kinds(other)),
1505 }
1506 }
1507
1508 /// The rule still holds over the hosts left: with jarrah down, karri and birch silent
1509 /// together is the watcher's link, judged not at all.
1510 #[test]
1511 fn with_one_peer_down_every_live_host_silent_is_still_the_watchers_link() {
1512 let cfg = watch_of(&["jarrah", "karri", "birch"], 2);
1513 let mut state = states_for(&cfg);
1514 let mut link_down = false;
1515 let t0 = Instant::now();
1516 let mut raised = Vec::new();
1517 for i in 0..2u64 {
1518 let (ev, _) = judge_round(&cfg, "conifer", &mut state, &mut link_down,
1519 round_of(vec![Probe::Silent, Probe::Up(None), Probe::Up(None)]),
1520 t0 + Duration::from_secs(60 * i), 1_000 + 60 * i);
1521 raised.extend(ev);
1522 }
1523 assert_eq!(kinds(&raised), vec!["down"]);
1524 for i in 2..10u64 {
1525 let (ev, samples) = judge_round(&cfg, "conifer", &mut state, &mut link_down,
1526 round_of(vec![Probe::Silent, Probe::Silent, Probe::Silent]),
1527 t0 + Duration::from_secs(60 * i), 1_000 + 60 * i);
1528 assert!(ev.is_empty() && samples.is_empty() && link_down, "round {} was judged", i);
1529 }
1530 assert_eq!(state[1].health, Health::Up { failures: 0 });
1531 assert_eq!(state[2].health, Health::Up { failures: 0 });
1532 }
1533
1534 /// Serve one canned reply to the first connection on a loopback port, returning the port.
1535 async fn serve_once(reply: Vec<u8>) -> u16 {
1536 use tokio::io::{
1537 AsyncReadExt,
1538 AsyncWriteExt,
1539 };
1540 let listener = match tokio::net::TcpListener::bind("127.0.0.1:0").await {
1541 Ok(l) => l,
1542 Err(e) => panic!("no loopback listener: {}", e),
1543 };
1544 let port = match listener.local_addr() {
1545 Ok(a) => a.port(),
1546 Err(e) => panic!("no local address: {}", e),
1547 };
1548 tokio::spawn(async move {
1549 if let Ok((mut sock, _)) = listener.accept().await {
1550 let mut buf = [0u8; 4096];
1551 let _ = sock.read(&mut buf).await;
1552 let _ = sock.write_all(&reply).await;
1553 let _ = sock.shutdown().await;
1554 }
1555 });
1556 port
1557 }
1558
1559 /// A `200` carrying a well-formed health body of `fields` fields.
1560 fn health_reply(fields: usize) -> Vec<u8> {
1561 let mut b = HealthBody::new();
1562 for i in 0..fields {
1563 b.set(&fmt!("k{}", i), i as i64);
1564 }
1565 let body = b.to_json();
1566 let mut out = fmt!("HTTP/1.1 200 OK\r\nContent-Type: application/json\r\n\
1567 Content-Length: {}\r\n\r\n", body.len()).into_bytes();
1568 out.extend_from_slice(body.as_bytes());
1569 out
1570 }
1571
1572 /// The same `200`, sent chunked and so declaring no length, as a proxy streaming the page
1573 /// answers.
1574 fn chunked_health_reply(fields: usize) -> Vec<u8> {
1575 let mut b = HealthBody::new();
1576 for i in 0..fields {
1577 b.set(&fmt!("k{}", i), i as i64);
1578 }
1579 let body = b.to_json();
1580 let mut out = b"HTTP/1.1 200 OK\r\nContent-Type: application/json\r\n\
1581 Transfer-Encoding: chunked\r\n\r\n".to_vec();
1582 for piece in body.as_bytes().chunks(4096) {
1583 out.extend_from_slice(fmt!("{:x}\r\n", piece.len()).as_bytes());
1584 out.extend_from_slice(piece);
1585 out.extend_from_slice(b"\r\n");
1586 }
1587 out.extend_from_slice(b"0\r\n\r\n");
1588 out
1589 }
1590
1591 /// A probe reads a health body of ordinary size and refuses one far larger, since each body
1592 /// read is kept in the dashboard's ring for an hour (D-06 audit D1).
1593 #[tokio::test]
1594 async fn a_probe_reads_a_health_body_but_not_one_of_any_size() {
1595 let tls = match oxedyne_fe2o3_net::tls::default_client_config() {
1596 Ok(c) => Arc::new(c),
1597 Err(e) => panic!("no TLS client config: {}", e),
1598 };
1599 let small = health_reply(20);
1600 let large = health_reply(20_000);
1601 assert!(small.len() < PROBE_BODY_MAX && large.len() > 2 * PROBE_BODY_MAX);
1602
1603 let port = serve_once(small).await;
1604 let mut peer = plain_peer("small", &fmt!("http://127.0.0.1:{}/_steel/health", port));
1605 peer.plain_ok = true;
1606 match probe(&peer, Duration::from_secs(5), tls.clone()).await {
1607 (Probe::Up(Some(b)), _) => assert_eq!(b.fields.len(), 20),
1608 _ => panic!("a health body of ordinary size was not read"),
1609 }
1610
1611 let port = serve_once(large).await;
1612 let mut peer = plain_peer("large", &fmt!("http://127.0.0.1:{}/_steel/health", port));
1613 peer.plain_ok = true;
1614 match probe(&peer, Duration::from_secs(5), tls.clone()).await {
1615 (Probe::Up(Some(b)), _) => panic!("a body of {} fields, over {} bytes, was read \
1616 whole and would be kept", b.fields.len(), PROBE_BODY_MAX),
1617 // An answer, so not silence: it must not count towards the watcher's own link.
1618 (Probe::Refused(0), _) => (),
1619 (other, _) => panic!("an oversized answer read as {:?}, not a refusal", other),
1620 }
1621
1622 // Chunked, the body declares no length, so it is cut off part way through decoding, and
1623 // the overflow reaches the probe from beneath a wrapping frame (re-check B-H1).
1624 let port = serve_once(chunked_health_reply(20)).await;
1625 let mut peer = plain_peer("chunked", &fmt!("http://127.0.0.1:{}/_steel/health", port));
1626 peer.plain_ok = true;
1627 match probe(&peer, Duration::from_secs(5), tls.clone()).await {
1628 (Probe::Up(Some(b)), _) => assert_eq!(b.fields.len(), 20),
1629 (other, _) => panic!("a chunked health body of ordinary size read as {:?}", other),
1630 }
1631
1632 let port = serve_once(chunked_health_reply(20_000)).await;
1633 let mut peer = plain_peer("chunked", &fmt!("http://127.0.0.1:{}/_steel/health", port));
1634 peer.plain_ok = true;
1635 match probe(&peer, Duration::from_secs(5), tls).await {
1636 (Probe::Refused(0), _) => (),
1637 (other, _) => panic!("an oversized chunked answer read as {:?}, not a refusal", other),
1638 }
1639 }
1640
1641 /// Distress is told by mail alone, so a host with no mail recipient names the peers whose
1642 /// distress it would keep to itself, and a host with one names none.
1643 #[test]
1644 fn distress_on_a_host_without_mail_is_named_at_start() {
1645 let peers = vec![
1646 plain_peer("karri", "https://karri.example.test/"),
1647 distress_peer("jarrah", &[("mem_pct", 85)], &[("mem_pct", 70)]),
1648 ];
1649 assert_eq!(distress_unheard(&peers, false), vec![fmt!("jarrah")]);
1650 assert!(distress_unheard(&peers, true).is_empty());
1651 }
1652}