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",
"serde_json", "serde_json",
"sha2 0.11.0", "sha2 0.11.0",
"subtle",
"swarm-queue-client", "swarm-queue-client",
"swarm-secret-client",
"tokio", "tokio",
"tracing", "tracing",
"tracing-subscriber", "tracing-subscriber",

View file

@ -230,6 +230,9 @@ matrix-sdk = { version = "0.18", default-features = false, features = [
futures-util = "0.3" futures-util = "0.3"
hmac = "0.13" hmac = "0.13"
sha2 = "0.11" 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. # The NATS protocol client, for the swarm queue's auth-callout responder.
# `default-features = false` because the default set is broad - jetstream, kv, # `default-features = false` because the default set is broad - jetstream, kv,
# object-store, websockets, service - and a callout responder speaks none of # 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. startup, so she needs her identity minted by hand, once.
```bash ```bash
# On the swarm-controller host: give ruth her store identity. <hive> is the # On the swarm-controller host: give ruth her store identity.
# name of the hive she runs on. swarmctl agent mint-identity ruth
swarmctl agent mint-identity ruth --hive <hive>
# On ruth's hive: re-apply her container config, which is when hive-c0re # On ruth's hive: re-apply her container config, which is when hive-c0re
# hands the new identity to the container. # 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. environment logs an error and the harness keeps serving.
What an agent does with that connection is publish its terminal. Every row its 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 own web UI renders also goes to `$SWARM.term.<agent>`, one subject per agent, so
agent, so a swarm-level terminal can follow one agent without subscribing to a swarm-level terminal can follow one agent without subscribing to the swarm's
the swarm's whole traffic. The `<hive>` is the one the agent's client id names, whole traffic. An agent connected with its own queue credential gets that
which is the same string the broker builds its grant from. Publishing only: an 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 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 subscriber that wasn't listening missed them, the same as on the agent's own
live stream. 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. 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 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 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 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), 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: it. Re-run the mint for one agent with:
```sh ```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 The queue secret half is idempotent — an agent that already has one keeps exactly the
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
value it holds, so running this against an already-migrated agent doesn't drop 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 its queue connection. The certificate half isn't: the agent gets a fresh leaf
and picks it up on its next boot. 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. 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:** ###### **Arguments:**
@ -100,9 +100,6 @@ Queues and returns, the same way `agent create` does — watch the swarm UI's jo
###### **Options:** ###### **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. * `--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`. 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 swarm_queue_client::wanted::AgentState;
use crate::events::{Bus, BusEvent, LiveEvent, TurnState}; use crate::events::{Bus, BusEvent, LiveEvent, TurnState};
use crate::swarm_queue::{Connection, Presented};
use crate::term_msg::iso8601_utc; use crate::term_msg::iso8601_utc;
/// Subject family carrying agent turn-state headers, the swarm-wide /// 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 /// The subject this agent's header goes to under the credential it connected
/// identity the subject can be derived from. /// with, or `None` when a hive client id names no hive.
/// fn subject(presented: &Presented, agent: &str) -> Option<String> {
/// Returns without spawning in every other case — no queue, an unparseable match presented {
/// client id, no label — each of which is a legal state for an agent rather Presented::Agent => Some(format!("{SUBJECT_PREFIX}.{agent}")),
/// than an error, and each logged once here rather than per transition. Presented::Hive { client_id } => {
pub fn spawn(bus: &Bus) { let Some(hive) = hive_from_client_id(client_id) else {
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 {
tracing::warn!( tracing::warn!(
client_id = %cfg.client_id, %client_id,
expected = format!("{CLIENT_ID_PREFIX}<hive>{CLIENT_ID_SUFFIX}"), expected = format!("{CLIENT_ID_PREFIX}<hive>{CLIENT_ID_SUFFIX}"),
"queue client id does not name a hive; not publishing turn state upward" "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(); let agent = crate::identity::label();
if agent.is_empty() { if agent.is_empty() {
tracing::warn!("this agent has no label; not publishing turn state upward"); tracing::warn!("this agent has no label; not publishing turn state upward");
return; return;
} }
let subject = format!("{SUBJECT_PREFIX}.{hive}.{agent}"); tokio::spawn(run(bus.subscribe(), bus.clone(), agent));
tokio::spawn(run(bus.subscribe(), bus.clone(), subject));
} }
/// Watch the bus and publish whenever the header actually changed. /// 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 /// and it carries nothing this header reads. Every other variant, including
/// any added later, funnels into the comparison and costs nothing when it /// any added later, funnels into the comparison and costs nothing when it
/// changes nothing. /// changes nothing.
async fn run(mut rx: broadcast::Receiver<BusEvent>, bus: Bus, subject: String) { async fn run(mut rx: broadcast::Receiver<BusEvent>, bus: Bus, agent: String) {
let Some(client) = crate::swarm_queue::client().await else { let Some(Connection { client, presented }) = crate::swarm_queue::client().await else {
return;
};
let Some(subject) = subject(&presented, &agent) else {
return; return;
}; };
tracing::info!(subject, "publishing agent turn state to the swarm queue"); tracing::info!(subject, "publishing agent turn state to the swarm queue");
@ -307,7 +321,7 @@ async fn publish(
#[cfg(test)] #[cfg(test)]
mod tests { 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 crate::events::TurnState;
use swarm_queue_client::wanted::AgentState; use swarm_queue_client::wanted::AgentState;
@ -393,10 +407,22 @@ mod tests {
/// this pins the half that lives here. /// this pins the half that lives here.
#[test] #[test]
fn the_subject_is_the_prefix_then_the_hive_then_the_agent() { 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!( assert_eq!(
format!("{SUBJECT_PREFIX}.{hive}.mara"), subject(&presented, "mara").as_deref(),
"$SWARM.agent-state.alpha.mara" 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 //! 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`). //! fetched by the container itself (`nix/agent-modules/queue-identity.nix`).
//! //!
//! It is reported here but not yet *presented*: the queue's auth-callout //! When it is present, [`client`] connects with it first, and the agent
//! responder (`swarm-nats-auth`) validates only the hive-scoped token, and an //! publishes on its own hive-free subjects. When it is absent, or the queue
//! agent offering a credential nothing on the other end reads back would be //! refuses it, the connect falls back to the hive's shared client and the
//! refused. Until that responder learns the same path, the connect path below //! hive-scoped subjects. [`Connection::presented`] says which, so a publisher
//! is unchanged and this is the fetching half. //! builds the subject that credential is granted.
use std::path::{Path, PathBuf}; use std::path::{Path, PathBuf};
use std::sync::OnceLock; use std::sync::OnceLock;
use anyhow::Context as _;
use tokio::sync::OnceCell; use tokio::sync::OnceCell;
use swarm_queue_client::QueueConfig; use swarm_queue_client::QueueConfig;
@ -42,10 +43,39 @@ const ENV_PREFIX: &str = "HIVE_AGENT";
/// one answer rather than re-deriving it per call. /// one answer rather than re-deriving it per call.
static CONFIG: OnceLock<Option<QueueConfig>> = OnceLock::new(); 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 /// The one connection every publisher in this process shares. Separate from
/// [`CONFIG`] because resolving the coordinates is synchronous boot work and /// [`CONFIG`] because resolving the coordinates is synchronous boot work and
/// connecting is not — see [`client`]. /// 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 /// 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 /// 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 /// Decide what this agent's queue configuration is, given the environment and
/// whatever the client-id file held. /// whatever the client-id file held.
/// ///
@ -216,15 +256,10 @@ pub fn init() {
let _ = CONFIG.set(resolved); let _ = CONFIG.set(resolved);
// Independent of everything above: this agent may hold its own secret on // 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 // a hive whose shared client is not published yet, or the shared client
// of its own yet. Reported either way, because "which credential is this // and no secret of its own.
// agent able to present" is a question only this process can answer, and let secret = decide_agent_secret(env.agent_secret_file.as_deref());
// it is the one the next slice's rollout will be asked repeatedly. if let Some(path) = &secret {
//
// 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()) {
// The path, never the bytes: the file holds the secret itself. // The path, never the bytes: the file holds the secret itself.
tracing::info!( tracing::info!(
path = %path.display(), path = %path.display(),
@ -236,6 +271,13 @@ pub fn init() {
by its hive's shared client" 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. /// 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 /// 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 /// there is nothing to publish onto. Connecting is lazy so that an agent on a
/// hive with no queue pays nothing at boot. /// 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() 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()?; let cfg = config()?;
match swarm_queue_client::connect(cfg.clone()).await { 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) => { Err(e) => {
// `chain`, not `{:#}`: this is `swarm_queue_client::Error`, whose // `chain`, not `{:#}`: this is `swarm_queue_client::Error`, whose
// `Display` ignores the alternate flag, so `{:#}` renders the // `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)] #[cfg(test)]
mod tests { 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 { fn env(parts: [Option<&str>; 4]) -> QueueEnv {
let [nats_url, token_endpoint, client_id_file, client_secret_file] = parts; 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!(matches!(decide(&e, None), Resolution::Absent(_)));
assert!(decide_agent_secret(path.to_str()).is_some()); 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 tokio::sync::broadcast;
use crate::events::BusEvent; use crate::events::BusEvent;
use crate::swarm_queue::{Connection, Presented};
use crate::term_msg::{ClassifyCtx, TermMsg, classify}; use crate::term_msg::{ClassifyCtx, TermMsg, classify};
/// Subject family carrying agent terminal rows, the swarm-wide agreement this /// 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()) serde_json::to_vec(msg).ok().map(|v| v.len())
} }
/// Start the publish task, if this agent has both queue coordinates and an /// The subject this agent's rows go to under the credential it connected
/// identity the subject can be derived from. /// with, or `None` when a hive client id names no hive.
/// fn subject(presented: &Presented, agent: &str) -> Option<String> {
/// Returns without spawning in every other case — no queue, an unparseable match presented {
/// client id, no label — each of which is a legal state for an agent rather Presented::Agent => Some(format!("{SUBJECT_PREFIX}.{agent}")),
/// than an error, and each logged once here rather than per row. Presented::Hive { client_id } => {
pub fn spawn(rx: broadcast::Receiver<BusEvent>) { let Some(hive) = hive_from_client_id(client_id) else {
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 {
tracing::warn!( tracing::warn!(
client_id = %cfg.client_id, %client_id,
expected = format!("{CLIENT_ID_PREFIX}<hive>{CLIENT_ID_SUFFIX}"), expected = format!("{CLIENT_ID_PREFIX}<hive>{CLIENT_ID_SUFFIX}"),
"queue client id does not name a hive; not publishing the terminal upward" "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(); let agent = crate::identity::label();
if agent.is_empty() { if agent.is_empty() {
tracing::warn!("this agent has no label; not publishing the terminal upward"); tracing::warn!("this agent has no label; not publishing the terminal upward");
return; return;
} }
let subject = format!("{SUBJECT_PREFIX}.{hive}.{agent}"); tokio::spawn(run(rx, agent));
tokio::spawn(run(rx, subject));
} }
async fn run(mut rx: broadcast::Receiver<BusEvent>, subject: String) { async fn run(mut rx: broadcast::Receiver<BusEvent>, agent: String) {
let Some(client) = crate::swarm_queue::client().await else { let Some(Connection { client, presented }) = crate::swarm_queue::client().await else {
return;
};
let Some(subject) = subject(&presented, &agent) else {
return; return;
}; };
tracing::info!(subject, "publishing the agent terminal to the swarm queue"); 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)] #[cfg(test)]
mod tests { 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::events::LiveEvent;
use crate::term_msg::{BodyFormat, ClassifyCtx, Level, TermMsg, classify}; 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] #[test]
fn a_row_that_already_fits_is_published_unchanged() { fn a_row_that_already_fits_is_published_unchanged() {
let msg = TermMsg::new(Level::Info, "turn ok") let msg = TermMsg::new(Level::Info, "turn ok")

View file

@ -32,6 +32,7 @@
./glue-grafana-oidc-client.nix ./glue-grafana-oidc-client.nix
./glue-matrix-bao-token.nix ./glue-matrix-bao-token.nix
./glue-matrix-ctl-bao-identity.nix ./glue-matrix-ctl-bao-identity.nix
./glue-nats-auth-bao-identity.nix
./glue-nats-bao-identity.nix ./glue-nats-bao-identity.nix
./glue-queue-agent-credential.nix ./glue-queue-agent-credential.nix
./glue-secret-publisher-bao-identity.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. # leaf logs in with. Stays on this host; see its default above.
[ -s ${pkiDir}/granter.pem ] || ${signLeaf} ${pkiDir} granter \ [ -s ${pkiDir}/granter.pem ] || ${signLeaf} ${pkiDir} granter \
${lib.escapeShellArg deployCfg.bao.granterCommonName} "" clientAuth ${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 # 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 # above: the role attaches the policy by spelling it identically, and one
# string for both objects removes the way they drift apart. # 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 { matrixCtlHiveName = lib.mkOption {
type = lib.types.str; type = lib.types.str;
default = toString hyperhiveCfg.hiveName; default = toString hyperhiveCfg.hiveName;
@ -2021,6 +2051,7 @@ in
"swarm-bao-services-issuer-policy" "swarm-bao-services-issuer-policy"
"swarm-bao-nats-tls-policy" "swarm-bao-nats-tls-policy"
"swarm-bao-agent-pki" "swarm-bao-agent-pki"
"swarm-bao-nats-auth-policy"
]; ];
# 🚫 No `swarm.otel.scrapeTargets.bao` entry any more, and its absence is # 🚫 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-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 # 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 # 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 # grant: a bao policy reaches nothing until a login role hands it to a

View file

@ -219,6 +219,25 @@ let
tlsKeyCredential = "tls-key"; tlsKeyCredential = "tls-key";
tlsKeyCredentialPath = "/run/credentials/nats.service/${tlsKeyCredential}"; 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` # 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 # grants it (./swarm-bao.nix), so the daily timer below has two weeks of
# retries before it lapses. # retries before it lapses.
@ -498,6 +517,34 @@ in
A path, never a value. 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 { config = lib.mkIf deployCfg.nats.enable {
@ -905,6 +952,11 @@ in
# an append-only row stream, a header is one current value # an append-only row stream, a header is one current value
# republished on change. # republished on change.
"--agent-publish-subject ${lib.escapeShellArg "\$\$SWARM.agent-state.{hive}.>"}" "--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 # Every credential arrives by `LoadCredential` and is named
# on the command line only as a **path** — `argv` is # on the command line only as a **path** — `argv` is
@ -914,12 +966,25 @@ in
"callout-user.seed:${inContainer "callout-user.seed"}" "callout-user.seed:${inContainer "callout-user.seed"}"
"issuer.seed:${inContainer "issuer.seed"}" "issuer.seed:${inContainer "issuer.seed"}"
"oidc-client.secret:${inContainer "oidc-client.secret"}" "oidc-client.secret:${inContainer "oidc-client.secret"}"
]; ]
++ lib.mapAttrsToList (id: _: "${id}:${inContainer id}") authStoreFiles;
DynamicUser = true; DynamicUser = true;
Restart = "on-failure"; Restart = "on-failure";
RestartSec = "5s"; RestartSec = "5s";
SyslogIdentifier = "swarm-nats-auth"; 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 # 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 # being up says nothing about whether its in-container secrets unit
# has finished. The wait in the script is what actually closes it; # has finished. The wait in the script is what actually closes it;
# this only stops us spinning for the full timeout on every boot. # 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"; requires = lib.optional deployCfg.nats.autoGenerateCallout "swarm-nats-callout-keys.service";
serviceConfig = { serviceConfig = {
Type = "oneshot"; Type = "oneshot";
@ -1006,7 +1074,22 @@ in
install -m 0400 "$secret" \ install -m 0400 "$secret" \
${lib.escapeShellArg (hostPath "oidc-client.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 # ⚠️ Minted on the HOST, not in the container, because the responder is

View file

@ -53,6 +53,7 @@ let
deployCfg.bao.servicesIssuerCommonName deployCfg.bao.servicesIssuerCommonName
deployCfg.bao.natsCommonName deployCfg.bao.natsCommonName
deployCfg.bao.granterCommonName deployCfg.bao.granterCommonName
deployCfg.bao.natsAuthCommonName
] ]
# The two per-hive readers' subjects, spelled out per hive rather than as the # 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 # 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 _: u: (u.environment.BAO_CLIENT_CERT or null) == granterCertFile
) baoGrantWithConsumers.systemd.services; ) 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. # control below.
grantingUnitNames = [ grantingUnitNames = [
"swarm-bao-controller-policy" "swarm-bao-controller-policy"
@ -165,6 +165,7 @@ let
"swarm-bao-services-issuer-policy" "swarm-bao-services-issuer-policy"
"swarm-bao-nats-tls-policy" "swarm-bao-nats-tls-policy"
"swarm-bao-agent-pki" "swarm-bao-agent-pki"
"swarm-bao-nats-auth-policy"
]; ];
# Comment lines dropped first: both the HCL and the scripts explain # Comment lines dropped first: both the HCL and the scripts explain
@ -561,6 +562,33 @@ let
&& !(lib.hasInfix "swarm-grafana" s) && !(lib.hasInfix "swarm-grafana" s)
&& !(lib.hasInfix "sys/policies/acl" 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 # 🩸 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 # 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 # 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 # `lib.mkIf haveGranter` could be dropped from any of them and every other
# case here would still pass. # 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 = ok =
let let
s = baoGranterOptOut.systemd.services; s = baoGranterOptOut.systemd.services;
in in
lib.all (unit: !(s ? ${unit})) (grantingUnitNames ++ [ "swarm-bao-granter-role" ]) 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; && 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. # 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 # 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. # 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 = ok =
let let
s = baoGranterNoToken.systemd.services; s = baoGranterNoToken.systemd.services;
@ -1129,8 +1157,8 @@ let
} }
{ {
# What makes the case above mean something: discovery by the granter's # What makes the case above mean something: discovery by the granter's
# certificate reaches all eleven units, and each yields calls. # certificate reaches all twelve units, and each yields calls.
name = "the granter-policy check sees all eleven granting units, and parses calls from each"; name = "the granter-policy check sees all twelve granting units, and parses calls from each";
ok = ok =
lib.sort lib.lessThan (lib.attrNames granterUnits) == lib.sort lib.lessThan grantingUnitNames lib.sort lib.lessThan (lib.attrNames granterUnits) == lib.sort lib.lessThan grantingUnitNames
&& lib.all (u: baoCalls u.script != [ ]) (lib.attrValues granterUnits) && lib.all (u: baoCalls u.script != [ ]) (lib.attrValues granterUnits)

View file

@ -60,6 +60,28 @@ let
swarm.nats.calloutIssuerSeedFile = "/run/secrets/nats-issuer.seed"; 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 # 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 # `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 # does *not* say: it is the shape whose IdP address used to be null, and
@ -117,6 +139,73 @@ let
in in
lib.hasInfix "--agent-publish-subject '$$SWARM.agent-state.{hive}.>'" exec; 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 # Not a rename test. `hostClientSecretDir` is `readOnly`, so the fixture
# cannot define it; what can break is a reader left pointing at the # 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 /// 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 /// failure here fails the job node and nothing else — the agent is still
/// created, without a store identity. /// 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 (mount, pki_role) = agent_pki(|k| std::env::var(k).ok())?;
let name = policy::agent_object_name(agent)?; let name = policy::agent_object_name(agent)?;
let path = mtls::identity_path(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) .read_optional(&queue_path)
.await .await
.with_context(|| format!("checking whether {queue_path} already holds a credential"))?; .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. // The secret survives a re-run. An object naming a different agent is
// An object whose `hive` disagrees with the hive this node was invoked // corrected by rewriting the name around the *same* `value`, which no live
// with would grant its holder subjects on the wrong hive, so it is // connection notices.
// corrected — but by rewriting the two name fields around the *same*
// `value`, which is a correction no live connection notices.
let wanted = queue::AgentCredential { let wanted = queue::AgentCredential {
value: match &existing { value: match &existing {
Some(existing) => existing.value.clone(), Some(existing) => existing.value.clone(),
None => generate_queue_secret()?, None => generate_queue_secret()?,
}, },
agent: agent.to_owned(), agent: agent.to_owned(),
hive: hive.to_owned(),
}; };
if existing.as_ref() == Some(&wanted) { if existing.as_ref() == Some(&wanted) {
tracing::info!( 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}"))?; .with_context(|| format!("publishing the agent queue credential at {queue_path}"))?;
tracing::info!( tracing::info!(
agent, agent,
hive,
%queue_path, %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. // around it moved.
corrected = existing.is_some(), corrected = existing.is_some(),
"agent queue credential published" "agent queue credential published"
@ -280,9 +276,9 @@ async fn read_back_as_agent(
.read(queue_path) .read(queue_path)
.await .await
.with_context(|| format!("reading {queue_path} back under {role}'s own token"))?; .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 agent is compared as well as the secret: the verifying end refuses
// the verifying end will grant subjects from, so a mismatch here is the // an object naming a different agent than the path, so a mismatch here is
// same class of fault as an unreadable path. // the same class of fault as an unreadable path.
if read_back != *queue_credential { if read_back != *queue_credential {
bail!("the store returned a different object at {queue_path} than the one just published"); 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 /// Subject family carrying agent turn-state headers — must match
/// `hive-agent::swarm_agent_state::SUBJECT_PREFIX` exactly, since the two /// `hive-agent::swarm_agent_state::SUBJECT_PREFIX` exactly, since the two
/// ends never see the constant together. `$SWARM.agent-state.<hive>.<agent>` /// ends never see the constant together. `crate::term_stream::agent_subjects`
/// is the full subject; as with the terminal relay's own copy, there is no /// 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 /// way to check the two sides agree short of this comment and the module
/// docs on both ends staying honest about it. /// docs on both ends staying honest about it.
const SUBJECT_PREFIX: &str = "$SWARM.agent-state"; 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 = crate::term_stream::subscribe_all(
let subscriber = client.subscribe(subject.clone()).await.map_err(|e| { &client,
tracing::warn!(%subject, error = %e, "state stream: subscribe failed"); &crate::term_stream::agent_subjects(SUBJECT_PREFIX, &hive, &agent),
crate::error_problem( "state",
axum::http::StatusCode::INTERNAL_SERVER_ERROR,
&format!("subscribing to {subject} failed: {e}"),
) )
})?; .await?;
tracing::info!(%subject, "state stream: client attached");
let stream = let stream =
subscriber.map(|msg| Ok(Event::default().data(String::from_utf8_lossy(&msg.payload)))); subscriber.map(|msg| Ok(Event::default().data(String::from_utf8_lossy(&msg.payload))));
Ok(Sse::new(stream).keep_alive(KeepAlive::default())) 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 /// `agent_identity::mint_and_verify` — including why this node does not
/// report success on a write. /// report success on a write.
/// ///
/// Carries the hive for a different reason than `TriggerDeploy` does: /// Carries no hive: neither the certificate, the queue credential nor
/// not as an address, but because the agent's queue credential names the /// the agent's policy names one, so an agent keeps one store identity
/// hive it may take subjects on, so it cannot be written without knowing /// whichever hive it runs on.
/// which hive the agent belongs to. MintAgentIdentity { agent: String },
MintAgentIdentity { hive: String, agent: String },
/// Make sure `agent` holds a live forge access token in the swarm secret /// 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 /// 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 /// `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 /// Declare `agent` on `hive` as `Paused` in the swarm's wanted-state
/// store, so a freshly created agent does not start driving turns the /// 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`. /// moment it's deployed — the operator has to explicitly flip it to `Up`.
/// Carries the hive for the same reason `TriggerDeploy`/`MintAgentIdentity` /// Carries the hive because the wanted-state bucket is keyed per hive.
/// do: the wanted-state bucket is keyed per hive.
SetAgentWanted { hive: String, agent: String }, SetAgentWanted { hive: String, agent: String },
/// Tell `hive` to rebuild `agent`, by publishing on the swarm's deploy /// Tell `hive` to rebuild `agent`, by publishing on the swarm's deploy
/// subject. The one node kind whose effect leaves this host. /// 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::CreateForgeUser { agent }
| SwarmNodeKind::AddRepoMember { agent } | SwarmNodeKind::AddRepoMember { agent }
| SwarmNodeKind::InitAgentConfigRepo { agent } | SwarmNodeKind::InitAgentConfigRepo { agent }
| SwarmNodeKind::MintAgentIdentity { agent }
| SwarmNodeKind::MintAgentForgeToken { agent } | SwarmNodeKind::MintAgentForgeToken { agent }
| SwarmNodeKind::MintAgentMatrixAccount { agent } => { | SwarmNodeKind::MintAgentMatrixAccount { agent } => {
serde_json::json!({ "agent": agent }) serde_json::json!({ "agent": agent })
} }
SwarmNodeKind::TriggerDeploy { hive, agent } SwarmNodeKind::TriggerDeploy { hive, agent }
| SwarmNodeKind::MintAgentIdentity { hive, agent }
| SwarmNodeKind::SetAgentWanted { hive, agent } => { | SwarmNodeKind::SetAgentWanted { hive, agent } => {
serde_json::json!({ "agent": agent, "hive": hive }) serde_json::json!({ "agent": agent, "hive": hive })
} }
@ -314,7 +312,7 @@ async fn run_swarm_node(
Err(e) => Outcome::Failed(format!("{e:#}")), 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::MintAgentForgeToken { agent } => mint_forge_token(deps.forge, &agent).await,
SwarmNodeKind::MintAgentMatrixAccount { agent } => { SwarmNodeKind::MintAgentMatrixAccount { agent } => {
mint_matrix_account(deps.matrix_homeserver.as_deref(), &agent).await 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 /// The `MintAgentIdentity` arm, lifted out so `run_swarm_node` stays under
/// `clippy::too_many_lines`. /// `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; 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, Ok(()) => Outcome::Done,
Err(e) => Outcome::Failed(format!("{e:#}")), Err(e) => Outcome::Failed(format!("{e:#}")),
} }
@ -1512,7 +1510,6 @@ fn declare_agent_job(
// authelia. // authelia.
let mint_identity = b let mint_identity = b
.node(SwarmNodeKind::MintAgentIdentity { .node(SwarmNodeKind::MintAgentIdentity {
hive: hive.to_owned(),
agent: agent.to_owned(), agent: agent.to_owned(),
}) })
.after_ok(create_identity); .after_ok(create_identity);
@ -1670,20 +1667,6 @@ async fn get_agent_config_pr(
Ok(Json(cache.get(&name))) 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`. /// Success body of `POST /api/agents/{name}/identity`.
#[derive(Clone, Debug, Serialize, ToSchema)] #[derive(Clone, Debug, Serialize, ToSchema)]
struct MintAgentIdentityResponse { struct MintAgentIdentityResponse {
@ -1712,10 +1695,9 @@ struct MintAgentIdentityResponse {
post, post,
path = "/api/agents/{name}/identity", path = "/api/agents/{name}/identity",
params(("name" = String, Path, description = "agent name")), params(("name" = String, Path, description = "agent name")),
request_body = MintAgentIdentityRequest,
responses( responses(
(status = 200, description = "mint queued", body = MintAgentIdentityResponse), (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), (status = 500, description = "the job could not be queued (problem+json)", body = String),
), ),
tag = "agents" tag = "agents"
@ -1723,28 +1705,12 @@ struct MintAgentIdentityResponse {
async fn mint_agent_identity( async fn mint_agent_identity(
State(state): State<AppState>, State(state): State<AppState>,
Path(name): Path<String>, Path(name): Path<String>,
Json(req): Json<MintAgentIdentityRequest>,
) -> Result<Json<MintAgentIdentityResponse>, problem_details::ProblemDetails> { ) -> Result<Json<MintAgentIdentityResponse>, problem_details::ProblemDetails> {
// Both names are interpolated into store paths and policy documents // The name is interpolated into store paths and policy documents
// downstream, so both are validated here as well as there. // downstream, so it is validated here as well as there.
let agent = hive_types::Ident::parse(&name) let agent = hive_types::Ident::parse(&name)
.map_err(|reason| error_problem(axum::http::StatusCode::BAD_REQUEST, reason))? .map_err(|reason| error_problem(axum::http::StatusCode::BAD_REQUEST, reason))?
.into_string(); .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`: // No reserved-name or collision warnings here, unlike `create_agent`:
// those answer "is this name available", and this route is only ever // 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| { .insert_job(None, |b| {
vec![ vec![
b.node(SwarmNodeKind::MintAgentIdentity { b.node(SwarmNodeKind::MintAgentIdentity {
hive: hive.clone(),
agent: agent.clone(), agent: agent.clone(),
}) })
.guid(), .guid(),
@ -2803,39 +2768,6 @@ mod tests {
assert_eq!(resp.change, super::ForgeAdminChange::AlreadyAdmin); 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 /// 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 /// 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 /// on the hive; the agent name has no roster to check against, so this
@ -2848,9 +2780,6 @@ mod tests {
super::mint_agent_identity( super::mint_agent_identity(
axum::extract::State(state), axum::extract::State(state),
axum::extract::Path("../beta".to_owned()), axum::extract::Path("../beta".to_owned()),
axum::Json(super::MintAgentIdentityRequest {
hive: "pr1ma".to_owned(),
}),
) )
.await .await
.expect_err("a traversal in the agent name must be refused"); .expect_err("a traversal in the agent name must be refused");
@ -2876,12 +2805,9 @@ mod tests {
let queued = super::mint_agent_identity( let queued = super::mint_agent_identity(
axum::extract::State(state), axum::extract::State(state),
axum::extract::Path("atlas".to_owned()), axum::extract::Path("atlas".to_owned()),
axum::Json(super::MintAgentIdentityRequest {
hive: "pr1ma".to_owned(),
}),
) )
.await .await
.expect("a hive in the roster must be accepted"); .expect("a valid agent name must be accepted");
let guard = sched let guard = sched
.lock() .lock()
@ -2893,10 +2819,9 @@ mod tests {
assert!( assert!(
matches!( matches!(
&node.payload, &node.payload,
SwarmNodeKind::MintAgentIdentity { hive, agent } SwarmNodeKind::MintAgentIdentity { agent } if agent == "atlas"
if hive == "pr1ma" && 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 node.payload
); );
// The id the operator is told to watch has to be the node that was // The id the operator is told to watch has to be the node that was
@ -3025,7 +2950,6 @@ mod tests {
let id = sched let id = sched
.append( .append(
SwarmNodeKind::MintAgentIdentity { SwarmNodeKind::MintAgentIdentity {
hive: "pr1ma".to_owned(),
agent: "atlas".to_owned(), agent: "atlas".to_owned(),
}, },
Vec::new(), 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 /// The ordering the operator's ruling requires: a hive cannot pass down
/// a certificate the swarm has not published, so the deploy message must /// a certificate the swarm has not published, so the deploy message must
/// not leave before the mint is terminal. /// not leave before the mint is terminal.

View file

@ -9,7 +9,8 @@
//! //!
//! **Live tail only, on purpose.** `hive-agent` publishes each //! **Live tail only, on purpose.** `hive-agent` publishes each
//! already-classified `TermMsg` row to the core subject //! 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 //! 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 //! listening missed the row, same as on the agent's own local SSE
//! stream. This handler relays exactly that: no replay, no last-N //! 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 /// Subject family carrying agent terminal rows — must match
/// `hive-agent::swarm_term::SUBJECT_PREFIX` exactly, since the two ends /// `hive-agent::swarm_term::SUBJECT_PREFIX` exactly, since the two ends
/// never see the constant together. `$SWARM.term.<hive>.<agent>` is the /// never see the constant together. [`agent_subjects`] spells the full
/// full subject; there is no way to check the two sides agree short of /// 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. /// this comment and the module docs on both ends staying honest about it.
const SUBJECT_PREFIX: &str = "$SWARM.term"; 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| { 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( crate::error_problem(
axum::http::StatusCode::INTERNAL_SERVER_ERROR, axum::http::StatusCode::INTERNAL_SERVER_ERROR,
&format!("subscribing to {subject} failed: {e}"), &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"); #[cfg(test)]
let stream = mod tests {
subscriber.map(|msg| Ok(Event::default().data(String::from_utf8_lossy(&msg.payload)))); use super::*;
Ok(Sse::new(stream).keep_alive(KeepAlive::default()))
/// 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 serde_json.workspace = true
# The jti digest: base32hex(sha256(claims)) over every JWT this crate signs. # The jti digest: base32hex(sha256(claims)) over every JWT this crate signs.
sha2.workspace = true 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 # For `status::BUCKET` and `notices::STREAM` - the subjects a hive may
# publish to are derived from these names, and every end that touches them # 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 # must agree on the same one. Deliberately WITHOUT the `kv` feature: this
# crate derives subject strings, it never opens the bucket. `notices` # crate derives subject strings, it never opens the bucket. `notices`
# is name-only too (no `jetstream`/`kv` surface), same reason. # is name-only too (no `jetstream`/`kv` surface), same reason.
swarm-queue-client = { workspace = true, features = ["notices"] } 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 tokio.workspace = true
tracing.workspace = true tracing.workspace = true
tracing-subscriber.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 //! 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 //! server refuses to start if that entry carries a username), subscribes to
//! `$SYS.REQ.USER.AUTH`, validates the presented bearer token against //! `$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 //! by the account key. A rejection is answered explicitly: silence is
//! indistinguishable from the responder being down, and the queue is the //! indistinguishable from the responder being down, and the queue is the
//! swarm's control path. //! swarm's control path.
@ -27,6 +28,7 @@ use anyhow::Context;
use clap::Parser; use clap::Parser;
use futures_util::StreamExt; use futures_util::StreamExt;
mod agent_token;
mod introspect; mod introspect;
mod policy; mod policy;
mod request; mod request;
@ -89,7 +91,7 @@ struct Args {
/// config change plus a reload. /// config change plus a reload.
/// ///
/// So this identity says which hive an agent belongs to and never which /// 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 /// 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 /// `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. /// is refused, which is loud rather than silently over-broad.
#[arg(long = "agent-publish-subject")] #[arg(long = "agent-publish-subject")]
agent_publish_subjects: Vec<String>, 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. /// Read a secret file and strip surrounding whitespace.
@ -177,7 +191,9 @@ async fn main() -> anyhow::Result<()> {
args.reader_clients.clone(), args.reader_clients.clone(),
args.hive_publish_subjects.clone(), args.hive_publish_subjects.clone(),
args.agent_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 http = reqwest::Client::new();
let issuer = nkeys::KeyPair::from_seed(&read_secret(&args.issuer_seed_file)?) let issuer = nkeys::KeyPair::from_seed(&read_secret(&args.issuer_seed_file)?)
.context("parse the account signing seed")?; .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 // 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 // normal thing for a client to attempt and an abnormal thing to
// grant. Introspection is only reached once something was presented. // grant. Neither the store nor introspection is reached until
// // something was presented.
// The caller is an identity or nothing — see `introspect`'s module let (caller, permissions) = match req
// docs. There is no "admitted, identity unknown" branch to write here .connect_opts
// because there is no such value to receive. .auth_token
let caller = match &req.connect_opts.auth_token { .as_deref()
Some(token) => introspect::identify_caller( .map(agent_token::classify)
&http, {
&args.introspection_url, // The caller is the name the token claims, logged whether or not
&args.client_id, // the secret proved it; `granted` says which.
&client_secret, Some(agent_token::Presented::Agent(token)) => (
token, Some(format!("agent:{}", token.agent)),
) agent_token::authorize(&policy, store.as_ref(), &token).await,
.await ),
// An introspection that could not be *made* is a denial too. The Some(agent_token::Presented::Malformed(e)) => {
// failure modes of an HTTP call are exactly the conditions under tracing::warn!(error = %e, "malformed agent token; denying");
// which an attacker would most like this to fall open. (None, None)
.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"
);
} }
// The client id is an identifier, not a credential, and it is the Some(agent_token::Presented::Bearer(token)) => {
// only thing tying a connection in this log to a hive. 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!( tracing::info!(
user_nkey = %req.user_nkey, user_nkey = %req.user_nkey,
server_id = %req.server_id.id, server_id = %req.server_id.id,
@ -281,6 +286,50 @@ async fn main() -> anyhow::Result<()> {
anyhow::bail!("subscription to {AUTH_SUBJECT} ended") 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)] #[cfg(test)]
mod tests { mod tests {
use super::*; use super::*;
@ -306,9 +355,16 @@ mod tests {
"hive-", "hive-",
"--agent-client-suffix", "--agent-client-suffix",
"-agent", "-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"); .expect("the unit's own argument vector must parse");
assert_eq!(args.agent_client_suffix, "-agent"); 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 /// The control for the case above: an ordinary value parses through the

View file

@ -45,12 +45,16 @@ pub struct Policy {
readers: Vec<String>, readers: Vec<String>,
extra_hive_subjects: Vec<String>, extra_hive_subjects: Vec<String>,
extra_agent_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 /// Placeholder replaced with the hive's own name in `extra_hive_subjects` and
/// `extra_agent_subjects`. /// `extra_agent_subjects`.
const HIVE_PLACEHOLDER: &str = "{hive}"; const HIVE_PLACEHOLDER: &str = "{hive}";
/// Placeholder replaced with the agent's own name in `agent_token_subjects`.
const AGENT_PLACEHOLDER: &str = "{agent}";
impl Policy { impl Policy {
/// `hive_prefix` is the client-id prefix that marks a hive and /// `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 — /// `agent_suffix` what a hive's agent containers carry **on top of** it —
@ -130,9 +134,43 @@ impl Policy {
readers, readers,
extra_hive_subjects, extra_hive_subjects,
extra_agent_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. /// The permissions for `client_id`, or `None` when no rule matches.
/// ///
/// `None` is a denial. It is not "grant nothing and let them connect": /// `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:?}" "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. /// inside it. See the module doc for why the two must not merge.
pub mod agent_status; 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 /// The subject the swarm controller publishes on when the hive-wide knowledge
/// repository has changed. One writer, many readers — every hive subscribes. /// 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. /// degrades every other client of it.
const MAX_RECONNECT_DELAY: std::time::Duration = std::time::Duration::from_mins(1); 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. /// Where the controller finds the queue and what it authenticates with.
/// ///
/// Every field comes from an environment variable the NixOS module sets, the /// 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 // which turns an unreachable queue into a permanent 4s poll — and, before
// the cache above, a permanent 4s token-request loop against authelia. // the cache above, a permanent 4s token-request loop against authelia.
// Exponential from 500ms so a momentary blip still reconnects promptly. // Exponential from 500ms so a momentary blip still reconnects promptly.
.reconnect_delay_callback(|attempts| { .reconnect_delay_callback(reconnect_delay)
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,
)
})
// The controller and the queue are separate units on (possibly) // The controller and the queue are separate units on (possibly)
// separate hosts, and nothing orders them. Without this, a queue that // separate hosts, and nothing orders them. Without this, a queue that
// comes up one second later leaves the controller permanently // 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) 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)] #[cfg(test)]
mod tests { mod tests {
use super::*; use super::*;

View file

@ -22,6 +22,9 @@
//! reaching the store and nothing else; deriving a queue identity from it would //! reaching the store and nothing else; deriving a queue identity from it would
//! couple the two credentials' lifetimes, so that renewing one would mean //! couple the two credentials' lifetimes, so that renewing one would mean
//! renewing the other. //! 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}; use serde::{Deserialize, Serialize};
@ -57,14 +60,11 @@ pub fn agent_queue_path(agent: &str) -> Result<String, Error> {
Ok(format!("{prefix}/queue")) 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 /// No hive: an agent's identity is not tied to one, and the subjects the
/// composite principal string. Hive and agent names draw from the same /// verifier grants are keyed on the agent alone. Objects written with a `hive`
/// alphabet (`hive_types::Ident`, `[a-z0-9-]`), so a principal spelled /// field still decode, because unknown fields are ignored.
/// `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.
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct AgentCredential { pub struct AgentCredential {
/// The secret itself. Named to match [`Credential::value`] and /// The secret itself. Named to match [`Credential::value`] and
@ -74,13 +74,6 @@ pub struct AgentCredential {
/// The agent this secret authenticates. /// The agent this secret authenticates.
pub agent: String, 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. /// What the path holds: the client secret, plus the client id it belongs to.
@ -198,7 +191,6 @@ mod tests {
let c = AgentCredential { let c = AgentCredential {
value: "s3cr3t".to_owned(), value: "s3cr3t".to_owned(),
agent: "atlas".to_owned(), agent: "atlas".to_owned(),
hive: "alpha".to_owned(),
}; };
let json = serde_json::to_string(&c).expect("serialises"); let json = serde_json::to_string(&c).expect("serialises");
assert_eq!( assert_eq!(
@ -212,29 +204,39 @@ mod tests {
let json = serde_json::to_value(AgentCredential { let json = serde_json::to_value(AgentCredential {
value: "s3cr3t".to_owned(), value: "s3cr3t".to_owned(),
agent: "atlas".to_owned(), agent: "atlas".to_owned(),
hive: "alpha".to_owned(),
}) })
.expect("serialises"); .expect("serialises");
assert_eq!(json["value"], "s3cr3t"); assert_eq!(json["value"], "s3cr3t");
assert_eq!(json["agent"], "atlas"); 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 /// A stored object may carry a `hive` field. It is ignored, and the object
/// credential — a verifier holding `None` for the hive can only guess at /// decodes.
/// the subjects to grant, and guessing is the failure this shape exists to
/// prevent.
#[test] #[test]
fn an_agent_object_missing_a_principal_does_not_decode() { fn a_stored_agent_object_that_still_names_a_hive_decodes() {
assert!( let c: AgentCredential =
serde_json::from_str::<AgentCredential>(r#"{"value":"s","agent":"atlas"}"#).is_err() 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!( assert!(
serde_json::from_str::<AgentCredential>(r#"{"value":"s","hive":"alpha"}"#).is_err() 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 /// decodes as the other — the property that keeps a reader from picking up
/// a hive-shared credential where a per-agent one was meant. /// a hive-shared credential where a per-agent one was meant.
#[test] #[test]
@ -249,7 +251,6 @@ mod tests {
let agent_json = serde_json::to_string(&AgentCredential { let agent_json = serde_json::to_string(&AgentCredential {
value: "s".to_owned(), value: "s".to_owned(),
agent: "atlas".to_owned(), agent: "atlas".to_owned(),
hive: "alpha".to_owned(),
}) })
.expect("serialises"); .expect("serialises");
assert!(serde_json::from_str::<Credential>(&agent_json).is_err()); assert!(serde_json::from_str::<Credential>(&agent_json).is_err());

View file

@ -51,13 +51,6 @@ struct CreateAgentResponse {
warnings: Vec<String>, 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`. /// Success body of `POST /api/agents/{name}/identity`.
#[derive(Deserialize)] #[derive(Deserialize)]
struct MintIdentityResponse { 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 /// Synchronous for the same reason [`create`] is, and built on the same
/// round trip. /// 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 // 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 name = parse_ident(name, "agent name")?;
let hive = parse_ident(hive, "hive")?;
let rt = tokio::runtime::Builder::new_current_thread() let rt = tokio::runtime::Builder::new_current_thread()
.enable_io() .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( let resp: MintIdentityResponse = rt.block_on(post(
socket, socket,
&format!("/api/agents/{name}/identity"), &format!("/api/agents/{name}/identity"),
&MintIdentityRequest { hive: &hive }, &serde_json::json!({}),
"mint-identity", "mint-identity",
))?; ))?;
println!("queued: job node {}", resp.node_id); println!("queued: job node {}", resp.node_id);
println!( println!(
"agent {name:?} will have its identity re-minted on hive {hive:?} once the job graph \ "agent {name:?} will have its identity re-minted once the job graph runs; `swarmctl` \
runs; `swarmctl` does not wait for it. An existing queue secret is kept as it is; the \ does not wait for it. An existing queue secret is kept as it is; the store \
store certificate is re-minted and the agent picks the new one up on its next boot" certificate is re-minted and the agent picks the new one up on its next boot"
); );
Ok(()) Ok(())
} }

View file

@ -242,16 +242,6 @@ struct AgentMintForgeTokenArgs {
struct AgentMintIdentityArgs { struct AgentMintIdentityArgs {
/// Name of an agent that already exists. /// Name of an agent that already exists.
name: String, 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. /// swarm-controller's unix socket.
/// ///
/// Supplied by the nix module that installs this binary, from the same /// Supplied by the nix module that installs this binary, from the same
@ -357,7 +347,7 @@ fn main() -> Result<()> {
command: AgentVerb::MintIdentity(args), command: AgentVerb::MintIdentity(args),
} => { } => {
let socket = path_from(args.controller_socket, "SWARM_CONTROLLER_SOCKET")?; 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. // Same socket-resolution reasoning as `Create` above.
Verb::Agent { Verb::Agent {
@ -732,19 +722,9 @@ mod tests {
assert!(args.controller_socket.is_none()); 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] #[test]
fn the_backfill_verb_takes_an_agent_and_a_hive() { fn the_backfill_verb_takes_an_agent_and_no_hive() {
let cli = Cli::try_parse_from([ let cli = Cli::try_parse_from(["swarmctl", "agent", "mint-identity", "scribe"])
"swarmctl",
"agent",
"mint-identity",
"scribe",
"--hive",
"alpha",
])
.expect("the minimal form parses"); .expect("the minimal form parses");
let Verb::Agent { let Verb::Agent {
command: AgentVerb::MintIdentity(args), command: AgentVerb::MintIdentity(args),
@ -753,12 +733,18 @@ mod tests {
panic!("expected `agent mint-identity`"); panic!("expected `agent mint-identity`");
}; };
assert_eq!(args.name, "scribe"); assert_eq!(args.name, "scribe");
assert_eq!(args.hive, "alpha");
assert!(args.controller_socket.is_none()); assert!(args.controller_socket.is_none());
assert!( assert!(
Cli::try_parse_from(["swarmctl", "agent", "mint-identity", "scribe"]).is_err(), Cli::try_parse_from([
"an omitted hive must not be defaulted" "swarmctl",
"agent",
"mint-identity",
"scribe",
"--hive",
"a"
])
.is_err(),
"the identity has no hive, so the verb must not take one"
); );
} }