Watch
0
0
Fork
You've already forked hyperhive
0

swarm-controller: read the queue client secret from the store, drop the file
Some checks were skipped
public bin cache / build + push to preem:grid (push) Has been skipped

The controller's OIDC client secret (client `swarm-controller`, used for
the queue connection, the auth-bridge bearer and the OTLP push) came from
an operator-placed file, `deploy.swarm-controller.queue.clientSecretFile`,
handed in by `LoadCredential=`.

Now `swarm-secret-publish`, which already copies authelia's minted OIDC
secrets into the store, also publishes this one, to
`swarm/controller/swarm-controller/oidc/client`. That path sits under
`controller/`, which no hive's policy reads. The controller reads it once
at start with its existing store certificate and holds it in memory, as
`swarm_queue_client::ClientSecret::Value`. If the store is down, it
retries for about a minute and then fails the start, so `Restart=` tries
again.

Policy delta: the controller gets `read` on that leaf, and the publisher
gets `create`/`update` on that leaf.

Removed: the `queue.clientSecretFile` option (both spellings, now removed
options with a message), its singleHostSwarm default, the credential and
placeholder, and the path watcher plus its restart oneshot. A controller
without a store identity is now an eval error, because it has no other
way to get the secret.
This commit is contained in:
atlas 2026-09-28 19:00:31 +02:00
commit e94406cdb9
22 changed files with 595 additions and 247 deletions

View file

@ -3,7 +3,7 @@
//! README for why the file cannot be touched from here).
//!
//! Authenticated with THIS daemon's own queue OIDC identity
//! (`SWARM_CONTROLLER_OIDC_*`, the same one `swarm-queue-client` mints for
//! ([`crate::queue_identity`], the same one `swarm-queue-client` mints for
//! the queue connection) — "one identity per principal" means a second
//! op that needs to prove who this process is reuses the identity it
//! already has rather than provisioning a new one. A fresh token is
@ -41,7 +41,7 @@ pub struct AuthBridge {
impl AuthBridge {
/// Read `SWARM_CONTROLLER_AUTH_BRIDGE_URL`; `Ok(None)` when unset.
///
/// The queue identity (`SWARM_CONTROLLER_OIDC_*`) is not optional once
/// The queue identity ([`crate::queue_identity`]) is not optional once
/// the bridge URL is set: the nix module sets `queueEnv` unconditionally
/// for every controller (the queue is required, not just co-located
/// service), so a bridge URL with no queue identity to authenticate
@ -51,15 +51,13 @@ impl AuthBridge {
let Ok(base_url) = std::env::var("SWARM_CONTROLLER_AUTH_BRIDGE_URL") else {
return Ok(None);
};
let queue_cfg = swarm_queue_client::QueueConfig::from_env("SWARM_CONTROLLER")
.context("reading the queue OIDC identity the auth bridge authenticates with")?
.ok_or_else(|| {
anyhow::anyhow!(
"SWARM_CONTROLLER_AUTH_BRIDGE_URL is set but SWARM_CONTROLLER_OIDC_* is \
let queue_cfg = crate::queue_identity::get().cloned().ok_or_else(|| {
anyhow::anyhow!(
"SWARM_CONTROLLER_AUTH_BRIDGE_URL is set but SWARM_CONTROLLER_OIDC_* is \
not — the bridge is authenticated with this daemon's queue identity, so \
that identity must exist first"
)
})?;
)
})?;
let http = reqwest::Client::builder()
.connect_timeout(HTTP_CONNECT_TIMEOUT)
.timeout(HTTP_TIMEOUT)

View file

@ -50,6 +50,7 @@ mod forge;
mod issue_report;
mod matrix_account;
mod otel_http_client;
mod queue_identity;
mod read_policy;
mod status;
mod store;
@ -2220,11 +2221,12 @@ fn register_swarm_webhooks(forge: Option<Arc<forge::Client>>, secret: Option<Arc
/// Deliberately NOT fatal on failure: the controller's HTTP surface is
/// useful without the queue, and a hive that cannot be read from renders
/// as `unknown` rather than as an outage of this daemon. What IS fatal is
/// a half-set environment — `QueueConfig::from_env` refuses that, because
/// silently behaving like an unconfigured host is how every hive ends up
/// reading `never_reported` with nothing to point at.
/// an identity that could not be resolved, which `queue_identity::init`
/// reports in `main` before this runs: silently behaving like an
/// unconfigured host is how every hive ends up reading `never_reported`
/// with nothing to point at.
async fn connect_status_reader() -> Result<Option<Arc<status::StatusReader>>> {
let Some(cfg) = swarm_queue_client::QueueConfig::from_env("SWARM_CONTROLLER")? else {
let Some(cfg) = queue_identity::get().cloned() else {
tracing::info!("no swarm queue configured; status aggregation is off");
return Ok(None);
};
@ -2304,6 +2306,9 @@ async fn main() -> Result<()> {
.with_context(|| format!("chmod {}", path.display()))?;
tracing::info!(socket = %path.display(), "swarm-controller listening");
// Before every user of the identity: the queue, the bridge and both OTLP
// exporters read it from `queue_identity::get`.
queue_identity::init().await?;
let status = connect_status_reader().await?;
// Same "not fatal, log and carry on" shape as the queue connect above:

View file

@ -0,0 +1,250 @@
//! This daemon's OIDC identity: the one it connects to the queue with and
//! mints the auth bridge's and the OTLP receiver's bearer tokens from.
//!
//! The queue's address, the token endpoint and the client id come from the
//! environment. The client secret comes from the swarm's secret store, at
//! `swarm_secret_client::queue::controller_client_path`, read once at start
//! under this daemon's own certificate and held only in memory — it is never
//! on this host's disk. `swarm-secret-publish`, on the host that runs
//! authelia, puts it there.
//!
//! [`init`] runs once in `main`, before anything calls [`get`].
use std::future::Future;
use std::path::PathBuf;
use std::sync::OnceLock;
use std::time::Duration;
use anyhow::{Context, Result, bail};
use swarm_queue_client::{ClientSecret, QueueConfig};
use swarm_secret_client::queue::{self, ControllerCredential};
/// The environment prefix the nix module's `queueEnv` sets.
const PREFIX: &str = "SWARM_CONTROLLER";
/// Seconds to wait after each failed fetch before the next one. About a minute
/// in all, which rides out a store restarting alongside this daemon. A longer
/// outage fails the start, and the unit's `Restart=` tries again from the top.
const RETRY_DELAYS_S: [u64; 6] = [1, 2, 4, 8, 16, 30];
static IDENTITY: OnceLock<Option<QueueConfig>> = OnceLock::new();
/// Everything in a [`QueueConfig`] except the secret.
#[derive(Debug, PartialEq, Eq)]
struct Coordinates {
url: String,
token_endpoint: String,
client_id: String,
ca_file: Option<PathBuf>,
}
impl Coordinates {
fn with_secret(self, secret: String) -> QueueConfig {
QueueConfig {
url: self.url,
token_endpoint: self.token_endpoint,
client_id: self.client_id,
client_secret: ClientSecret::Value(secret),
ca_file: self.ca_file,
}
}
}
/// Resolve the identity and hold it for [`get`].
///
/// # Errors
/// A half-set environment, which is a deployment bug and not an absent queue;
/// a secret the store still would not give up after [`RETRY_DELAYS_S`]; or a
/// second call.
pub async fn init() -> Result<()> {
let identity = match coordinates_from(|name| std::env::var(name).ok())? {
None => None,
Some(coordinates) => {
let delays = RETRY_DELAYS_S.map(Duration::from_secs);
Some(coordinates.with_secret(fetch_with_retry(&delays, fetch_secret).await?))
}
};
IDENTITY
.set(identity)
.map_err(|_| anyhow::anyhow!("the queue identity was initialised twice"))
}
/// The identity, or `None` when this deployment wired no queue up.
pub fn get() -> Option<&'static QueueConfig> {
IDENTITY.get().and_then(Option::as_ref)
}
/// `<PREFIX>_NATS_URL`, `_OIDC_TOKEN_ENDPOINT` and `_OIDC_CLIENT_ID`, all or
/// none; `_OIDC_CA_FILE` optionally.
fn coordinates_from(var: impl Fn(&str) -> Option<String>) -> Result<Option<Coordinates>> {
let url = var(&format!("{PREFIX}_NATS_URL"));
let token_endpoint = var(&format!("{PREFIX}_OIDC_TOKEN_ENDPOINT"));
let client_id = var(&format!("{PREFIX}_OIDC_CLIENT_ID"));
let ca_file = var(&format!("{PREFIX}_OIDC_CA_FILE")).map(PathBuf::from);
match (url, token_endpoint, client_id) {
(None, None, None) => Ok(None),
(Some(url), Some(token_endpoint), Some(client_id)) => Ok(Some(Coordinates {
url,
token_endpoint,
client_id,
ca_file,
})),
_ => bail!(
"swarm queue is half-configured: {PREFIX}_NATS_URL, {PREFIX}_OIDC_TOKEN_ENDPOINT \
and {PREFIX}_OIDC_CLIENT_ID must be set together or not at all"
),
}
}
/// One login and one read. An empty value is refused like an absent one: it
/// would reach authelia as the secret and come back as `invalid_client`, which
/// names neither the store nor the publisher.
async fn fetch_secret() -> Result<String> {
let path = queue::controller_client_path()?;
let store = crate::store::connect()
.await
.context("logging in to the swarm secret store")?;
let stored: Option<ControllerCredential> = store
.read_optional(&path)
.await
.with_context(|| format!("reading {path}"))?;
match stored {
Some(c) if !c.value.trim().is_empty() => Ok(c.value.trim().to_owned()),
_ => bail!(
"nothing at {path}: swarm-secret-publish, on the host that runs authelia, has not \
published this daemon's client secret"
),
}
}
/// Run `attempt`, waiting each of `delays` after a failure, then once more.
/// Answers with the first success, or the last attempt's error.
async fn fetch_with_retry<F, Fut>(delays: &[Duration], mut attempt: F) -> Result<String>
where
F: FnMut() -> Fut,
Fut: Future<Output = Result<String>>,
{
for delay in delays {
match attempt().await {
Ok(secret) => return Ok(secret),
Err(e) => {
tracing::warn!(
error = %format!("{e:#}"),
retry_in = ?delay,
"fetching the queue client secret from the swarm secret store failed"
);
tokio::time::sleep(*delay).await;
}
}
}
attempt().await.with_context(|| {
format!(
"fetching the queue client secret from the swarm secret store, after {} attempts",
delays.len() + 1
)
})
}
#[cfg(test)]
mod tests {
use std::collections::HashMap;
use super::*;
fn env(pairs: &[(&str, &str)]) -> impl Fn(&str) -> Option<String> {
let map: HashMap<String, String> = pairs
.iter()
.map(|(k, v)| ((*k).to_owned(), (*v).to_owned()))
.collect();
move |name| map.get(name).cloned()
}
const FULL: [(&str, &str); 3] = [
("SWARM_CONTROLLER_NATS_URL", "tls://nats.example:4222"),
(
"SWARM_CONTROLLER_OIDC_TOKEN_ENDPOINT",
"https://auth.example/api/oidc/token",
),
("SWARM_CONTROLLER_OIDC_CLIENT_ID", "swarm-controller"),
];
#[test]
fn no_coordinates_is_no_queue() {
assert_eq!(
coordinates_from(env(&[])).expect("absent is not an error"),
None
);
}
#[test]
fn a_half_set_environment_is_an_error_naming_the_variables() {
let err = coordinates_from(env(&FULL[..2])).expect_err("a partial set is not absent");
assert!(
format!("{err}").contains("SWARM_CONTROLLER_OIDC_CLIENT_ID"),
"{err}"
);
}
#[test]
fn a_full_environment_resolves_with_its_optional_ca() {
let mut with_ca = FULL.to_vec();
with_ca.push(("SWARM_CONTROLLER_OIDC_CA_FILE", "/etc/ca.pem"));
let c = coordinates_from(env(&with_ca))
.expect("complete")
.expect("configured");
assert_eq!(c.client_id, "swarm-controller");
assert_eq!(c.ca_file, Some(PathBuf::from("/etc/ca.pem")));
let c = coordinates_from(env(&FULL))
.expect("complete")
.expect("configured");
assert_eq!(c.ca_file, None);
}
#[test]
fn the_fetched_secret_is_held_as_a_value_not_a_path() {
let cfg = coordinates_from(env(&FULL))
.expect("complete")
.expect("configured")
.with_secret("s3cret".to_owned());
assert!(
matches!(&cfg.client_secret, ClientSecret::Value(v) if v == "s3cret"),
"{:?}",
cfg.client_secret
);
}
#[tokio::test]
async fn a_store_that_comes_back_within_the_window_is_waited_for() {
let mut calls = 0;
let got = fetch_with_retry(&[Duration::ZERO; 3], || {
calls += 1;
let n = calls;
async move {
if n < 3 {
bail!("store sealed")
}
Ok("s3cret".to_owned())
}
})
.await
.expect("the third attempt succeeds");
assert_eq!(got, "s3cret");
assert_eq!(calls, 3);
}
#[tokio::test]
async fn a_store_that_stays_down_fails_after_one_attempt_per_delay_plus_one() {
let mut calls = 0;
let err = fetch_with_retry(&[Duration::ZERO; 3], || {
calls += 1;
async { bail!("store sealed") }
})
.await
.expect_err("every attempt fails");
assert_eq!(calls, 4);
let msg = format!("{err:#}");
assert!(msg.contains("after 4 attempts"), "{msg}");
assert!(msg.contains("store sealed"), "{msg}");
}
}

View file

@ -138,7 +138,7 @@ fn build_provider(interval: Duration) -> Result<SdkMeterProvider> {
/// Build the [`crate::otel_http_client::AuthenticatedHttpClient`] this
/// exporter pushes through.
///
/// The queue identity (`SWARM_CONTROLLER_OIDC_*`) is read fresh here rather
/// The queue identity is read from [`crate::queue_identity::get`] rather
/// than threaded in from a caller — by the time this runs, [`endpoint`] has
/// already confirmed OTEL is configured for this host, and `queueEnv` is set
/// **unconditionally** for every `swarm-controller` (the nix module's own
@ -162,15 +162,13 @@ fn build_provider(interval: Duration) -> Result<SdkMeterProvider> {
/// the same fact stated twice.
pub(crate) fn authenticated_http_client() -> Result<crate::otel_http_client::AuthenticatedHttpClient>
{
let cfg = swarm_queue_client::QueueConfig::from_env("SWARM_CONTROLLER")
.context("reading the queue identity this daemon authenticates its OTLP push with")?
.ok_or_else(|| {
anyhow::anyhow!(
"SWARM_CONTROLLER_NATS_URL and friends are unset, but OTEL_EXPORTER_OTLP_ENDPOINT \
let cfg = crate::queue_identity::get().cloned().ok_or_else(|| {
anyhow::anyhow!(
"SWARM_CONTROLLER_NATS_URL and friends are unset, but OTEL_EXPORTER_OTLP_ENDPOINT \
is — the queue is required for every controller, so this combination is a \
deployment bug, not a supported partial config"
)
})?;
)
})?;
let audience = std::env::var("SWARM_CONTROLLER_OTEL_AUDIENCE").context(
"SWARM_CONTROLLER_OTEL_AUDIENCE is unset, but OTEL_EXPORTER_OTLP_ENDPOINT is — the nix \
module sets both together",