Watch
0
0
Fork
You've already forked hyperhive
0

swarm-controller: re-issue agent certificates and re-mint queue secrets at half-life

A five-minute pass over every agent some hive's wanted state declares as
anything but destroyed queues, per agent:

- `MintAgentIdentity` (the node agent creation uses) when the stored
  certificate at swarm/agents/<agent>/bao-mtls is past half its validity,
  read from its own notBefore/notAfter: day 45 of the role's 90;
- the new `RenewAgentQueueCredential` node when the queue secret at
  swarm/agents/<agent>/queue is 45 days old or has no mint time. The node
  re-decides, writes a fresh value with `minted_at`, reads it back, and logs
  the agent and the old age.

When both are due the secret node runs after_any the certificate node,
because mint_and_verify compares the queue secret it read with the one it
reads back. A credential that is not stored is never created here.

`queue::AgentCredential` gains an optional `minted_at` (unix seconds);
agent creation now sets it. Stored objects without it decode unchanged and
count as due, so every existing queue secret is re-minted on the first pass.

Both replacements reach the agent at its next start. The old certificate
stays valid until it expires; the old queue secret does not, so a queue
reconnect before that restart is denied.

Adds x509-cert 0.2 (with der_derive and flagset) to read the validity.

docs/swarm/credentials.md: the renewal column splits into automatic re-mint
and automatic re-pull, filled from the code as it stands.
This commit is contained in:
atlas 2026-09-28 21:35:01 +02:00
commit 2115ec2bb3
10 changed files with 763 additions and 29 deletions

View file

@ -117,6 +117,8 @@ tracing.workspace = true
url.workspace = true
utoipa.workspace = true
utoipa-axum.workspace = true
# `agent_renewal.rs` reads an agent certificate's validity window.
x509-cert.workspace = true
[lints]
workspace = true

View file

@ -25,8 +25,8 @@
//! The authority is the store's own agent CA, generated inside its agent PKI
//! mount; its key never leaves the store. It signs no host leaf, and no host
//! role pins it, so an agent's certificate satisfies only that agent's role.
//! Nothing re-issues a leaf before it expires (the role's `ttl`); until a
//! renewal path exists, an operator re-runs agent creation.
//! [`crate::agent_renewal`] re-issues a live agent's leaf once it is past half
//! its validity (the role's `ttl`), by queueing the same node as creation.
use anyhow::{Context, Result, bail};
use swarm_secret_client::{
@ -98,7 +98,7 @@ const QUEUE_SECRET_BYTES: usize = 32;
/// When the kernel will not supply randomness. Bubbled rather than panicked
/// on: the caller is a job node that reports a named failure, and a secret
/// from a degraded source is worse than no secret.
fn generate_queue_secret() -> Result<String> {
pub(crate) fn generate_queue_secret() -> Result<String> {
let mut bytes = [0u8; QUEUE_SECRET_BYTES];
getrandom::fill(&mut bytes).context("drawing a queue secret from the kernel's CSPRNG")?;
Ok(base64::Engine::encode(
@ -130,7 +130,8 @@ fn generate_queue_secret() -> Result<String> {
/// ⚠️ **Step 3 is idempotent and steps 1–2 are not.** Re-running issues a
/// fresh leaf, picked up on the agent's next boot, but leaves an existing
/// queue secret alone: this is re-run against running agents, which hold that
/// secret in a live connection. Revoking one means deleting the path.
/// secret in a live connection. Revoking one means deleting the path;
/// replacing one by age is [`crate::agent_renewal`]'s.
///
/// # Errors
/// Anything that stops one of those five steps, with the step named. A
@ -161,15 +162,15 @@ pub async fn mint_and_verify(agent: &str) -> Result<()> {
.read_optional(&queue_path)
.await
.with_context(|| format!("checking whether {queue_path} already holds a credential"))?;
// The secret survives a re-run. An object naming a different agent is
// corrected by rewriting the name around the *same* `value`, which no live
// connection notices.
let wanted = queue::AgentCredential {
value: match &existing {
Some(existing) => existing.value.clone(),
None => generate_queue_secret()?,
// The secret and its mint time survive a re-run. An object naming a
// different agent is corrected by rewriting the name around the *same*
// `value`, which no live connection notices.
let wanted = match &existing {
Some(existing) => queue::AgentCredential {
agent: agent.to_owned(),
..existing.clone()
},
agent: agent.to_owned(),
None => crate::agent_renewal::fresh(agent, crate::agent_renewal::unix_now())?,
};
if existing.as_ref() == Some(&wanted) {
tracing::info!(

View file

@ -0,0 +1,496 @@
//! Renewing each live agent's two store credentials by age: its mTLS
//! certificate (`swarm/agents/<agent>/bao-mtls`) once it is past half its
//! validity, and its queue secret (`swarm/agents/<agent>/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 due.
//!
//! 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, 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_is_due`], [`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.
const RECONCILE_INTERVAL: std::time::Duration = std::time::Duration::from_mins(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))
}
/// Whether a queue secret is due for a re-mint at `now`: it has no mint time,
/// or it is at least [`SECRET_RENEW_AFTER`] old.
pub fn secret_is_due(stored: &queue::AgentCredential, now: i64) -> bool {
let renew_after = i64::try_from(SECRET_RENEW_AFTER.as_secs()).unwrap_or(i64::MAX);
stored
.minted_at
.is_none_or(|at| now.saturating_sub(at) >= renew_after)
}
/// 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<queue::AgentCredential> {
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<String> {
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, with its age in seconds when that is known.
Due(Option<i64>),
/// 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 at least one
/// [`Observed::Due`] credential.
pub fn plan(observed: &[(String, AgentObserved)]) -> Vec<Renewal> {
observed
.iter()
.map(|(agent, o)| Renewal {
agent: agent.clone(),
cert: matches!(o.cert, Observed::Due(_)),
secret: matches!(o.secret, Observed::Due(_)),
})
.filter(|r| r.cert || r.secret)
.collect()
}
/// An age for a log line: whole days, or `unrecorded`.
struct Age(Option<i64>);
impl std::fmt::Display for Age {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self.0 {
Some(secs) => write!(f, "{}d", secs / 86_400),
None => f.write_str("unrecorded"),
}
}
}
/// Re-mint `agent`'s queue secret if it is still stored and still due. The
/// whole job of the `RenewAgentQueueCredential` node.
///
/// Re-decides rather than trusting the pass that queued it, so a node queued
/// twice renews 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::<queue::AgentCredential>(&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();
if !secret_is_due(&stored, now) {
tracing::debug!(agent, "agent queue credential is current; left as it is");
return Ok(());
}
let renewed = fresh(agent, now)?;
store
.write(&path, &renewed)
.await
.with_context(|| format!("writing the re-minted agent queue credential at {path}"))?;
let read_back: queue::AgentCredential = store
.read(&path)
.await
.with_context(|| format!("reading {path} back"))?;
if read_back != renewed {
bail!("the store returned a different object at {path} than the one just written");
}
tracing::info!(
agent,
%path,
old_age = %Age(stored.minted_at.map(|at| now.saturating_sub(at))),
"agent queue credential re-minted; the agent presents it from its next restart"
);
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::<mtls::Credential>(&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(Some(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::<queue::AgentCredential>(&path).await,
Err(e) => Err(e),
};
match stored {
Ok(Some(stored)) if secret_is_due(&stored, now) => {
Observed::Due(stored.minted_at.map(|at| now.saturating_sub(at)))
}
Ok(Some(_)) => 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<Vec<(String, AgentObserved)>> {
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");
}
observed.push((agent, o));
}
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.
///
/// The first tick fires immediately. 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"
),
}
}
});
}
#[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<i64>) -> 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_due() {
assert!(secret_is_due(&minted(None), NOW));
}
#[test]
fn a_secret_is_due_from_forty_five_days_and_not_before() {
assert!(!secret_is_due(&minted(Some(NOW)), NOW));
assert!(!secret_is_due(&minted(Some(NOW - 45 * DAY + 1)), NOW));
assert!(secret_is_due(&minted(Some(NOW - 45 * DAY)), NOW));
assert!(secret_is_due(&minted(Some(NOW - 90 * DAY)), NOW));
}
/// A mint time ahead of the clock is not due, rather than overflowing
/// into a large age.
#[test]
fn a_mint_time_in_the_future_is_not_due() {
assert!(!secret_is_due(&minted(Some(NOW + DAY)), NOW));
}
#[test]
fn a_re_mint_is_a_new_value_with_the_mint_time_for_the_same_agent() {
let old = minted(None);
let renewed = fresh("atlas", NOW).expect("the kernel supplies randomness");
assert_ne!(renewed.value, old.value);
assert_eq!(renewed.minted_at, Some(NOW));
assert_eq!(renewed.agent, "atlas");
assert!(!secret_is_due(&renewed, NOW), "a fresh secret is not due");
let again = fresh("atlas", NOW).expect("twice");
assert_ne!(again.value, renewed.value, "two re-mints must not agree");
}
#[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::<Vec<_>>(),
["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};
let o = |cert, secret| AgentObserved { cert, secret };
let observed = [
("both".to_owned(), o(Due(Some(DAY)), Due(None))),
("cert".to_owned(), o(Due(None), Current)),
("secret".to_owned(), o(Absent, Due(None))),
("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),
]
);
}
#[test]
fn the_age_logged_is_whole_days_or_unrecorded() {
assert_eq!(Age(Some(46 * DAY + 5)).to_string(), "46d");
assert_eq!(Age(None).to_string(), "unrecorded");
}
}

View file

@ -42,6 +42,7 @@ use utoipa_axum::{router::OpenApiRouter, routes};
mod agent_icon;
mod agent_identity;
mod agent_renewal;
mod agent_state_stream;
mod agent_status;
mod auth;
@ -122,6 +123,11 @@ enum SwarmNodeKind {
/// Carries no hive: the path it writes has no hive segment, so an agent
/// that moves between hives keeps one matrix identity.
MintAgentMatrixAccount { agent: String },
/// Re-mint `agent`'s queue secret if it is still stored and old enough.
/// See `agent_renewal`.
///
/// Carries no hive: the secret's store path has none.
RenewAgentQueueCredential { agent: String },
/// Declare `agent` on `hive` as `Paused` in the swarm's wanted-state
/// store, so a freshly created agent does not start driving turns the
/// moment it's deployed — the operator has to explicitly flip it to `Up`.
@ -148,6 +154,9 @@ impl hive_jobq_wire::WireNode for SwarmNodeKind {
SwarmNodeKind::MintAgentIdentity { .. } => "mint_agent_identity".to_owned(),
SwarmNodeKind::MintAgentForgeToken { .. } => "mint_agent_forge_token".to_owned(),
SwarmNodeKind::MintAgentMatrixAccount { .. } => "mint_agent_matrix_account".to_owned(),
SwarmNodeKind::RenewAgentQueueCredential { .. } => {
"renew_agent_queue_credential".to_owned()
}
SwarmNodeKind::SetAgentWanted { .. } => "set_agent_wanted".to_owned(),
SwarmNodeKind::TriggerDeploy { .. } => "trigger_deploy".to_owned(),
}
@ -168,7 +177,8 @@ impl hive_jobq_wire::WireNode for SwarmNodeKind {
| SwarmNodeKind::InitAgentConfigRepo { agent }
| SwarmNodeKind::MintAgentIdentity { agent }
| SwarmNodeKind::MintAgentForgeToken { agent }
| SwarmNodeKind::MintAgentMatrixAccount { agent } => {
| SwarmNodeKind::MintAgentMatrixAccount { agent }
| SwarmNodeKind::RenewAgentQueueCredential { agent } => {
serde_json::json!({ "agent": agent })
}
SwarmNodeKind::TriggerDeploy { hive, agent }
@ -319,6 +329,12 @@ async fn run_swarm_node(
SwarmNodeKind::MintAgentMatrixAccount { agent } => {
mint_matrix_account(deps.matrix_homeserver.as_deref(), &agent).await
}
SwarmNodeKind::RenewAgentQueueCredential { agent } => {
match agent_renewal::renew(&agent).await {
Ok(()) => Outcome::Done,
Err(e) => Outcome::Failed(format!("{e:#}")),
}
}
SwarmNodeKind::SetAgentWanted { hive, agent } => match deps.wanted {
None => Outcome::Failed(
"no swarm queue is configured on this host, so no wanted-state \
@ -1875,6 +1891,73 @@ fn spawn_forge_workers(
});
}
/// Insert one job per renewal and return the ids of its nodes:
/// `MintAgentIdentity` for a certificate, `RenewAgentQueueCredential` for a
/// queue secret, the second `after_any` the first when an agent needs both.
/// `agent_renewal::spawn`'s periodic pass comes through here.
///
/// Chained rather than parallel because `mint_and_verify` reads the queue
/// secret back and compares it with the one it read first; a re-mint landing
/// in between fails that node. `after_any`, so a certificate that could not be
/// re-issued does not hold up the secret.
fn queue_agent_renewals(
sched: &Mutex<hive_jobq::scheduler::Scheduler<SwarmNodeKind, SwarmResourceKind>>,
renewals: Vec<agent_renewal::Renewal>,
) -> Result<Vec<hive_jobq::NodeId>> {
let mut sched = sched
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
let mut ids = Vec::with_capacity(renewals.len());
for renewal in renewals {
let queued = sched
.insert_job(None, |b| {
let mut asked = Vec::new();
let cert = renewal.cert.then(|| {
b.node(SwarmNodeKind::MintAgentIdentity {
agent: renewal.agent.clone(),
})
.guid()
});
asked.extend(cert);
if renewal.secret {
let secret = b.node(SwarmNodeKind::RenewAgentQueueCredential {
agent: renewal.agent.clone(),
});
let secret = match cert {
Some(cert) => secret.after_any(cert),
None => secret,
};
asked.push(secret.guid());
}
asked
})
.map_err(|e| anyhow::anyhow!("{e}"))?;
ids.extend(queued);
}
Ok(ids)
}
/// Start the agent credential renewal pass when a swarm queue is wired up:
/// the pass reads the wanted-state declarations to tell live agents from
/// destroyed ones, and without a queue it has none. Lifted out of `main` for
/// `clippy::too_many_lines`.
fn spawn_agent_renewal(
jobq: &Arc<Mutex<hive_jobq::scheduler::Scheduler<SwarmNodeKind, SwarmResourceKind>>>,
wanted: Option<Arc<wanted::WantedWriter>>,
hives: &[HiveEntry],
) {
let Some(wanted) = wanted else {
return;
};
let hives = hives.iter().map(|h| h.name.clone()).collect();
let sched = Arc::clone(jobq);
agent_renewal::spawn(wanted, hives, move |renewals| {
if let Err(e) = queue_agent_renewals(&sched, renewals) {
tracing::warn!(error = %format!("{e:#}"), "agent credential renewal: queueing failed");
}
});
}
/// Check an agent's forge token now, and mint one if it is missing or stale.
///
/// The periodic pass (`forge::agent_token::spawn`) does the same every five
@ -2402,6 +2485,7 @@ async fn main() -> Result<()> {
let state_forge = keep_forge_for_state(forge_client, webhook_secret.clone());
let hives = load_hives();
spawn_agent_renewal(&jobq, wanted_writer(status.as_ref()), &hives);
// Before serving, because a hive whose role does not exist cannot log in,
// and one whose policy does not exist logs in able to read nothing —
// either way it cannot collect what this daemon writes for it. A store
@ -3387,6 +3471,88 @@ mod tests {
assert_eq!(kind.data(1)["agent"], "atlas");
}
/// The periodic pass's queueing: a certificate re-issue is the node agent
/// creation mints with, a secret re-mint its own node, and an agent due
/// both gets the secret `after_any` the certificate.
#[test]
fn a_queued_renewal_chains_the_secret_after_the_certificate() {
use crate::agent_renewal::Renewal;
use hive_jobq_wire::WireNode as _;
let sched = std::sync::Mutex::new(hive_jobq::scheduler::Scheduler::new(
hive_jobq::Graph::new(),
hive_jobq::resources::ResourceTable::new(),
));
let r = |agent: &str, cert, secret| Renewal {
agent: agent.to_owned(),
cert,
secret,
};
let ids = super::queue_agent_renewals(
&sched,
vec![
r("both", true, true),
r("cert", true, false),
r("secret", false, true),
],
)
.expect("three jobs insert");
assert_eq!(ids.len(), 4, "one returned id per node");
let guard = sched
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner);
let graph = guard.graph();
let agent_of = |n: &hive_jobq::Node<SwarmNodeKind, super::SwarmResourceKind>| {
n.payload.data(n.id.get())["agent"]
.as_str()
.expect("agent is a string")
.to_owned()
};
let mut nodes: Vec<(String, String)> = graph
.nodes()
.map(|n| (agent_of(n), n.payload.label()))
.collect();
nodes.sort();
assert_eq!(
nodes,
[
("both".to_owned(), "mint_agent_identity".to_owned()),
("both".to_owned(), "renew_agent_queue_credential".to_owned()),
("cert".to_owned(), "mint_agent_identity".to_owned()),
(
"secret".to_owned(),
"renew_agent_queue_credential".to_owned()
),
]
);
let find = |agent: &str, label: &str| {
graph
.nodes()
.find(|n| n.payload.label() == label && agent_of(n) == agent)
.unwrap_or_else(|| panic!("{agent} has a {label} node"))
};
let cert = find("both", "mint_agent_identity").id;
let secret = find("both", "renew_agent_queue_credential");
let when = secret
.deps
.iter()
.find_map(|d| match d {
hive_jobq::Dep::Node { id, when } if *id == cert => Some(*when),
_ => None,
})
.expect("the secret waits for the certificate");
assert!(
when.accepts(hive_jobq::TerminalState::Failed),
"a certificate that could not be re-issued must not hold up the secret"
);
assert!(
find("secret", "renew_agent_queue_credential")
.deps
.is_empty(),
"a secret on its own waits for nothing"
);
}
/// Agent creation mints the matrix account as a root of its own, and the
/// deploy waits for it without being cancelled by it: a host with no
/// homeserver configured must still deploy the agent.