From 9e06e191a357adf2340f81873689a790c2bbc423 Mon Sep 17 00:00:00 2001 From: atlas Date: Wed, 30 Sep 2026 15:05:25 +0200 Subject: [PATCH] swarm-controller: first credential-renewal pass waits for the queue connection The renewal pass ran the moment swarm-controller started, before its queue client had connected, so `WantedWriter::view` refused with "client state Pending" and the pass logged `agent credential renewal: pass failed; retrying next tick`. That fired twice in 24h on muede-lpt2, both under a second after start, and would trip a Grafana rule on that WARN on ordinary restarts. The first pass now waits until the queue client is connected, polling `swarm_queue_client::ensure_connected` every 5s the way `AgentIconReader::create_when_connected` does. The wait is bounded by one RECONCILE_INTERVAL (5 min): the old code's retry after a false start also came one interval later, so a queue that never connects gets its first pass no later than before. Hitting the bound logs one WARN and runs the pass anyway. Only startup waits; a later disconnect still fails a pass and logs the WARN. Refs #4717 --- swarm-controller/Cargo.toml | 5 + swarm-controller/src/agent_renewal.rs | 179 ++++++++++++++++++++------ swarm-controller/src/wanted.rs | 6 + 3 files changed, 151 insertions(+), 39 deletions(-) diff --git a/swarm-controller/Cargo.toml b/swarm-controller/Cargo.toml index e0ad9742..d5929e7b 100644 --- a/swarm-controller/Cargo.toml +++ b/swarm-controller/Cargo.toml @@ -120,5 +120,10 @@ utoipa-axum.workspace = true # `agent_renewal.rs` reads an agent certificate's validity window. x509-cert.workspace = true +# `test-util` for `#[tokio::test(start_paused = true)]`: `agent_renewal`'s +# tests wait out its five-minute bound on paused time. +[dev-dependencies] +tokio = { workspace = true, features = ["test-util"] } + [lints] workspace = true diff --git a/swarm-controller/src/agent_renewal.rs b/swarm-controller/src/agent_renewal.rs index f315322d..6097fc98 100644 --- a/swarm-controller/src/agent_renewal.rs +++ b/swarm-controller/src/agent_renewal.rs @@ -15,13 +15,13 @@ //! an open connection is unaffected, but a reconnect before that restart //! presents the old secret and is denied. //! -//! [`spawn`]'s pass, at start 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. +//! [`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 @@ -44,9 +44,13 @@ use crate::wanted::WantedWriter; /// cadence. pub const SECRET_RENEW_AFTER: std::time::Duration = std::time::Duration::from_hours(45 * 24); -/// How often [`spawn`] re-checks every agent. +/// 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() @@ -344,44 +348,81 @@ async fn observe_all( Ok(observed) } -/// Check every live agent now and every [`RECONCILE_INTERVAL`] after, and hand -/// the renewals due to `enqueue`, which inserts a job for each. +/// 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. /// -/// The first tick fires immediately. A pass that fails is logged and retried -/// on the next tick; it never stops the daemon. +/// 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 { - let mut ticker = tokio::time::interval(RECONCILE_INTERVAL); - loop { - ticker.tick().await; - match observe_all(&wanted, &hives).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" - ), - } - } + run( + || wanted.is_connected(), + || observe_all(&wanted, &hives), + enqueue, + ) + .await; }); } @@ -558,4 +599,64 @@ hWx3sOwmUgjkRQXoxY+p 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(); + } } diff --git a/swarm-controller/src/wanted.rs b/swarm-controller/src/wanted.rs index 254ba5cd..b6a9ffb6 100644 --- a/swarm-controller/src/wanted.rs +++ b/swarm-controller/src/wanted.rs @@ -69,6 +69,12 @@ impl WantedWriter { Self { client } } + /// Whether the queue is connected, which [`Self::view`] and [`Self::set`] + /// require. + pub fn is_connected(&self) -> bool { + swarm_queue_client::ensure_connected(&self.client).is_ok() + } + /// The declaration currently published for `hive`, or `None`. pub async fn view(&self, hive: &str) -> Result> { // An unconnected client does not fail a JetStream request, it hangs