Watch
0
0
Fork
You've already forked hyperhive
0

Compare commits

..
Author SHA1 Message Date
atlas
9cbee27e6d docs: reword vale-flagged prose in docs/swarm/README.md 2026-09-28 08:24:52 +02:00
atlas
727065c960 hive-agent: present this agent's own queue credential, then fall back
When the per-agent secret `queue-identity.nix` fetched is present, the
harness connects with `swarm-agent.<agent>.<secret>` as a static token
and publishes on `$SWARM.term.<agent>` and `$SWARM.agent-state.<agent>`.
When it is absent, or that first connect fails for any reason, a refusal
from a responder that does not verify agent tokens included, it connects
with the hive's shared OIDC client and publishes on the hive-scoped
subjects as before. Which one it took is logged once per connect.

`swarm_queue_client::connect_with_token` is the static-token connect: no
retry on the initial attempt, so the caller sees the refusal and can
fall back. Reconnects share the existing backoff, now a named function.

Closes #4630
2026-09-28 08:24:52 +02:00
atlas
0c1fb44a4f swarm-controller: relay an agent's hive-free subjects beside the old ones
The terminal and turn-state relays now subscribe to `$SWARM.term.<agent>`
and `$SWARM.agent-state.<agent>`, the subjects a verified agent token is
granted, as well as the hive-scoped `<prefix>.<hive>.<agent>` an agent
still on its hive's shared credential publishes to. Both at once, so a
swarm-ui terminal keeps working whichever credential an agent connected
with and in whatever order hosts deploy.
2026-09-28 08:24:52 +02:00
atlas
9bad58d86d 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.
2026-09-28 08:24:52 +02:00
atlas
fc97c237dc swarm-queue-client: one agent-token spelling, and no hive in AgentCredential
`swarm_queue_client::agent_token::format_agent_token` / `parse_agent_token`
are the spelling an agent presents its own queue secret in,
`swarm-agent.<agent>.<secret>`, and the one the auth-callout responder
reads back. The prefix is what separates it from an OIDC access token,
which may itself contain `.`. Parsing distinguishes "not an agent token"
(no prefix) from "a malformed one"; the error names the problem and never
the value. The module is store-free, so the agent formats its token
without linking the secret-store client.

`swarm_secret_client::queue::AgentCredential` loses `hive`: an agent's
identity is not tied to a hive, and nothing reads the field. Objects
already in the store carry it and still decode, since unknown fields are
ignored; a test parses one. The controller stops writing it.

With the credential no longer naming a hive, and the agent's policy
naming none since #4762, nothing in the mint consumes one. `hive` goes
from `mint_and_verify`, from the `MintAgentIdentity` node, and from
`POST /api/agents/{name}/identity`, which now takes no body and no longer
checks a hive against the roster; a caller that still sends one is not
refused, the body is ignored. `swarmctl agent mint-identity` loses
`--hive`, so passing it is now a usage error.
2026-09-28 08:24:52 +02:00
30 changed files with 1435 additions and 353 deletions

2
Cargo.lock generated
View file

@ -4828,7 +4828,9 @@ dependencies = [
"serde",
"serde_json",
"sha2 0.11.0",
"subtle",
"swarm-queue-client",
"swarm-secret-client",
"tokio",
"tracing",
"tracing-subscriber",

View file

@ -230,6 +230,9 @@ matrix-sdk = { version = "0.18", default-features = false, features = [
futures-util = "0.3"
hmac = "0.13"
sha2 = "0.11"
# Constant-time comparison of a presented secret against the stored one
# (`swarm-nats-auth::agent_token`). Already in the tree through the TLS stack.
subtle = "2.6"
# The NATS protocol client, for the swarm queue's auth-callout responder.
# `default-features = false` because the default set is broad - jetstream, kv,
# object-store, websockets, service - and a callout responder speaks none of

View file

@ -33,9 +33,8 @@ its identity then. Ruth doesn't: hive-c0re creates her on its own at
startup, so she needs her identity minted by hand, once.
```bash
# On the swarm-controller host: give ruth her store identity. <hive> is the
# name of the hive she runs on.
swarmctl agent mint-identity ruth --hive <hive>
# On the swarm-controller host: give ruth her store identity.
swarmctl agent mint-identity ruth
# On ruth's hive: re-apply her container config, which is when hive-c0re
# hands the new identity to the container.

View file

@ -410,10 +410,13 @@ agents sets none of the four and each agent logs that it has none; a half-set
environment logs an error and the harness keeps serving.
What an agent does with that connection is publish its terminal. Every row its
own web UI renders also goes to `$SWARM.term.<hive>.<agent>`, one subject per
agent, so a swarm-level terminal can follow one agent without subscribing to
the swarm's whole traffic. The `<hive>` is the one the agent's client id names,
which is the same string the broker builds its grant from. Publishing only: an
own web UI renders also goes to `$SWARM.term.<agent>`, one subject per agent, so
a swarm-level terminal can follow one agent without subscribing to the swarm's
whole traffic. An agent connected with its own queue credential gets that
subject. An agent without one, or whose own credential the queue
refused, connects with its hive's shared client and publishes to
`$SWARM.term.<hive>.<agent>` instead, the `<hive>` being the one that client id
names. The swarm controller relays both. Publishing only: an
agent talks about itself here and reads nothing. Rows aren't retained — a
subscriber that wasn't listening missed them, the same as on the agent's own
live stream.
@ -424,7 +427,7 @@ sending and leaves a marker in its place; the summary, level and icon still
arrive. The harness logs and skips a row that's too large even without its body.
The second thing an agent publishes is its **turn-state header**, on
`$SWARM.agent-state.<hive>.<agent>` — same shape of subject, same grant
`$SWARM.agent-state.<agent>` (or `$SWARM.agent-state.<hive>.<agent>`) — same shape of subject, same grant
mechanics, same lack of retention. It carries what a header bar wants: what the
turn loop is doing (`turn_state`, plus `turn_state_since` as an ISO 8601 UTC
stamp), which model (`model` and the resolved id the last turn actually ran on),

View file

@ -101,12 +101,10 @@ mint therefore never receives one, and nothing will ever come back around to
it. Re-run the mint for one agent with:
```sh
swarmctl agent mint-identity <agent> --hive <hive>
swarmctl agent mint-identity <agent>
```
`--hive` has no default: neither the CLI nor the controller keeps a roster of
which agent runs where, and the credentials this mints name a hive. The queue
secret half is idempotent — an agent that already has one keeps exactly the
The queue secret half is idempotent — an agent that already has one keeps exactly the
value it holds, so running this against an already-migrated agent doesn't drop
its queue connection. The certificate half isn't: the agent gets a fresh leaf
and picks it up on its next boot.

View file

@ -92,7 +92,7 @@ Queue a re-mint of an existing agent's identity at the swarm's secret store.
Queues and returns, the same way `agent create` does — watch the swarm UI's job view for the outcome.
**Usage:** `swarmctl agent mint-identity [OPTIONS] --hive <HIVE> <NAME>`
**Usage:** `swarmctl agent mint-identity [OPTIONS] <NAME>`
###### **Arguments:**
@ -100,9 +100,6 @@ Queues and returns, the same way `agent create` does — watch the swarm UI's jo
###### **Options:**
* `--hive <HIVE>` — The hive that agent runs on.
Required, and deliberately not defaulted: the credentials this mints name a hive, and neither this CLI nor the controller keeps a roster of which agent is on which hive. Naming the wrong one gives the agent an identity scoped to a hive it doesn't run on. The controller checks the value against the swarm's hive roster and names the known hives if it misses.
* `--controller-socket <PATH>` — swarm-controller's unix socket.
Supplied by the nix module that installs this binary, from the same `socketPath` option the daemon binds; falls back to `SWARM_CONTROLLER_SOCKET`.

View file

@ -35,6 +35,7 @@ use tokio::sync::broadcast;
use swarm_queue_client::wanted::AgentState;
use crate::events::{Bus, BusEvent, LiveEvent, TurnState};
use crate::swarm_queue::{Connection, Presented};
use crate::term_msg::iso8601_utc;
/// Subject family carrying agent turn-state headers, the swarm-wide
@ -171,33 +172,43 @@ fn snapshot(bus: &Bus) -> AgentStateMsg {
}
}
/// Start the publish task, if this agent has both queue coordinates and an
/// identity the subject can be derived from.
///
/// Returns without spawning in every other case — no queue, an unparseable
/// client id, no label — each of which is a legal state for an agent rather
/// than an error, and each logged once here rather than per transition.
pub fn spawn(bus: &Bus) {
let Some(cfg) = crate::swarm_queue::config() else {
// `swarm_queue::init` already said why at boot; repeating it here
// would be the same fact logged twice.
return;
};
let Some(hive) = hive_from_client_id(&cfg.client_id) else {
/// The subject this agent's header goes to under the credential it connected
/// with, or `None` when a hive client id names no hive.
fn subject(presented: &Presented, agent: &str) -> Option<String> {
match presented {
Presented::Agent => Some(format!("{SUBJECT_PREFIX}.{agent}")),
Presented::Hive { client_id } => {
let Some(hive) = hive_from_client_id(client_id) else {
tracing::warn!(
client_id = %cfg.client_id,
%client_id,
expected = format!("{CLIENT_ID_PREFIX}<hive>{CLIENT_ID_SUFFIX}"),
"queue client id does not name a hive; not publishing turn state upward"
);
return;
return None;
};
Some(format!("{SUBJECT_PREFIX}.{hive}.{agent}"))
}
}
}
/// Start the publish task, if this agent has a queue credential and a label
/// to name its subject with.
///
/// Returns without spawning in every other case — no queue or no label — each
/// of which is a legal state for an agent rather than an error, and each
/// logged once here rather than per transition.
pub fn spawn(bus: &Bus) {
if !crate::swarm_queue::configured() {
// `swarm_queue::init` already said why at boot; repeating it here
// would be the same fact logged twice.
return;
}
let agent = crate::identity::label();
if agent.is_empty() {
tracing::warn!("this agent has no label; not publishing turn state upward");
return;
}
let subject = format!("{SUBJECT_PREFIX}.{hive}.{agent}");
tokio::spawn(run(bus.subscribe(), bus.clone(), subject));
tokio::spawn(run(bus.subscribe(), bus.clone(), agent));
}
/// Watch the bus and publish whenever the header actually changed.
@ -218,8 +229,11 @@ pub fn spawn(bus: &Bus) {
/// and it carries nothing this header reads. Every other variant, including
/// any added later, funnels into the comparison and costs nothing when it
/// changes nothing.
async fn run(mut rx: broadcast::Receiver<BusEvent>, bus: Bus, subject: String) {
let Some(client) = crate::swarm_queue::client().await else {
async fn run(mut rx: broadcast::Receiver<BusEvent>, bus: Bus, agent: String) {
let Some(Connection { client, presented }) = crate::swarm_queue::client().await else {
return;
};
let Some(subject) = subject(&presented, &agent) else {
return;
};
tracing::info!(subject, "publishing agent turn state to the swarm queue");
@ -307,7 +321,7 @@ async fn publish(
#[cfg(test)]
mod tests {
use super::{AgentStateMsg, SUBJECT_PREFIX, hive_from_client_id};
use super::{AgentStateMsg, Presented, hive_from_client_id, subject};
use crate::events::TurnState;
use swarm_queue_client::wanted::AgentState;
@ -393,10 +407,22 @@ mod tests {
/// this pins the half that lives here.
#[test]
fn the_subject_is_the_prefix_then_the_hive_then_the_agent() {
let hive = hive_from_client_id("hive-alpha-agent").expect("names a hive");
let presented = Presented::Hive {
client_id: "hive-alpha-agent".to_owned(),
};
assert_eq!(
format!("{SUBJECT_PREFIX}.{hive}.mara"),
"$SWARM.agent-state.alpha.mara"
subject(&presented, "mara").as_deref(),
Some("$SWARM.agent-state.alpha.mara")
);
}
/// An agent that connected with its own credential publishes on the
/// subject that credential is granted, which names no hive.
#[test]
fn an_agent_on_its_own_credential_publishes_on_its_hive_free_subject() {
assert_eq!(
subject(&Presented::Agent, "mara").as_deref(),
Some("$SWARM.agent-state.mara")
);
}

View file

@ -21,15 +21,16 @@
//! agent. The fifth is this agent's own, minted per agent at swarm level and
//! fetched by the container itself (`nix/agent-modules/queue-identity.nix`).
//!
//! It is reported here but not yet *presented*: the queue's auth-callout
//! responder (`swarm-nats-auth`) validates only the hive-scoped token, and an
//! agent offering a credential nothing on the other end reads back would be
//! refused. Until that responder learns the same path, the connect path below
//! is unchanged and this is the fetching half.
//! When it is present, [`client`] connects with it first, and the agent
//! publishes on its own hive-free subjects. When it is absent, or the queue
//! refuses it, the connect falls back to the hive's shared client and the
//! hive-scoped subjects. [`Connection::presented`] says which, so a publisher
//! builds the subject that credential is granted.
use std::path::{Path, PathBuf};
use std::sync::OnceLock;
use anyhow::Context as _;
use tokio::sync::OnceCell;
use swarm_queue_client::QueueConfig;
@ -42,10 +43,39 @@ const ENV_PREFIX: &str = "HIVE_AGENT";
/// one answer rather than re-deriving it per call.
static CONFIG: OnceLock<Option<QueueConfig>> = OnceLock::new();
/// Where to present this agent's own credential, resolved once at boot.
static AGENT: OnceLock<Option<AgentPath>> = OnceLock::new();
/// The one connection every publisher in this process shares. Separate from
/// [`CONFIG`] because resolving the coordinates is synchronous boot work and
/// connecting is not — see [`client`].
static CLIENT: OnceCell<Option<async_nats::Client>> = OnceCell::const_new();
static CLIENT: OnceCell<Option<Connection>> = OnceCell::const_new();
/// What connecting with this agent's own credential needs.
#[derive(Debug, PartialEq, Eq)]
struct AgentPath {
url: String,
ca_file: Option<PathBuf>,
/// The fetched secret. A path; the bytes are read at connect.
secret_file: PathBuf,
}
/// Which credential the shared connection presented.
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum Presented {
/// This agent's own, granted `<prefix>.<agent>`.
Agent,
/// The hive's shared client, granted `<prefix>.<hive>.>` for the hive the
/// client id names.
Hive { client_id: String },
}
/// The shared queue connection and the credential it was made with.
#[derive(Clone)]
pub struct Connection {
pub client: async_nats::Client,
pub presented: Presented,
}
/// The four variables the harness unit sets, before the client-id file is
/// read. Collected into a struct so [`decide`] is pure over them and the
@ -147,6 +177,16 @@ fn decide_agent_secret(path: Option<&str>) -> Option<PathBuf> {
}
}
/// Decide whether this agent can connect with its own credential: it needs the
/// queue's address and a fetched secret, and nothing of the hive's client.
fn decide_agent_path(env: &QueueEnv, secret_file: Option<PathBuf>) -> Option<AgentPath> {
Some(AgentPath {
url: env.nats_url.clone()?,
ca_file: env.ca_file.as_ref().map(Into::into),
secret_file: secret_file?,
})
}
/// Decide what this agent's queue configuration is, given the environment and
/// whatever the client-id file held.
///
@ -216,15 +256,10 @@ pub fn init() {
let _ = CONFIG.set(resolved);
// Independent of everything above: this agent may hold its own secret on
// a hive with no queue coordinates, or hold the coordinates and no secret
// of its own yet. Reported either way, because "which credential is this
// agent able to present" is a question only this process can answer, and
// it is the one the next slice's rollout will be asked repeatedly.
//
// The answer is only logged here. Presenting it needs the queue's
// auth-callout responder to verify it, which is the next slice — see this
// module's header.
if let Some(path) = decide_agent_secret(env.agent_secret_file.as_deref()) {
// a hive whose shared client is not published yet, or the shared client
// and no secret of its own.
let secret = decide_agent_secret(env.agent_secret_file.as_deref());
if let Some(path) = &secret {
// The path, never the bytes: the file holds the secret itself.
tracing::info!(
path = %path.display(),
@ -236,6 +271,13 @@ pub fn init() {
by its hive's shared client"
);
}
let _ = AGENT.set(decide_agent_path(&env, secret));
}
/// Whether this agent has any credential to reach the queue with. `false`
/// before [`init`] has run.
pub fn configured() -> bool {
config().is_some() || AGENT.get().is_some_and(Option::is_some)
}
/// What [`init`] resolved, or `None` when this agent has no queue.
@ -263,14 +305,48 @@ pub fn config() -> Option<&'static QueueConfig> {
/// failed" — a caller does nothing differently between them, since either way
/// there is nothing to publish onto. Connecting is lazy so that an agent on a
/// hive with no queue pays nothing at boot.
pub async fn client() -> Option<async_nats::Client> {
pub async fn client() -> Option<Connection> {
CLIENT.get_or_init(connect_once).await.clone()
}
async fn connect_once() -> Option<async_nats::Client> {
/// This agent's own credential first, then the hive's. Each path taken is
/// logged once, here.
async fn connect_once() -> Option<Connection> {
if let Some(agent) = AGENT.get().and_then(Option::as_ref) {
match connect_as_agent(agent).await {
Ok(client) => {
tracing::info!(
url = %agent.url,
"connected to the swarm queue with this agent's own credential"
);
return Some(Connection {
client,
presented: Presented::Agent,
});
}
Err(e) => tracing::warn!(
error = format!("{e:#}"),
"connecting with this agent's own credential failed; falling back to \
its hive's shared client"
),
}
}
let cfg = config()?;
match swarm_queue_client::connect(cfg.clone()).await {
Ok(client) => Some(client),
Ok(client) => {
tracing::info!(
url = %cfg.url,
client_id = %cfg.client_id,
"connecting to the swarm queue with the hive's shared client; the \
client retries in the background until the queue accepts it"
);
Some(Connection {
client,
presented: Presented::Hive {
client_id: cfg.client_id.clone(),
},
})
}
Err(e) => {
// `chain`, not `{:#}`: this is `swarm_queue_client::Error`, whose
// `Display` ignores the alternate flag, so `{:#}` renders the
@ -284,9 +360,26 @@ async fn connect_once() -> Option<async_nats::Client> {
}
}
/// Connect presenting this agent's own token. Any failure, a refusal
/// included, is returned for [`connect_once`] to fall back on.
async fn connect_as_agent(agent: &AgentPath) -> anyhow::Result<async_nats::Client> {
let name = crate::identity::label();
let secret = std::fs::read_to_string(&agent.secret_file)
.with_context(|| format!("reading {}", agent.secret_file.display()))?;
let token = swarm_queue_client::agent_token::format_agent_token(&name, secret.trim())?;
swarm_queue_client::connect_with_token(&agent.url, agent.ca_file.as_deref(), token)
.await
.map_err(|e| anyhow::anyhow!(swarm_queue_client::chain(&e)))
}
#[cfg(test)]
mod tests {
use super::{QueueEnv, Resolution, decide, decide_agent_secret, read_client_id};
use std::path::PathBuf;
use super::{
AgentPath, QueueEnv, Resolution, decide, decide_agent_path, decide_agent_secret,
read_client_id,
};
fn env(parts: [Option<&str>; 4]) -> QueueEnv {
let [nats_url, token_endpoint, client_id_file, client_secret_file] = parts;
@ -435,4 +528,34 @@ mod tests {
assert!(matches!(decide(&e, None), Resolution::Absent(_)));
assert!(decide_agent_secret(path.to_str()).is_some());
}
/// The agent's own path needs the queue's address and its own secret, and
/// not the hive's client id: that is what lets it connect on a hive whose
/// shared client is not published.
#[test]
fn an_agent_with_its_own_secret_and_the_queue_address_connects_as_itself() {
let secret = PathBuf::from("/run/queue-identity/secret");
assert_eq!(
decide_agent_path(&full(), Some(secret.clone())),
Some(AgentPath {
url: "nats://10.42.0.1:4222".to_owned(),
ca_file: None,
secret_file: secret.clone(),
})
);
let url_only = env([Some("nats://10.42.0.1:4222"), None, None, None]);
assert!(decide_agent_path(&url_only, Some(secret)).is_some());
}
#[test]
fn without_its_own_secret_or_the_queue_address_an_agent_does_not_connect_as_itself() {
assert_eq!(decide_agent_path(&full(), None), None);
assert_eq!(
decide_agent_path(
&env([None, None, None, None]),
Some(PathBuf::from("/run/queue-identity/secret"))
),
None
);
}
}

View file

@ -27,6 +27,7 @@
use tokio::sync::broadcast;
use crate::events::BusEvent;
use crate::swarm_queue::{Connection, Presented};
use crate::term_msg::{ClassifyCtx, TermMsg, classify};
/// Subject family carrying agent terminal rows, the swarm-wide agreement this
@ -112,37 +113,50 @@ fn serialized_len(msg: &TermMsg) -> Option<usize> {
serde_json::to_vec(msg).ok().map(|v| v.len())
}
/// Start the publish task, if this agent has both queue coordinates and an
/// identity the subject can be derived from.
///
/// Returns without spawning in every other case — no queue, an unparseable
/// client id, no label — each of which is a legal state for an agent rather
/// than an error, and each logged once here rather than per row.
pub fn spawn(rx: broadcast::Receiver<BusEvent>) {
let Some(cfg) = crate::swarm_queue::config() else {
// `swarm_queue::init` already said why at boot; repeating it here
// would be the same fact logged twice.
return;
};
let Some(hive) = hive_from_client_id(&cfg.client_id) else {
/// The subject this agent's rows go to under the credential it connected
/// with, or `None` when a hive client id names no hive.
fn subject(presented: &Presented, agent: &str) -> Option<String> {
match presented {
Presented::Agent => Some(format!("{SUBJECT_PREFIX}.{agent}")),
Presented::Hive { client_id } => {
let Some(hive) = hive_from_client_id(client_id) else {
tracing::warn!(
client_id = %cfg.client_id,
%client_id,
expected = format!("{CLIENT_ID_PREFIX}<hive>{CLIENT_ID_SUFFIX}"),
"queue client id does not name a hive; not publishing the terminal upward"
);
return;
return None;
};
Some(format!("{SUBJECT_PREFIX}.{hive}.{agent}"))
}
}
}
/// Start the publish task, if this agent has a queue credential and a label
/// to name its subject with.
///
/// Returns without spawning in every other case — no queue or no label — each
/// of which is a legal state for an agent rather than an error, and each
/// logged once here rather than per row.
pub fn spawn(rx: broadcast::Receiver<BusEvent>) {
if !crate::swarm_queue::configured() {
// `swarm_queue::init` already said why at boot; repeating it here
// would be the same fact logged twice.
return;
}
let agent = crate::identity::label();
if agent.is_empty() {
tracing::warn!("this agent has no label; not publishing the terminal upward");
return;
}
let subject = format!("{SUBJECT_PREFIX}.{hive}.{agent}");
tokio::spawn(run(rx, subject));
tokio::spawn(run(rx, agent));
}
async fn run(mut rx: broadcast::Receiver<BusEvent>, subject: String) {
let Some(client) = crate::swarm_queue::client().await else {
async fn run(mut rx: broadcast::Receiver<BusEvent>, agent: String) {
let Some(Connection { client, presented }) = crate::swarm_queue::client().await else {
return;
};
let Some(subject) = subject(&presented, &agent) else {
return;
};
tracing::info!(subject, "publishing the agent terminal to the swarm queue");
@ -203,7 +217,7 @@ async fn publish(client: &async_nats::Client, subject: &str, msg: TermMsg) {
#[cfg(test)]
mod tests {
use super::{DROPPED_BODY, fit, hive_from_client_id};
use super::{DROPPED_BODY, Presented, fit, hive_from_client_id, subject};
use crate::events::LiveEvent;
use crate::term_msg::{BodyFormat, ClassifyCtx, Level, TermMsg, classify};
@ -243,6 +257,27 @@ mod tests {
}
}
/// Each credential publishes on the subject it is granted: its own under
/// the agent's credential, its hive's under the shared one.
#[test]
fn the_subject_follows_the_credential_the_agent_connected_with() {
assert_eq!(
subject(&Presented::Agent, "mara").as_deref(),
Some("$SWARM.term.mara")
);
let hive = Presented::Hive {
client_id: "hive-alpha-agent".to_owned(),
};
assert_eq!(
subject(&hive, "mara").as_deref(),
Some("$SWARM.term.alpha.mara")
);
let unparseable = Presented::Hive {
client_id: "hive-alpha".to_owned(),
};
assert_eq!(subject(&unparseable, "mara"), None);
}
#[test]
fn a_row_that_already_fits_is_published_unchanged() {
let msg = TermMsg::new(Level::Info, "turn ok")

View file

@ -32,6 +32,7 @@
./glue-grafana-oidc-client.nix
./glue-matrix-bao-token.nix
./glue-matrix-ctl-bao-identity.nix
./glue-nats-auth-bao-identity.nix
./glue-nats-bao-identity.nix
./glue-queue-agent-credential.nix
./glue-secret-publisher-bao-identity.nix

View file

@ -260,6 +260,12 @@ in
# leaf logs in with. Stays on this host; see its default above.
[ -s ${pkiDir}/granter.pem ] || ${signLeaf} ${pkiDir} granter \
${lib.escapeShellArg deployCfg.bao.granterCommonName} "" clientAuth
# The queue's auth-callout responder, which reads agent queue
# credentials. Minted here for the matrix-ctl leaf's reason: the queue
# is a swarm singleton, so elsewhere this is the file an operator copies.
[ -s ${pkiDir}/nats-auth.pem ] || ${signLeaf} ${pkiDir} nats-auth \
${lib.escapeShellArg deployCfg.bao.natsAuthCommonName} "" clientAuth
'';
};
};

View file

@ -0,0 +1,37 @@
# Glue: point the queue's auth-callout responder at the bao leaf minted for it.
#
# ONE PAIRING PER FILE — the swarm-nats-auth principal ← bao, and nothing else.
# Deleting this leaves a responder with no store identity unless the operator
# names one: every agent token is then denied, and OIDC clients are unaffected.
#
# ⚠️ The minting is NOT here. ./glue-bao-tls.nix holds the CA and signs the
# leaf. What belongs here is the pairing: which paths the responder presents.
#
# ⚠️ Gated on the leaf existing, not on the store being enabled, for the reason
# ./glue-nats-bao-identity.nix states: the hive hosting the queue need not be
# the hive hosting the store.
#
# Everything is `mkDefault`. An operator naming their own paths wins.
{
lib,
config,
...
}:
let
hyperhiveCfg = config.services.hyperhive;
deployCfg = hyperhiveCfg.deploy;
baoDeploy = deployCfg.bao;
# Where ./glue-bao-tls.nix puts the leaves, derived from the reader's own path
# rather than repeating that file's directory literal.
haveMintedPki = baoDeploy.clientCertFile != null;
pkiDir = if haveMintedPki then builtins.dirOf baoDeploy.clientCertFile else null;
in
{
config = lib.mkIf (hyperhiveCfg.enable && deployCfg.nats.enable && haveMintedPki) {
services.hyperhive.deploy.nats = {
authBaoClientCertFile = lib.mkDefault "${pkiDir}/nats-auth.pem";
authBaoClientKeyFile = lib.mkDefault "${pkiDir}/nats-auth-key.pem";
};
};
}

View file

@ -754,6 +754,18 @@ let
}
];
# The queue's auth-callout responder, which checks the secret an agent
# presents against the one stored for it. One leaf under every agent: `+` is
# one path segment, where `*` globs only at the end and would reach every
# other credential an agent holds.
natsAuthReaders = [
{
name = "swarm-nats-auth";
cn = baoDeploy.natsAuthCommonName;
policyText = readStanza "${credentialMountPath}/data/swarm/agents/+/queue";
}
];
# The role name IS the policy name, as for the three service principals
# above: the role attaches the policy by spelling it identically, and one
# string for both objects removes the way they drift apart.
@ -1438,6 +1450,24 @@ in
'';
};
natsAuthCommonName = lib.mkOption {
type = lib.types.str;
default = "swarm-nats-auth";
description = ''
Subject the store's `swarm-nats-auth` cert-auth role accepts: the
identity the queue's auth-callout responder presents to read agent
queue credentials, `swarm/agents/<agent>/queue` and nothing else under
an agent.
Its own principal rather than
{option}`services.hyperhive.deploy.bao.natsCommonName`: that one
issues the queue's TLS leaf from a host unit, this one reads secrets
from inside the queue's container.
⚠️ Reserved as a hive name by ./swarm.nix, like its siblings.
'';
};
matrixCtlHiveName = lib.mkOption {
type = lib.types.str;
default = toString hyperhiveCfg.hiveName;
@ -2021,6 +2051,7 @@ in
"swarm-bao-services-issuer-policy"
"swarm-bao-nats-tls-policy"
"swarm-bao-agent-pki"
"swarm-bao-nats-auth-policy"
];
# 🚫 No `swarm.otel.scrapeTargets.bao` entry any more, and its absence is
@ -2723,6 +2754,8 @@ in
systemd.services.swarm-bao-forwarder-oidc-policy = readerPolicyUnit "write the store forwarder's OIDC-secret-reader bao policy and cert-auth role" forwarderOidcReaders;
systemd.services.swarm-bao-nats-auth-policy = readerPolicyUnit "write the queue responder's agent-credential-reader bao policy and cert-auth role" natsAuthReaders;
# A FOURTH sibling, same shape and same reasons as the two above. This
# one is what turns `swarm-services-issuer` from a declaration into a
# grant: a bao policy reaches nothing until a login role hands it to a

View file

@ -219,6 +219,25 @@ let
tlsKeyCredential = "tls-key";
tlsKeyCredentialPath = "/run/credentials/nats.service/${tlsKeyCredential}";
# The responder's own store identity, for reading agent queue credentials.
# Without it the responder denies every agent token; agents then fall back to
# their hive's OIDC client.
authStoreActive =
deployCfg.nats.authBaoClientCertFile != null && deployCfg.nats.authBaoClientKeyFile != null;
# The role ./swarm-bao.nix writes on the store's `cert` mount, by the same name.
authCertRole = "swarm-nats-auth";
# Store identity files the copy unit delivers, as credential id → host source.
authStoreFiles = lib.optionalAttrs authStoreActive (
{
"bao-client.pem" = deployCfg.nats.authBaoClientCertFile;
"bao-client-key.pem" = deployCfg.nats.authBaoClientKeyFile;
}
// lib.optionalAttrs (baoDeploy.serverCaFile != null) {
"bao-ca.pem" = baoDeploy.serverCaFile;
}
);
authCredential = id: "/run/credentials/swarm-nats-auth.service/${id}";
# Re-issue once the leaf is past half of the 720h `pki/roles/swarm-nats`
# grants it (./swarm-bao.nix), so the daily timer below has two weeks of
# retries before it lapses.
@ -498,6 +517,34 @@ in
A path, never a value.
'';
};
authBaoClientCertFile = lib.mkOption {
type = lib.types.nullOr lib.types.str;
default = null;
example = "/var/lib/swarm-bao-pki/nats-auth.pem";
description = ''
Client certificate the auth-callout responder presents to the swarm's
secret store to read agent queue credentials. Its subject must be
{option}`services.hyperhive.deploy.bao.natsAuthCommonName`.
No default. ./glue-nats-auth-bao-identity.nix points it at the leaf
./glue-bao-tls.nix mints, where this host mints one. Unset, or set to
a file that does not exist, the responder denies every agent token and
admits OIDC clients as before.
A path, never a value.
'';
};
authBaoClientKeyFile = lib.mkOption {
type = lib.types.nullOr lib.types.str;
default = null;
example = "/var/lib/swarm-bao-pki/nats-auth-key.pem";
description = ''
Private key for {option}`services.hyperhive.deploy.nats.authBaoClientCertFile`.
A path, never a value.
'';
};
};
config = lib.mkIf deployCfg.nats.enable {
@ -905,6 +952,11 @@ in
# an append-only row stream, a header is one current value
# republished on change.
"--agent-publish-subject ${lib.escapeShellArg "\$\$SWARM.agent-state.{hive}.>"}"
# What an agent that proved its own credential may publish to:
# the same two streams, keyed on the agent alone.
"--agent-token-publish-subject ${lib.escapeShellArg "\$\$SWARM.term.{agent}"}"
"--agent-token-publish-subject ${lib.escapeShellArg "\$\$SWARM.agent-state.{agent}"}"
"--store-cert-role ${lib.escapeShellArg authCertRole}"
];
# Every credential arrives by `LoadCredential` and is named
# on the command line only as a **path** — `argv` is
@ -914,12 +966,25 @@ in
"callout-user.seed:${inContainer "callout-user.seed"}"
"issuer.seed:${inContainer "issuer.seed"}"
"oidc-client.secret:${inContainer "oidc-client.secret"}"
];
]
++ lib.mapAttrsToList (id: _: "${id}:${inContainer id}") authStoreFiles;
DynamicUser = true;
Restart = "on-failure";
RestartSec = "5s";
SyslogIdentifier = "swarm-nats-auth";
};
# The store the responder reads agent credentials from. Absent
# without an identity, which the responder reads as "no store".
environment = lib.optionalAttrs authStoreActive (
{
BAO_ADDR = "https://${baoCfg.domain}:${toString baoCfg.port}";
BAO_CLIENT_CERT = authCredential "bao-client.pem";
BAO_CLIENT_KEY = authCredential "bao-client-key.pem";
}
// lib.optionalAttrs (authStoreFiles ? "bao-ca.pem") {
BAO_CACERT = authCredential "bao-ca.pem";
}
);
};
# The server binary, so an operator with a shell in here can
@ -957,7 +1022,10 @@ in
# being up says nothing about whether its in-container secrets unit
# has finished. The wait in the script is what actually closes it;
# this only stops us spinning for the full timeout on every boot.
++ lib.optional deployCfg.authelia.enable "container@${autheliaCfg.machine}.service";
++ lib.optional deployCfg.authelia.enable "container@${autheliaCfg.machine}.service"
# Where the store's PKI is minted on this host, the responder's leaf is
# one of its files. Ordering only: elsewhere the unit does not exist.
++ lib.optional authStoreActive "swarm-bao-pki.service";
requires = lib.optional deployCfg.nats.autoGenerateCallout "swarm-nats-callout-keys.service";
serviceConfig = {
Type = "oneshot";
@ -1006,7 +1074,22 @@ in
install -m 0400 "$secret" \
${lib.escapeShellArg (hostPath "oidc-client.secret")}
'';
''
# The store identity is optional to the responder, so an absent source
# is delivered as an empty file rather than failing this unit: a missing
# `LoadCredential` source stops the responder starting, and a responder
# that does not start denies every client. An empty identity fails only
# the store login, which denies agent tokens.
+ lib.concatStrings (
lib.mapAttrsToList (id: source: ''
if [ -s ${lib.escapeShellArg source} ]; then
install -m 0400 ${lib.escapeShellArg source} ${lib.escapeShellArg (hostPath id)}
else
echo "${source} is absent; the queue responder denies agent tokens until it exists" >&2
install -m 0400 /dev/null ${lib.escapeShellArg (hostPath id)}
fi
'') authStoreFiles
);
};
# ⚠️ Minted on the HOST, not in the container, because the responder is

View file

@ -53,6 +53,7 @@ let
deployCfg.bao.servicesIssuerCommonName
deployCfg.bao.natsCommonName
deployCfg.bao.granterCommonName
deployCfg.bao.natsAuthCommonName
]
# The two per-hive readers' subjects, spelled out per hive rather than as the
# prefix. The prefix alone would reserve the wrong string: the role for hive

View file

@ -151,7 +151,7 @@ let
_: u: (u.environment.BAO_CLIENT_CERT or null) == granterCertFile
) baoGrantWithConsumers.systemd.services;
# The eleven units that write a `swarm-*` grant, by name, for the discovery
# The twelve units that write a `swarm-*` grant, by name, for the discovery
# control below.
grantingUnitNames = [
"swarm-bao-controller-policy"
@ -165,6 +165,7 @@ let
"swarm-bao-services-issuer-policy"
"swarm-bao-nats-tls-policy"
"swarm-bao-agent-pki"
"swarm-bao-nats-auth-policy"
];
# Comment lines dropped first: both the HCL and the scripts explain
@ -561,6 +562,33 @@ let
&& !(lib.hasInfix "swarm-grafana" s)
&& !(lib.hasInfix "sys/policies/acl" s);
}
{
# `+` is one path segment, so this reaches `swarm/agents/<agent>/queue`
# and no other credential an agent holds; `swarm/agents/*` would reach
# all of them.
name = "the queue responder's grant is every agent's queue credential and nothing else";
ok =
let
s = baoGrantHere.systemd.services.swarm-bao-nats-auth-policy.script;
in
lib.hasInfix "path \"secret/data/swarm/agents/+/queue\" {" s
&& lib.hasInfix "capabilities = [\"read\"]" s
&& lib.length (lib.filter lib.isList (builtins.split "path \"" s)) == 1
&& !(lib.hasInfix "secret/data/swarm/agents/*" s)
&& !(lib.hasInfix "secret/data/swarm/hives" s)
&& !(lib.hasInfix "secret/data/swarm/services" s)
&& lib.hasInfix "auth/cert/certs/swarm-nats-auth" s
&& lib.hasInfix "allowed_common_names=swarm-nats-auth" s
&& lib.hasInfix "token_policies=swarm-nats-auth" s;
}
{
name = "the PKI unit signs the queue responder's leaf under its own subject";
ok =
let
s = baoGrantHere.systemd.services.swarm-bao-pki.script;
in
lib.hasInfix "/nats-auth.pem ]" s && lib.hasInfix "swarm-nats-auth \"\" clientAuth" s;
}
{
# 🩸 The half that makes the policies above bind: a policy grants only
# through a token that carries it, and a token is minted by a cert-auth
@ -748,24 +776,24 @@ let
}
{
# A store host without the granter's pair writes its grants some other
# way, so none of the eleven units may exist. Without this arm
# way, so none of the twelve units may exist. Without this arm
# `lib.mkIf haveGranter` could be dropped from any of them and every other
# case here would still pass.
name = "without the granter's pair none of the eleven granting units render";
name = "without the granter's pair none of the twelve granting units render";
ok =
let
s = baoGranterOptOut.systemd.services;
in
lib.all (unit: !(s ? ${unit})) (grantingUnitNames ++ [ "swarm-bao-granter-role" ])
# The control: the same store with the pair renders all eleven.
# The control: the same store with the pair renders all twelve.
&& lib.all (unit: baoGrantHere.systemd.services ? ${unit}) grantingUnitNames;
}
{
# 🩸 What replaced the silent skip. With no bootstrap token the eleven still
# 🩸 What replaced the silent skip. With no bootstrap token the twelve still
# render, and a refused granter fails them with the step that fixes it.
# A store host that never named a token is told to name one, since the
# unit that sets the granter up renders only where it has.
name = "a store host without a bootstrap token renders the eleven, each failing loudly with the one-time step";
name = "a store host without a bootstrap token renders the twelve, each failing loudly with the one-time step";
ok =
let
s = baoGranterNoToken.systemd.services;
@ -1129,8 +1157,8 @@ let
}
{
# What makes the case above mean something: discovery by the granter's
# certificate reaches all eleven units, and each yields calls.
name = "the granter-policy check sees all eleven granting units, and parses calls from each";
# certificate reaches all twelve units, and each yields calls.
name = "the granter-policy check sees all twelve granting units, and parses calls from each";
ok =
lib.sort lib.lessThan (lib.attrNames granterUnits) == lib.sort lib.lessThan grantingUnitNames
&& lib.all (u: baoCalls u.script != [ ]) (lib.attrValues granterUnits)

View file

@ -60,6 +60,28 @@ let
swarm.nats.calloutIssuerSeedFile = "/run/secrets/nats-issuer.seed";
};
# The queue with a store identity for its responder, placed by hand: the
# queue host that is not the store's.
natsWithStoreIdentity = hive {
deploy.nats.enable = true;
deploy.nats.autoGenerateCallout = false;
deploy.nats.calloutUserPublicKey = "UTESTUSERPUBKEYAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAA";
deploy.nats.calloutIssuerPublicKey = "ATESTISSUERPUBKEYAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAA";
deploy.nats.calloutUserSeedFile = "/run/secrets/nats-user.seed";
deploy.nats.calloutIssuerSeedFile = "/run/secrets/nats-issuer.seed";
deploy.nats.authPackage = pkgs.emptyDirectory;
deploy.nats.authBaoClientCertFile = "/etc/pki/nats-auth.pem";
deploy.nats.authBaoClientKeyFile = "/etc/pki/nats-auth-key.pem";
};
# The queue on the store's own host, where the PKI glue mints every leaf.
natsOnStoreHost = hive {
deploy.bao.enable = true;
deploy.nats.enable = true;
};
responderOf = h: h.containers.swarm-nats.config.systemd.services.swarm-nats-auth;
# A hive running NOTHING of the swarm's own services — no IdP here, no
# `swarm.authelia.url` set by hand. The whole point of the fixture is what it
# does *not* say: it is the shape whose IdP address used to be null, and
@ -117,6 +139,73 @@ let
in
lib.hasInfix "--agent-publish-subject '$$SWARM.agent-state.{hive}.>'" exec;
}
{
# The per-agent grant: the same two streams keyed on the agent alone,
# with the same doubled dollar, and the role the store writes for it.
name = "the responder grants a verified agent its own hive-free subjects";
ok =
let
exec = (responderOf natsOldPath).serviceConfig.ExecStart;
in
lib.hasInfix "--agent-token-publish-subject '$$SWARM.term.{agent}'" exec
&& lib.hasInfix "--agent-token-publish-subject '$$SWARM.agent-state.{agent}'" exec
&& lib.hasInfix "--store-cert-role swarm-nats-auth" exec;
}
{
# Each credential by `LoadCredential`, and the environment naming where
# the unit sees it: a DynamicUser cannot read the copies directly.
name = "a responder with a store identity loads it and is pointed at the store";
ok =
let
u = responderOf natsWithStoreIdentity;
creds = u.serviceConfig.LoadCredential;
in
lib.elem "bao-client.pem:/var/lib/swarm-nats-auth/bao-client.pem" creds
&& lib.elem "bao-client-key.pem:/var/lib/swarm-nats-auth/bao-client-key.pem" creds
&& u.environment.BAO_CLIENT_CERT == "/run/credentials/swarm-nats-auth.service/bao-client.pem"
&& u.environment.BAO_CLIENT_KEY == "/run/credentials/swarm-nats-auth.service/bao-client-key.pem"
&& lib.hasPrefix "https://" u.environment.BAO_ADDR;
}
{
# An absent leaf must not stop the responder, which would deny every
# client: it is delivered empty, and the three credentials the responder
# cannot run without are delivered too.
name = "the copy unit delivers the store identity, and an absent one as an empty file";
ok =
let
s = natsWithStoreIdentity.systemd.services.swarm-nats-auth-secrets.script;
in
lib.hasInfix "install -m 0400 /etc/pki/nats-auth.pem " s
&& lib.hasInfix "install -m 0400 /etc/pki/nats-auth-key.pem " s
&& lib.hasInfix "install -m 0400 /dev/null " s
&& lib.hasInfix "oidc-client.secret" s;
}
{
# The control: no identity, no store wiring, and the responder's
# credentials are exactly the three it cannot run without.
name = "a responder with no store identity is not pointed at a store";
ok =
let
u = responderOf natsOldPath;
in
!(u.environment ? BAO_ADDR)
&& !(u.environment ? BAO_CLIENT_CERT)
&& lib.length u.serviceConfig.LoadCredential == 3;
}
{
# On the store's host the glue pairs the responder with a leaf of its
# own, never the queue's TLS-issuing one or the hive's.
name = "on the store's host the responder presents its own leaf";
ok =
let
n = natsOnStoreHost.services.hyperhive.deploy.nats;
b = natsOnStoreHost.services.hyperhive.deploy.bao;
in
n.authBaoClientCertFile == "/var/lib/swarm-bao-pki/nats-auth.pem"
&& n.authBaoClientKeyFile == "/var/lib/swarm-bao-pki/nats-auth-key.pem"
&& n.authBaoClientCertFile != n.baoClientCertFile
&& n.authBaoClientCertFile != b.clientCertFile;
}
{
# Not a rename test. `hostClientSecretDir` is `readOnly`, so the fixture
# cannot define it; what can break is a reader left pointing at the

View file

@ -136,7 +136,7 @@ fn generate_queue_secret() -> Result<String> {
/// Anything that stops one of those five steps, with the step named. A
/// failure here fails the job node and nothing else — the agent is still
/// created, without a store identity.
pub async fn mint_and_verify(agent: &str, hive: &str) -> Result<()> {
pub async fn mint_and_verify(agent: &str) -> Result<()> {
let (mount, pki_role) = agent_pki(|k| std::env::var(k).ok())?;
let name = policy::agent_object_name(agent)?;
let path = mtls::identity_path(agent)?;
@ -161,18 +161,15 @@ pub async fn mint_and_verify(agent: &str, hive: &str) -> Result<()> {
.read_optional(&queue_path)
.await
.with_context(|| format!("checking whether {queue_path} already holds a credential"))?;
// The secret survives a re-run; the principal it names does not get to.
// An object whose `hive` disagrees with the hive this node was invoked
// with would grant its holder subjects on the wrong hive, so it is
// corrected — but by rewriting the two name fields around the *same*
// `value`, which is a correction no live connection notices.
// 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()?,
},
agent: agent.to_owned(),
hive: hive.to_owned(),
};
if existing.as_ref() == Some(&wanted) {
tracing::info!(
@ -187,9 +184,8 @@ pub async fn mint_and_verify(agent: &str, hive: &str) -> Result<()> {
.with_context(|| format!("publishing the agent queue credential at {queue_path}"))?;
tracing::info!(
agent,
hive,
%queue_path,
// Never "rotated": the secret is the same one, only the names
// Never "rotated": the secret is the same one, only the name
// around it moved.
corrected = existing.is_some(),
"agent queue credential published"
@ -280,9 +276,9 @@ async fn read_back_as_agent(
.read(queue_path)
.await
.with_context(|| format!("reading {queue_path} back under {role}'s own token"))?;
// The two name fields are compared as well as the secret: they are what
// the verifying end will grant subjects from, so a mismatch here is the
// same class of fault as an unreadable path.
// The agent is compared as well as the secret: the verifying end refuses
// an object naming a different agent than the path, so a mismatch here is
// the same class of fault as an unreadable path.
if read_back != *queue_credential {
bail!("the store returned a different object at {queue_path} than the one just published");
}

View file

@ -33,8 +33,8 @@ use crate::AppState;
/// Subject family carrying agent turn-state headers — must match
/// `hive-agent::swarm_agent_state::SUBJECT_PREFIX` exactly, since the two
/// ends never see the constant together. `$SWARM.agent-state.<hive>.<agent>`
/// is the full subject; as with the terminal relay's own copy, there is no
/// ends never see the constant together. `crate::term_stream::agent_subjects`
/// spells the full subjects; as with the terminal relay's own copy, there is no
/// way to check the two sides agree short of this comment and the module
/// docs on both ends staying honest about it.
const SUBJECT_PREFIX: &str = "$SWARM.agent-state";
@ -117,16 +117,12 @@ pub(crate) async fn stream_agent_state(
)
})?;
let subject = format!("{SUBJECT_PREFIX}.{hive}.{agent}");
let subscriber = client.subscribe(subject.clone()).await.map_err(|e| {
tracing::warn!(%subject, error = %e, "state stream: subscribe failed");
crate::error_problem(
axum::http::StatusCode::INTERNAL_SERVER_ERROR,
&format!("subscribing to {subject} failed: {e}"),
let subscriber = crate::term_stream::subscribe_all(
&client,
&crate::term_stream::agent_subjects(SUBJECT_PREFIX, &hive, &agent),
"state",
)
})?;
tracing::info!(%subject, "state stream: client attached");
.await?;
let stream =
subscriber.map(|msg| Ok(Event::default().data(String::from_utf8_lossy(&msg.payload))));
Ok(Sse::new(stream).keep_alive(KeepAlive::default()))

View file

@ -100,11 +100,10 @@ enum SwarmNodeKind {
/// `agent_identity::mint_and_verify` — including why this node does not
/// report success on a write.
///
/// Carries the hive for a different reason than `TriggerDeploy` does:
/// not as an address, but because the agent's queue credential names the
/// hive it may take subjects on, so it cannot be written without knowing
/// which hive the agent belongs to.
MintAgentIdentity { hive: String, agent: String },
/// Carries no hive: neither the certificate, the queue credential nor
/// the agent's policy names one, so an agent keeps one store identity
/// whichever hive it runs on.
MintAgentIdentity { agent: String },
/// Make sure `agent` holds a live forge access token in the swarm secret
/// store, minting one with the forge's admin API when it does not. See
/// `forge::agent_token` — including why a rotation is a delete then a
@ -124,8 +123,7 @@ enum SwarmNodeKind {
/// 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`.
/// Carries the hive for the same reason `TriggerDeploy`/`MintAgentIdentity`
/// do: the wanted-state bucket is keyed per hive.
/// Carries the hive because the wanted-state bucket is keyed per hive.
SetAgentWanted { hive: String, agent: String },
/// Tell `hive` to rebuild `agent`, by publishing on the swarm's deploy
/// subject. The one node kind whose effect leaves this host.
@ -166,12 +164,12 @@ impl hive_jobq_wire::WireNode for SwarmNodeKind {
| SwarmNodeKind::CreateForgeUser { agent }
| SwarmNodeKind::AddRepoMember { agent }
| SwarmNodeKind::InitAgentConfigRepo { agent }
| SwarmNodeKind::MintAgentIdentity { agent }
| SwarmNodeKind::MintAgentForgeToken { agent }
| SwarmNodeKind::MintAgentMatrixAccount { agent } => {
serde_json::json!({ "agent": agent })
}
SwarmNodeKind::TriggerDeploy { hive, agent }
| SwarmNodeKind::MintAgentIdentity { hive, agent }
| SwarmNodeKind::SetAgentWanted { hive, agent } => {
serde_json::json!({ "agent": agent, "hive": hive })
}
@ -314,7 +312,7 @@ async fn run_swarm_node(
Err(e) => Outcome::Failed(format!("{e:#}")),
},
},
SwarmNodeKind::MintAgentIdentity { hive, agent } => mint_identity(&agent, &hive).await,
SwarmNodeKind::MintAgentIdentity { agent } => mint_identity(&agent).await,
SwarmNodeKind::MintAgentForgeToken { agent } => mint_forge_token(deps.forge, &agent).await,
SwarmNodeKind::MintAgentMatrixAccount { agent } => {
mint_matrix_account(deps.matrix_homeserver.as_deref(), &agent).await
@ -340,10 +338,10 @@ async fn run_swarm_node(
/// The `MintAgentIdentity` arm, lifted out so `run_swarm_node` stays under
/// `clippy::too_many_lines`.
async fn mint_identity(agent: &str, hive: &str) -> hive_jobq::scheduler::Outcome {
async fn mint_identity(agent: &str) -> hive_jobq::scheduler::Outcome {
use hive_jobq::scheduler::Outcome;
match agent_identity::mint_and_verify(agent, hive).await {
match agent_identity::mint_and_verify(agent).await {
Ok(()) => Outcome::Done,
Err(e) => Outcome::Failed(format!("{e:#}")),
}
@ -1512,7 +1510,6 @@ fn declare_agent_job(
// authelia.
let mint_identity = b
.node(SwarmNodeKind::MintAgentIdentity {
hive: hive.to_owned(),
agent: agent.to_owned(),
})
.after_ok(create_identity);
@ -1670,20 +1667,6 @@ async fn get_agent_config_pr(
Ok(Json(cache.get(&name)))
}
/// Body of `POST /api/agents/{name}/identity`.
#[derive(Deserialize, ToSchema)]
struct MintAgentIdentityRequest {
/// The hive this agent belongs to.
///
/// Required, for the same reason [`CreateAgentRequest`]'s is: the
/// credentials this mints name a hive, and the controller has nowhere to
/// look one up — agents are created on hives at runtime and this daemon
/// keeps no roster of which agent is where. An operator naming the wrong
/// one would hand the agent subjects on a hive it does not run on, so it
/// is asked for rather than guessed at.
hive: String,
}
/// Success body of `POST /api/agents/{name}/identity`.
#[derive(Clone, Debug, Serialize, ToSchema)]
struct MintAgentIdentityResponse {
@ -1712,10 +1695,9 @@ struct MintAgentIdentityResponse {
post,
path = "/api/agents/{name}/identity",
params(("name" = String, Path, description = "agent name")),
request_body = MintAgentIdentityRequest,
responses(
(status = 200, description = "mint queued", body = MintAgentIdentityResponse),
(status = 400, description = "`name` or `hive` is not a valid identifier, or `hive` is not in this swarm (problem+json)", body = String),
(status = 400, description = "`name` is not a valid identifier (problem+json)", body = String),
(status = 500, description = "the job could not be queued (problem+json)", body = String),
),
tag = "agents"
@ -1723,28 +1705,12 @@ struct MintAgentIdentityResponse {
async fn mint_agent_identity(
State(state): State<AppState>,
Path(name): Path<String>,
Json(req): Json<MintAgentIdentityRequest>,
) -> Result<Json<MintAgentIdentityResponse>, problem_details::ProblemDetails> {
// Both names are interpolated into store paths and policy documents
// downstream, so both are validated here as well as there.
// The name is interpolated into store paths and policy documents
// downstream, so it is validated here as well as there.
let agent = hive_types::Ident::parse(&name)
.map_err(|reason| error_problem(axum::http::StatusCode::BAD_REQUEST, reason))?
.into_string();
let hive = hive_types::Ident::parse(&req.hive)
.map_err(|reason| error_problem(axum::http::StatusCode::BAD_REQUEST, reason))?
.into_string();
if !state.hives.iter().any(|h| h.name == hive) {
let known: Vec<&str> = state.hives.iter().map(|h| h.name.as_str()).collect();
let known = if known.is_empty() {
"(none configured)".to_owned()
} else {
known.join(", ")
};
return Err(error_problem(
axum::http::StatusCode::BAD_REQUEST,
&format!("hive {hive:?} is not in this swarm — known hives: {known}"),
));
}
// No reserved-name or collision warnings here, unlike `create_agent`:
// those answer "is this name available", and this route is only ever
@ -1757,7 +1723,6 @@ async fn mint_agent_identity(
.insert_job(None, |b| {
vec![
b.node(SwarmNodeKind::MintAgentIdentity {
hive: hive.clone(),
agent: agent.clone(),
})
.guid(),
@ -2803,39 +2768,6 @@ mod tests {
assert_eq!(resp.change, super::ForgeAdminChange::AlreadyAdmin);
}
/// The backfill route's roster check, asserted by effect for the same
/// reason its sibling above is: a refusal that queued first would still
/// re-mint the agent's certificate, which every running agent on the
/// named hive picks up on its next boot.
#[tokio::test]
async fn a_backfill_for_a_hive_outside_the_roster_queues_nothing() {
let (state, sched) = state_with_roster();
let err = super::mint_agent_identity(
axum::extract::State(state),
axum::extract::Path("atlas".to_owned()),
axum::Json(super::MintAgentIdentityRequest {
hive: "pr1maa".to_owned(),
}),
)
.await
.expect_err("a hive outside the roster must be refused");
let rendered = format!("{err:?}");
assert!(
rendered.contains("pr1ma"),
"the refusal should name the known hives, got: {rendered}"
);
let queued = sched
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
.graph()
.nodes()
.count();
assert_eq!(queued, 0, "a refused backfill must queue no work");
}
/// A name that is not an identifier is refused before it can reach a
/// store path. `create_agent` gets this for free from the roster check
/// on the hive; the agent name has no roster to check against, so this
@ -2848,9 +2780,6 @@ mod tests {
super::mint_agent_identity(
axum::extract::State(state),
axum::extract::Path("../beta".to_owned()),
axum::Json(super::MintAgentIdentityRequest {
hive: "pr1ma".to_owned(),
}),
)
.await
.expect_err("a traversal in the agent name must be refused");
@ -2876,12 +2805,9 @@ mod tests {
let queued = super::mint_agent_identity(
axum::extract::State(state),
axum::extract::Path("atlas".to_owned()),
axum::Json(super::MintAgentIdentityRequest {
hive: "pr1ma".to_owned(),
}),
)
.await
.expect("a hive in the roster must be accepted");
.expect("a valid agent name must be accepted");
let guard = sched
.lock()
@ -2893,10 +2819,9 @@ mod tests {
assert!(
matches!(
&node.payload,
SwarmNodeKind::MintAgentIdentity { hive, agent }
if hive == "pr1ma" && agent == "atlas"
SwarmNodeKind::MintAgentIdentity { agent } if agent == "atlas"
),
"the one node must be the mint, carrying both names: {:?}",
"the one node must be the mint, for the named agent: {:?}",
node.payload
);
// The id the operator is told to watch has to be the node that was
@ -3025,7 +2950,6 @@ mod tests {
let id = sched
.append(
SwarmNodeKind::MintAgentIdentity {
hive: "pr1ma".to_owned(),
agent: "atlas".to_owned(),
},
Vec::new(),
@ -3192,25 +3116,6 @@ mod tests {
);
}
/// The node carries the hive, and the viewer has to see it. The `data`
/// match is an or-pattern on purpose (see its own comment), and this is
/// the assertion that the new variant joined the two-field arm rather
/// than the agent-only one — a viewer silently missing the hive is the
/// failure that comment describes having already happened once.
#[test]
fn a_mint_node_renders_both_the_agent_and_the_hive() {
use hive_jobq_wire::WireNode as _;
let kind = SwarmNodeKind::MintAgentIdentity {
hive: "pr1ma".to_owned(),
agent: "atlas".to_owned(),
};
assert_eq!(kind.label(), "mint_agent_identity");
let data = kind.data(1);
assert_eq!(data["agent"], "atlas");
assert_eq!(data["hive"], "pr1ma");
}
/// The ordering the operator's ruling requires: a hive cannot pass down
/// a certificate the swarm has not published, so the deploy message must
/// not leave before the mint is terminal.

View file

@ -9,7 +9,8 @@
//!
//! **Live tail only, on purpose.** `hive-agent` publishes each
//! already-classified `TermMsg` row to the core subject
//! `$SWARM.term.{hive}.{agent}` (see `hive-agent::swarm_term`'s module
//! `$SWARM.term.{agent}` (`$SWARM.term.{hive}.{agent}` from an agent on its
//! hive's shared credential; see `hive-agent::swarm_term`'s module
//! doc) — a core subject, not `JetStream`, so a subscriber who was not
//! listening missed the row, same as on the agent's own local SSE
//! stream. This handler relays exactly that: no replay, no last-N
@ -34,8 +35,8 @@ use crate::AppState;
/// Subject family carrying agent terminal rows — must match
/// `hive-agent::swarm_term::SUBJECT_PREFIX` exactly, since the two ends
/// never see the constant together. `$SWARM.term.<hive>.<agent>` is the
/// full subject; there is no way to check the two sides agree short of
/// never see the constant together. [`agent_subjects`] spells the full
/// subjects; there is no way to check the two sides agree short of
/// this comment and the module docs on both ends staying honest about it.
const SUBJECT_PREFIX: &str = "$SWARM.term";
@ -116,17 +117,63 @@ pub(crate) async fn stream_agent_term(
)
})?;
let subject = format!("{SUBJECT_PREFIX}.{hive}.{agent}");
let subscriber = subscribe_all(
&client,
&agent_subjects(SUBJECT_PREFIX, &hive, &agent),
"term",
)
.await?;
let stream =
subscriber.map(|msg| Ok(Event::default().data(String::from_utf8_lossy(&msg.payload))));
Ok(Sse::new(stream).keep_alive(KeepAlive::default()))
}
/// The subjects one agent publishes on under `prefix`: its own
/// `<prefix>.<agent>`, and `<prefix>.<hive>.<agent>`, which an agent still on
/// its hive's shared queue credential publishes to.
pub(crate) fn agent_subjects(prefix: &str, hive: &str, agent: &str) -> [String; 2] {
[
format!("{prefix}.{agent}"),
format!("{prefix}.{hive}.{agent}"),
]
}
/// Subscribe to every one of `subjects`, merged into one stream. `what` names
/// the relay in logs and errors.
pub(crate) async fn subscribe_all(
client: &async_nats::Client,
subjects: &[String],
what: &str,
) -> Result<futures_util::stream::SelectAll<async_nats::Subscriber>, problem_details::ProblemDetails>
{
let mut subscribers = Vec::with_capacity(subjects.len());
for subject in subjects {
let subscriber = client.subscribe(subject.clone()).await.map_err(|e| {
tracing::warn!(%subject, error = %e, "term stream: subscribe failed");
tracing::warn!(%subject, error = %e, "{what} stream: subscribe failed");
crate::error_problem(
axum::http::StatusCode::INTERNAL_SERVER_ERROR,
&format!("subscribing to {subject} failed: {e}"),
)
})?;
subscribers.push(subscriber);
}
tracing::info!(?subjects, "{what} stream: client attached");
Ok(futures_util::stream::select_all(subscribers))
}
tracing::info!(%subject, "term stream: client attached");
let stream =
subscriber.map(|msg| Ok(Event::default().data(String::from_utf8_lossy(&msg.payload))));
Ok(Sse::new(stream).keep_alive(KeepAlive::default()))
#[cfg(test)]
mod tests {
use super::*;
/// Both shapes, so an agent on either credential is streamed.
#[test]
fn a_relay_subscribes_to_the_agents_own_subject_and_its_hive_scoped_one() {
assert_eq!(
agent_subjects(SUBJECT_PREFIX, "h1", "atlas"),
[
"$SWARM.term.atlas".to_owned(),
"$SWARM.term.h1.atlas".to_owned(),
]
);
}
}

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,
};
// 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"
);
// 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)
}
// The client id is an identifier, not a credential, and it is the
// only thing tying a connection in this log to a hive.
Some(agent_token::Presented::Bearer(token)) => {
introspected(&policy, &http, &args, &client_secret, token).await
}
None => (None, None),
};
// 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());
}
}

View file

@ -0,0 +1,155 @@
//! The one spelling of the token an agent presents its own queue secret in,
//! shared by the agent that formats it and the auth-callout responder that
//! parses it back:
//! [`AGENT_TOKEN_PREFIX`](crate::agent_token::AGENT_TOKEN_PREFIX), the agent's
//! name, `.`, the secret.
//!
//! The secret itself lives at `swarm/agents/<agent>/queue` in the swarm's
//! secret store (`swarm_secret_client::queue`). This module knows nothing of
//! the store, so an agent formats its token without linking a store client.
/// Marks an `auth_token` as an agent's own credential.
///
/// An OIDC access token may itself contain `.`, so the `<agent>.<secret>`
/// shape alone does not tell the two apart. Authelia's access tokens start
/// `authelia_at_`.
pub const AGENT_TOKEN_PREFIX: &str = "swarm-agent.";
/// An agent's own credential as presented at the queue.
///
/// No `Debug`: `secret` is the credential itself.
pub struct AgentToken<'a> {
/// The agent the presenter claims to be. Unproven until `secret` is
/// checked against that agent's stored credential.
pub agent: &'a str,
/// The presented secret.
pub secret: &'a str,
}
/// A token that is not `<agent>.<secret>` after its prefix, or parts that
/// would not make one. Carries the reason only: a token holds a secret.
#[derive(Debug, thiserror::Error)]
#[error("malformed agent token: {0}")]
pub struct Malformed(pub &'static str);
/// Spell `agent`'s credential as the token it presents at the queue.
///
/// # Errors
/// [`Malformed`] when `agent` or `secret` is empty or holds anything outside
/// `[A-Za-z0-9_-]`. Either would make [`parse_agent_token`] split the token
/// differently than it was joined.
pub fn format_agent_token(agent: &str, secret: &str) -> Result<String, Malformed> {
check_agent(agent)?;
check_secret(secret)?;
Ok(format!("{AGENT_TOKEN_PREFIX}{agent}.{secret}"))
}
/// Read a presented token back into an [`AgentToken`].
///
/// `None` when `token` does not start with [`AGENT_TOKEN_PREFIX`]: it is some
/// other kind of token, not a malformed one of these. `Some(Err(_))` when it
/// does and is not exactly `<agent>.<secret>` after it.
#[must_use]
pub fn parse_agent_token(token: &str) -> Option<Result<AgentToken<'_>, Malformed>> {
let rest = token.strip_prefix(AGENT_TOKEN_PREFIX)?;
Some(
rest.split_once('.')
.ok_or(Malformed("no `.` between the agent and the secret"))
.and_then(|(agent, secret)| {
check_agent(agent)?;
check_secret(secret)?;
Ok(AgentToken { agent, secret })
}),
)
}
fn is_segment(s: &str) -> bool {
!s.is_empty()
&& s.bytes()
.all(|b| b.is_ascii_alphanumeric() || b == b'-' || b == b'_')
}
/// The alphabet a store path segment allows, so the name cannot widen a
/// subject with `.`, `*` or `>`, or address another agent's path.
fn check_agent(agent: &str) -> Result<(), Malformed> {
if is_segment(agent) {
Ok(())
} else {
Err(Malformed("the agent is not a single [A-Za-z0-9_-] segment"))
}
}
/// The alphabet the controller mints secrets in, base64url without padding.
/// The error never carries the value.
fn check_secret(secret: &str) -> Result<(), Malformed> {
if is_segment(secret) {
Ok(())
} else {
Err(Malformed(
"the secret is empty or holds a byte outside [A-Za-z0-9_-]",
))
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn an_agent_token_round_trips() {
let token = format_agent_token("atlas", "Ab9_-z").expect("legal");
assert_eq!(token, "swarm-agent.atlas.Ab9_-z");
let parsed = parse_agent_token(&token)
.expect("carries the prefix")
.expect("well-formed");
assert_eq!(parsed.agent, "atlas");
assert_eq!(parsed.secret, "Ab9_-z");
}
/// Anything without the prefix belongs to the OIDC path, including a
/// token that happens to look like `<agent>.<secret>`.
#[test]
fn a_token_without_the_prefix_is_not_an_agent_token() {
for token in ["authelia_at_abc.def", "atlas.s3cr3t", "", "swarm-agent"] {
assert!(parse_agent_token(token).is_none(), "{token:?}");
}
}
#[test]
fn a_malformed_agent_token_is_refused() {
for token in [
"swarm-agent.",
"swarm-agent.atlas",
"swarm-agent.atlas.",
"swarm-agent..s3cr3t",
"swarm-agent.atlas.s3.cr3t",
"swarm-agent.at*las.s3cr3t",
"swarm-agent.at>las.s3cr3t",
"swarm-agent.at/las.s3cr3t",
"swarm-agent.atlas.s3cr3t=",
"swarm-agent.atlas.s3 cr3t",
] {
assert!(
matches!(parse_agent_token(token), Some(Err(_))),
"{token:?} must be refused"
);
}
}
/// What formats always parses back into the same two parts.
#[test]
fn formatting_refuses_what_parsing_would_split_differently() {
assert!(format_agent_token("at.las", "s3cr3t").is_err());
assert!(format_agent_token("atlas", "s3.cr3t").is_err());
assert!(format_agent_token("atlas", "").is_err());
assert!(format_agent_token("", "s3cr3t").is_err());
}
#[test]
fn a_malformed_secret_is_not_echoed_in_the_error() {
let Some(Err(e)) = parse_agent_token("swarm-agent.atlas.hunter2!") else {
panic!("must be refused");
};
assert!(!e.to_string().contains("hunter2"), "{e}");
}
}

View file

@ -178,6 +178,10 @@ pub mod wanted;
/// inside it. See the module doc for why the two must not merge.
pub mod agent_status;
/// The token an agent presents its own queue secret in. Store-free, so an agent
/// formats it without linking the secret-store client.
pub mod agent_token;
/// The subject the swarm controller publishes on when the hive-wide knowledge
/// repository has changed. One writer, many readers — every hive subscribes.
///
@ -316,6 +320,15 @@ const TOKEN_REFRESH_SKEW: std::time::Duration = std::time::Duration::from_mins(2
/// degrades every other client of it.
const MAX_RECONNECT_DELAY: std::time::Duration = std::time::Duration::from_mins(1);
/// Exponential from 500ms, capped at [`MAX_RECONNECT_DELAY`].
fn reconnect_delay(attempts: usize) -> std::time::Duration {
let exp = u32::try_from(attempts.saturating_sub(1)).unwrap_or(u32::MAX);
std::cmp::min(
std::time::Duration::from_millis(500u64.saturating_mul(2u64.saturating_pow(exp.min(8)))),
MAX_RECONNECT_DELAY,
)
}
/// Where the controller finds the queue and what it authenticates with.
///
/// Every field comes from an environment variable the NixOS module sets, the
@ -738,15 +751,7 @@ pub async fn connect(cfg: QueueConfig) -> Result<async_nats::Client, Error> {
// which turns an unreachable queue into a permanent 4s poll — and, before
// the cache above, a permanent 4s token-request loop against authelia.
// Exponential from 500ms so a momentary blip still reconnects promptly.
.reconnect_delay_callback(|attempts| {
let exp = u32::try_from(attempts.saturating_sub(1)).unwrap_or(u32::MAX);
std::cmp::min(
std::time::Duration::from_millis(
500u64.saturating_mul(2u64.saturating_pow(exp.min(8))),
),
MAX_RECONNECT_DELAY,
)
})
.reconnect_delay_callback(reconnect_delay)
// 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
@ -780,6 +785,29 @@ pub async fn connect(cfg: QueueConfig) -> Result<async_nats::Client, Error> {
Ok(client)
}
/// Connect presenting a fixed `token`, as an agent presents its own queue
/// credential. Reconnects present the same token.
///
/// Unlike [`connect`], the first attempt must succeed: a refused or failed
/// connect is returned rather than retried in the background, so the caller
/// can fall back to another credential. `ca_file` is the queue's trust anchor,
/// as [`QueueConfig::ca_file`].
pub async fn connect_with_token(
url: &str,
ca_file: Option<&std::path::Path>,
token: String,
) -> Result<async_nats::Client, Error> {
let mut options =
async_nats::ConnectOptions::with_token(token).reconnect_delay_callback(reconnect_delay);
if let Some(path) = ca_file {
options = options.add_root_certificates(path.to_path_buf());
}
options.connect(url).await.map_err(|source| Error::Connect {
url: url.to_owned(),
source,
})
}
#[cfg(test)]
mod tests {
use super::*;

View file

@ -22,6 +22,9 @@
//! reaching the store and nothing else; deriving a queue identity from it would
//! couple the two credentials' lifetimes, so that renewing one would mean
//! renewing the other.
//!
//! The token an agent presents that secret in at the queue is spelled by
//! `swarm_queue_client::agent_token`, which needs no store client.
use serde::{Deserialize, Serialize};
@ -57,14 +60,11 @@ pub fn agent_queue_path(agent: &str) -> Result<String, Error> {
Ok(format!("{prefix}/queue"))
}
/// What [`agent_queue_path`] holds: the secret, and the principal it proves.
/// What [`agent_queue_path`] holds: the secret, and the agent it proves.
///
/// Both names ride **in the object** rather than being parsed back out of a
/// composite principal string. Hive and agent names draw from the same
/// alphabet (`hive_types::Ident`, `[a-z0-9-]`), so a principal spelled
/// `hive-<hive>-agent-<agent>` parses two ways for a name containing `-agent-`
/// — and an ambiguous principal parse in an authorisation path is a caller that
/// authenticates fine and is handed somebody else's grant.
/// No hive: an agent's identity is not tied to one, and the subjects the
/// verifier grants are keyed on the agent alone. Objects written with a `hive`
/// field still decode, because unknown fields are ignored.
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct AgentCredential {
/// The secret itself. Named to match [`Credential::value`] and
@ -74,13 +74,6 @@ pub struct AgentCredential {
/// The agent this secret authenticates.
pub agent: String,
/// The hive that agent belongs to.
///
/// Here because the verifying end has no roster to look it up in, and
/// because the subjects an agent is granted are hive-templated — without
/// this field the verifier would know *who* is connecting and not *where*.
pub hive: String,
}
/// What the path holds: the client secret, plus the client id it belongs to.
@ -198,7 +191,6 @@ mod tests {
let c = AgentCredential {
value: "s3cr3t".to_owned(),
agent: "atlas".to_owned(),
hive: "alpha".to_owned(),
};
let json = serde_json::to_string(&c).expect("serialises");
assert_eq!(
@ -212,29 +204,39 @@ mod tests {
let json = serde_json::to_value(AgentCredential {
value: "s3cr3t".to_owned(),
agent: "atlas".to_owned(),
hive: "alpha".to_owned(),
})
.expect("serialises");
assert_eq!(json["value"], "s3cr3t");
assert_eq!(json["agent"], "atlas");
assert_eq!(json["hive"], "alpha");
assert!(json.get("hive").is_none(), "{json}");
}
/// Neither name is optional. An object missing one is not a usable
/// credential — a verifier holding `None` for the hive can only guess at
/// the subjects to grant, and guessing is the failure this shape exists to
/// prevent.
/// A stored object may carry a `hive` field. It is ignored, and the object
/// decodes.
#[test]
fn an_agent_object_missing_a_principal_does_not_decode() {
assert!(
serde_json::from_str::<AgentCredential>(r#"{"value":"s","agent":"atlas"}"#).is_err()
fn a_stored_agent_object_that_still_names_a_hive_decodes() {
let c: AgentCredential =
serde_json::from_str(r#"{"value":"s3cr3t","agent":"atlas","hive":"alpha"}"#)
.expect("a stored object carrying `hive` decodes");
assert_eq!(
c,
AgentCredential {
value: "s3cr3t".to_owned(),
agent: "atlas".to_owned(),
}
);
}
/// The agent is required: it is what the verifier checks the presented
/// name against.
#[test]
fn an_agent_object_missing_its_agent_does_not_decode() {
assert!(serde_json::from_str::<AgentCredential>(r#"{"value":"s"}"#).is_err());
assert!(
serde_json::from_str::<AgentCredential>(r#"{"value":"s","hive":"alpha"}"#).is_err()
);
}
/// The two kinds are different objects at different paths, and neither
/// decodes as the other — the property that keeps a reader from picking up
/// a hive-shared credential where a per-agent one was meant.
#[test]
@ -249,7 +251,6 @@ mod tests {
let agent_json = serde_json::to_string(&AgentCredential {
value: "s".to_owned(),
agent: "atlas".to_owned(),
hive: "alpha".to_owned(),
})
.expect("serialises");
assert!(serde_json::from_str::<Credential>(&agent_json).is_err());

View file

@ -51,13 +51,6 @@ struct CreateAgentResponse {
warnings: Vec<String>,
}
/// Body of `POST /api/agents/{name}/identity`. Mirrors the controller's own
/// `MintAgentIdentityRequest` — see this module's doc comment.
#[derive(Serialize)]
struct MintIdentityRequest<'a> {
hive: &'a str,
}
/// Success body of `POST /api/agents/{name}/identity`.
#[derive(Deserialize)]
struct MintIdentityResponse {
@ -117,11 +110,10 @@ pub(crate) fn parse_ident(value: &str, what: &str) -> Result<String> {
///
/// Synchronous for the same reason [`create`] is, and built on the same
/// round trip.
pub(crate) fn mint_identity(socket: &Path, name: &str, hive: &str) -> Result<()> {
pub(crate) fn mint_identity(socket: &Path, name: &str) -> Result<()> {
// Client-side first, so a typo is a local error rather than a 400 the
// operator waits for. The controller validates both again.
// operator waits for. The controller validates it again.
let name = parse_ident(name, "agent name")?;
let hive = parse_ident(hive, "hive")?;
let rt = tokio::runtime::Builder::new_current_thread()
.enable_io()
@ -130,15 +122,15 @@ pub(crate) fn mint_identity(socket: &Path, name: &str, hive: &str) -> Result<()>
let resp: MintIdentityResponse = rt.block_on(post(
socket,
&format!("/api/agents/{name}/identity"),
&MintIdentityRequest { hive: &hive },
&serde_json::json!({}),
"mint-identity",
))?;
println!("queued: job node {}", resp.node_id);
println!(
"agent {name:?} will have its identity re-minted on hive {hive:?} once the job graph \
runs; `swarmctl` does not wait for it. An existing queue secret is kept as it is; the \
store certificate is re-minted and the agent picks the new one up on its next boot"
"agent {name:?} will have its identity re-minted once the job graph runs; `swarmctl` \
does not wait for it. An existing queue secret is kept as it is; the store \
certificate is re-minted and the agent picks the new one up on its next boot"
);
Ok(())
}

View file

@ -242,16 +242,6 @@ struct AgentMintForgeTokenArgs {
struct AgentMintIdentityArgs {
/// Name of an agent that already exists.
name: String,
/// The hive that agent runs on.
///
/// Required, and deliberately not defaulted: the credentials this mints
/// name a hive, and neither this CLI nor the controller keeps a roster of
/// which agent is on which hive. Naming the wrong one gives the agent an
/// identity scoped to a hive it doesn't run on. The controller checks
/// the value against the swarm's hive roster and names the known hives if
/// it misses.
#[arg(long, value_name = "HIVE")]
hive: String,
/// swarm-controller's unix socket.
///
/// Supplied by the nix module that installs this binary, from the same
@ -357,7 +347,7 @@ fn main() -> Result<()> {
command: AgentVerb::MintIdentity(args),
} => {
let socket = path_from(args.controller_socket, "SWARM_CONTROLLER_SOCKET")?;
agent::mint_identity(&socket, &args.name, &args.hive)
agent::mint_identity(&socket, &args.name)
}
// Same socket-resolution reasoning as `Create` above.
Verb::Agent {
@ -732,19 +722,9 @@ mod tests {
assert!(args.controller_socket.is_none());
}
/// The backfill verb takes the same two names as `create`, and `--hive`
/// is required on it for the same reason: it is an address nobody can
/// infer.
#[test]
fn the_backfill_verb_takes_an_agent_and_a_hive() {
let cli = Cli::try_parse_from([
"swarmctl",
"agent",
"mint-identity",
"scribe",
"--hive",
"alpha",
])
fn the_backfill_verb_takes_an_agent_and_no_hive() {
let cli = Cli::try_parse_from(["swarmctl", "agent", "mint-identity", "scribe"])
.expect("the minimal form parses");
let Verb::Agent {
command: AgentVerb::MintIdentity(args),
@ -753,12 +733,18 @@ mod tests {
panic!("expected `agent mint-identity`");
};
assert_eq!(args.name, "scribe");
assert_eq!(args.hive, "alpha");
assert!(args.controller_socket.is_none());
assert!(
Cli::try_parse_from(["swarmctl", "agent", "mint-identity", "scribe"]).is_err(),
"an omitted hive must not be defaulted"
Cli::try_parse_from([
"swarmctl",
"agent",
"mint-identity",
"scribe",
"--hive",
"a"
])
.is_err(),
"the identity has no hive, so the verb must not take one"
);
}