//! Offering this hive's readiness upward to the swarm. //! //! A hive **offers**; the swarm controller never reaches down to collect. //! That direction is the design and not an implementation detail: the //! 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 go dark exactly when it is needed. A hive computes //! its own status locally whether or not the swarm can be reached, and //! this task is only the part that carries it. //! //! **It publishes what the hive already says about itself.** //! [`crate::warnings::readiness`] is the same value `/health/ready` //! serves — not a swarm-specific recomputation. Two producers of "is this //! hive OK" would be free to disagree, and the disagreement would surface //! as the dashboard and the swarm view contradicting each other about the //! same host, each internally consistent. //! //! **The key is this hive's `hiveName`.** Nothing here has to arrange for //! that to match the controller's roster: `swarm.nix` asserts that //! `services.hyperhive.hiveName` is a key of `services.hyperhive.swarm.hives`, //! so a hive that evaluates at all publishes under a name the roster //! knows. The alternative — a key the controller has never heard of — //! renders as `unknown` rather than being silently dropped, but the //! assertion means it should not arise. //! //! **Nothing here stamps a time.** Freshness is derived by the reader //! from when the value landed in the bucket, so a hive cannot make itself //! look fresher than it is, and a hive with a wrong clock skews only its //! own payload. use std::time::Duration; use anyhow::{Context, Result}; // `Subscriber` is a `Stream`, so reading the next event needs the extension // trait — there is no inherent `next()` on it. use futures_util::StreamExt as _; use crate::stats::sweep_health::{self, SweepHealth}; /// How often this hive offers a snapshot. /// /// **One decision with the controller's `staleAfterSeconds` (default /// 120s), not two.** The ratio is what either option means: at 2, a /// single lost publish still reads `fresh` and two consecutive misses /// read `stale`. Set them equal and any one hiccup alarms; open the gap /// wide and `stale` stops meaning anything. Changing one without the /// other silently re-tunes the swarm's definition of "quiet". pub const PUBLISH_INTERVAL: Duration = Duration::from_mins(1); /// Consecutive failed publishes before the dashboard banners. /// /// At [`PUBLISH_INTERVAL`] this is ~3 minutes of genuine failure, so a /// blip does not flap a banner an operator learns to ignore. const FAILURES_BEFORE_BANNER: u32 = 3; /// Start the publish loop, if this deployment wired up a swarm queue. /// /// The connect itself is shared with every other swarm-queue consumer in /// this process — see [`crate::swarm_queue`] for why one connection and /// not one per consumer, and for where "no queue configured" vs. "queue /// configured but unreachable" gets bannered. This function only decides /// whether *status* has anything to offer once a client exists. pub fn spawn( coord: std::sync::Arc, mut shutdown: tokio::sync::watch::Receiver, ) { let Some(hive) = crate::container_view::hive_swarm_names().0 else { crate::warnings::set_boot_warning( "swarm_status_config", "warn", "swarm status publishing is off: HYPERHIVE_HIVE_NAME is unset, so this \ hive has no key to publish under", ); return; }; tokio::spawn(async move { let Some(client) = crate::swarm_queue::client().await else { // Absent or failed — either way already handled (an `info` // log or a `swarm_queue_config` banner) by the shared // connector; nothing left to report here. return; }; // The hive's ONE queue connection, now serving both directions: // status goes up, swarm events come down. A second `connect` would // double the auth-callout traffic against authelia and give the two // paths independent reconnect state, so one could be serving while // the other was still down. `async_nats::Client` is a handle, so the // clone is cheap. tokio::spawn(drain_swarm_events( client.clone(), coord, hive.clone(), shutdown.clone(), )); let mut health = SweepHealth::new("swarm_status_publish", "warn", FAILURES_BEFORE_BANNER); loop { match publish(&client, &hive).await { Ok(()) => health.record_ok(), Err(e) => { tracing::warn!(error = ?e, "swarm status: publish failed"); let err = format!("{e:#}"); health.record_err(|ctx| { let age = ctx.since_last_ok.map_or_else( || "no success this session".to_owned(), |d| format!("last ok {} ago", sweep_health::fmt_age(d)), ); format!( "swarm status publishing is failing ({} consecutive, {age}) \ — the swarm sees this hive as stale, the hive itself is \ unaffected: {err}", ctx.consecutive ) }); } } // Publish first, then wait: a hive that has just come up is // exactly the one whose status someone is looking at, and // sleeping first would make every restart read `stale` for a // full interval. The shutdown arm means a stop is not spent // waiting out that interval either. tokio::select! { () = tokio::time::sleep(PUBLISH_INTERVAL) => {} _ = shutdown.changed() => { tracing::info!("swarm status: shutdown signal received"); break; } } } }); } /// Listen on the swarm's event subjects and act on what arrives. /// /// The controller decides *what a forge delivery means* and addresses the /// result here; this end does not know a forge exists. Two events today: /// **the knowledge repository changed**, answered by the pull this daemon /// already runs at boot, and **deploy this agent**, answered by the same /// rebuild insert the operator's own verb makes. /// /// Only the deploy event carries a payload, and only the agent name: its /// subject already names the hive, so this one listens on its own rather than /// filtering a swarm-wide feed. Even then it is a trigger, never the config, /// which git owns. /// /// # A missed message costs the two events very differently /// /// Core NATS, so delivery is at-most-once. Harmless for knowledge: this /// daemon pulls at startup regardless, and the webhook it replaced was lost /// identically when a hive was down. /// /// ⚠️ **The deploy event has no such second path** — nothing else would ever /// tell this hive to build that agent. Its backstop is the hive-side reconcile /// loop ("hives pull and self-update"); until that exists this is a nudge with /// no safety net. Not papered over with `JetStream`: durability on one subject /// looks like a fix while the desired state still lives only in a message. /// /// ⚠️ **A refused subscription is indistinguishable from a quiet one.** NATS /// reports an authorization violation asynchronously on the connection, not as /// an error from `subscribe`, so this task cannot tell "no events published" /// from "not allowed to hear them". If a hive stops picking up either event, /// the server log is the thing that knows why — nothing here will say. async fn drain_swarm_events( client: async_nats::Client, coord: std::sync::Arc, hive: String, mut shutdown: tokio::sync::watch::Receiver, ) { // One subject for the whole swarm, so this hive's own name never enters // it: the controller publishes once and core NATS fans out to whoever is // subscribed. let subject = swarm_queue_client::KNOWLEDGE_SUBJECT; let mut sub = match client.subscribe(subject).await { Ok(sub) => sub, Err(e) => { // Warn rather than a boot banner: the hive is fully functional // without this, it just falls back to learning about knowledge // changes at its next boot. tracing::warn!( %subject, error = %e, "swarm events: subscribe failed; this hive will not hear knowledge changes" ); return; } }; // This hive's own deploy subject, not a swarm-wide one: the controller // addresses each hive, so nothing arrives here that is not for us. let deploy_subject = swarm_queue_client::deploy_subject(&hive); let mut deploy_sub = match client.subscribe(deploy_subject.clone()).await { Ok(sub) => sub, Err(e) => { tracing::warn!( subject = %deploy_subject, error = %e, "swarm events: subscribe failed; this hive will not hear deploy requests" ); return; } }; tracing::info!(%subject, %deploy_subject, "swarm events: listening"); loop { tokio::select! { msg = sub.next() => { if msg.is_none() { // The subscription ended — the connection went away for // good. Returning is right: `async-nats` reconnects // underneath a live subscription, so a closed stream is // not a blip this should spin on. tracing::warn!(%subject, "swarm events: subscription closed"); return; } tracing::info!(%subject, "swarm events: knowledge change announced, pulling"); if let Err(e) = crate::workers::knowledge::pull(&coord).await { tracing::warn!(error = ?e, "swarm events: knowledge pull failed"); } } msg = deploy_sub.next() => { let Some(msg) = msg else { tracing::warn!(subject = %deploy_subject, "swarm events: deploy subscription closed"); return; }; handle_deploy_request(&coord, &msg.payload); } _ = shutdown.changed() => { tracing::info!("swarm events: shutdown signal received"); return; } } } } /// Act on one deploy request. Everything that arrives on this hive's own /// deploy subject is for this hive, so there is nothing to filter. /// /// A payload that will not decode is worth a `warn`: the controller and this /// end share one type, so a decode failure means they disagree about it. fn handle_deploy_request(coord: &std::sync::Arc, payload: &[u8]) { let request: swarm_queue_client::DeployRequest = match serde_json::from_slice(payload) { Ok(request) => request, Err(e) => { tracing::warn!(error = %e, "swarm events: undecodable deploy request"); return; } }; // The same insert the operator's own `rebuild` verb makes, relock and // all: "deploy this agent" means here exactly what it already meant, and // a swarm-triggered rebuild that quietly did something narrower would be // a second definition of the word. let agent = request.agent; let inserted = coord .job_queue .insert_job(|b| crate::job_queue::templates::rebuild(b, &agent, true)); match inserted { Ok(_) => { tracing::info!(%agent, "swarm events: deploy requested by the swarm, rebuild queued"); coord.emit_rebuild_queue_snapshot(); } Err(e) => tracing::warn!(%agent, error = %e, "swarm events: queueing the deploy failed"), } } /// Offer one snapshot: this hive's current readiness, under its own key. async fn publish(client: &async_nats::Client, hive: &str) -> Result<()> { // An unconnected client does not fail a JetStream request, it hangs // on it — which here would hang the loop, shutdown arm included. swarm_queue_client::ensure_connected(client)?; let store = swarm_queue_client::status::open_or_create(client).await?; let payload = serde_json::to_vec(&crate::warnings::readiness()).context("serialising the readiness")?; store .put(hive, payload.into()) .await .with_context(|| format!("publishing the status snapshot for {hive}"))?; Ok(()) }