Watch
0
0
Fork
You've already forked hyperhive
0

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
This commit is contained in:
atlas 2026-09-30 15:05:25 +02:00 • committed by mara
commit 9e06e191a3
3 changed files with 151 additions and 39 deletions

View file

@ -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<P, F>(connected: impl Fn() -> bool, mut pass: P, enqueue: impl Fn(Vec<Renewal>))
where
P: FnMut() -> F,
F: std::future::Future<Output = Result<Vec<(String, AgentObserved)>>>,
{
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<WantedWriter>,
hives: Vec<String>,
enqueue: impl Fn(Vec<Renewal>) + 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<std::sync::atomic::AtomicBool>,
Arc<std::sync::atomic::AtomicUsize>,
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();
}
}

View file

@ -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<Option<HiveWanted>> {
// An unconnected client does not fail a JetStream request, it hangs