//! Renewing each live agent's two store credentials by age: its mTLS //! certificate (`swarm/agents//bao-mtls`) once it is past half its //! validity, and its queue secret (`swarm/agents//queue`) once it is //! [`SECRET_RENEW_AFTER`] old. //! //! The certificate's age is its own `notBefore`/`notAfter`. The queue secret //! has no expiry and `swarm-nats-auth` checks nothing about time, so its age is //! the controller's own record, [`queue::AgentCredential::minted_at`]; a secret //! without one is stamped, value unchanged ([`SecretStep::Backfill`]). //! //! Either replacement reaches the agent at its next start, when the container //! is handed the certificate and fetches the secret; nothing pulls either into //! a running container. The old certificate stays valid until it expires. The //! old secret stops working at once: the queue checks it only at `CONNECT`, so //! an open connection is unaffected, but a reconnect before that restart //! presents the old secret and is denied. //! //! [`spawn`]'s pass, once the queue is up and every [`RECONCILE_INTERVAL`], //! inserts a job per agent with anything due: `MintAgentIdentity` for the //! certificate ([`crate::agent_identity::mint_and_verify`], which keeps the //! queue secret as it is) and `RenewAgentQueueCredential` ([`renew`]) for the //! secret, the second after the first when both are due. The decisions are //! pure ([`live_agents`], [`cert_is_due`], [`secret_step`], [`plan`]) so the //! tests pin them; the IO on either side only reads or acts. //! //! **Only agents some hive is declared to run are renewed** ([`live_agents`]): //! the wanted-state declarations are where a destroy is recorded, and an agent //! declared nowhere, or only as `Destroyed`, is skipped. A credential that is //! not stored is skipped too — creating one is agent creation's job, and a //! destroy deletes it. use std::collections::BTreeSet; use std::sync::Arc; use anyhow::{Context, Result, bail}; use swarm_queue_client::wanted::{AgentState, HiveWanted}; use swarm_secret_client::{SecretStore, mtls, queue}; use x509_cert::der::DecodePem as _; use crate::wanted::WantedWriter; /// Age at which a queue secret is re-minted: half the agent certificate's /// 90-day life (`agentPkiLeafTtl` in `swarm-bao.nix`), so the two renew on one /// cadence. pub const SECRET_RENEW_AFTER: std::time::Duration = std::time::Duration::from_hours(45 * 24); /// How often [`spawn`] re-checks every agent. Also the longest it waits for /// the queue connection before its first pass. const RECONCILE_INTERVAL: std::time::Duration = std::time::Duration::from_mins(5); /// How often [`spawn`] checks the queue connection before its first pass. const CONNECT_POLL: std::time::Duration = std::time::Duration::from_secs(5); /// Now, in the unix seconds [`queue::AgentCredential::minted_at`] holds. pub fn unix_now() -> i64 { std::time::SystemTime::now() .duration_since(std::time::UNIX_EPOCH) .map_or(0, |d| i64::try_from(d.as_secs()).unwrap_or(i64::MAX)) } /// What a stored queue secret needs at `now`. #[derive(Debug, Clone, Copy, PartialEq, Eq)] pub enum SecretStep { /// Younger than [`SECRET_RENEW_AFTER`]. Keep, /// No mint time recorded: stamp `now` beside the unchanged value, so its /// clock starts here and no holder of it is cut off. Backfill, /// At least [`SECRET_RENEW_AFTER`] old: replace the value. Remint, } /// Decide [`SecretStep`] for `stored` at `now`. pub fn secret_step(stored: &queue::AgentCredential, now: i64) -> SecretStep { let renew_after = i64::try_from(SECRET_RENEW_AFTER.as_secs()).unwrap_or(i64::MAX); match stored.minted_at { None => SecretStep::Backfill, Some(at) if now.saturating_sub(at) >= renew_after => SecretStep::Remint, Some(_) => SecretStep::Keep, } } /// The object `step` writes for `agent` at `now`, or `None` for /// [`SecretStep::Keep`]. /// /// # Errors /// When a re-mint cannot draw randomness from the kernel. pub fn apply_secret_step( step: SecretStep, stored: &queue::AgentCredential, agent: &str, now: i64, ) -> Result> { Ok(match step { SecretStep::Keep => None, SecretStep::Backfill => Some(queue::AgentCredential { minted_at: Some(now), ..stored.clone() }), SecretStep::Remint => Some(fresh(agent, now)?), }) } /// A certificate's validity window, `(notBefore, notAfter)` in unix seconds. /// /// # Errors /// When `pem` is not one PEM-encoded X.509 certificate. pub fn cert_validity(pem: &str) -> Result<(i64, i64)> { let cert = x509_cert::Certificate::from_pem(pem.as_bytes()) .context("decoding the stored certificate")?; let validity = &cert.tbs_certificate.validity; let secs = |t: x509_cert::time::Time| { i64::try_from(t.to_unix_duration().as_secs()).unwrap_or(i64::MAX) }; Ok((secs(validity.not_before), secs(validity.not_after))) } /// Whether a certificate valid from `not_before` to `not_after` is past half /// its life at `now`. pub fn cert_is_due(not_before: i64, not_after: i64, now: i64) -> bool { let half = not_after.saturating_sub(not_before) / 2; now >= not_before.saturating_add(half) } /// A new queue credential for `agent`, minted at `now`. Shared with agent /// creation, so every secret this daemon writes carries a mint time. /// /// # Errors /// When the kernel will not supply randomness. pub fn fresh(agent: &str, now: i64) -> Result { Ok(queue::AgentCredential { value: crate::agent_identity::generate_queue_secret()?, agent: agent.to_owned(), minted_at: Some(now), }) } /// Every agent declared on at least one hive in a state other than /// `Destroyed`. pub fn live_agents(declarations: &[HiveWanted]) -> BTreeSet { declarations .iter() .flat_map(|d| &d.agents) .filter(|(_, wanted)| wanted.state != AgentState::Destroyed) .map(|(agent, _)| agent.clone()) .collect() } /// What one pass found for one credential of one live agent. #[derive(Debug, Clone, Copy, PartialEq, Eq)] pub enum Observed { /// Stored and due for replacement, with its age in seconds. Due(i64), /// A queue secret stored with no mint time: due for /// [`SecretStep::Backfill`], never for a new value. Unstamped, /// Stored and not yet due. Current, /// Nothing stored. Never renewed: see the module doc. Absent, /// The read or the decode failed. Nothing is known, so nothing is done. Unknown, } /// What one pass found for one live agent. #[derive(Debug, Clone, Copy, PartialEq, Eq)] pub struct AgentObserved { pub cert: Observed, pub secret: Observed, } /// One agent's share of a pass: which of its credentials to renew. #[derive(Debug, Clone, PartialEq, Eq)] pub struct Renewal { pub agent: String, pub cert: bool, pub secret: bool, } /// The renewals a pass queues: one per agent with a certificate that is /// [`Observed::Due`], or a queue secret that is [`Observed::Due`] or /// [`Observed::Unstamped`]. pub fn plan(observed: &[(String, AgentObserved)]) -> Vec { observed .iter() .map(|(agent, o)| Renewal { agent: agent.clone(), cert: matches!(o.cert, Observed::Due(_)), secret: matches!(o.secret, Observed::Due(_) | Observed::Unstamped), }) .filter(|r| r.cert || r.secret) .collect() } /// An age in seconds, for a log line: whole days. struct Age(i64); impl std::fmt::Display for Age { fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { write!(f, "{}d", self.0 / 86_400) } } /// Backfill or re-mint `agent`'s queue secret if it is still stored and still /// needs it ([`secret_step`]). The whole job of the `RenewAgentQueueCredential` /// node. /// /// Re-decides rather than trusting the pass that queued it, so a node queued /// twice acts once. Reads the write back before reporting success. /// /// # Errors /// When the store refuses a step, or returns a different object than the one /// just written. pub async fn renew(agent: &str) -> Result<()> { let path = queue::agent_queue_path(agent)?; let store = crate::store::connect() .await .context("logging in to the swarm secret store")?; let Some(stored) = store .read_optional::(&path) .await .with_context(|| format!("reading {path}"))? else { tracing::info!(agent, %path, "agent queue credential: nothing stored, nothing to renew"); return Ok(()); }; let now = unix_now(); let step = secret_step(&stored, now); let Some(written) = apply_secret_step(step, &stored, agent, now)? else { tracing::debug!(agent, "agent queue credential is current; left as it is"); return Ok(()); }; store .write(&path, &written) .await .with_context(|| format!("writing the agent queue credential at {path}"))?; let read_back: queue::AgentCredential = store .read(&path) .await .with_context(|| format!("reading {path} back"))?; if read_back != written { bail!("the store returned a different object at {path} than the one just written"); } if let Some(at) = stored.minted_at { tracing::info!( agent, %path, old_age = %Age(now.saturating_sub(at)), "agent queue credential re-minted; the agent presents it from its next restart" ); } else { tracing::info!( agent, %path, "agent queue credential had no mint time; stamped now, value unchanged" ); } Ok(()) } /// What the store says about `agent`'s certificate. async fn observe_cert(store: &SecretStore, agent: &str, now: i64) -> Observed { let stored = match mtls::identity_path(agent) { Ok(path) => store.read_optional::(&path).await, Err(e) => Err(e), }; let credential = match stored { Ok(Some(credential)) => credential, Ok(None) => return Observed::Absent, Err(e) => { tracing::warn!(agent, error = %e, "agent certificate: store read failed"); return Observed::Unknown; } }; match cert_validity(&credential.cert) { Ok((not_before, not_after)) if cert_is_due(not_before, not_after, now) => { Observed::Due(now.saturating_sub(not_before)) } Ok(_) => Observed::Current, Err(e) => { tracing::warn!(agent, error = %format!("{e:#}"), "agent certificate: unreadable"); Observed::Unknown } } } /// What the store says about `agent`'s queue secret. async fn observe_secret(store: &SecretStore, agent: &str, now: i64) -> Observed { let stored = match queue::agent_queue_path(agent) { Ok(path) => store.read_optional::(&path).await, Err(e) => Err(e), }; match stored { Ok(Some(stored)) => match (secret_step(&stored, now), stored.minted_at) { (SecretStep::Remint, Some(at)) => Observed::Due(now.saturating_sub(at)), (SecretStep::Backfill, _) => Observed::Unstamped, _ => Observed::Current, }, Ok(None) => Observed::Absent, Err(e) => { tracing::warn!(agent, error = %e, "agent queue credential: store read failed"); Observed::Unknown } } } /// One pass: every live agent, observed. Fails as a whole when any hive's /// declaration cannot be read, since that hive may be the one declaring an /// agent destroyed. async fn observe_all( wanted: &WantedWriter, hives: &[String], ) -> Result> { let mut declarations = Vec::with_capacity(hives.len()); for hive in hives { if let Some(d) = wanted .view(hive) .await .with_context(|| format!("reading {hive}'s wanted-state declaration"))? { declarations.push(d); } } let store = crate::store::connect() .await .context("logging in to the swarm secret store")?; let now = unix_now(); let mut observed = Vec::new(); for agent in live_agents(&declarations) { let o = AgentObserved { cert: observe_cert(&store, &agent, now).await, secret: observe_secret(&store, &agent, now).await, }; if let Observed::Due(age) = o.cert { tracing::info!(agent, age = %Age(age), "agent certificate past half its life; re-issuing"); } if let Observed::Due(age) = o.secret { tracing::info!(agent, age = %Age(age), "agent queue credential due; re-minting"); } if o.secret == Observed::Unstamped { tracing::info!( agent, "agent queue credential has no mint time; stamping it" ); } observed.push((agent, o)); } Ok(observed) } /// Wait until `connected` holds, checking every [`CONNECT_POLL`] for at most /// [`RECONCILE_INTERVAL`]. Returns whether it held. async fn wait_for_queue(connected: impl Fn() -> bool) -> bool { let deadline = tokio::time::Instant::now() + RECONCILE_INTERVAL; loop { if connected() { return true; } if tokio::time::Instant::now() >= deadline { return false; } tokio::time::sleep(CONNECT_POLL).await; } } /// [`spawn`]'s loop, with the connection check and the pass handed in. async fn run(connected: impl Fn() -> bool, mut pass: P, enqueue: impl Fn(Vec)) where P: FnMut() -> F, F: std::future::Future>>, { if !wait_for_queue(connected).await { tracing::warn!( waited_s = RECONCILE_INTERVAL.as_secs(), "agent credential renewal: queue still not connected; running the first pass anyway" ); } let mut ticker = tokio::time::interval(RECONCILE_INTERVAL); loop { ticker.tick().await; match pass().await { Ok(observed) => { let renewals = plan(&observed); if renewals.is_empty() { tracing::debug!( checked = observed.len(), "agent credential renewal: none due" ); } else { tracing::info!( checked = observed.len(), renewing = renewals.len(), "agent credential renewal: queueing" ); enqueue(renewals); } } Err(e) => tracing::warn!( error = %format!("{e:#}"), retry_in_s = RECONCILE_INTERVAL.as_secs(), "agent credential renewal: pass failed; retrying next tick" ), } } } /// Check every live agent once the queue is connected and every /// [`RECONCILE_INTERVAL`] after, and hand the renewals due to `enqueue`, which /// inserts a job for each. /// /// Only the first pass waits for the connection, and for at most one /// [`RECONCILE_INTERVAL`]. A pass that fails is logged and retried on the next /// tick; it never stops the daemon. pub fn spawn( wanted: Arc, hives: Vec, enqueue: impl Fn(Vec) + Send + 'static, ) { tokio::spawn(async move { run( || wanted.is_connected(), || observe_all(&wanted, &hives), enqueue, ) .await; }); } #[cfg(test)] mod tests { use super::*; use swarm_queue_client::wanted::AgentWanted; const NOW: i64 = 1_790_000_000; const DAY: i64 = 86_400; /// Self-signed, CN `hive-agent-atlas`, valid 2026-01-01T00:00:00Z to /// 2026-04-01T00:00:00Z: 90 days, the agent role's `ttl`. Its key was /// discarded. const CERT_PEM: &str = "-----BEGIN CERTIFICATE----- MIIBizCCATGgAwIBAgIUSFpuYmy1UoHvq/MYiyrAudRaWXcwCgYIKoZIzj0EAwIw GzEZMBcGA1UEAwwQaGl2ZS1hZ2VudC1hdGxhczAeFw0yNjAxMDEwMDAwMDBaFw0y NjA0MDEwMDAwMDBaMBsxGTAXBgNVBAMMEGhpdmUtYWdlbnQtYXRsYXMwWTATBgcq hkjOPQIBBggqhkjOPQMBBwNCAAR7p2eZAxjiE1RdsuSHE3CcTjbnd4swzfnmqUPW 9NJN7mvOtqEzxCKfOSFjaIhhyShIQ41/V6vGi80EPDB+uR56o1MwUTAdBgNVHQ4E FgQUxgZFsN1VY3ODrPS1NUfY/x3EnmswHwYDVR0jBBgwFoAUxgZFsN1VY3ODrPS1 NUfY/x3EnmswDwYDVR0TAQH/BAUwAwEB/zAKBggqhkjOPQQDAgNIADBFAiEAnU9e WOioNTUNK9b6FignejpK5FpyrODUuH/iL4TqCh0CIAz6GzJeLebhRPvYlfVUofdV hWx3sOwmUgjkRQXoxY+p -----END CERTIFICATE----- "; const CERT_NOT_BEFORE: i64 = 1_767_225_600; const CERT_NOT_AFTER: i64 = 1_775_001_600; fn minted(at: Option) -> queue::AgentCredential { queue::AgentCredential { value: "Ab9_-zSECRET".to_owned(), agent: "atlas".to_owned(), minted_at: at, } } #[test] fn a_secret_without_a_mint_time_is_backfilled_with_its_value_unchanged() { let stored = minted(None); assert_eq!(secret_step(&stored, NOW), SecretStep::Backfill); let written = apply_secret_step(SecretStep::Backfill, &stored, "atlas", NOW) .expect("no randomness needed") .expect("a backfill writes"); assert_eq!(written.value, stored.value, "no holder of it is cut off"); assert_eq!(written.agent, stored.agent); assert_eq!(written.minted_at, Some(NOW), "its clock starts now"); assert_eq!(secret_step(&written, NOW), SecretStep::Keep); } #[test] fn a_secret_is_re_minted_from_forty_five_days_with_a_new_value_and_mint_time() { for age in [45 * DAY, 90 * DAY] { let stored = minted(Some(NOW - age)); assert_eq!(secret_step(&stored, NOW), SecretStep::Remint, "{age}"); let written = apply_secret_step(SecretStep::Remint, &stored, "atlas", NOW) .expect("the kernel supplies randomness") .expect("a re-mint writes"); assert_ne!(written.value, stored.value); assert_eq!(written.minted_at, Some(NOW)); assert_eq!(written.agent, "atlas"); assert_eq!(secret_step(&written, NOW), SecretStep::Keep); } let a = fresh("atlas", NOW).expect("randomness"); let b = fresh("atlas", NOW).expect("randomness"); assert_ne!(a.value, b.value, "two re-mints must not agree"); } /// Includes a mint time ahead of the clock, which must not overflow into a /// large age. #[test] fn a_secret_younger_than_forty_five_days_is_left_alone() { for at in [NOW, NOW - 45 * DAY + 1, NOW + DAY] { let stored = minted(Some(at)); assert_eq!(secret_step(&stored, NOW), SecretStep::Keep, "{at}"); assert_eq!( apply_secret_step(SecretStep::Keep, &stored, "atlas", NOW).expect("no IO"), None ); } } #[test] fn a_certificates_validity_is_read_from_the_certificate() { assert_eq!( cert_validity(CERT_PEM).expect("a well-formed certificate"), (CERT_NOT_BEFORE, CERT_NOT_AFTER) ); } #[test] fn something_that_is_not_a_certificate_is_an_error() { assert!(cert_validity("not a certificate").is_err()); assert!(cert_validity("").is_err()); } #[test] fn a_ninety_day_certificate_is_due_from_day_forty_five() { let (nb, na) = (CERT_NOT_BEFORE, CERT_NOT_AFTER); assert!(!cert_is_due(nb, na, nb)); assert!(!cert_is_due(nb, na, nb + 45 * DAY - 1)); assert!(cert_is_due(nb, na, nb + 45 * DAY)); assert!(cert_is_due(nb, na, na + DAY), "an expired one too"); } fn declared(agents: &[(&str, AgentState)]) -> HiveWanted { let mut d = HiveWanted::default(); for (agent, state) in agents { d.agents .insert((*agent).to_owned(), AgentWanted { state: *state }); } d } #[test] fn a_destroyed_agent_is_not_live_and_every_other_state_is() { let live = live_agents(&[declared(&[ ("up", AgentState::Up), ("offline", AgentState::Offline), ("paused", AgentState::Paused), ("gone", AgentState::Destroyed), ])]); assert_eq!( live.into_iter().collect::>(), ["offline", "paused", "up"] ); } /// An agent destroyed on one hive and declared on another is live: it /// still runs somewhere. #[test] fn an_agent_live_on_any_hive_is_live() { let live = live_agents(&[ declared(&[("atlas", AgentState::Destroyed)]), declared(&[("atlas", AgentState::Up)]), ]); assert!(live.contains("atlas"), "{live:?}"); } #[test] fn no_declaration_is_no_live_agent() { assert!(live_agents(&[]).is_empty()); } #[test] fn only_due_credentials_are_planned() { use Observed::{Absent, Current, Due, Unknown, Unstamped}; let o = |cert, secret| AgentObserved { cert, secret }; let observed = [ ("both".to_owned(), o(Due(46 * DAY), Due(45 * DAY))), ("cert".to_owned(), o(Due(46 * DAY), Current)), ("secret".to_owned(), o(Absent, Due(45 * DAY))), ("unstamped".to_owned(), o(Current, Unstamped)), ("neither".to_owned(), o(Current, Absent)), ("unknown".to_owned(), o(Unknown, Unknown)), ]; let r = |agent: &str, cert, secret| Renewal { agent: agent.to_owned(), cert, secret, }; assert_eq!( plan(&observed), [ r("both", true, true), r("cert", true, false), r("secret", false, true), r("unstamped", false, true), ] ); } #[test] fn the_age_logged_is_whole_days() { assert_eq!(Age(46 * DAY + 5).to_string(), "46d"); } /// [`run`] over a connection the test flips and a pass that only counts /// itself. fn counting_run() -> ( Arc, Arc, tokio::task::JoinHandle<()>, ) { use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering}; let connected = Arc::new(AtomicBool::new(false)); let passes = Arc::new(AtomicUsize::new(0)); let (c, p) = (Arc::clone(&connected), Arc::clone(&passes)); let task = tokio::spawn(run( move || c.load(Ordering::SeqCst), move || { p.fetch_add(1, Ordering::SeqCst); std::future::ready(Ok(Vec::new())) }, |_| {}, )); (connected, passes, task) } /// The sleeps end between [`CONNECT_POLL`] boundaries, so a check and an /// assertion never fall on the same instant. #[tokio::test(start_paused = true)] async fn the_first_pass_waits_for_the_queue_and_later_passes_do_not() { use std::sync::atomic::Ordering; let (connected, passes, task) = counting_run(); tokio::time::sleep(std::time::Duration::from_secs(62)).await; assert_eq!(passes.load(Ordering::SeqCst), 0, "no pass while Pending"); connected.store(true, Ordering::SeqCst); tokio::time::sleep(CONNECT_POLL).await; assert_eq!(passes.load(Ordering::SeqCst), 1, "a pass once connected"); connected.store(false, Ordering::SeqCst); tokio::time::sleep(RECONCILE_INTERVAL).await; assert_eq!( passes.load(Ordering::SeqCst), 2, "a later disconnect does not hold back the next pass" ); task.abort(); } #[tokio::test(start_paused = true)] async fn a_queue_that_never_connects_gets_a_pass_after_one_interval() { use std::sync::atomic::Ordering; let (_connected, passes, task) = counting_run(); tokio::time::sleep(RECONCILE_INTERVAL.saturating_sub(std::time::Duration::from_secs(2))) .await; assert_eq!(passes.load(Ordering::SeqCst), 0); tokio::time::sleep(std::time::Duration::from_secs(4)).await; assert_eq!(passes.load(Ordering::SeqCst), 1); task.abort(); } }