The controller verified a knowledge delivery, logged it, and returned OK. Nothing downstream ever heard about it, so a hive learned the repository had changed only by registering its own webhook — which is the last-writer-wins contention this issue is about. The event carries no payload. The hive-side handler this replaces reads two fields from Forgejo's push webhook and uses neither — both are filters — then runs `git pull`, which re-derives everything from the repository. What crosses the queue is an edge trigger, and fields would invent a contract nobody reads. One subject per hive, so the callout policy can express "this hive may hear its own events" at all; a subject with no hive component is the same subject for every hive. `ConfigPr` deliveries are deliberately not forwarded. A hive does not want to hear that a config PR opened — it wants to be told when to rebuild from main, which the controller decides after a merge rather than by relaying this delivery. That is deploy coordination's job, and the empty arm is there so the omission reads as scoped rather than forgotten. Fails soft: a missed announcement costs a hive stale knowledge until its next boot pull, which is the same cost as a webhook delivery to a hive that was down — what this replaces. A permission failure cannot be observed at the call site (a NATS authorization violation is reported asynchronously on the connection, reaching a client as a timeout or as nothing), so the doc says the flush proves only that the bytes left this process and points at the server log.
530 lines
20 KiB
Rust
530 lines
20 KiB
Rust
//! Swarm-wide view of what each hive last said about itself.
|
|
//!
|
|
//! Hives **offer** a snapshot upward; this daemon never reaches down to
|
|
//! collect one. That direction is the design, not an implementation
|
|
//! detail: the hive gateway has gone down in a way where every recovery
|
|
//! channel ran through the one broken thing, so a status path that
|
|
//! depended on the controller would have gone dark exactly when it was
|
|
//! needed to diagnose the controller's own network. A hive computes its
|
|
//! status locally regardless of whether the swarm can be reached.
|
|
//!
|
|
//! **The queue is the store.** A hive publishes into a `JetStream` KV
|
|
//! bucket, which retains the last value per key; this daemon reads that
|
|
//! bucket per request and keeps no copy. A cache here would be a second
|
|
//! answer to the same question, free to disagree with the first — and
|
|
//! the disagreement would surface as a hive reading healthy on a
|
|
//! dashboard while the bucket says otherwise.
|
|
//!
|
|
//! **Absence is the case this is built around** — the freshness states
|
|
//! and the reasoning behind each are in `docs/swarm/README.md`. The two
|
|
//! properties that constrain the code rather than describe it:
|
|
//! freshness is **derived at read time**, never stored (a stored
|
|
//! `healthy: bool` goes stale silently the moment nothing arrives), and
|
|
//! rows come from the **roster**, not the bucket, so an empty bucket
|
|
//! cannot render as a healthy swarm.
|
|
//!
|
|
//! One consequence worth stating because it is the opposite of what a
|
|
//! cache would give: losing the bucket degrades **to honesty**. Every
|
|
//! hive reads `never_reported` until its next publish, which is the true
|
|
//! answer — not a remembered "healthy" from before the loss.
|
|
|
|
use std::collections::BTreeMap;
|
|
use std::time::{Duration, SystemTime, UNIX_EPOCH};
|
|
|
|
use anyhow::{Context, Result};
|
|
use futures_util::TryStreamExt as _;
|
|
use serde::{Deserialize, Serialize};
|
|
use utoipa::ToSchema;
|
|
|
|
use crate::HiveEntry;
|
|
|
|
/// Default age past which a snapshot is reported stale.
|
|
///
|
|
/// A threshold is a statement about how often hives offer, and that
|
|
/// cadence is decided by the publisher (a later slice), so this is a
|
|
/// default to be overridden rather than a constant to be relied on.
|
|
pub const DEFAULT_STALE_AFTER: Duration = Duration::from_mins(2);
|
|
|
|
/// Env var the NixOS module sets from
|
|
/// `services.hyperhive.swarm.controller.staleAfterSeconds`.
|
|
pub const STALE_AFTER_ENV: &str = "SWARM_CONTROLLER_STALE_AFTER_SECS";
|
|
|
|
/// How a hive's last report reads *now* — a function of the clock, not a
|
|
/// property of the report.
|
|
#[derive(Clone, Copy, Debug, PartialEq, Eq, Serialize, Deserialize, ToSchema)]
|
|
#[serde(rename_all = "snake_case")]
|
|
pub enum Freshness {
|
|
/// Reported within the staleness threshold.
|
|
Fresh,
|
|
/// Reported, but longer ago than the threshold. The payload is still
|
|
/// rendered: "old" and "absent" are different answers and a consumer
|
|
/// may want the last thing a hive managed to say.
|
|
Stale,
|
|
/// In the roster, has never offered a snapshot. Distinct from
|
|
/// `Stale` because it separates "went quiet" from "never spoke" —
|
|
/// the first is a fault, the second is usually a deployment that
|
|
/// hasn't happened yet.
|
|
NeverReported,
|
|
/// Offered a snapshot but is not in the roster. Not an error this
|
|
/// daemon can resolve, and not one it should hide.
|
|
Unknown,
|
|
}
|
|
|
|
/// One row of the aggregate.
|
|
#[derive(Clone, Debug, PartialEq, Serialize, Deserialize, ToSchema)]
|
|
pub struct HiveStatus {
|
|
pub name: String,
|
|
/// From the roster; `None` for a hive the roster doesn't list.
|
|
pub domain: Option<String>,
|
|
pub freshness: Freshness,
|
|
/// When the snapshot was stored, unix seconds. `None` when nothing
|
|
/// has been.
|
|
///
|
|
/// This is the **bucket's** stamp, applied by the NATS server when
|
|
/// the value landed — not a field inside the payload. A publisher
|
|
/// therefore cannot make itself look fresher than it is, and a hive
|
|
/// with a wrong clock skews its own payload rather than this.
|
|
pub last_seen_unix: Option<i64>,
|
|
/// Age at render time. Carried alongside `last_seen_unix` so a
|
|
/// consumer with a different threshold need not re-derive it from a
|
|
/// clock that may not match this host's.
|
|
pub age_seconds: Option<u64>,
|
|
/// Whatever the hive published, unopened — stored opaquely so the
|
|
/// snapshot's contents stay settleable later without reworking the
|
|
/// aggregate.
|
|
///
|
|
/// `None` in two cases the consumer can tell apart by `freshness`:
|
|
/// nothing has ever been published (`never_reported`), or something
|
|
/// was published that is not JSON (any other freshness — the row
|
|
/// still reports *when* the hive last spoke, and the read logs a
|
|
/// warning naming it).
|
|
pub snapshot: Option<serde_json::Value>,
|
|
}
|
|
|
|
/// A snapshot as retained by the bucket.
|
|
#[derive(Clone, Debug)]
|
|
struct Offered {
|
|
received_at: SystemTime,
|
|
/// `None` when the stored bytes are not JSON — see [`HiveStatus::snapshot`].
|
|
payload: Option<serde_json::Value>,
|
|
}
|
|
|
|
/// Reads the aggregate out of the KV bucket.
|
|
///
|
|
/// Holds a NATS client rather than a bucket handle: the bucket is
|
|
/// resolved on first use and cached, so a controller that starts before
|
|
/// the bucket exists picks it up without a restart. Resolution failures
|
|
/// are not cached — [`tokio::sync::OnceCell::get_or_try_init`] retries —
|
|
/// which is what makes the queue coming up *after* this daemon a
|
|
/// non-event rather than a permanent degradation.
|
|
pub struct StatusReader {
|
|
client: async_nats::Client,
|
|
store: tokio::sync::OnceCell<async_nats::jetstream::kv::Store>,
|
|
stale_after: Duration,
|
|
}
|
|
|
|
impl StatusReader {
|
|
#[must_use]
|
|
pub fn new(client: async_nats::Client, stale_after: Duration) -> Self {
|
|
Self {
|
|
client,
|
|
store: tokio::sync::OnceCell::new(),
|
|
stale_after,
|
|
}
|
|
}
|
|
|
|
/// A handle on the queue connection this reader holds.
|
|
///
|
|
/// The controller has exactly **one** connection to the swarm queue and
|
|
/// more than one thing to do with it: status is read out of a KV bucket,
|
|
/// swarm events are published on a subject. Handing out a clone is cheap —
|
|
/// `async_nats::Client` is a handle, not a socket — and is strictly better
|
|
/// than opening a second connection, which would double the auth-callout
|
|
/// traffic and give the two paths independent reconnect state, so one could
|
|
/// be serving while the other was still down.
|
|
///
|
|
/// That this lives on the *status* reader is an accident of who constructed
|
|
/// the connection first, not a claim that events are a kind of status.
|
|
#[must_use]
|
|
pub fn queue_client(&self) -> async_nats::Client {
|
|
self.client.clone()
|
|
}
|
|
|
|
/// Reads [`STALE_AFTER_ENV`], falling back to
|
|
/// [`DEFAULT_STALE_AFTER`]. A zero or unparseable value takes the
|
|
/// default rather than failing startup — same rule as `load_hives`:
|
|
/// a controller whose own config is wrong must still serve.
|
|
#[must_use]
|
|
pub fn stale_after_from_env() -> Duration {
|
|
std::env::var(STALE_AFTER_ENV)
|
|
.ok()
|
|
.and_then(|raw| raw.trim().parse::<u64>().ok())
|
|
.filter(|secs| *secs > 0)
|
|
.map_or(DEFAULT_STALE_AFTER, Duration::from_secs)
|
|
}
|
|
|
|
/// The bucket handle, created on first use if nothing has made it yet.
|
|
///
|
|
/// The name and the creation config come from
|
|
/// [`swarm_queue_client::status`] rather than from here: the hive that
|
|
/// writes this bucket opens it with the same call, and a bucket both
|
|
/// ends may create is one both ends have to describe identically.
|
|
///
|
|
/// Whichever side arrives first creates it, and both sides ask for the
|
|
/// same shape, so this is a race with one outcome.
|
|
///
|
|
/// Returns the queue client's own error rather than an `anyhow::Error`:
|
|
/// `OnceCell::get_or_try_init` takes its error type from the closure,
|
|
/// so widening here would mean converting *inside* the closure for no
|
|
/// gain. `view` below `?`s it and anyhow converts there — which is the
|
|
/// whole point of the library keeping a typed error while the binary
|
|
/// keeps anyhow.
|
|
async fn store(
|
|
&self,
|
|
) -> std::result::Result<&async_nats::jetstream::kv::Store, swarm_queue_client::Error> {
|
|
self.store
|
|
.get_or_try_init(|| swarm_queue_client::status::open_or_create(&self.client))
|
|
.await
|
|
}
|
|
|
|
/// The aggregate, rendered against `now`.
|
|
///
|
|
/// Every roster hive produces a row whether or not it has ever
|
|
/// reported; a reporting hive outside the roster produces one too.
|
|
pub async fn view(&self, roster: &[HiveEntry], now: SystemTime) -> Result<Vec<HiveStatus>> {
|
|
// Before anything JetStream: an unconnected client does not fail a
|
|
// request, it hangs on it. The rule and the measurement behind it live
|
|
// in `swarm_queue_client::ensure_connected` — every consumer of the
|
|
// queue needs it, so it is not this daemon's to keep.
|
|
swarm_queue_client::ensure_connected(&self.client)?;
|
|
|
|
let store = self.store().await?;
|
|
|
|
// Keys first, then a fetch per key. The roster is a handful of
|
|
// hives, so the round-trip count is not worth trading for a
|
|
// watcher whose "I have seen everything current" condition is
|
|
// one more thing to get right on a read path.
|
|
let mut keys = store.keys().await.context("listing status bucket keys")?;
|
|
let mut entries: BTreeMap<String, Offered> = BTreeMap::new();
|
|
while let Some(key) = keys
|
|
.try_next()
|
|
.await
|
|
.context("reading the status bucket's key list")?
|
|
{
|
|
let Some(entry) = store
|
|
.entry(&key)
|
|
.await
|
|
.with_context(|| format!("reading status entry {key}"))?
|
|
else {
|
|
// Deleted between listing and fetching. Not an error:
|
|
// the next read simply won't list it.
|
|
continue;
|
|
};
|
|
let payload = match serde_json::from_slice(&entry.value) {
|
|
Ok(value) => Some(value),
|
|
Err(e) => {
|
|
tracing::warn!(
|
|
hive = %key,
|
|
error = %e,
|
|
"status snapshot is not JSON; reporting the timestamp without it"
|
|
);
|
|
None
|
|
}
|
|
};
|
|
entries.insert(
|
|
key,
|
|
Offered {
|
|
received_at: to_system_time(entry.created.unix_timestamp()),
|
|
payload,
|
|
},
|
|
);
|
|
}
|
|
|
|
Ok(render(roster, &entries, now, self.stale_after))
|
|
}
|
|
}
|
|
|
|
/// A bucket timestamp as a [`SystemTime`].
|
|
///
|
|
/// A pre-epoch stamp is not representable here and is not a thing a NATS
|
|
/// server produces; treating it as the epoch renders the row as
|
|
/// extremely stale, which is the safe direction — a nonsense timestamp
|
|
/// must never read as fresh.
|
|
fn to_system_time(unix_seconds: i64) -> SystemTime {
|
|
u64::try_from(unix_seconds).map_or(UNIX_EPOCH, |secs| UNIX_EPOCH + Duration::from_secs(secs))
|
|
}
|
|
|
|
/// Turn a roster plus whatever the bucket held into the rendered rows.
|
|
///
|
|
/// Split out of [`StatusReader::view`] deliberately: this is where every
|
|
/// rule the acceptance criterion cares about lives, and keeping it a
|
|
/// pure function means those rules are tested against a table rather
|
|
/// than against a running NATS server.
|
|
fn render(
|
|
roster: &[HiveEntry],
|
|
entries: &BTreeMap<String, Offered>,
|
|
now: SystemTime,
|
|
stale_after: Duration,
|
|
) -> Vec<HiveStatus> {
|
|
let mut rows: Vec<HiveStatus> = roster
|
|
.iter()
|
|
.map(|hive| match entries.get(&hive.name) {
|
|
Some(offered) => row(
|
|
hive.name.clone(),
|
|
Some(hive.domain.clone()),
|
|
offered,
|
|
now,
|
|
stale_after,
|
|
),
|
|
None => HiveStatus {
|
|
name: hive.name.clone(),
|
|
domain: Some(hive.domain.clone()),
|
|
freshness: Freshness::NeverReported,
|
|
last_seen_unix: None,
|
|
age_seconds: None,
|
|
snapshot: None,
|
|
},
|
|
})
|
|
.collect();
|
|
|
|
rows.extend(
|
|
entries
|
|
.iter()
|
|
.filter(|(name, _)| !roster.iter().any(|hive| &&hive.name == name))
|
|
.map(|(name, offered)| {
|
|
let mut unknown = row(name.clone(), None, offered, now, stale_after);
|
|
unknown.freshness = Freshness::Unknown;
|
|
unknown
|
|
}),
|
|
);
|
|
rows
|
|
}
|
|
|
|
fn row(
|
|
name: String,
|
|
domain: Option<String>,
|
|
offered: &Offered,
|
|
now: SystemTime,
|
|
stale_after: Duration,
|
|
) -> HiveStatus {
|
|
// A snapshot stamped in the future (clock skew between the NATS
|
|
// server and this host) yields no age rather than a negative one,
|
|
// and is treated as fresh — the honest reading of "this arrived, I
|
|
// cannot tell how long ago".
|
|
let age = now.duration_since(offered.received_at).ok();
|
|
let freshness = match age {
|
|
Some(age) if age > stale_after => Freshness::Stale,
|
|
_ => Freshness::Fresh,
|
|
};
|
|
HiveStatus {
|
|
name,
|
|
domain,
|
|
freshness,
|
|
last_seen_unix: offered
|
|
.received_at
|
|
.duration_since(UNIX_EPOCH)
|
|
.ok()
|
|
.and_then(|d| i64::try_from(d.as_secs()).ok()),
|
|
age_seconds: age.map(|age| age.as_secs()),
|
|
snapshot: offered.payload.clone(),
|
|
}
|
|
}
|
|
|
|
#[cfg(test)]
|
|
mod tests {
|
|
use super::{DEFAULT_STALE_AFTER, Freshness, Offered, STALE_AFTER_ENV, StatusReader, render};
|
|
use crate::HiveEntry;
|
|
use std::collections::BTreeMap;
|
|
use std::time::{Duration, SystemTime};
|
|
|
|
fn roster() -> Vec<HiveEntry> {
|
|
vec![
|
|
HiveEntry {
|
|
name: "pr1ma".to_owned(),
|
|
domain: "pr1ma.example.com".to_owned(),
|
|
},
|
|
HiveEntry {
|
|
name: "umbra".to_owned(),
|
|
domain: "umbra.example.com".to_owned(),
|
|
},
|
|
]
|
|
}
|
|
|
|
fn entries(rows: &[(&str, SystemTime)]) -> BTreeMap<String, Offered> {
|
|
rows.iter()
|
|
.map(|(name, at)| {
|
|
(
|
|
(*name).to_owned(),
|
|
Offered {
|
|
received_at: *at,
|
|
payload: Some(serde_json::json!({ "ok": true })),
|
|
},
|
|
)
|
|
})
|
|
.collect()
|
|
}
|
|
|
|
fn t0() -> SystemTime {
|
|
SystemTime::UNIX_EPOCH + Duration::from_secs(1_700_000_000)
|
|
}
|
|
|
|
/// The whole point of the module: an empty bucket must not render as
|
|
/// a healthy swarm. Rows come from the roster, so silence is visible.
|
|
#[test]
|
|
fn an_empty_bucket_renders_every_hive_as_never_reported() {
|
|
let rows = render(&roster(), &BTreeMap::new(), t0(), DEFAULT_STALE_AFTER);
|
|
assert_eq!(rows.len(), 2, "a row per roster hive, not per report");
|
|
assert!(
|
|
rows.iter().all(|r| r.freshness == Freshness::NeverReported),
|
|
"nothing published means nothing known — not healthy"
|
|
);
|
|
assert!(rows.iter().all(|r| r.snapshot.is_none()));
|
|
}
|
|
|
|
/// Freshness is derived from the clock at read time, so the same
|
|
/// stored value reads differently as it ages. The boundary is where
|
|
/// an off-by-one would hide, so it is pinned in both directions.
|
|
#[test]
|
|
fn the_threshold_boundary_is_inclusive() {
|
|
let stale_after = Duration::from_secs(90);
|
|
let stored = entries(&[("pr1ma", t0())]);
|
|
|
|
assert_eq!(
|
|
render(
|
|
&roster(),
|
|
&stored,
|
|
t0() + Duration::from_secs(90),
|
|
stale_after
|
|
)[0]
|
|
.freshness,
|
|
Freshness::Fresh,
|
|
"exactly at the threshold is inside it"
|
|
);
|
|
assert_eq!(
|
|
render(
|
|
&roster(),
|
|
&stored,
|
|
t0() + Duration::from_secs(91),
|
|
stale_after
|
|
)[0]
|
|
.freshness,
|
|
Freshness::Stale,
|
|
"one second past is outside it"
|
|
);
|
|
}
|
|
|
|
/// One hive reporting must not make its silent neighbour look
|
|
/// healthy — the failure mode of any aggregate that renders only
|
|
/// what it has.
|
|
#[test]
|
|
fn a_reporting_hive_does_not_vouch_for_a_silent_one() {
|
|
let rows = render(
|
|
&roster(),
|
|
&entries(&[("pr1ma", t0())]),
|
|
t0(),
|
|
DEFAULT_STALE_AFTER,
|
|
);
|
|
assert_eq!(rows[0].name, "pr1ma");
|
|
assert_eq!(rows[0].freshness, Freshness::Fresh);
|
|
assert_eq!(rows[1].name, "umbra");
|
|
assert_eq!(rows[1].freshness, Freshness::NeverReported);
|
|
}
|
|
|
|
/// An observation the daemon cannot explain is surfaced, not dropped.
|
|
#[test]
|
|
fn a_hive_outside_the_roster_is_surfaced_as_unknown() {
|
|
let rows = render(
|
|
&roster(),
|
|
&entries(&[("ghost", t0())]),
|
|
t0(),
|
|
DEFAULT_STALE_AFTER,
|
|
);
|
|
assert_eq!(rows.len(), 3, "two roster hives plus the stranger");
|
|
let ghost = rows.last().expect("rows is non-empty");
|
|
assert_eq!(ghost.name, "ghost");
|
|
assert_eq!(ghost.freshness, Freshness::Unknown);
|
|
assert!(
|
|
ghost.domain.is_none(),
|
|
"the roster is where a domain comes from, and this hive isn't in it"
|
|
);
|
|
}
|
|
|
|
/// Clock skew must not produce a negative age or a panic. A snapshot
|
|
/// stamped in the future reads fresh with no age — "it arrived, I
|
|
/// cannot tell how long ago".
|
|
#[test]
|
|
fn a_future_timestamp_yields_no_age_rather_than_a_wrong_one() {
|
|
let rows = render(
|
|
&roster(),
|
|
&entries(&[("pr1ma", t0() + Duration::from_secs(30))]),
|
|
t0(),
|
|
Duration::from_secs(90),
|
|
);
|
|
assert_eq!(rows[0].freshness, Freshness::Fresh);
|
|
assert_eq!(rows[0].age_seconds, None);
|
|
}
|
|
|
|
/// A hive that published something unreadable still gets its
|
|
/// timestamp reported: *when* it last spoke is exactly what this
|
|
/// aggregate is for, and dropping the row would read as silence.
|
|
#[test]
|
|
fn an_unparseable_payload_still_reports_when_it_arrived() {
|
|
let mut stored = BTreeMap::new();
|
|
stored.insert(
|
|
"pr1ma".to_owned(),
|
|
Offered {
|
|
received_at: t0(),
|
|
payload: None,
|
|
},
|
|
);
|
|
|
|
let row = &render(&roster(), &stored, t0(), DEFAULT_STALE_AFTER)[0];
|
|
assert_eq!(
|
|
row.freshness,
|
|
Freshness::Fresh,
|
|
"unreadable is not the same as absent — freshness is what \
|
|
separates them on the wire"
|
|
);
|
|
assert!(row.snapshot.is_none());
|
|
assert_eq!(row.last_seen_unix, Some(1_700_000_000));
|
|
}
|
|
|
|
/// SAFETY: single-threaded mutation of a process env var no other
|
|
/// test in this crate reads; restored before returning. One test
|
|
/// rather than four for the same reason `load_hives`'s is — the
|
|
/// parallel runner would race them.
|
|
#[test]
|
|
fn stale_after_from_env_covers_missing_bogus_zero_and_valid() {
|
|
unsafe {
|
|
std::env::remove_var(STALE_AFTER_ENV);
|
|
}
|
|
assert_eq!(StatusReader::stale_after_from_env(), DEFAULT_STALE_AFTER);
|
|
|
|
unsafe {
|
|
std::env::set_var(STALE_AFTER_ENV, "not a number");
|
|
}
|
|
assert_eq!(
|
|
StatusReader::stale_after_from_env(),
|
|
DEFAULT_STALE_AFTER,
|
|
"a controller whose own config is wrong must still serve"
|
|
);
|
|
|
|
unsafe {
|
|
std::env::set_var(STALE_AFTER_ENV, "0");
|
|
}
|
|
assert_eq!(
|
|
StatusReader::stale_after_from_env(),
|
|
DEFAULT_STALE_AFTER,
|
|
"zero would make every snapshot instantly stale — take the default"
|
|
);
|
|
|
|
unsafe {
|
|
std::env::set_var(STALE_AFTER_ENV, "300");
|
|
}
|
|
assert_eq!(StatusReader::stale_after_from_env(), Duration::from_mins(5));
|
|
|
|
unsafe {
|
|
std::env::remove_var(STALE_AFTER_ENV);
|
|
}
|
|
}
|
|
}
|