Watch
0
0
Fork
You've already forked hyperhive
0

swarm-nats-auth: verify an agent's own token against the store

An `auth_token` spelled `swarm-agent.<agent>.<secret>` is no longer sent
to introspection. The responder reads `swarm/agents/<agent>/queue` with
an identity of its own, checks that the stored object names the same
agent, compares the secret in constant time, and grants the subjects
`--agent-token-publish-subject` lists with `{agent}` expanded. Every
other outcome denies: a malformed token, no store identity, nothing
stored, a failed or slow lookup, a different secret. A token without
the prefix takes the OIDC path unchanged.

The journal's `auth request` line names such a caller `agent:<agent>`;
the hive-shared credential keeps `hive-<h>-agent`.

The new principal: a `swarm-nats-auth` cert-auth role and policy with
read on `secret/data/swarm/agents/+/queue` alone, a leaf signed by the
store's PKI glue, and `glue-nats-auth-bao-identity.nix` pairing the two.
The copy unit delivers the identity into the queue's container, and an
absent leaf is delivered empty so the responder still starts and only
agent tokens are refused.

The policy and role are written by `swarm-bao-nats-auth-policy`, logged in
as the bao granter: both names fall under its `swarm-*` globs, so the
deploy writes them with no operator step. module-eval counts it among the
granting units, so every generic granting-unit case covers it.

The secret compare uses `subtle`, already in the lock file through the
TLS stack; no workspace crate offered one directly.
This commit is contained in:
atlas 2026-09-27 20:38:34 +02:00
commit 9bad58d86d
14 changed files with 846 additions and 51 deletions

View file

@ -22,12 +22,17 @@ serde.workspace = true
serde_json.workspace = true
# The jti digest: base32hex(sha256(claims)) over every JWT this crate signs.
sha2.workspace = true
# Comparing an agent's presented secret with the stored one.
subtle.workspace = true
# For `status::BUCKET` and `notices::STREAM` - the subjects a hive may
# publish to are derived from these names, and every end that touches them
# must agree on the same one. Deliberately WITHOUT the `kv` feature: this
# crate derives subject strings, it never opens the bucket. `notices`
# is name-only too (no `jetstream`/`kv` surface), same reason.
swarm-queue-client = { workspace = true, features = ["notices"] }
# The agent-token spelling the agent also uses, and the read of the stored
# credential it is checked against.
swarm-secret-client.workspace = true
tokio.workspace = true
tracing.workspace = true
tracing-subscriber.workspace = true

View file

@ -0,0 +1,382 @@
//! Verifying an agent's own credential, presented as `auth_token`.
//!
//! The spelling is `swarm_queue_client::agent_token`'s, shared with the agent
//! that presents it: a prefix, the agent's name, and its secret. The name is only a
//! claim. It becomes an identity when the secret equals the one stored at
//! `swarm/agents/<agent>/queue`, and nothing in this module grants on the name
//! alone: every path that does not reach that comparison, and every lookup
//! that fails, is a denial.
//!
//! A token without the prefix is not this module's: it goes to introspection
//! unchanged.
use std::future::Future;
use anyhow::Context;
use subtle::ConstantTimeEq;
use swarm_queue_client::agent_token::{AgentToken, Malformed, parse_agent_token};
use swarm_secret_client::client::{DEFAULT_CERT_MOUNT, Settings};
use swarm_secret_client::queue::{self, AgentCredential};
use crate::policy::{Permissions, Policy};
/// What a presented `auth_token` is.
pub enum Presented<'a> {
/// An agent's own credential, to be checked against the store.
Agent(AgentToken<'a>),
/// Carries the agent-token prefix but is not well-formed. Denied without a
/// lookup, and never handed to introspection: it holds a secret meant for
/// the store, not for the `IdP`.
Malformed(Malformed),
/// Anything else, which is an OIDC access token for introspection.
Bearer(&'a str),
}
/// Sort `token` onto the agent path or the introspection path.
pub fn classify(token: &str) -> Presented<'_> {
match parse_agent_token(token) {
Some(Ok(agent)) => Presented::Agent(agent),
Some(Err(e)) => Presented::Malformed(e),
None => Presented::Bearer(token),
}
}
/// Where the stored credential for an agent comes from.
pub trait CredentialSource {
/// The object at `agent`'s queue path, `None` when nothing is stored
/// there, or an error when the store could not answer.
fn lookup(
&self,
agent: &str,
) -> impl Future<Output = anyhow::Result<Option<AgentCredential>>> + Send;
}
/// The swarm secret store, reached with this responder's own certificate.
pub struct Store {
settings: Settings,
cert_role: String,
}
impl Store {
/// The store named by the `BAO_*` environment, or `None` when that is
/// unset: a responder with no store identity denies every agent token and
/// serves the OIDC path unchanged.
pub fn from_env(cert_role: String) -> Option<Self> {
match Settings::from_env() {
Ok(settings) => Some(Self {
settings,
cert_role,
}),
Err(e) => {
tracing::info!(reason = %e, "no secret store configured; agent tokens are denied");
None
}
}
}
}
impl CredentialSource for Store {
/// Logs in per lookup. An agent connects once per boot and per reconnect,
/// so a login each time costs little, and it leaves no token to expire
/// inside a long-lived process.
async fn lookup(&self, agent: &str) -> anyhow::Result<Option<AgentCredential>> {
let path = queue::agent_queue_path(agent)?;
let store = swarm_secret_client::SecretStore::connect(
&self.settings,
&self.cert_role,
DEFAULT_CERT_MOUNT,
)
.await
.context("logging in to the secret store")?;
store
.read_optional(&path)
.await
.with_context(|| format!("reading {path}"))
}
}
/// Whether `token`'s secret is the one stored for the agent it names.
///
/// The lookup is bounded by introspection's budget, under the server's
/// `authorization.timeout`. The callout loop answers one request at a time, so
/// a store that does not answer holds every request behind it, OIDC ones
/// included, for at most that long, and this then denies.
async fn verify(source: &impl CredentialSource, token: &AgentToken<'_>) -> bool {
let agent = token.agent;
let stored = match tokio::time::timeout(
crate::introspect::INTROSPECTION_TIMEOUT,
source.lookup(agent),
)
.await
{
Ok(Ok(Some(stored))) => stored,
Ok(Ok(None)) => {
tracing::warn!(
agent,
"no queue credential is stored for this agent; denying"
);
return false;
}
Ok(Err(e)) => {
tracing::warn!(
agent,
error = format!("{e:#}"),
"credential lookup failed; denying"
);
return false;
}
Err(_) => {
tracing::warn!(agent, "credential lookup timed out; denying");
return false;
}
};
if stored.agent != agent {
tracing::warn!(
agent,
stored_agent = %stored.agent,
"the stored credential names a different agent than its path; denying"
);
return false;
}
// Constant time, so the time to refuse says nothing about how much of the
// secret was right. A length mismatch returns early; the length is not
// secret, every minted secret has the same one.
stored
.value
.as_bytes()
.ct_eq(token.secret.as_bytes())
.into()
}
/// The grant for an agent token, or `None` for a denial.
///
/// `source` is `None` when this responder has no store identity, which denies.
pub async fn authorize<S: CredentialSource>(
policy: &Policy,
source: Option<&S>,
token: &AgentToken<'_>,
) -> Option<Permissions> {
let Some(source) = source else {
tracing::warn!(
agent = token.agent,
"agent token presented, but this responder has no secret store; denying"
);
return None;
};
if !verify(source, token).await {
return None;
}
let permissions = policy.agent_token_permissions(token.agent);
if permissions.is_none() {
tracing::warn!(
agent = token.agent,
"verified agent token, but no --agent-token-publish-subject is configured; denying"
);
}
permissions
}
#[cfg(test)]
mod tests {
use super::*;
/// A store holding at most one credential, or failing outright.
enum Fake {
/// `credential` is stored at `path_agent`'s queue path.
Holds {
path_agent: &'static str,
credential: AgentCredential,
},
Fails,
/// A store that accepts the request and never answers.
Hangs,
}
impl CredentialSource for Fake {
async fn lookup(&self, agent: &str) -> anyhow::Result<Option<AgentCredential>> {
match self {
Self::Holds {
path_agent,
credential,
} => Ok((agent == *path_agent).then(|| credential.clone())),
Self::Fails => anyhow::bail!("store unreachable"),
Self::Hangs => std::future::pending().await,
}
}
}
fn stored(path_agent: &'static str, named: &str) -> Fake {
Fake::Holds {
path_agent,
credential: AgentCredential {
value: "Ab9_-zSECRET".to_owned(),
agent: named.to_owned(),
},
}
}
fn atlas() -> Fake {
stored("atlas", "atlas")
}
fn policy() -> Policy {
Policy::new(
"hive-".to_owned(),
"-agent".to_owned(),
"hive-status".to_owned(),
vec!["swarm-controller".to_owned()],
vec![],
vec!["$SWARM.term.{hive}.>".to_owned()],
)
.expect("valid")
.with_agent_token_subjects(vec![
"$SWARM.term.{agent}".to_owned(),
"$SWARM.agent-state.{agent}".to_owned(),
])
.expect("valid")
}
async fn grant(source: Option<&Fake>, token: &str) -> Option<Permissions> {
let Presented::Agent(token) = classify(token) else {
panic!("{token:?} must classify as an agent token");
};
authorize(&policy(), source, &token).await
}
#[tokio::test]
async fn a_valid_agent_token_is_granted_exactly_its_own_subjects() {
let g = grant(Some(&atlas()), "swarm-agent.atlas.Ab9_-zSECRET")
.await
.expect("granted");
assert_eq!(
g.publish,
vec![
"$SWARM.term.atlas".to_owned(),
"$SWARM.agent-state.atlas".to_owned(),
]
);
}
#[tokio::test]
async fn a_wrong_secret_is_denied() {
assert!(
grant(Some(&atlas()), "swarm-agent.atlas.Ab9_-zSECREX")
.await
.is_none()
);
assert!(
grant(Some(&atlas()), "swarm-agent.atlas.Ab9_-zSECRE")
.await
.is_none()
);
}
/// The secret check is what stops one agent claiming another's name: the
/// right secret under someone else's name finds that agent's credential,
/// or none.
#[tokio::test]
async fn another_agents_name_with_this_agents_secret_is_denied() {
assert!(
grant(Some(&atlas()), "swarm-agent.argus.Ab9_-zSECRET")
.await
.is_none()
);
}
#[tokio::test]
async fn an_agent_with_nothing_stored_is_denied() {
assert!(
grant(
Some(&stored("argus", "argus")),
"swarm-agent.atlas.Ab9_-zSECRET"
)
.await
.is_none()
);
}
#[tokio::test]
async fn a_failed_lookup_is_denied() {
assert!(
grant(Some(&Fake::Fails), "swarm-agent.atlas.Ab9_-zSECRET")
.await
.is_none()
);
}
/// The callout loop answers one request at a time, so a store that never
/// answers must cost at most the bound, and then deny.
#[tokio::test]
async fn a_store_that_never_answers_is_denied_within_the_bound() {
let bound = crate::introspect::INTROSPECTION_TIMEOUT;
let started = std::time::Instant::now();
assert!(
grant(Some(&Fake::Hangs), "swarm-agent.atlas.Ab9_-zSECRET")
.await
.is_none()
);
let took = started.elapsed();
assert!(took >= bound, "denied before the bound: {took:?}");
assert!(
took < bound + std::time::Duration::from_millis(250),
"denied well after the bound: {took:?}"
);
}
#[tokio::test]
async fn no_store_is_denied() {
assert!(
grant(None, "swarm-agent.atlas.Ab9_-zSECRET")
.await
.is_none()
);
}
#[tokio::test]
async fn a_stored_object_naming_another_agent_is_denied() {
assert!(
grant(
Some(&stored("argus", "atlas")),
"swarm-agent.argus.Ab9_-zSECRET"
)
.await
.is_none()
);
}
#[tokio::test]
async fn a_verified_agent_with_no_subjects_configured_is_denied() {
let bare = Policy::new(
"hive-".to_owned(),
"-agent".to_owned(),
"hive-status".to_owned(),
vec![],
vec![],
vec![],
)
.expect("valid");
let Presented::Agent(token) = classify("swarm-agent.atlas.Ab9_-zSECRET") else {
panic!("an agent token");
};
assert!(authorize(&bare, Some(&atlas()), &token).await.is_none());
}
/// Everything without the prefix reaches introspection as it arrived.
#[test]
fn a_token_without_the_prefix_goes_to_introspection_unchanged() {
for token in ["authelia_at_abc.def", "atlas.Ab9_-zSECRET"] {
let Presented::Bearer(t) = classify(token) else {
panic!("{token:?} is not an agent token");
};
assert_eq!(t, token);
}
}
#[test]
fn a_malformed_agent_token_goes_nowhere() {
assert!(matches!(
classify("swarm-agent.atlas"),
Presented::Malformed(_)
));
}
}

View file

@ -8,7 +8,8 @@
//! It connects as the one callout-exempt user (by nkey, never by name — the
//! server refuses to start if that entry carries a username), subscribes to
//! `$SYS.REQ.USER.AUTH`, validates the presented bearer token against
//! authelia's introspection endpoint, and answers with a NATS user JWT signed
//! authelia's introspection endpoint (or, for an agent's own token, against the
//! secret store: see `agent_token`), and answers with a NATS user JWT signed
//! by the account key. A rejection is answered explicitly: silence is
//! indistinguishable from the responder being down, and the queue is the
//! swarm's control path.
@ -27,6 +28,7 @@ use anyhow::Context;
use clap::Parser;
use futures_util::StreamExt;
mod agent_token;
mod introspect;
mod policy;
mod request;
@ -89,7 +91,7 @@ struct Args {
/// config change plus a reload.
///
/// So this identity says which hive an agent belongs to and never which
/// agent: two agents on one hive are indistinguishable to this responder.
/// agent: two agents on one hive are indistinguishable under it.
///
/// A suffix on the hive's id rather than a prefix of its own, because
/// `agent-<name>` reads as *the agent called `<name>`* — the one thing
@ -127,6 +129,18 @@ struct Args {
/// is refused, which is loud rather than silently over-broad.
#[arg(long = "agent-publish-subject")]
agent_publish_subjects: Vec<String>,
/// Subjects an agent that presented its own credential may publish to,
/// with `{agent}` standing for its name. Repeatable, empty by default,
/// which denies every agent token.
#[arg(long = "agent-token-publish-subject")]
agent_token_publish_subjects: Vec<String>,
/// Role on the secret store's `cert` auth mount this responder logs in
/// as to read agent credentials. The store is found through the `BAO_*`
/// environment; with none, every agent token is denied.
#[arg(long, default_value = "swarm-nats-auth")]
store_cert_role: String,
}
/// Read a secret file and strip surrounding whitespace.
@ -177,7 +191,9 @@ async fn main() -> anyhow::Result<()> {
args.reader_clients.clone(),
args.hive_publish_subjects.clone(),
args.agent_publish_subjects.clone(),
)?;
)?
.with_agent_token_subjects(args.agent_token_publish_subjects.clone())?;
let store = agent_token::Store::from_env(args.store_cert_role.clone());
let http = reqwest::Client::new();
let issuer = nkeys::KeyPair::from_seed(&read_secret(&args.issuer_seed_file)?)
.context("parse the account signing seed")?;
@ -210,44 +226,33 @@ async fn main() -> anyhow::Result<()> {
};
// No token is a denial, not an error: an anonymous connect is a
// normal thing for a client to attempt and an abnormal thing to
// grant. Introspection is only reached once something was presented.
//
// The caller is an identity or nothing — see `introspect`'s module
// docs. There is no "admitted, identity unknown" branch to write here
// because there is no such value to receive.
let caller = match &req.connect_opts.auth_token {
Some(token) => introspect::identify_caller(
&http,
&args.introspection_url,
&args.client_id,
&client_secret,
token,
)
.await
// An introspection that could not be *made* is a denial too. The
// failure modes of an HTTP call are exactly the conditions under
// which an attacker would most like this to fall open.
.unwrap_or_else(|e| {
tracing::warn!(error = ?e, "introspection failed; denying");
None
}),
None => None,
// grant. Neither the store nor introspection is reached until
// something was presented.
let (caller, permissions) = match req
.connect_opts
.auth_token
.as_deref()
.map(agent_token::classify)
{
// The caller is the name the token claims, logged whether or not
// the secret proved it; `granted` says which.
Some(agent_token::Presented::Agent(token)) => (
Some(format!("agent:{}", token.agent)),
agent_token::authorize(&policy, store.as_ref(), &token).await,
),
Some(agent_token::Presented::Malformed(e)) => {
tracing::warn!(error = %e, "malformed agent token; denying");
(None, None)
}
Some(agent_token::Presented::Bearer(token)) => {
introspected(&policy, &http, &args, &client_secret, token).await
}
None => (None, None),
};
// Admission said who; the policy says what. A caller the `IdP`
// vouches for but no rule matches is denied — see `policy`'s module
// docs for why that is deny and not "connect with nothing".
let permissions = caller.as_deref().and_then(|id| policy.permissions(id));
if let (Some(id), None) = (caller.as_deref(), permissions.as_ref()) {
// Loud, and the one case an operator has to be able to find: a
// valid credential refused by our own policy. The alternative is
// a client that authenticates fine and mysteriously cannot work.
tracing::warn!(
caller = %id,
"authenticated client matches no policy rule; denying"
);
}
// The client id is an identifier, not a credential, and it is the
// only thing tying a connection in this log to a hive.
// The caller is an identifier, not a credential: a client id, or
// `agent:<name>` for an agent token. One line per auth request, a connect
// or a reconnect, so `hive-<h>-agent` lines count requests made on a
// hive's shared credential, not agents.
tracing::info!(
user_nkey = %req.user_nkey,
server_id = %req.server_id.id,
@ -281,6 +286,50 @@ async fn main() -> anyhow::Result<()> {
anyhow::bail!("subscription to {AUTH_SUBJECT} ended")
}
/// The OIDC path: who the `IdP` says presented `token`, and what the policy
/// grants that client. `None` permissions is a denial.
///
/// The caller is an identity or nothing — see `introspect`'s module docs.
/// There is no "admitted, identity unknown" branch to write here because
/// there is no such value to receive.
async fn introspected(
policy: &policy::Policy,
http: &reqwest::Client,
args: &Args,
client_secret: &str,
token: &str,
) -> (Option<String>, Option<policy::Permissions>) {
let caller = introspect::identify_caller(
http,
&args.introspection_url,
&args.client_id,
client_secret,
token,
)
.await
// An introspection that could not be *made* is a denial too. The
// failure modes of an HTTP call are exactly the conditions under
// which an attacker would most like this to fall open.
.unwrap_or_else(|e| {
tracing::warn!(error = ?e, "introspection failed; denying");
None
});
// Admission said who; the policy says what. A caller the `IdP`
// vouches for but no rule matches is denied — see `policy`'s module
// docs for why that is deny and not "connect with nothing".
let permissions = caller.as_deref().and_then(|id| policy.permissions(id));
if let (Some(id), None) = (caller.as_deref(), permissions.as_ref()) {
// Loud, and the one case an operator has to be able to find: a
// valid credential refused by our own policy. The alternative is
// a client that authenticates fine and mysteriously cannot work.
tracing::warn!(
caller = %id,
"authenticated client matches no policy rule; denying"
);
}
(caller, permissions)
}
#[cfg(test)]
mod tests {
use super::*;
@ -306,9 +355,16 @@ mod tests {
"hive-",
"--agent-client-suffix",
"-agent",
"--agent-token-publish-subject",
"$SWARM.term.{agent}",
"--agent-token-publish-subject",
"$SWARM.agent-state.{agent}",
"--store-cert-role",
"swarm-nats-auth",
])
.expect("the unit's own argument vector must parse");
assert_eq!(args.agent_client_suffix, "-agent");
assert_eq!(args.agent_token_publish_subjects.len(), 2);
}
/// The control for the case above: an ordinary value parses through the

View file

@ -45,12 +45,16 @@ pub struct Policy {
readers: Vec<String>,
extra_hive_subjects: Vec<String>,
extra_agent_subjects: Vec<String>,
agent_token_subjects: Vec<String>,
}
/// Placeholder replaced with the hive's own name in `extra_hive_subjects` and
/// `extra_agent_subjects`.
const HIVE_PLACEHOLDER: &str = "{hive}";
/// Placeholder replaced with the agent's own name in `agent_token_subjects`.
const AGENT_PLACEHOLDER: &str = "{agent}";
impl Policy {
/// `hive_prefix` is the client-id prefix that marks a hive and
/// `agent_suffix` what a hive's agent containers carry **on top of** it —
@ -130,9 +134,43 @@ impl Policy {
readers,
extra_hive_subjects,
extra_agent_subjects,
agent_token_subjects: Vec::new(),
})
}
/// Set the subjects an agent that proved its own credential may publish
/// to, with `{agent}` standing for its name.
///
/// # Errors
///
/// A template with no `{agent}` in it is refused: it would be one subject
/// shared by every agent in the swarm rather than the agent's own.
pub fn with_agent_token_subjects(mut self, subjects: Vec<String>) -> anyhow::Result<Self> {
if let Some(bad) = subjects.iter().find(|s| !s.contains(AGENT_PLACEHOLDER)) {
anyhow::bail!(
"--agent-token-publish-subject {bad:?} contains no {AGENT_PLACEHOLDER}: every \
agent in the swarm would be granted that exact subject"
);
}
self.agent_token_subjects = subjects;
Ok(self)
}
/// The permissions for an agent whose own credential was verified, or
/// `None` when none are configured.
///
/// Keyed on the agent alone: its identity is not tied to a hive. `agent`
/// must already be a single `[A-Za-z0-9_-]` segment, which the token parse
/// guarantees, so it cannot widen a subject with `.`, `*` or `>`.
pub fn agent_token_permissions(&self, agent: &str) -> Option<Permissions> {
let publish: Vec<String> = self
.agent_token_subjects
.iter()
.map(|s| s.replace(AGENT_PLACEHOLDER, agent))
.collect();
(!publish.is_empty()).then_some(Permissions { publish })
}
/// The permissions for `client_id`, or `None` when no rule matches.
///
/// `None` is a denial. It is not "grant nothing and let them connect":
@ -1175,4 +1213,35 @@ mod tests {
"an empty hive name must never be expanded into a subject: {g:?}"
);
}
#[test]
fn an_agent_token_grant_is_the_agents_own_subjects_and_nothing_else() {
let p = policy_with_agent_subject()
.with_agent_token_subjects(vec![
"$SWARM.term.{agent}".to_owned(),
"$SWARM.agent-state.{agent}".to_owned(),
])
.expect("per-agent templates are valid");
let g = p.agent_token_permissions("atlas").expect("configured");
assert_eq!(
g.publish,
vec![
"$SWARM.term.atlas".to_owned(),
"$SWARM.agent-state.atlas".to_owned(),
]
);
}
#[test]
fn an_agent_token_subject_without_the_placeholder_is_refused() {
let err = policy()
.with_agent_token_subjects(vec!["$SWARM.term.all".to_owned()])
.expect_err("a subject shared by every agent is not the agent's own");
assert!(format!("{err}").contains("$SWARM.term.all"), "{err}");
}
#[test]
fn with_no_agent_token_subject_configured_an_agent_token_is_refused() {
assert!(policy().agent_token_permissions("atlas").is_none());
}
}