//! 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> = 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, } 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) } /// `_NATS_URL`, `_OIDC_TOKEN_ENDPOINT` and `_OIDC_CLIENT_ID`, all or /// none; `_OIDC_CA_FILE` optionally. fn coordinates_from(var: impl Fn(&str) -> Option) -> Result> { 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 { let path = queue::controller_client_path()?; let store = crate::store::connect() .await .context("logging in to the swarm secret store")?; let stored: Option = 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(delays: &[Duration], mut attempt: F) -> Result where F: FnMut() -> Fut, Fut: Future>, { 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 { let map: HashMap = 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}"); } }