Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
79bc198165 | ||
|
|
d71222c206 | ||
|
|
62b9c76d69 | ||
|
|
a9603214c2 |
8 changed files with 424 additions and 212 deletions
15
Cargo.lock
generated
15
Cargo.lock
generated
|
|
@ -4561,9 +4561,9 @@ dependencies = [
|
|||
"async-nats",
|
||||
"axum",
|
||||
"futures-util",
|
||||
"reqwest 0.13.1",
|
||||
"serde",
|
||||
"serde_json",
|
||||
"swarm-queue-client",
|
||||
"tokio",
|
||||
"tracing",
|
||||
"tracing-subscriber",
|
||||
|
|
@ -4591,6 +4591,19 @@ dependencies = [
|
|||
"tracing-subscriber",
|
||||
]
|
||||
|
||||
[[package]]
|
||||
name = "swarm-queue-client"
|
||||
version = "0.1.0"
|
||||
dependencies = [
|
||||
"async-nats",
|
||||
"reqwest 0.13.1",
|
||||
"serde",
|
||||
"serde_json",
|
||||
"thiserror 2.0.18",
|
||||
"tokio",
|
||||
"tracing",
|
||||
]
|
||||
|
||||
[[package]]
|
||||
name = "swarmctl"
|
||||
version = "0.1.0"
|
||||
|
|
|
|||
|
|
@ -23,6 +23,7 @@ members = [
|
|||
"hivectl",
|
||||
"swarm-controller",
|
||||
"swarm-nats-auth",
|
||||
"swarm-queue-client",
|
||||
"swarmctl",
|
||||
]
|
||||
|
||||
|
|
@ -84,6 +85,7 @@ hive-host-sock = { path = "hive-host-sock" }
|
|||
hive-priv-sock = { path = "hive-priv-sock" }
|
||||
hive-sock-client = { path = "hive-sock-client" }
|
||||
hive-types = { path = "hive-types" }
|
||||
swarm-queue-client = { path = "swarm-queue-client" }
|
||||
thiserror = "2"
|
||||
tower-http = { version = "0.7", features = ["fs"] }
|
||||
uuid = { version = "1", features = ["v4"] }
|
||||
|
|
|
|||
|
|
@ -18,9 +18,13 @@ anyhow.workspace = true
|
|||
async-nats = { workspace = true, features = ["kv"] }
|
||||
axum.workspace = true
|
||||
futures-util.workspace = true
|
||||
reqwest.workspace = true
|
||||
serde.workspace = true
|
||||
serde_json.workspace = true
|
||||
# The queue connect (token mint + auth callback + reconnect) is shared with
|
||||
# every other participant - a hive publishing its own status runs the same
|
||||
# code with a different client id. Two copies of credential handling is one
|
||||
# token-refresh fix that has to be found twice.
|
||||
swarm-queue-client.workspace = true
|
||||
tokio.workspace = true
|
||||
tracing.workspace = true
|
||||
tracing-subscriber.workspace = true
|
||||
|
|
|
|||
|
|
@ -31,7 +31,6 @@ use serde::{Deserialize, Serialize};
|
|||
use utoipa::{OpenApi, ToSchema};
|
||||
use utoipa_axum::{router::OpenApiRouter, routes};
|
||||
|
||||
mod queue;
|
||||
mod status;
|
||||
|
||||
/// Where the daemon binds, overridable via `SWARM_CONTROLLER_SOCKET`.
|
||||
|
|
@ -298,12 +297,12 @@ async fn main() -> Result<()> {
|
|||
// 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.
|
||||
let status = match queue::QueueConfig::from_env()? {
|
||||
let status = match swarm_queue_client::QueueConfig::from_env("SWARM_CONTROLLER")? {
|
||||
None => {
|
||||
tracing::info!("no swarm queue configured; status aggregation is off");
|
||||
None
|
||||
}
|
||||
Some(cfg) => match queue::connect(cfg).await {
|
||||
Some(cfg) => match swarm_queue_client::connect(cfg).await {
|
||||
Ok(client) => {
|
||||
// NOT "connected": `retry_on_initial_connect` returns a client
|
||||
// before any connection has been established, so claiming a
|
||||
|
|
@ -319,7 +318,13 @@ async fn main() -> Result<()> {
|
|||
)))
|
||||
}
|
||||
Err(e) => {
|
||||
tracing::warn!(error = format!("{e:#}"), "swarm queue unreachable");
|
||||
// `chain`, not `{:#}`: this is the queue client's own
|
||||
// error type, and thiserror's Display ignores the
|
||||
// alternate flag — the source would be dropped silently.
|
||||
tracing::warn!(
|
||||
error = swarm_queue_client::chain(&e),
|
||||
"swarm queue unreachable"
|
||||
);
|
||||
None
|
||||
}
|
||||
},
|
||||
|
|
|
|||
|
|
@ -1,206 +0,0 @@
|
|||
//! The controller's client end of the swarm message queue.
|
||||
//!
|
||||
//! The queue admits every non-responder client through `auth_callout`: a
|
||||
//! client presents a token at CONNECT, the callout responder introspects it
|
||||
//! against authelia and mints a user JWT if it is good. So the controller is
|
||||
//! an ordinary client and needs an identity of its own — it is not a hive, and
|
||||
//! the per-hive clients issued from the roster are not its to use.
|
||||
//!
|
||||
//! Two things about that shape drive everything here:
|
||||
//!
|
||||
//! - **A token expires.** Authelia issues `client_credentials` access tokens
|
||||
//! with `expires_in: 3599`. Authentication happens at CONNECT, so a
|
||||
//! long-lived connection is fine — but a *reconnect* an hour later needs a
|
||||
//! token that was minted an hour later.
|
||||
//! - **`async-nats` re-runs an auth callback per connection attempt** (it is
|
||||
//! handed that attempt's nonce). So the refresh belongs in the callback and
|
||||
//! not in a timer: there is no window in which the client holds a token it
|
||||
//! minted for a previous connection.
|
||||
//!
|
||||
//! The alternative — mint once, pass a static `auth_token`, own the reconnect
|
||||
//! loop — fails in the way this subsystem exists to prevent: the controller
|
||||
//! keeps serving, its status data quietly stops updating, and nothing says so
|
||||
//! until someone reads a dashboard.
|
||||
|
||||
use std::path::PathBuf;
|
||||
|
||||
use anyhow::{Context, Result, bail};
|
||||
|
||||
/// Only the one field this needs; authelia returns several.
|
||||
#[derive(serde::Deserialize)]
|
||||
struct TokenResponse {
|
||||
access_token: String,
|
||||
}
|
||||
|
||||
/// Where the controller finds the queue and what it authenticates with.
|
||||
///
|
||||
/// Every field comes from an environment variable the NixOS module sets, the
|
||||
/// same way `load_hives` takes the roster — a config change is a redeploy, and
|
||||
/// this process reads no file it was not pointed at.
|
||||
#[derive(Debug, Clone)]
|
||||
pub struct QueueConfig {
|
||||
/// `nats://host:port` for the swarm queue.
|
||||
pub url: String,
|
||||
/// Authelia's token endpoint, e.g. `https://auth.<swarm>/api/oidc/token`.
|
||||
pub token_endpoint: String,
|
||||
/// The controller's own `OAuth2` client id.
|
||||
pub client_id: String,
|
||||
/// File holding the client secret's PLAINTEXT.
|
||||
///
|
||||
/// A path and not a value: the secret is minted on the authelia host and
|
||||
/// read here, and putting it in the environment would publish it to
|
||||
/// anything that can read `/proc/<pid>/environ`.
|
||||
pub client_secret_file: PathBuf,
|
||||
}
|
||||
|
||||
impl QueueConfig {
|
||||
/// Read the config from the environment, or `None` when the queue was not
|
||||
/// wired up for this deployment.
|
||||
///
|
||||
/// `None` rather than an error on purpose: the controller serves its HTTP
|
||||
/// surface on hosts where the queue is not enabled, and refusing to start
|
||||
/// there would trade a missing feature for a dead daemon. What must NOT
|
||||
/// happen is a *half* configuration silently behaving like an absent one —
|
||||
/// hence the explicit partial check below.
|
||||
pub fn from_env() -> Result<Option<Self>> {
|
||||
let url = std::env::var("SWARM_CONTROLLER_NATS_URL").ok();
|
||||
let token_endpoint = std::env::var("SWARM_CONTROLLER_OIDC_TOKEN_ENDPOINT").ok();
|
||||
let client_id = std::env::var("SWARM_CONTROLLER_OIDC_CLIENT_ID").ok();
|
||||
let secret = std::env::var("SWARM_CONTROLLER_OIDC_CLIENT_SECRET_FILE").ok();
|
||||
|
||||
match (url, token_endpoint, client_id, secret) {
|
||||
(None, None, None, None) => Ok(None),
|
||||
(Some(url), Some(token_endpoint), Some(client_id), Some(secret)) => Ok(Some(Self {
|
||||
url,
|
||||
token_endpoint,
|
||||
client_id,
|
||||
client_secret_file: PathBuf::from(secret),
|
||||
})),
|
||||
// A partially-set environment is a deployment bug, and the failure
|
||||
// it would otherwise produce is the expensive kind: the controller
|
||||
// comes up "fine", never connects, and every hive reads as having
|
||||
// never reported. Naming the missing variables costs one line.
|
||||
_ => bail!(
|
||||
"swarm queue is half-configured: SWARM_CONTROLLER_NATS_URL, \
|
||||
_OIDC_TOKEN_ENDPOINT, _OIDC_CLIENT_ID and \
|
||||
_OIDC_CLIENT_SECRET_FILE must be set together or not at all"
|
||||
),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// Mint a fresh access token for the controller's own client.
|
||||
///
|
||||
/// `client_credentials`, because there is no user here: the controller
|
||||
/// authenticates as itself. Authelia refuses the `openid` scope for this grant
|
||||
/// (a machine client receives an access token and never an id-token), so no
|
||||
/// scope is requested.
|
||||
async fn mint_token(http: &reqwest::Client, cfg: &QueueConfig) -> Result<String> {
|
||||
// Read per call rather than caching: the file is small, and a cached
|
||||
// secret would survive a rotation that the operator believes took effect.
|
||||
let secret = tokio::fs::read_to_string(&cfg.client_secret_file)
|
||||
.await
|
||||
.with_context(|| {
|
||||
format!(
|
||||
"reading the queue client secret from {}",
|
||||
cfg.client_secret_file.display()
|
||||
)
|
||||
})?;
|
||||
|
||||
let response = http
|
||||
.post(&cfg.token_endpoint)
|
||||
.form(&[
|
||||
("grant_type", "client_credentials"),
|
||||
("client_id", cfg.client_id.as_str()),
|
||||
("client_secret", secret.trim()),
|
||||
])
|
||||
.send()
|
||||
.await
|
||||
.context("requesting an access token from authelia")?;
|
||||
|
||||
// The body carries authelia's own error description, and it is far more
|
||||
// useful than the status alone: a wrong grant says `unauthorized_client`,
|
||||
// a wrong secret says `invalid_client`, and those point at different
|
||||
// config.
|
||||
let status = response.status();
|
||||
let body = response.text().await.unwrap_or_default();
|
||||
if !status.is_success() {
|
||||
bail!("authelia refused the controller's token request ({status}): {body}");
|
||||
}
|
||||
|
||||
let parsed: TokenResponse =
|
||||
serde_json::from_str(&body).context("parsing authelia's token response")?;
|
||||
Ok(parsed.access_token)
|
||||
}
|
||||
|
||||
/// Connect to the swarm queue, minting a token for each connection attempt.
|
||||
pub async fn connect(cfg: QueueConfig) -> Result<async_nats::Client> {
|
||||
// A timeout, because this client runs INSIDE the auth callback: a token
|
||||
// endpoint that accepts the connection and then never answers would hang
|
||||
// the callback, and with it the connection attempt that invoked it, with
|
||||
// no retry and nothing in the log to say why. Failing fast lets
|
||||
// `async-nats` do what it already does well — back off and try again.
|
||||
// 10s is generous for a form POST to a local IdP.
|
||||
let http = reqwest::Client::builder()
|
||||
.timeout(std::time::Duration::from_secs(10))
|
||||
.build()
|
||||
.context("building the token-endpoint HTTP client")?;
|
||||
let url = cfg.url.clone();
|
||||
|
||||
let client = async_nats::ConnectOptions::with_auth_callback(move |_nonce| {
|
||||
let http = http.clone();
|
||||
let cfg = cfg.clone();
|
||||
async move {
|
||||
let token = mint_token(&http, &cfg)
|
||||
.await
|
||||
// The callback's error type carries a string, so the context
|
||||
// chain would be lost; flatten it rather than dropping it.
|
||||
.map_err(|e| async_nats::AuthError::new(format!("{e:#}")))?;
|
||||
let mut auth = async_nats::Auth::new();
|
||||
auth.token = Some(token);
|
||||
Ok(auth)
|
||||
}
|
||||
})
|
||||
// The controller and the queue are separate units on (possibly)
|
||||
// separate hosts, and nothing orders them. Without this, a queue that
|
||||
// comes up one second later leaves the controller permanently
|
||||
// queue-less until someone restarts it — a boot-order race that
|
||||
// presents as "status has been unavailable since Tuesday".
|
||||
//
|
||||
// It also composes with the callback above rather than fighting it:
|
||||
// each background attempt is a connection attempt, so each one mints
|
||||
// its own token instead of retrying a stale one.
|
||||
.retry_on_initial_connect()
|
||||
.connect(&url)
|
||||
.await
|
||||
.with_context(|| format!("connecting to the swarm queue at {url}"))?;
|
||||
|
||||
Ok(client)
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
|
||||
/// The all-unset case is the common one — most hosts do not run the queue.
|
||||
#[test]
|
||||
fn an_absent_environment_is_not_an_error() {
|
||||
// Guard: this test would pass vacuously inside a configured
|
||||
// environment, so it asserts the variables really are unset first.
|
||||
for k in [
|
||||
"SWARM_CONTROLLER_NATS_URL",
|
||||
"SWARM_CONTROLLER_OIDC_TOKEN_ENDPOINT",
|
||||
"SWARM_CONTROLLER_OIDC_CLIENT_ID",
|
||||
"SWARM_CONTROLLER_OIDC_CLIENT_SECRET_FILE",
|
||||
] {
|
||||
if std::env::var(k).is_ok() {
|
||||
return;
|
||||
}
|
||||
}
|
||||
assert!(
|
||||
QueueConfig::from_env()
|
||||
.expect("absent is not an error")
|
||||
.is_none()
|
||||
);
|
||||
}
|
||||
}
|
||||
23
swarm-queue-client/Cargo.toml
Normal file
23
swarm-queue-client/Cargo.toml
Normal file
|
|
@ -0,0 +1,23 @@
|
|||
[package]
|
||||
name = "swarm-queue-client"
|
||||
version.workspace = true
|
||||
readme = "README.md"
|
||||
edition.workspace = true
|
||||
|
||||
[dependencies]
|
||||
# No `kv`/`jetstream` feature here on purpose: this crate's job ends at a
|
||||
# connected client. What a consumer does with it - KV for the controller and
|
||||
# the hive, plain messaging for anything later - is the consumer's business,
|
||||
# and its Cargo.toml is where that requirement should be visible.
|
||||
async-nats.workspace = true
|
||||
reqwest.workspace = true
|
||||
serde.workspace = true
|
||||
serde_json.workspace = true
|
||||
# A library, so its errors are a matchable enum rather than an opaque
|
||||
# `anyhow::Error`. The binaries that consume this keep anyhow; `?` converts.
|
||||
thiserror.workspace = true
|
||||
tokio.workspace = true
|
||||
tracing.workspace = true
|
||||
|
||||
[lints]
|
||||
workspace = true
|
||||
55
swarm-queue-client/README.md
Normal file
55
swarm-queue-client/README.md
Normal file
|
|
@ -0,0 +1,55 @@
|
|||
# swarm-queue-client
|
||||
|
||||
Connecting to the swarm message queue as an authenticated client. Shared by
|
||||
every process that participates: the swarm controller reads hive status out of
|
||||
the queue, a hive publishes its own status into it.
|
||||
|
||||
## Why a crate and not a module per binary
|
||||
|
||||
The *connect* is identical for every participant — mint an authelia token,
|
||||
present it at CONNECT for the `auth_callout` responder to introspect, let
|
||||
`async-nats` re-run the callback on each connection attempt. Only the **use**
|
||||
differs.
|
||||
|
||||
Two copies of that would be two copies of credential handling, and a
|
||||
token-refresh fix would have to be found twice. The same reasoning already put
|
||||
`hive-sock-client` in its own crate rather than in each daemon that speaks to a
|
||||
unix socket.
|
||||
|
||||
## The two properties that constrain the code
|
||||
|
||||
**A token expires.** Authelia issues `client_credentials` access tokens with
|
||||
`expires_in: 3599`. Authentication happens at CONNECT, so a long-lived
|
||||
connection is fine — but a *reconnect* an hour later needs a token minted an
|
||||
hour later.
|
||||
|
||||
**The refresh therefore lives in the auth callback, not in a timer.**
|
||||
`async-nats` invokes it per connection attempt, so there is no window in which
|
||||
the client holds a token it minted for a previous connection. The alternative —
|
||||
mint once, own the reconnect loop — fails in the way this subsystem exists to
|
||||
prevent: the process keeps serving while its data quietly stops moving, and
|
||||
nothing says so until someone reads a dashboard.
|
||||
|
||||
## Configuration
|
||||
|
||||
`QueueConfig::from_env(prefix)` reads `<prefix>_NATS_URL`,
|
||||
`<prefix>_OIDC_TOKEN_ENDPOINT`, `<prefix>_OIDC_CLIENT_ID` and
|
||||
`<prefix>_OIDC_CLIENT_SECRET_FILE`.
|
||||
|
||||
The prefix is a parameter because the variables belong to the consuming unit —
|
||||
a NixOS module sets them alongside its other options. What is shared is the
|
||||
rule, not the spelling: **all four together or none at all.** A half-set
|
||||
environment is a hard error, because the failure it would otherwise produce is
|
||||
the expensive kind — the process comes up "fine", never connects, and the data
|
||||
it was supposed to move silently stops.
|
||||
|
||||
The client secret is a **path, not a value**: putting it in the environment
|
||||
would publish it to anything that can read `/proc/<pid>/environ`. It is read
|
||||
per token request rather than cached, so a rotation the operator believes took
|
||||
effect actually did.
|
||||
|
||||
## What this crate does not do
|
||||
|
||||
It ends at a connected client. No `jetstream`/`kv` feature is enabled here —
|
||||
what a consumer does with the connection is its own business, and its
|
||||
`Cargo.toml` is where that requirement should be visible.
|
||||
316
swarm-queue-client/src/lib.rs
Normal file
316
swarm-queue-client/src/lib.rs
Normal file
|
|
@ -0,0 +1,316 @@
|
|||
//! Connecting to the swarm message queue as an authenticated client.
|
||||
//!
|
||||
//! Shared by every process that needs the queue — the swarm controller reads
|
||||
//! hive status out of it, a hive publishes its own status into it — because
|
||||
//! the *connect* is identical for all of them and only the use differs.
|
||||
//! Duplicating it per binary would put credential handling in two places, and
|
||||
//! a token-refresh fix would then have to be found twice.
|
||||
//!
|
||||
//! The queue admits every non-responder client through `auth_callout`: a
|
||||
//! client presents a token at CONNECT, the callout responder introspects it
|
||||
//! against authelia and mints a user JWT if it is good. So each participant is
|
||||
//! an ordinary client that needs an identity of its own — the controller is
|
||||
//! not a hive, and the per-hive clients issued from the roster are not its to
|
||||
//! use.
|
||||
//!
|
||||
//! Two things about that shape drive everything here:
|
||||
//!
|
||||
//! - **A token expires.** Authelia issues `client_credentials` access tokens
|
||||
//! with `expires_in: 3599`. Authentication happens at CONNECT, so a
|
||||
//! long-lived connection is fine — but a *reconnect* an hour later needs a
|
||||
//! token that was minted an hour later.
|
||||
//! - **`async-nats` re-runs an auth callback per connection attempt** (it is
|
||||
//! handed that attempt's nonce). So the refresh belongs in the callback and
|
||||
//! not in a timer: there is no window in which the client holds a token it
|
||||
//! minted for a previous connection.
|
||||
//!
|
||||
//! The alternative — mint once, pass a static `auth_token`, own the reconnect
|
||||
//! loop — fails in the way this subsystem exists to prevent: the controller
|
||||
//! keeps serving, its status data quietly stops updating, and nothing says so
|
||||
//! until someone reads a dashboard.
|
||||
|
||||
use std::path::PathBuf;
|
||||
|
||||
/// Everything that can go wrong reaching the swarm queue.
|
||||
///
|
||||
/// `thiserror` and not `anyhow` because this is a library: a caller gets a
|
||||
/// type it can match on, and the binaries that consume it keep using
|
||||
/// `anyhow` — `?` converts for free, so nothing downstream is more verbose
|
||||
/// for it. The same split `hive-claude` uses.
|
||||
///
|
||||
/// The variants are the failures an operator acts on differently: a
|
||||
/// half-configured environment is a deployment bug, a refused token is an
|
||||
/// identity-provider config problem, an unreachable queue is a network
|
||||
/// one. Collapsing them into one string would make that distinction a
|
||||
/// matter of reading prose.
|
||||
#[derive(Debug, thiserror::Error)]
|
||||
pub enum Error {
|
||||
#[error(
|
||||
"swarm queue is half-configured: {prefix}_NATS_URL, \
|
||||
{prefix}_OIDC_TOKEN_ENDPOINT, {prefix}_OIDC_CLIENT_ID and \
|
||||
{prefix}_OIDC_CLIENT_SECRET_FILE must be set together or not at all"
|
||||
)]
|
||||
PartialConfig { prefix: String },
|
||||
|
||||
#[error("reading the queue client secret from {path}")]
|
||||
ClientSecret {
|
||||
path: String,
|
||||
#[source]
|
||||
source: std::io::Error,
|
||||
},
|
||||
|
||||
#[error("building the token-endpoint HTTP client")]
|
||||
HttpClient(#[source] reqwest::Error),
|
||||
|
||||
#[error("requesting an access token from authelia")]
|
||||
TokenRequest(#[source] reqwest::Error),
|
||||
|
||||
/// The body carries authelia's own `error_description`, and it is far
|
||||
/// more useful than the status alone: a wrong grant says
|
||||
/// `unauthorized_client`, a wrong secret says `invalid_client`, and
|
||||
/// those point at different config.
|
||||
#[error("authelia refused the token request ({status}): {body}")]
|
||||
TokenRefused {
|
||||
status: reqwest::StatusCode,
|
||||
body: String,
|
||||
},
|
||||
|
||||
#[error("parsing authelia's token response")]
|
||||
TokenResponse(#[source] serde_json::Error),
|
||||
|
||||
#[error("connecting to the swarm queue at {url}")]
|
||||
Connect {
|
||||
url: String,
|
||||
#[source]
|
||||
source: async_nats::ConnectError,
|
||||
},
|
||||
}
|
||||
|
||||
/// Render an error and its source chain on one line.
|
||||
///
|
||||
/// **Use this instead of `{:#}` anywhere an [`Error`] is rendered without
|
||||
/// first being `?`-converted into an `anyhow::Error`.** `anyhow`'s
|
||||
/// `Display` special-cases `f.alternate()` to walk the source chain;
|
||||
/// `thiserror`'s derive does not, so `{e:#}` and `{e}` render
|
||||
/// identically for this type. A call site that formatted an
|
||||
/// `anyhow::Error` with `{:#}` and now holds an [`Error`] therefore keeps
|
||||
/// compiling, keeps looking right, and silently drops the cause — which
|
||||
/// is the half that says *why*, and is exactly what a one-shot boot
|
||||
/// warning with no retry needs most.
|
||||
///
|
||||
/// Public for that reason: the fix cannot live only inside this crate's
|
||||
/// own auth callback while the callers it was written for reach for
|
||||
/// `{:#}` and get nothing.
|
||||
pub fn chain(error: &dyn std::error::Error) -> String {
|
||||
let mut rendered = error.to_string();
|
||||
let mut source = error.source();
|
||||
while let Some(cause) = source {
|
||||
rendered.push_str(": ");
|
||||
rendered.push_str(&cause.to_string());
|
||||
source = cause.source();
|
||||
}
|
||||
rendered
|
||||
}
|
||||
|
||||
/// Only the one field this needs; authelia returns several.
|
||||
#[derive(serde::Deserialize)]
|
||||
struct TokenResponse {
|
||||
access_token: String,
|
||||
}
|
||||
|
||||
/// Where the controller finds the queue and what it authenticates with.
|
||||
///
|
||||
/// Every field comes from an environment variable the NixOS module sets, the
|
||||
/// same way `load_hives` takes the roster — a config change is a redeploy, and
|
||||
/// this process reads no file it was not pointed at.
|
||||
#[derive(Debug, Clone)]
|
||||
pub struct QueueConfig {
|
||||
/// `nats://host:port` for the swarm queue.
|
||||
pub url: String,
|
||||
/// Authelia's token endpoint, e.g. `https://auth.<swarm>/api/oidc/token`.
|
||||
pub token_endpoint: String,
|
||||
/// The controller's own `OAuth2` client id.
|
||||
pub client_id: String,
|
||||
/// File holding the client secret's PLAINTEXT.
|
||||
///
|
||||
/// A path and not a value: the secret is minted on the authelia host and
|
||||
/// read here, and putting it in the environment would publish it to
|
||||
/// anything that can read `/proc/<pid>/environ`.
|
||||
pub client_secret_file: PathBuf,
|
||||
}
|
||||
|
||||
impl QueueConfig {
|
||||
/// Read the config from `<prefix>_NATS_URL`, `<prefix>_OIDC_TOKEN_ENDPOINT`,
|
||||
/// `<prefix>_OIDC_CLIENT_ID` and `<prefix>_OIDC_CLIENT_SECRET_FILE`, or
|
||||
/// `None` when the queue was not wired up for this deployment.
|
||||
///
|
||||
/// The prefix is a parameter rather than a constant because the variables
|
||||
/// belong to the *consuming unit* — a NixOS module sets them alongside its
|
||||
/// other options, and two daemons sharing one name would be a worse
|
||||
/// coupling than passing four characters. What is shared is the RULE
|
||||
/// below, not the spelling.
|
||||
///
|
||||
/// `None` rather than an error on purpose: a daemon serves its other
|
||||
/// surfaces on hosts where the queue is not enabled, and refusing to start
|
||||
/// there would trade a missing feature for a dead process. What must NOT
|
||||
/// happen is a *half* configuration silently behaving like an absent one —
|
||||
/// hence the explicit partial check below.
|
||||
pub fn from_env(prefix: &str) -> Result<Option<Self>, Error> {
|
||||
let url = std::env::var(format!("{prefix}_NATS_URL")).ok();
|
||||
let token_endpoint = std::env::var(format!("{prefix}_OIDC_TOKEN_ENDPOINT")).ok();
|
||||
let client_id = std::env::var(format!("{prefix}_OIDC_CLIENT_ID")).ok();
|
||||
let secret = std::env::var(format!("{prefix}_OIDC_CLIENT_SECRET_FILE")).ok();
|
||||
|
||||
match (url, token_endpoint, client_id, secret) {
|
||||
(None, None, None, None) => Ok(None),
|
||||
(Some(url), Some(token_endpoint), Some(client_id), Some(secret)) => Ok(Some(Self {
|
||||
url,
|
||||
token_endpoint,
|
||||
client_id,
|
||||
client_secret_file: PathBuf::from(secret),
|
||||
})),
|
||||
// A partially-set environment is a deployment bug, and the failure
|
||||
// it would otherwise produce is the expensive kind: the process
|
||||
// comes up "fine", never connects, and the data it was supposed to
|
||||
// move silently stops moving. Naming the variables costs one line.
|
||||
_ => Err(Error::PartialConfig {
|
||||
prefix: prefix.to_owned(),
|
||||
}),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// Mint a fresh access token for the controller's own client.
|
||||
///
|
||||
/// `client_credentials`, because there is no user here: the controller
|
||||
/// authenticates as itself. Authelia refuses the `openid` scope for this grant
|
||||
/// (a machine client receives an access token and never an id-token), so no
|
||||
/// scope is requested.
|
||||
async fn mint_token(http: &reqwest::Client, cfg: &QueueConfig) -> Result<String, Error> {
|
||||
// Read per call rather than caching: the file is small, and a cached
|
||||
// secret would survive a rotation that the operator believes took effect.
|
||||
let secret = tokio::fs::read_to_string(&cfg.client_secret_file)
|
||||
.await
|
||||
.map_err(|source| Error::ClientSecret {
|
||||
path: cfg.client_secret_file.display().to_string(),
|
||||
source,
|
||||
})?;
|
||||
|
||||
let response = http
|
||||
.post(&cfg.token_endpoint)
|
||||
.form(&[
|
||||
("grant_type", "client_credentials"),
|
||||
("client_id", cfg.client_id.as_str()),
|
||||
("client_secret", secret.trim()),
|
||||
])
|
||||
.send()
|
||||
.await
|
||||
.map_err(Error::TokenRequest)?;
|
||||
|
||||
// Why the body and not just the status: see `Error::TokenRefused`.
|
||||
let status = response.status();
|
||||
let body = response.text().await.unwrap_or_default();
|
||||
if !status.is_success() {
|
||||
return Err(Error::TokenRefused { status, body });
|
||||
}
|
||||
|
||||
let parsed: TokenResponse = serde_json::from_str(&body).map_err(Error::TokenResponse)?;
|
||||
Ok(parsed.access_token)
|
||||
}
|
||||
|
||||
/// Connect to the swarm queue, minting a token for each connection attempt.
|
||||
pub async fn connect(cfg: QueueConfig) -> Result<async_nats::Client, Error> {
|
||||
// A timeout, because this client runs INSIDE the auth callback: a token
|
||||
// endpoint that accepts the connection and then never answers would hang
|
||||
// the callback, and with it the connection attempt that invoked it, with
|
||||
// no retry and nothing in the log to say why. Failing fast lets
|
||||
// `async-nats` do what it already does well — back off and try again.
|
||||
// 10s is generous for a form POST to a local IdP.
|
||||
let http = reqwest::Client::builder()
|
||||
.timeout(std::time::Duration::from_secs(10))
|
||||
.build()
|
||||
.map_err(Error::HttpClient)?;
|
||||
let url = cfg.url.clone();
|
||||
|
||||
let client = async_nats::ConnectOptions::with_auth_callback(move |_nonce| {
|
||||
let http = http.clone();
|
||||
let cfg = cfg.clone();
|
||||
async move {
|
||||
let token = mint_token(&http, &cfg)
|
||||
.await
|
||||
// The callback's error type carries a string, so the source
|
||||
// chain would be lost; flatten it rather than dropping it.
|
||||
.map_err(|e| async_nats::AuthError::new(chain(&e)))?;
|
||||
let mut auth = async_nats::Auth::new();
|
||||
auth.token = Some(token);
|
||||
Ok(auth)
|
||||
}
|
||||
})
|
||||
// The controller and the queue are separate units on (possibly)
|
||||
// separate hosts, and nothing orders them. Without this, a queue that
|
||||
// comes up one second later leaves the controller permanently
|
||||
// queue-less until someone restarts it — a boot-order race that
|
||||
// presents as "status has been unavailable since Tuesday".
|
||||
//
|
||||
// It also composes with the callback above rather than fighting it:
|
||||
// each background attempt is a connection attempt, so each one mints
|
||||
// its own token instead of retrying a stale one.
|
||||
.retry_on_initial_connect()
|
||||
.connect(&url)
|
||||
.await
|
||||
.map_err(|source| Error::Connect {
|
||||
url: url.clone(),
|
||||
source,
|
||||
})?;
|
||||
|
||||
Ok(client)
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
|
||||
/// The all-unset case is the common one — most hosts do not run the queue.
|
||||
///
|
||||
/// Uses a prefix no deployment sets, so it cannot pass vacuously by
|
||||
/// running inside a configured environment (the guard below covers the
|
||||
/// same ground, and both are cheap).
|
||||
#[test]
|
||||
fn an_absent_environment_is_not_an_error() {
|
||||
for k in [
|
||||
"SWARM_QUEUE_TEST_NATS_URL",
|
||||
"SWARM_QUEUE_TEST_OIDC_TOKEN_ENDPOINT",
|
||||
"SWARM_QUEUE_TEST_OIDC_CLIENT_ID",
|
||||
"SWARM_QUEUE_TEST_OIDC_CLIENT_SECRET_FILE",
|
||||
] {
|
||||
assert!(std::env::var(k).is_err(), "{k} must be unset for this test");
|
||||
}
|
||||
assert!(
|
||||
QueueConfig::from_env("SWARM_QUEUE_TEST")
|
||||
.expect("absent is not an error")
|
||||
.is_none()
|
||||
);
|
||||
}
|
||||
|
||||
/// The half-set case is the one the rule exists for: a deployment bug that
|
||||
/// would otherwise look exactly like "no queue configured".
|
||||
///
|
||||
/// SAFETY: single-threaded mutation of a process env var under a prefix no
|
||||
/// other test or deployment uses; removed before returning.
|
||||
#[test]
|
||||
fn a_half_set_environment_is_a_hard_error() {
|
||||
unsafe {
|
||||
std::env::set_var("SWARM_QUEUE_HALF_NATS_URL", "nats://127.0.0.1:4222");
|
||||
}
|
||||
let err = QueueConfig::from_env("SWARM_QUEUE_HALF")
|
||||
.expect_err("a partial set must not read as absent");
|
||||
let msg = format!("{err}");
|
||||
assert!(
|
||||
msg.contains("SWARM_QUEUE_HALF_OIDC_CLIENT_ID"),
|
||||
"the error must name the missing variables, got: {msg}"
|
||||
);
|
||||
unsafe {
|
||||
std::env::remove_var("SWARM_QUEUE_HALF_NATS_URL");
|
||||
}
|
||||
}
|
||||
}
|
||||
Loading…
Reference in a new issue