diff --git a/Cargo.lock b/Cargo.lock index 11a6dd49..41b34768 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -4828,7 +4828,9 @@ dependencies = [ "serde", "serde_json", "sha2 0.11.0", + "subtle", "swarm-queue-client", + "swarm-secret-client", "tokio", "tracing", "tracing-subscriber", diff --git a/Cargo.toml b/Cargo.toml index 22eb6dda..7368ead1 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -230,6 +230,9 @@ matrix-sdk = { version = "0.18", default-features = false, features = [ futures-util = "0.3" hmac = "0.13" sha2 = "0.11" +# Constant-time comparison of a presented secret against the stored one +# (`swarm-nats-auth::agent_token`). Already in the tree through the TLS stack. +subtle = "2.6" # The NATS protocol client, for the swarm queue's auth-callout responder. # `default-features = false` because the default set is broad - jetstream, kv, # object-store, websockets, service - and a callout responder speaks none of diff --git a/docs/getting-started/setup.md b/docs/getting-started/setup.md index b00365c6..5522b5ec 100644 --- a/docs/getting-started/setup.md +++ b/docs/getting-started/setup.md @@ -33,9 +33,8 @@ its identity then. Ruth doesn't: hive-c0re creates her on its own at startup, so she needs her identity minted by hand, once. ```bash -# On the swarm-controller host: give ruth her store identity. is the -# name of the hive she runs on. -swarmctl agent mint-identity ruth --hive +# On the swarm-controller host: give ruth her store identity. +swarmctl agent mint-identity ruth # On ruth's hive: re-apply her container config, which is when hive-c0re # hands the new identity to the container. diff --git a/docs/swarm/README.md b/docs/swarm/README.md index 37808df0..8dbd3f25 100644 --- a/docs/swarm/README.md +++ b/docs/swarm/README.md @@ -410,10 +410,13 @@ agents sets none of the four and each agent logs that it has none; a half-set environment logs an error and the harness keeps serving. What an agent does with that connection is publish its terminal. Every row its -own web UI renders also goes to `$SWARM.term..`, one subject per -agent, so a swarm-level terminal can follow one agent without subscribing to -the swarm's whole traffic. The `` is the one the agent's client id names, -which is the same string the broker builds its grant from. Publishing only: an +own web UI renders also goes to `$SWARM.term.`, one subject per agent, so +a swarm-level terminal can follow one agent without subscribing to the swarm's +whole traffic. An agent connected with its own queue credential gets that +subject. An agent without one, or whose own credential the queue +refused, connects with its hive's shared client and publishes to +`$SWARM.term..` instead, the `` being the one that client id +names. The swarm controller relays both. Publishing only: an agent talks about itself here and reads nothing. Rows aren't retained — a subscriber that wasn't listening missed them, the same as on the agent's own live stream. @@ -424,7 +427,7 @@ sending and leaves a marker in its place; the summary, level and icon still arrive. The harness logs and skips a row that's too large even without its body. The second thing an agent publishes is its **turn-state header**, on -`$SWARM.agent-state..` — same shape of subject, same grant +`$SWARM.agent-state.` (or `$SWARM.agent-state..`) — same shape of subject, same grant mechanics, same lack of retention. It carries what a header bar wants: what the turn loop is doing (`turn_state`, plus `turn_state_since` as an ISO 8601 UTC stamp), which model (`model` and the resolved id the last turn actually ran on), diff --git a/docs/swarm/credentials.md b/docs/swarm/credentials.md index d816e502..0ec0292b 100644 --- a/docs/swarm/credentials.md +++ b/docs/swarm/credentials.md @@ -101,12 +101,10 @@ mint therefore never receives one, and nothing will ever come back around to it. Re-run the mint for one agent with: ```sh -swarmctl agent mint-identity --hive +swarmctl agent mint-identity ``` -`--hive` has no default: neither the CLI nor the controller keeps a roster of -which agent runs where, and the credentials this mints name a hive. The queue -secret half is idempotent — an agent that already has one keeps exactly the +The queue secret half is idempotent — an agent that already has one keeps exactly the value it holds, so running this against an already-migrated agent doesn't drop its queue connection. The certificate half isn't: the agent gets a fresh leaf and picks it up on its next boot. diff --git a/docs/tools/swarmctl-cli.md b/docs/tools/swarmctl-cli.md index 031d536f..39d69995 100644 --- a/docs/tools/swarmctl-cli.md +++ b/docs/tools/swarmctl-cli.md @@ -92,7 +92,7 @@ Queue a re-mint of an existing agent's identity at the swarm's secret store. Queues and returns, the same way `agent create` does — watch the swarm UI's job view for the outcome. -**Usage:** `swarmctl agent mint-identity [OPTIONS] --hive ` +**Usage:** `swarmctl agent mint-identity [OPTIONS] ` ###### **Arguments:** @@ -100,9 +100,6 @@ Queues and returns, the same way `agent create` does — watch the swarm UI's jo ###### **Options:** -* `--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 ` — 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`. diff --git a/hive-agent/src/swarm_agent_state.rs b/hive-agent/src/swarm_agent_state.rs index 4c7a988e..f4d48d77 100644 --- a/hive-agent/src/swarm_agent_state.rs +++ b/hive-agent/src/swarm_agent_state.rs @@ -35,6 +35,7 @@ use tokio::sync::broadcast; use swarm_queue_client::wanted::AgentState; use crate::events::{Bus, BusEvent, LiveEvent, TurnState}; +use crate::swarm_queue::{Connection, Presented}; use crate::term_msg::iso8601_utc; /// Subject family carrying agent turn-state headers, the swarm-wide @@ -171,33 +172,43 @@ fn snapshot(bus: &Bus) -> AgentStateMsg { } } -/// Start the publish task, if this agent has both queue coordinates and an -/// identity the subject can be derived from. +/// The subject this agent's header goes to under the credential it connected +/// with, or `None` when a hive client id names no hive. +fn subject(presented: &Presented, agent: &str) -> Option { + match presented { + Presented::Agent => Some(format!("{SUBJECT_PREFIX}.{agent}")), + Presented::Hive { client_id } => { + let Some(hive) = hive_from_client_id(client_id) else { + tracing::warn!( + %client_id, + expected = format!("{CLIENT_ID_PREFIX}{CLIENT_ID_SUFFIX}"), + "queue client id does not name a hive; not publishing turn state upward" + ); + 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, an unparseable -/// client id, no label — each of which is a legal state for an agent rather -/// than an error, and each logged once here rather than per transition. +/// 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) { - let Some(cfg) = crate::swarm_queue::config() else { + 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 Some(hive) = hive_from_client_id(&cfg.client_id) else { - tracing::warn!( - client_id = %cfg.client_id, - expected = format!("{CLIENT_ID_PREFIX}{CLIENT_ID_SUFFIX}"), - "queue client id does not name a hive; not publishing turn state upward" - ); - return; - }; + } let agent = crate::identity::label(); if agent.is_empty() { tracing::warn!("this agent has no label; not publishing turn state upward"); return; } - let subject = format!("{SUBJECT_PREFIX}.{hive}.{agent}"); - tokio::spawn(run(bus.subscribe(), bus.clone(), subject)); + tokio::spawn(run(bus.subscribe(), bus.clone(), agent)); } /// Watch the bus and publish whenever the header actually changed. @@ -218,8 +229,11 @@ pub fn spawn(bus: &Bus) { /// and it carries nothing this header reads. Every other variant, including /// any added later, funnels into the comparison and costs nothing when it /// changes nothing. -async fn run(mut rx: broadcast::Receiver, bus: Bus, subject: String) { - let Some(client) = crate::swarm_queue::client().await else { +async fn run(mut rx: broadcast::Receiver, bus: Bus, agent: String) { + let Some(Connection { client, presented }) = crate::swarm_queue::client().await else { + return; + }; + let Some(subject) = subject(&presented, &agent) else { return; }; tracing::info!(subject, "publishing agent turn state to the swarm queue"); @@ -307,7 +321,7 @@ async fn publish( #[cfg(test)] mod tests { - use super::{AgentStateMsg, SUBJECT_PREFIX, hive_from_client_id}; + use super::{AgentStateMsg, Presented, hive_from_client_id, subject}; use crate::events::TurnState; use swarm_queue_client::wanted::AgentState; @@ -393,10 +407,22 @@ mod tests { /// this pins the half that lives here. #[test] fn the_subject_is_the_prefix_then_the_hive_then_the_agent() { - let hive = hive_from_client_id("hive-alpha-agent").expect("names a hive"); + let presented = Presented::Hive { + client_id: "hive-alpha-agent".to_owned(), + }; assert_eq!( - format!("{SUBJECT_PREFIX}.{hive}.mara"), - "$SWARM.agent-state.alpha.mara" + subject(&presented, "mara").as_deref(), + Some("$SWARM.agent-state.alpha.mara") + ); + } + + /// An agent that connected with its own credential publishes on the + /// subject that credential is granted, which names no hive. + #[test] + fn an_agent_on_its_own_credential_publishes_on_its_hive_free_subject() { + assert_eq!( + subject(&Presented::Agent, "mara").as_deref(), + Some("$SWARM.agent-state.mara") ); } diff --git a/hive-agent/src/swarm_queue.rs b/hive-agent/src/swarm_queue.rs index b8d267e9..3613e02a 100644 --- a/hive-agent/src/swarm_queue.rs +++ b/hive-agent/src/swarm_queue.rs @@ -21,15 +21,16 @@ //! agent. The fifth is this agent's own, minted per agent at swarm level and //! fetched by the container itself (`nix/agent-modules/queue-identity.nix`). //! -//! It is reported here but not yet *presented*: the queue's auth-callout -//! responder (`swarm-nats-auth`) validates only the hive-scoped token, and an -//! agent offering a credential nothing on the other end reads back would be -//! refused. Until that responder learns the same path, the connect path below -//! is unchanged and this is the fetching half. +//! When it is present, [`client`] connects with it first, and the agent +//! publishes on its own hive-free subjects. When it is absent, or the queue +//! refuses it, the connect falls back to the hive's shared client and the +//! hive-scoped subjects. [`Connection::presented`] says which, so a publisher +//! builds the subject that credential is granted. use std::path::{Path, PathBuf}; use std::sync::OnceLock; +use anyhow::Context as _; use tokio::sync::OnceCell; use swarm_queue_client::QueueConfig; @@ -42,10 +43,39 @@ const ENV_PREFIX: &str = "HIVE_AGENT"; /// one answer rather than re-deriving it per call. static CONFIG: OnceLock> = OnceLock::new(); +/// Where to present this agent's own credential, resolved once at boot. +static AGENT: OnceLock> = OnceLock::new(); + /// The one connection every publisher in this process shares. Separate from /// [`CONFIG`] because resolving the coordinates is synchronous boot work and /// connecting is not — see [`client`]. -static CLIENT: OnceCell> = OnceCell::const_new(); +static CLIENT: OnceCell> = OnceCell::const_new(); + +/// What connecting with this agent's own credential needs. +#[derive(Debug, PartialEq, Eq)] +struct AgentPath { + url: String, + ca_file: Option, + /// 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 `.`. + Agent, + /// The hive's shared client, granted `..>` for the hive the + /// client id names. + Hive { client_id: String }, +} + +/// The shared queue connection and the credential it was made with. +#[derive(Clone)] +pub struct Connection { + pub client: async_nats::Client, + pub presented: Presented, +} /// The four variables the harness unit sets, before the client-id file is /// read. Collected into a struct so [`decide`] is pure over them and the @@ -147,6 +177,16 @@ fn decide_agent_secret(path: Option<&str>) -> Option { } } +/// 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) -> Option { + Some(AgentPath { + url: env.nats_url.clone()?, + ca_file: env.ca_file.as_ref().map(Into::into), + secret_file: secret_file?, + }) +} + /// Decide what this agent's queue configuration is, given the environment and /// whatever the client-id file held. /// @@ -216,15 +256,10 @@ pub fn init() { let _ = CONFIG.set(resolved); // Independent of everything above: this agent may hold its own secret on - // a hive with no queue coordinates, or hold the coordinates and no secret - // of its own yet. Reported either way, because "which credential is this - // agent able to present" is a question only this process can answer, and - // it is the one the next slice's rollout will be asked repeatedly. - // - // The answer is only logged here. Presenting it needs the queue's - // auth-callout responder to verify it, which is the next slice — see this - // module's header. - if let Some(path) = decide_agent_secret(env.agent_secret_file.as_deref()) { + // a hive whose shared client is not published yet, or the shared client + // and no secret of its own. + let secret = decide_agent_secret(env.agent_secret_file.as_deref()); + if let Some(path) = &secret { // The path, never the bytes: the file holds the secret itself. tracing::info!( path = %path.display(), @@ -236,6 +271,13 @@ pub fn init() { by its hive's shared client" ); } + let _ = AGENT.set(decide_agent_path(&env, secret)); +} + +/// Whether this agent has any credential to reach the queue with. `false` +/// before [`init`] has run. +pub fn configured() -> bool { + config().is_some() || AGENT.get().is_some_and(Option::is_some) } /// What [`init`] resolved, or `None` when this agent has no queue. @@ -263,14 +305,48 @@ pub fn config() -> Option<&'static QueueConfig> { /// failed" — a caller does nothing differently between them, since either way /// there is nothing to publish onto. Connecting is lazy so that an agent on a /// hive with no queue pays nothing at boot. -pub async fn client() -> Option { +pub async fn client() -> Option { CLIENT.get_or_init(connect_once).await.clone() } -async fn connect_once() -> Option { +/// This agent's own credential first, then the hive's. Each path taken is +/// logged once, here. +async fn connect_once() -> Option { + if let Some(agent) = AGENT.get().and_then(Option::as_ref) { + match connect_as_agent(agent).await { + Ok(client) => { + tracing::info!( + url = %agent.url, + "connected to the swarm queue with this agent's own credential" + ); + return Some(Connection { + client, + presented: Presented::Agent, + }); + } + Err(e) => tracing::warn!( + error = format!("{e:#}"), + "connecting with this agent's own credential failed; falling back to \ + its hive's shared client" + ), + } + } let cfg = config()?; match swarm_queue_client::connect(cfg.clone()).await { - Ok(client) => Some(client), + Ok(client) => { + tracing::info!( + url = %cfg.url, + client_id = %cfg.client_id, + "connecting to the swarm queue with the hive's shared client; the \ + client retries in the background until the queue accepts it" + ); + Some(Connection { + client, + presented: Presented::Hive { + client_id: cfg.client_id.clone(), + }, + }) + } Err(e) => { // `chain`, not `{:#}`: this is `swarm_queue_client::Error`, whose // `Display` ignores the alternate flag, so `{:#}` renders the @@ -284,9 +360,26 @@ async fn connect_once() -> Option { } } +/// 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 { + let name = crate::identity::label(); + let secret = std::fs::read_to_string(&agent.secret_file) + .with_context(|| format!("reading {}", agent.secret_file.display()))?; + let token = swarm_queue_client::agent_token::format_agent_token(&name, secret.trim())?; + swarm_queue_client::connect_with_token(&agent.url, agent.ca_file.as_deref(), token) + .await + .map_err(|e| anyhow::anyhow!(swarm_queue_client::chain(&e))) +} + #[cfg(test)] mod tests { - use super::{QueueEnv, Resolution, decide, decide_agent_secret, read_client_id}; + use std::path::PathBuf; + + use super::{ + AgentPath, QueueEnv, Resolution, decide, decide_agent_path, decide_agent_secret, + read_client_id, + }; fn env(parts: [Option<&str>; 4]) -> QueueEnv { let [nats_url, token_endpoint, client_id_file, client_secret_file] = parts; @@ -435,4 +528,34 @@ mod tests { assert!(matches!(decide(&e, None), Resolution::Absent(_))); assert!(decide_agent_secret(path.to_str()).is_some()); } + + /// The agent's own path needs the queue's address and its own secret, and + /// not the hive's client id: that is what lets it connect on a hive whose + /// shared client is not published. + #[test] + fn an_agent_with_its_own_secret_and_the_queue_address_connects_as_itself() { + let secret = PathBuf::from("/run/queue-identity/secret"); + assert_eq!( + decide_agent_path(&full(), Some(secret.clone())), + Some(AgentPath { + url: "nats://10.42.0.1:4222".to_owned(), + ca_file: None, + secret_file: secret.clone(), + }) + ); + let url_only = env([Some("nats://10.42.0.1:4222"), None, None, None]); + assert!(decide_agent_path(&url_only, Some(secret)).is_some()); + } + + #[test] + fn without_its_own_secret_or_the_queue_address_an_agent_does_not_connect_as_itself() { + assert_eq!(decide_agent_path(&full(), None), None); + assert_eq!( + decide_agent_path( + &env([None, None, None, None]), + Some(PathBuf::from("/run/queue-identity/secret")) + ), + None + ); + } } diff --git a/hive-agent/src/swarm_term.rs b/hive-agent/src/swarm_term.rs index 249e2b57..620b81d4 100644 --- a/hive-agent/src/swarm_term.rs +++ b/hive-agent/src/swarm_term.rs @@ -27,6 +27,7 @@ use tokio::sync::broadcast; use crate::events::BusEvent; +use crate::swarm_queue::{Connection, Presented}; use crate::term_msg::{ClassifyCtx, TermMsg, classify}; /// Subject family carrying agent terminal rows, the swarm-wide agreement this @@ -112,37 +113,50 @@ fn serialized_len(msg: &TermMsg) -> Option { serde_json::to_vec(msg).ok().map(|v| v.len()) } -/// Start the publish task, if this agent has both queue coordinates and an -/// identity the subject can be derived from. +/// The subject this agent's rows go to under the credential it connected +/// with, or `None` when a hive client id names no hive. +fn subject(presented: &Presented, agent: &str) -> Option { + match presented { + Presented::Agent => Some(format!("{SUBJECT_PREFIX}.{agent}")), + Presented::Hive { client_id } => { + let Some(hive) = hive_from_client_id(client_id) else { + tracing::warn!( + %client_id, + expected = format!("{CLIENT_ID_PREFIX}{CLIENT_ID_SUFFIX}"), + "queue client id does not name a hive; not publishing the terminal upward" + ); + 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, an unparseable -/// client id, no label — each of which is a legal state for an agent rather -/// than an error, and each logged once here rather than per row. +/// 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) { - let Some(cfg) = crate::swarm_queue::config() else { + 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 Some(hive) = hive_from_client_id(&cfg.client_id) else { - tracing::warn!( - client_id = %cfg.client_id, - expected = format!("{CLIENT_ID_PREFIX}{CLIENT_ID_SUFFIX}"), - "queue client id does not name a hive; not publishing the terminal upward" - ); - return; - }; + } let agent = crate::identity::label(); if agent.is_empty() { tracing::warn!("this agent has no label; not publishing the terminal upward"); return; } - let subject = format!("{SUBJECT_PREFIX}.{hive}.{agent}"); - tokio::spawn(run(rx, subject)); + tokio::spawn(run(rx, agent)); } -async fn run(mut rx: broadcast::Receiver, subject: String) { - let Some(client) = crate::swarm_queue::client().await else { +async fn run(mut rx: broadcast::Receiver, agent: String) { + let Some(Connection { client, presented }) = crate::swarm_queue::client().await else { + return; + }; + let Some(subject) = subject(&presented, &agent) else { return; }; tracing::info!(subject, "publishing the agent terminal to the swarm queue"); @@ -203,7 +217,7 @@ async fn publish(client: &async_nats::Client, subject: &str, msg: TermMsg) { #[cfg(test)] mod tests { - use super::{DROPPED_BODY, fit, hive_from_client_id}; + use super::{DROPPED_BODY, Presented, fit, hive_from_client_id, subject}; use crate::events::LiveEvent; use crate::term_msg::{BodyFormat, ClassifyCtx, Level, TermMsg, classify}; @@ -243,6 +257,27 @@ mod tests { } } + /// Each credential publishes on the subject it is granted: its own under + /// the agent's credential, its hive's under the shared one. + #[test] + fn the_subject_follows_the_credential_the_agent_connected_with() { + assert_eq!( + subject(&Presented::Agent, "mara").as_deref(), + Some("$SWARM.term.mara") + ); + let hive = Presented::Hive { + client_id: "hive-alpha-agent".to_owned(), + }; + assert_eq!( + subject(&hive, "mara").as_deref(), + Some("$SWARM.term.alpha.mara") + ); + let unparseable = Presented::Hive { + client_id: "hive-alpha".to_owned(), + }; + assert_eq!(subject(&unparseable, "mara"), None); + } + #[test] fn a_row_that_already_fits_is_published_unchanged() { let msg = TermMsg::new(Level::Info, "turn ok") diff --git a/nix/host-modules/default.nix b/nix/host-modules/default.nix index 4bcb348d..a02f8a82 100644 --- a/nix/host-modules/default.nix +++ b/nix/host-modules/default.nix @@ -32,6 +32,7 @@ ./glue-grafana-oidc-client.nix ./glue-matrix-bao-token.nix ./glue-matrix-ctl-bao-identity.nix + ./glue-nats-auth-bao-identity.nix ./glue-nats-bao-identity.nix ./glue-queue-agent-credential.nix ./glue-secret-publisher-bao-identity.nix diff --git a/nix/host-modules/glue-bao-tls.nix b/nix/host-modules/glue-bao-tls.nix index d410d9f9..179e22fc 100644 --- a/nix/host-modules/glue-bao-tls.nix +++ b/nix/host-modules/glue-bao-tls.nix @@ -260,6 +260,12 @@ in # leaf logs in with. Stays on this host; see its default above. [ -s ${pkiDir}/granter.pem ] || ${signLeaf} ${pkiDir} granter \ ${lib.escapeShellArg deployCfg.bao.granterCommonName} "" clientAuth + + # The queue's auth-callout responder, which reads agent queue + # credentials. Minted here for the matrix-ctl leaf's reason: the queue + # is a swarm singleton, so elsewhere this is the file an operator copies. + [ -s ${pkiDir}/nats-auth.pem ] || ${signLeaf} ${pkiDir} nats-auth \ + ${lib.escapeShellArg deployCfg.bao.natsAuthCommonName} "" clientAuth ''; }; }; diff --git a/nix/host-modules/glue-nats-auth-bao-identity.nix b/nix/host-modules/glue-nats-auth-bao-identity.nix new file mode 100644 index 00000000..7c6cd7cb --- /dev/null +++ b/nix/host-modules/glue-nats-auth-bao-identity.nix @@ -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"; + }; + }; +} diff --git a/nix/host-modules/swarm-bao.nix b/nix/host-modules/swarm-bao.nix index 988520b0..bd2ee06e 100644 --- a/nix/host-modules/swarm-bao.nix +++ b/nix/host-modules/swarm-bao.nix @@ -754,6 +754,18 @@ let } ]; + # The queue's auth-callout responder, which checks the secret an agent + # presents against the one stored for it. One leaf under every agent: `+` is + # one path segment, where `*` globs only at the end and would reach every + # other credential an agent holds. + natsAuthReaders = [ + { + name = "swarm-nats-auth"; + cn = baoDeploy.natsAuthCommonName; + policyText = readStanza "${credentialMountPath}/data/swarm/agents/+/queue"; + } + ]; + # The role name IS the policy name, as for the three service principals # above: the role attaches the policy by spelling it identically, and one # string for both objects removes the way they drift apart. @@ -1438,6 +1450,24 @@ in ''; }; + natsAuthCommonName = lib.mkOption { + type = lib.types.str; + default = "swarm-nats-auth"; + description = '' + Subject the store's `swarm-nats-auth` cert-auth role accepts: the + identity the queue's auth-callout responder presents to read agent + queue credentials, `swarm/agents//queue` and nothing else under + an agent. + + Its own principal rather than + {option}`services.hyperhive.deploy.bao.natsCommonName`: that one + issues the queue's TLS leaf from a host unit, this one reads secrets + from inside the queue's container. + + ⚠️ Reserved as a hive name by ./swarm.nix, like its siblings. + ''; + }; + matrixCtlHiveName = lib.mkOption { type = lib.types.str; default = toString hyperhiveCfg.hiveName; @@ -2021,6 +2051,7 @@ in "swarm-bao-services-issuer-policy" "swarm-bao-nats-tls-policy" "swarm-bao-agent-pki" + "swarm-bao-nats-auth-policy" ]; # 🚫 No `swarm.otel.scrapeTargets.bao` entry any more, and its absence is @@ -2723,6 +2754,8 @@ in systemd.services.swarm-bao-forwarder-oidc-policy = readerPolicyUnit "write the store forwarder's OIDC-secret-reader bao policy and cert-auth role" forwarderOidcReaders; + systemd.services.swarm-bao-nats-auth-policy = readerPolicyUnit "write the queue responder's agent-credential-reader bao policy and cert-auth role" natsAuthReaders; + # A FOURTH sibling, same shape and same reasons as the two above. This # one is what turns `swarm-services-issuer` from a declaration into a # grant: a bao policy reaches nothing until a login role hands it to a diff --git a/nix/host-modules/swarm-nats.nix b/nix/host-modules/swarm-nats.nix index 91cde3b3..fe1bff9b 100644 --- a/nix/host-modules/swarm-nats.nix +++ b/nix/host-modules/swarm-nats.nix @@ -219,6 +219,25 @@ let tlsKeyCredential = "tls-key"; tlsKeyCredentialPath = "/run/credentials/nats.service/${tlsKeyCredential}"; + # The responder's own store identity, for reading agent queue credentials. + # Without it the responder denies every agent token; agents then fall back to + # their hive's OIDC client. + authStoreActive = + deployCfg.nats.authBaoClientCertFile != null && deployCfg.nats.authBaoClientKeyFile != null; + # The role ./swarm-bao.nix writes on the store's `cert` mount, by the same name. + authCertRole = "swarm-nats-auth"; + # Store identity files the copy unit delivers, as credential id → host source. + authStoreFiles = lib.optionalAttrs authStoreActive ( + { + "bao-client.pem" = deployCfg.nats.authBaoClientCertFile; + "bao-client-key.pem" = deployCfg.nats.authBaoClientKeyFile; + } + // lib.optionalAttrs (baoDeploy.serverCaFile != null) { + "bao-ca.pem" = baoDeploy.serverCaFile; + } + ); + authCredential = id: "/run/credentials/swarm-nats-auth.service/${id}"; + # Re-issue once the leaf is past half of the 720h `pki/roles/swarm-nats` # grants it (./swarm-bao.nix), so the daily timer below has two weeks of # retries before it lapses. @@ -498,6 +517,34 @@ in A path, never a value. ''; }; + + authBaoClientCertFile = lib.mkOption { + type = lib.types.nullOr lib.types.str; + default = null; + example = "/var/lib/swarm-bao-pki/nats-auth.pem"; + description = '' + Client certificate the auth-callout responder presents to the swarm's + secret store to read agent queue credentials. Its subject must be + {option}`services.hyperhive.deploy.bao.natsAuthCommonName`. + + No default. ./glue-nats-auth-bao-identity.nix points it at the leaf + ./glue-bao-tls.nix mints, where this host mints one. Unset, or set to + a file that does not exist, the responder denies every agent token and + admits OIDC clients as before. + + A path, never a value. + ''; + }; + + authBaoClientKeyFile = lib.mkOption { + type = lib.types.nullOr lib.types.str; + default = null; + example = "/var/lib/swarm-bao-pki/nats-auth-key.pem"; + description = '' + Private key for {option}`services.hyperhive.deploy.nats.authBaoClientCertFile`. + A path, never a value. + ''; + }; }; config = lib.mkIf deployCfg.nats.enable { @@ -905,6 +952,11 @@ in # an append-only row stream, a header is one current value # republished on change. "--agent-publish-subject ${lib.escapeShellArg "\$\$SWARM.agent-state.{hive}.>"}" + # What an agent that proved its own credential may publish to: + # the same two streams, keyed on the agent alone. + "--agent-token-publish-subject ${lib.escapeShellArg "\$\$SWARM.term.{agent}"}" + "--agent-token-publish-subject ${lib.escapeShellArg "\$\$SWARM.agent-state.{agent}"}" + "--store-cert-role ${lib.escapeShellArg authCertRole}" ]; # Every credential arrives by `LoadCredential` and is named # on the command line only as a **path** — `argv` is @@ -914,12 +966,25 @@ in "callout-user.seed:${inContainer "callout-user.seed"}" "issuer.seed:${inContainer "issuer.seed"}" "oidc-client.secret:${inContainer "oidc-client.secret"}" - ]; + ] + ++ lib.mapAttrsToList (id: _: "${id}:${inContainer id}") authStoreFiles; DynamicUser = true; Restart = "on-failure"; RestartSec = "5s"; SyslogIdentifier = "swarm-nats-auth"; }; + # The store the responder reads agent credentials from. Absent + # without an identity, which the responder reads as "no store". + environment = lib.optionalAttrs authStoreActive ( + { + BAO_ADDR = "https://${baoCfg.domain}:${toString baoCfg.port}"; + BAO_CLIENT_CERT = authCredential "bao-client.pem"; + BAO_CLIENT_KEY = authCredential "bao-client-key.pem"; + } + // lib.optionalAttrs (authStoreFiles ? "bao-ca.pem") { + BAO_CACERT = authCredential "bao-ca.pem"; + } + ); }; # The server binary, so an operator with a shell in here can @@ -957,7 +1022,10 @@ in # being up says nothing about whether its in-container secrets unit # has finished. The wait in the script is what actually closes it; # this only stops us spinning for the full timeout on every boot. - ++ lib.optional deployCfg.authelia.enable "container@${autheliaCfg.machine}.service"; + ++ lib.optional deployCfg.authelia.enable "container@${autheliaCfg.machine}.service" + # Where the store's PKI is minted on this host, the responder's leaf is + # one of its files. Ordering only: elsewhere the unit does not exist. + ++ lib.optional authStoreActive "swarm-bao-pki.service"; requires = lib.optional deployCfg.nats.autoGenerateCallout "swarm-nats-callout-keys.service"; serviceConfig = { Type = "oneshot"; @@ -1006,7 +1074,22 @@ in install -m 0400 "$secret" \ ${lib.escapeShellArg (hostPath "oidc-client.secret")} - ''; + '' + # The store identity is optional to the responder, so an absent source + # is delivered as an empty file rather than failing this unit: a missing + # `LoadCredential` source stops the responder starting, and a responder + # that does not start denies every client. An empty identity fails only + # the store login, which denies agent tokens. + + lib.concatStrings ( + lib.mapAttrsToList (id: source: '' + if [ -s ${lib.escapeShellArg source} ]; then + install -m 0400 ${lib.escapeShellArg source} ${lib.escapeShellArg (hostPath id)} + else + echo "${source} is absent; the queue responder denies agent tokens until it exists" >&2 + install -m 0400 /dev/null ${lib.escapeShellArg (hostPath id)} + fi + '') authStoreFiles + ); }; # ⚠️ Minted on the HOST, not in the container, because the responder is diff --git a/nix/host-modules/swarm.nix b/nix/host-modules/swarm.nix index b7ac396e..86d90dd5 100644 --- a/nix/host-modules/swarm.nix +++ b/nix/host-modules/swarm.nix @@ -53,6 +53,7 @@ let deployCfg.bao.servicesIssuerCommonName deployCfg.bao.natsCommonName deployCfg.bao.granterCommonName + deployCfg.bao.natsAuthCommonName ] # The two per-hive readers' subjects, spelled out per hive rather than as the # prefix. The prefix alone would reserve the wrong string: the role for hive diff --git a/nix/module-eval/bao-grants.nix b/nix/module-eval/bao-grants.nix index 4f175725..a32e0348 100644 --- a/nix/module-eval/bao-grants.nix +++ b/nix/module-eval/bao-grants.nix @@ -151,7 +151,7 @@ let _: u: (u.environment.BAO_CLIENT_CERT or null) == granterCertFile ) baoGrantWithConsumers.systemd.services; - # The eleven units that write a `swarm-*` grant, by name, for the discovery + # The twelve units that write a `swarm-*` grant, by name, for the discovery # control below. grantingUnitNames = [ "swarm-bao-controller-policy" @@ -165,6 +165,7 @@ let "swarm-bao-services-issuer-policy" "swarm-bao-nats-tls-policy" "swarm-bao-agent-pki" + "swarm-bao-nats-auth-policy" ]; # Comment lines dropped first: both the HCL and the scripts explain @@ -561,6 +562,33 @@ let && !(lib.hasInfix "swarm-grafana" s) && !(lib.hasInfix "sys/policies/acl" s); } + { + # `+` is one path segment, so this reaches `swarm/agents//queue` + # and no other credential an agent holds; `swarm/agents/*` would reach + # all of them. + name = "the queue responder's grant is every agent's queue credential and nothing else"; + ok = + let + s = baoGrantHere.systemd.services.swarm-bao-nats-auth-policy.script; + in + lib.hasInfix "path \"secret/data/swarm/agents/+/queue\" {" s + && lib.hasInfix "capabilities = [\"read\"]" s + && lib.length (lib.filter lib.isList (builtins.split "path \"" s)) == 1 + && !(lib.hasInfix "secret/data/swarm/agents/*" s) + && !(lib.hasInfix "secret/data/swarm/hives" s) + && !(lib.hasInfix "secret/data/swarm/services" s) + && lib.hasInfix "auth/cert/certs/swarm-nats-auth" s + && lib.hasInfix "allowed_common_names=swarm-nats-auth" s + && lib.hasInfix "token_policies=swarm-nats-auth" s; + } + { + name = "the PKI unit signs the queue responder's leaf under its own subject"; + ok = + let + s = baoGrantHere.systemd.services.swarm-bao-pki.script; + in + lib.hasInfix "/nats-auth.pem ]" s && lib.hasInfix "swarm-nats-auth \"\" clientAuth" s; + } { # 🩸 The half that makes the policies above bind: a policy grants only # through a token that carries it, and a token is minted by a cert-auth @@ -748,24 +776,24 @@ let } { # A store host without the granter's pair writes its grants some other - # way, so none of the eleven units may exist. Without this arm + # way, so none of the twelve units may exist. Without this arm # `lib.mkIf haveGranter` could be dropped from any of them and every other # case here would still pass. - name = "without the granter's pair none of the eleven granting units render"; + name = "without the granter's pair none of the twelve granting units render"; ok = let s = baoGranterOptOut.systemd.services; in lib.all (unit: !(s ? ${unit})) (grantingUnitNames ++ [ "swarm-bao-granter-role" ]) - # The control: the same store with the pair renders all eleven. + # The control: the same store with the pair renders all twelve. && lib.all (unit: baoGrantHere.systemd.services ? ${unit}) grantingUnitNames; } { - # 🩸 What replaced the silent skip. With no bootstrap token the eleven still + # 🩸 What replaced the silent skip. With no bootstrap token the twelve still # render, and a refused granter fails them with the step that fixes it. # A store host that never named a token is told to name one, since the # unit that sets the granter up renders only where it has. - name = "a store host without a bootstrap token renders the eleven, each failing loudly with the one-time step"; + name = "a store host without a bootstrap token renders the twelve, each failing loudly with the one-time step"; ok = let s = baoGranterNoToken.systemd.services; @@ -1129,8 +1157,8 @@ let } { # What makes the case above mean something: discovery by the granter's - # certificate reaches all eleven units, and each yields calls. - name = "the granter-policy check sees all eleven granting units, and parses calls from each"; + # certificate reaches all twelve units, and each yields calls. + name = "the granter-policy check sees all twelve granting units, and parses calls from each"; ok = lib.sort lib.lessThan (lib.attrNames granterUnits) == lib.sort lib.lessThan grantingUnitNames && lib.all (u: baoCalls u.script != [ ]) (lib.attrValues granterUnits) diff --git a/nix/module-eval/nats-authelia.nix b/nix/module-eval/nats-authelia.nix index cd72aa1b..cce9c1aa 100644 --- a/nix/module-eval/nats-authelia.nix +++ b/nix/module-eval/nats-authelia.nix @@ -60,6 +60,28 @@ let swarm.nats.calloutIssuerSeedFile = "/run/secrets/nats-issuer.seed"; }; + # The queue with a store identity for its responder, placed by hand: the + # queue host that is not the store's. + natsWithStoreIdentity = hive { + deploy.nats.enable = true; + deploy.nats.autoGenerateCallout = false; + deploy.nats.calloutUserPublicKey = "UTESTUSERPUBKEYAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAA"; + deploy.nats.calloutIssuerPublicKey = "ATESTISSUERPUBKEYAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAAA"; + deploy.nats.calloutUserSeedFile = "/run/secrets/nats-user.seed"; + deploy.nats.calloutIssuerSeedFile = "/run/secrets/nats-issuer.seed"; + deploy.nats.authPackage = pkgs.emptyDirectory; + deploy.nats.authBaoClientCertFile = "/etc/pki/nats-auth.pem"; + deploy.nats.authBaoClientKeyFile = "/etc/pki/nats-auth-key.pem"; + }; + + # The queue on the store's own host, where the PKI glue mints every leaf. + natsOnStoreHost = hive { + deploy.bao.enable = true; + deploy.nats.enable = true; + }; + + responderOf = h: h.containers.swarm-nats.config.systemd.services.swarm-nats-auth; + # A hive running NOTHING of the swarm's own services — no IdP here, no # `swarm.authelia.url` set by hand. The whole point of the fixture is what it # does *not* say: it is the shape whose IdP address used to be null, and @@ -117,6 +139,73 @@ let in lib.hasInfix "--agent-publish-subject '$$SWARM.agent-state.{hive}.>'" exec; } + { + # The per-agent grant: the same two streams keyed on the agent alone, + # with the same doubled dollar, and the role the store writes for it. + name = "the responder grants a verified agent its own hive-free subjects"; + ok = + let + exec = (responderOf natsOldPath).serviceConfig.ExecStart; + in + lib.hasInfix "--agent-token-publish-subject '$$SWARM.term.{agent}'" exec + && lib.hasInfix "--agent-token-publish-subject '$$SWARM.agent-state.{agent}'" exec + && lib.hasInfix "--store-cert-role swarm-nats-auth" exec; + } + { + # Each credential by `LoadCredential`, and the environment naming where + # the unit sees it: a DynamicUser cannot read the copies directly. + name = "a responder with a store identity loads it and is pointed at the store"; + ok = + let + u = responderOf natsWithStoreIdentity; + creds = u.serviceConfig.LoadCredential; + in + lib.elem "bao-client.pem:/var/lib/swarm-nats-auth/bao-client.pem" creds + && lib.elem "bao-client-key.pem:/var/lib/swarm-nats-auth/bao-client-key.pem" creds + && u.environment.BAO_CLIENT_CERT == "/run/credentials/swarm-nats-auth.service/bao-client.pem" + && u.environment.BAO_CLIENT_KEY == "/run/credentials/swarm-nats-auth.service/bao-client-key.pem" + && lib.hasPrefix "https://" u.environment.BAO_ADDR; + } + { + # An absent leaf must not stop the responder, which would deny every + # client: it is delivered empty, and the three credentials the responder + # cannot run without are delivered too. + name = "the copy unit delivers the store identity, and an absent one as an empty file"; + ok = + let + s = natsWithStoreIdentity.systemd.services.swarm-nats-auth-secrets.script; + in + lib.hasInfix "install -m 0400 /etc/pki/nats-auth.pem " s + && lib.hasInfix "install -m 0400 /etc/pki/nats-auth-key.pem " s + && lib.hasInfix "install -m 0400 /dev/null " s + && lib.hasInfix "oidc-client.secret" s; + } + { + # The control: no identity, no store wiring, and the responder's + # credentials are exactly the three it cannot run without. + name = "a responder with no store identity is not pointed at a store"; + ok = + let + u = responderOf natsOldPath; + in + !(u.environment ? BAO_ADDR) + && !(u.environment ? BAO_CLIENT_CERT) + && lib.length u.serviceConfig.LoadCredential == 3; + } + { + # On the store's host the glue pairs the responder with a leaf of its + # own, never the queue's TLS-issuing one or the hive's. + name = "on the store's host the responder presents its own leaf"; + ok = + let + n = natsOnStoreHost.services.hyperhive.deploy.nats; + b = natsOnStoreHost.services.hyperhive.deploy.bao; + in + n.authBaoClientCertFile == "/var/lib/swarm-bao-pki/nats-auth.pem" + && n.authBaoClientKeyFile == "/var/lib/swarm-bao-pki/nats-auth-key.pem" + && n.authBaoClientCertFile != n.baoClientCertFile + && n.authBaoClientCertFile != b.clientCertFile; + } { # Not a rename test. `hostClientSecretDir` is `readOnly`, so the fixture # cannot define it; what can break is a reader left pointing at the diff --git a/swarm-controller/src/agent_identity.rs b/swarm-controller/src/agent_identity.rs index 6b240973..1d31e019 100644 --- a/swarm-controller/src/agent_identity.rs +++ b/swarm-controller/src/agent_identity.rs @@ -136,7 +136,7 @@ fn generate_queue_secret() -> Result { /// Anything that stops one of those five steps, with the step named. A /// failure here fails the job node and nothing else — the agent is still /// created, without a store identity. -pub async fn mint_and_verify(agent: &str, hive: &str) -> Result<()> { +pub async fn mint_and_verify(agent: &str) -> Result<()> { let (mount, pki_role) = agent_pki(|k| std::env::var(k).ok())?; let name = policy::agent_object_name(agent)?; let path = mtls::identity_path(agent)?; @@ -161,18 +161,15 @@ pub async fn mint_and_verify(agent: &str, hive: &str) -> Result<()> { .read_optional(&queue_path) .await .with_context(|| format!("checking whether {queue_path} already holds a credential"))?; - // The secret survives a re-run; the principal it names does not get to. - // An object whose `hive` disagrees with the hive this node was invoked - // with would grant its holder subjects on the wrong hive, so it is - // corrected — but by rewriting the two name fields around the *same* - // `value`, which is a correction no live connection notices. + // The secret survives a re-run. An object naming a different agent is + // corrected by rewriting the name around the *same* `value`, which no live + // connection notices. let wanted = queue::AgentCredential { value: match &existing { Some(existing) => existing.value.clone(), None => generate_queue_secret()?, }, agent: agent.to_owned(), - hive: hive.to_owned(), }; if existing.as_ref() == Some(&wanted) { tracing::info!( @@ -187,9 +184,8 @@ pub async fn mint_and_verify(agent: &str, hive: &str) -> Result<()> { .with_context(|| format!("publishing the agent queue credential at {queue_path}"))?; tracing::info!( agent, - hive, %queue_path, - // Never "rotated": the secret is the same one, only the names + // Never "rotated": the secret is the same one, only the name // around it moved. corrected = existing.is_some(), "agent queue credential published" @@ -280,9 +276,9 @@ async fn read_back_as_agent( .read(queue_path) .await .with_context(|| format!("reading {queue_path} back under {role}'s own token"))?; - // The two name fields are compared as well as the secret: they are what - // the verifying end will grant subjects from, so a mismatch here is the - // same class of fault as an unreadable path. + // The agent is compared as well as the secret: the verifying end refuses + // an object naming a different agent than the path, so a mismatch here is + // the same class of fault as an unreadable path. if read_back != *queue_credential { bail!("the store returned a different object at {queue_path} than the one just published"); } diff --git a/swarm-controller/src/agent_state_stream.rs b/swarm-controller/src/agent_state_stream.rs index a9adc67c..7eca9f93 100644 --- a/swarm-controller/src/agent_state_stream.rs +++ b/swarm-controller/src/agent_state_stream.rs @@ -33,8 +33,8 @@ use crate::AppState; /// Subject family carrying agent turn-state headers — must match /// `hive-agent::swarm_agent_state::SUBJECT_PREFIX` exactly, since the two -/// ends never see the constant together. `$SWARM.agent-state..` -/// is the full subject; as with the terminal relay's own copy, there is no +/// ends never see the constant together. `crate::term_stream::agent_subjects` +/// spells the full subjects; as with the terminal relay's own copy, there is no /// way to check the two sides agree short of this comment and the module /// docs on both ends staying honest about it. const SUBJECT_PREFIX: &str = "$SWARM.agent-state"; @@ -117,16 +117,12 @@ pub(crate) async fn stream_agent_state( ) })?; - let subject = format!("{SUBJECT_PREFIX}.{hive}.{agent}"); - let subscriber = client.subscribe(subject.clone()).await.map_err(|e| { - tracing::warn!(%subject, error = %e, "state stream: subscribe failed"); - crate::error_problem( - axum::http::StatusCode::INTERNAL_SERVER_ERROR, - &format!("subscribing to {subject} failed: {e}"), - ) - })?; - - tracing::info!(%subject, "state stream: client attached"); + let subscriber = crate::term_stream::subscribe_all( + &client, + &crate::term_stream::agent_subjects(SUBJECT_PREFIX, &hive, &agent), + "state", + ) + .await?; let stream = subscriber.map(|msg| Ok(Event::default().data(String::from_utf8_lossy(&msg.payload)))); Ok(Sse::new(stream).keep_alive(KeepAlive::default())) diff --git a/swarm-controller/src/main.rs b/swarm-controller/src/main.rs index d20d50df..a256aa1d 100644 --- a/swarm-controller/src/main.rs +++ b/swarm-controller/src/main.rs @@ -100,11 +100,10 @@ enum SwarmNodeKind { /// `agent_identity::mint_and_verify` — including why this node does not /// report success on a write. /// - /// Carries the hive for a different reason than `TriggerDeploy` does: - /// not as an address, but because the agent's queue credential names the - /// hive it may take subjects on, so it cannot be written without knowing - /// which hive the agent belongs to. - MintAgentIdentity { hive: String, agent: String }, + /// Carries no hive: neither the certificate, the queue credential nor + /// the agent's policy names one, so an agent keeps one store identity + /// whichever hive it runs on. + MintAgentIdentity { agent: String }, /// Make sure `agent` holds a live forge access token in the swarm secret /// store, minting one with the forge's admin API when it does not. See /// `forge::agent_token` — including why a rotation is a delete then a @@ -124,8 +123,7 @@ enum SwarmNodeKind { /// Declare `agent` on `hive` as `Paused` in the swarm's wanted-state /// store, so a freshly created agent does not start driving turns the /// moment it's deployed — the operator has to explicitly flip it to `Up`. - /// Carries the hive for the same reason `TriggerDeploy`/`MintAgentIdentity` - /// do: the wanted-state bucket is keyed per hive. + /// Carries the hive because the wanted-state bucket is keyed per hive. SetAgentWanted { hive: String, agent: String }, /// Tell `hive` to rebuild `agent`, by publishing on the swarm's deploy /// subject. The one node kind whose effect leaves this host. @@ -166,12 +164,12 @@ impl hive_jobq_wire::WireNode for SwarmNodeKind { | SwarmNodeKind::CreateForgeUser { agent } | SwarmNodeKind::AddRepoMember { agent } | SwarmNodeKind::InitAgentConfigRepo { agent } + | SwarmNodeKind::MintAgentIdentity { agent } | SwarmNodeKind::MintAgentForgeToken { agent } | SwarmNodeKind::MintAgentMatrixAccount { agent } => { serde_json::json!({ "agent": agent }) } SwarmNodeKind::TriggerDeploy { hive, agent } - | SwarmNodeKind::MintAgentIdentity { hive, agent } | SwarmNodeKind::SetAgentWanted { hive, agent } => { serde_json::json!({ "agent": agent, "hive": hive }) } @@ -314,7 +312,7 @@ async fn run_swarm_node( Err(e) => Outcome::Failed(format!("{e:#}")), }, }, - SwarmNodeKind::MintAgentIdentity { hive, agent } => mint_identity(&agent, &hive).await, + SwarmNodeKind::MintAgentIdentity { agent } => mint_identity(&agent).await, SwarmNodeKind::MintAgentForgeToken { agent } => mint_forge_token(deps.forge, &agent).await, SwarmNodeKind::MintAgentMatrixAccount { agent } => { mint_matrix_account(deps.matrix_homeserver.as_deref(), &agent).await @@ -340,10 +338,10 @@ async fn run_swarm_node( /// The `MintAgentIdentity` arm, lifted out so `run_swarm_node` stays under /// `clippy::too_many_lines`. -async fn mint_identity(agent: &str, hive: &str) -> hive_jobq::scheduler::Outcome { +async fn mint_identity(agent: &str) -> hive_jobq::scheduler::Outcome { use hive_jobq::scheduler::Outcome; - match agent_identity::mint_and_verify(agent, hive).await { + match agent_identity::mint_and_verify(agent).await { Ok(()) => Outcome::Done, Err(e) => Outcome::Failed(format!("{e:#}")), } @@ -1512,7 +1510,6 @@ fn declare_agent_job( // authelia. let mint_identity = b .node(SwarmNodeKind::MintAgentIdentity { - hive: hive.to_owned(), agent: agent.to_owned(), }) .after_ok(create_identity); @@ -1670,20 +1667,6 @@ async fn get_agent_config_pr( Ok(Json(cache.get(&name))) } -/// Body of `POST /api/agents/{name}/identity`. -#[derive(Deserialize, ToSchema)] -struct MintAgentIdentityRequest { - /// The hive this agent belongs to. - /// - /// Required, for the same reason [`CreateAgentRequest`]'s is: the - /// credentials this mints name a hive, and the controller has nowhere to - /// look one up — agents are created on hives at runtime and this daemon - /// keeps no roster of which agent is where. An operator naming the wrong - /// one would hand the agent subjects on a hive it does not run on, so it - /// is asked for rather than guessed at. - hive: String, -} - /// Success body of `POST /api/agents/{name}/identity`. #[derive(Clone, Debug, Serialize, ToSchema)] struct MintAgentIdentityResponse { @@ -1712,10 +1695,9 @@ struct MintAgentIdentityResponse { post, path = "/api/agents/{name}/identity", params(("name" = String, Path, description = "agent name")), - request_body = MintAgentIdentityRequest, responses( (status = 200, description = "mint queued", body = MintAgentIdentityResponse), - (status = 400, description = "`name` or `hive` is not a valid identifier, or `hive` is not in this swarm (problem+json)", body = String), + (status = 400, description = "`name` is not a valid identifier (problem+json)", body = String), (status = 500, description = "the job could not be queued (problem+json)", body = String), ), tag = "agents" @@ -1723,28 +1705,12 @@ struct MintAgentIdentityResponse { async fn mint_agent_identity( State(state): State, Path(name): Path, - Json(req): Json, ) -> Result, problem_details::ProblemDetails> { - // Both names are interpolated into store paths and policy documents - // downstream, so both are validated here as well as there. + // The name is interpolated into store paths and policy documents + // downstream, so it is validated here as well as there. let agent = hive_types::Ident::parse(&name) .map_err(|reason| error_problem(axum::http::StatusCode::BAD_REQUEST, reason))? .into_string(); - let hive = hive_types::Ident::parse(&req.hive) - .map_err(|reason| error_problem(axum::http::StatusCode::BAD_REQUEST, reason))? - .into_string(); - if !state.hives.iter().any(|h| h.name == hive) { - let known: Vec<&str> = state.hives.iter().map(|h| h.name.as_str()).collect(); - let known = if known.is_empty() { - "(none configured)".to_owned() - } else { - known.join(", ") - }; - return Err(error_problem( - axum::http::StatusCode::BAD_REQUEST, - &format!("hive {hive:?} is not in this swarm — known hives: {known}"), - )); - } // No reserved-name or collision warnings here, unlike `create_agent`: // those answer "is this name available", and this route is only ever @@ -1757,7 +1723,6 @@ async fn mint_agent_identity( .insert_job(None, |b| { vec![ b.node(SwarmNodeKind::MintAgentIdentity { - hive: hive.clone(), agent: agent.clone(), }) .guid(), @@ -2803,39 +2768,6 @@ mod tests { assert_eq!(resp.change, super::ForgeAdminChange::AlreadyAdmin); } - /// The backfill route's roster check, asserted by effect for the same - /// reason its sibling above is: a refusal that queued first would still - /// re-mint the agent's certificate, which every running agent on the - /// named hive picks up on its next boot. - #[tokio::test] - async fn a_backfill_for_a_hive_outside_the_roster_queues_nothing() { - let (state, sched) = state_with_roster(); - - let err = super::mint_agent_identity( - axum::extract::State(state), - axum::extract::Path("atlas".to_owned()), - axum::Json(super::MintAgentIdentityRequest { - hive: "pr1maa".to_owned(), - }), - ) - .await - .expect_err("a hive outside the roster must be refused"); - - let rendered = format!("{err:?}"); - assert!( - rendered.contains("pr1ma"), - "the refusal should name the known hives, got: {rendered}" - ); - - let queued = sched - .lock() - .unwrap_or_else(std::sync::PoisonError::into_inner) - .graph() - .nodes() - .count(); - assert_eq!(queued, 0, "a refused backfill must queue no work"); - } - /// A name that is not an identifier is refused before it can reach a /// store path. `create_agent` gets this for free from the roster check /// on the hive; the agent name has no roster to check against, so this @@ -2848,9 +2780,6 @@ mod tests { super::mint_agent_identity( axum::extract::State(state), axum::extract::Path("../beta".to_owned()), - axum::Json(super::MintAgentIdentityRequest { - hive: "pr1ma".to_owned(), - }), ) .await .expect_err("a traversal in the agent name must be refused"); @@ -2876,12 +2805,9 @@ mod tests { let queued = super::mint_agent_identity( axum::extract::State(state), axum::extract::Path("atlas".to_owned()), - axum::Json(super::MintAgentIdentityRequest { - hive: "pr1ma".to_owned(), - }), ) .await - .expect("a hive in the roster must be accepted"); + .expect("a valid agent name must be accepted"); let guard = sched .lock() @@ -2893,10 +2819,9 @@ mod tests { assert!( matches!( &node.payload, - SwarmNodeKind::MintAgentIdentity { hive, agent } - if hive == "pr1ma" && agent == "atlas" + SwarmNodeKind::MintAgentIdentity { agent } if agent == "atlas" ), - "the one node must be the mint, carrying both names: {:?}", + "the one node must be the mint, for the named agent: {:?}", node.payload ); // The id the operator is told to watch has to be the node that was @@ -3025,7 +2950,6 @@ mod tests { let id = sched .append( SwarmNodeKind::MintAgentIdentity { - hive: "pr1ma".to_owned(), agent: "atlas".to_owned(), }, Vec::new(), @@ -3192,25 +3116,6 @@ mod tests { ); } - /// The node carries the hive, and the viewer has to see it. The `data` - /// match is an or-pattern on purpose (see its own comment), and this is - /// the assertion that the new variant joined the two-field arm rather - /// than the agent-only one — a viewer silently missing the hive is the - /// failure that comment describes having already happened once. - #[test] - fn a_mint_node_renders_both_the_agent_and_the_hive() { - use hive_jobq_wire::WireNode as _; - - let kind = SwarmNodeKind::MintAgentIdentity { - hive: "pr1ma".to_owned(), - agent: "atlas".to_owned(), - }; - assert_eq!(kind.label(), "mint_agent_identity"); - let data = kind.data(1); - assert_eq!(data["agent"], "atlas"); - assert_eq!(data["hive"], "pr1ma"); - } - /// The ordering the operator's ruling requires: a hive cannot pass down /// a certificate the swarm has not published, so the deploy message must /// not leave before the mint is terminal. diff --git a/swarm-controller/src/term_stream.rs b/swarm-controller/src/term_stream.rs index 628f364f..3cd6a878 100644 --- a/swarm-controller/src/term_stream.rs +++ b/swarm-controller/src/term_stream.rs @@ -9,7 +9,8 @@ //! //! **Live tail only, on purpose.** `hive-agent` publishes each //! already-classified `TermMsg` row to the core subject -//! `$SWARM.term.{hive}.{agent}` (see `hive-agent::swarm_term`'s module +//! `$SWARM.term.{agent}` (`$SWARM.term.{hive}.{agent}` from an agent on its +//! hive's shared credential; see `hive-agent::swarm_term`'s module //! doc) — a core subject, not `JetStream`, so a subscriber who was not //! listening missed the row, same as on the agent's own local SSE //! stream. This handler relays exactly that: no replay, no last-N @@ -34,8 +35,8 @@ use crate::AppState; /// Subject family carrying agent terminal rows — must match /// `hive-agent::swarm_term::SUBJECT_PREFIX` exactly, since the two ends -/// never see the constant together. `$SWARM.term..` is the -/// full subject; there is no way to check the two sides agree short of +/// never see the constant together. [`agent_subjects`] spells the full +/// subjects; there is no way to check the two sides agree short of /// this comment and the module docs on both ends staying honest about it. const SUBJECT_PREFIX: &str = "$SWARM.term"; @@ -116,17 +117,63 @@ pub(crate) async fn stream_agent_term( ) })?; - let subject = format!("{SUBJECT_PREFIX}.{hive}.{agent}"); - let subscriber = client.subscribe(subject.clone()).await.map_err(|e| { - tracing::warn!(%subject, error = %e, "term stream: subscribe failed"); - crate::error_problem( - axum::http::StatusCode::INTERNAL_SERVER_ERROR, - &format!("subscribing to {subject} failed: {e}"), - ) - })?; - - tracing::info!(%subject, "term stream: client attached"); + 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 +/// `.`, and `..`, 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, problem_details::ProblemDetails> +{ + let mut subscribers = Vec::with_capacity(subjects.len()); + for subject in subjects { + let subscriber = client.subscribe(subject.clone()).await.map_err(|e| { + tracing::warn!(%subject, error = %e, "{what} stream: subscribe failed"); + crate::error_problem( + axum::http::StatusCode::INTERNAL_SERVER_ERROR, + &format!("subscribing to {subject} failed: {e}"), + ) + })?; + subscribers.push(subscriber); + } + tracing::info!(?subjects, "{what} stream: client attached"); + Ok(futures_util::stream::select_all(subscribers)) +} + +#[cfg(test)] +mod tests { + use super::*; + + /// Both shapes, so an agent on either credential is streamed. + #[test] + fn a_relay_subscribes_to_the_agents_own_subject_and_its_hive_scoped_one() { + assert_eq!( + agent_subjects(SUBJECT_PREFIX, "h1", "atlas"), + [ + "$SWARM.term.atlas".to_owned(), + "$SWARM.term.h1.atlas".to_owned(), + ] + ); + } +} diff --git a/swarm-nats-auth/Cargo.toml b/swarm-nats-auth/Cargo.toml index 20a52c4f..a921fb20 100644 --- a/swarm-nats-auth/Cargo.toml +++ b/swarm-nats-auth/Cargo.toml @@ -22,12 +22,17 @@ serde.workspace = true serde_json.workspace = true # The jti digest: base32hex(sha256(claims)) over every JWT this crate signs. sha2.workspace = true +# Comparing an agent's presented secret with the stored one. +subtle.workspace = true # For `status::BUCKET` and `notices::STREAM` - the subjects a hive may # publish to are derived from these names, and every end that touches them # must agree on the same one. Deliberately WITHOUT the `kv` feature: this # crate derives subject strings, it never opens the bucket. `notices` # is name-only too (no `jetstream`/`kv` surface), same reason. swarm-queue-client = { workspace = true, features = ["notices"] } +# The agent-token spelling the agent also uses, and the read of the stored +# credential it is checked against. +swarm-secret-client.workspace = true tokio.workspace = true tracing.workspace = true tracing-subscriber.workspace = true diff --git a/swarm-nats-auth/src/agent_token.rs b/swarm-nats-auth/src/agent_token.rs new file mode 100644 index 00000000..4a929df4 --- /dev/null +++ b/swarm-nats-auth/src/agent_token.rs @@ -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//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>> + 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 { + 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> { + 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( + policy: &Policy, + source: Option<&S>, + token: &AgentToken<'_>, +) -> Option { + 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> { + 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 { + 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(_) + )); + } +} diff --git a/swarm-nats-auth/src/main.rs b/swarm-nats-auth/src/main.rs index bfee3d5e..d88a88ed 100644 --- a/swarm-nats-auth/src/main.rs +++ b/swarm-nats-auth/src/main.rs @@ -8,7 +8,8 @@ //! It connects as the one callout-exempt user (by nkey, never by name — the //! server refuses to start if that entry carries a username), subscribes to //! `$SYS.REQ.USER.AUTH`, validates the presented bearer token against -//! authelia's introspection endpoint, and answers with a NATS user JWT signed +//! authelia's introspection endpoint (or, for an agent's own token, against the +//! secret store: see `agent_token`), and answers with a NATS user JWT signed //! by the account key. A rejection is answered explicitly: silence is //! indistinguishable from the responder being down, and the queue is the //! swarm's control path. @@ -27,6 +28,7 @@ use anyhow::Context; use clap::Parser; use futures_util::StreamExt; +mod agent_token; mod introspect; mod policy; mod request; @@ -89,7 +91,7 @@ struct Args { /// config change plus a reload. /// /// So this identity says which hive an agent belongs to and never which - /// agent: two agents on one hive are indistinguishable to this responder. + /// agent: two agents on one hive are indistinguishable under it. /// /// A suffix on the hive's id rather than a prefix of its own, because /// `agent-` reads as *the agent called ``* — the one thing @@ -127,6 +129,18 @@ struct Args { /// is refused, which is loud rather than silently over-broad. #[arg(long = "agent-publish-subject")] agent_publish_subjects: Vec, + + /// 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, + + /// Role on the secret store's `cert` auth mount this responder logs in + /// as to read agent credentials. The store is found through the `BAO_*` + /// environment; with none, every agent token is denied. + #[arg(long, default_value = "swarm-nats-auth")] + store_cert_role: String, } /// Read a secret file and strip surrounding whitespace. @@ -177,7 +191,9 @@ async fn main() -> anyhow::Result<()> { args.reader_clients.clone(), args.hive_publish_subjects.clone(), args.agent_publish_subjects.clone(), - )?; + )? + .with_agent_token_subjects(args.agent_token_publish_subjects.clone())?; + let store = agent_token::Store::from_env(args.store_cert_role.clone()); let http = reqwest::Client::new(); let issuer = nkeys::KeyPair::from_seed(&read_secret(&args.issuer_seed_file)?) .context("parse the account signing seed")?; @@ -210,44 +226,33 @@ async fn main() -> anyhow::Result<()> { }; // No token is a denial, not an error: an anonymous connect is a // normal thing for a client to attempt and an abnormal thing to - // grant. Introspection is only reached once something was presented. - // - // The caller is an identity or nothing — see `introspect`'s module - // docs. There is no "admitted, identity unknown" branch to write here - // because there is no such value to receive. - let caller = match &req.connect_opts.auth_token { - Some(token) => introspect::identify_caller( - &http, - &args.introspection_url, - &args.client_id, - &client_secret, - token, - ) - .await - // An introspection that could not be *made* is a denial too. The - // failure modes of an HTTP call are exactly the conditions under - // which an attacker would most like this to fall open. - .unwrap_or_else(|e| { - tracing::warn!(error = ?e, "introspection failed; denying"); - None - }), - None => None, + // grant. Neither the store nor introspection is reached until + // something was presented. + let (caller, permissions) = match req + .connect_opts + .auth_token + .as_deref() + .map(agent_token::classify) + { + // The caller is the name the token claims, logged whether or not + // the secret proved it; `granted` says which. + Some(agent_token::Presented::Agent(token)) => ( + Some(format!("agent:{}", token.agent)), + agent_token::authorize(&policy, store.as_ref(), &token).await, + ), + Some(agent_token::Presented::Malformed(e)) => { + tracing::warn!(error = %e, "malformed agent token; denying"); + (None, None) + } + Some(agent_token::Presented::Bearer(token)) => { + introspected(&policy, &http, &args, &client_secret, token).await + } + None => (None, None), }; - // Admission said who; the policy says what. A caller the `IdP` - // vouches for but no rule matches is denied — see `policy`'s module - // docs for why that is deny and not "connect with nothing". - let permissions = caller.as_deref().and_then(|id| policy.permissions(id)); - if let (Some(id), None) = (caller.as_deref(), permissions.as_ref()) { - // Loud, and the one case an operator has to be able to find: a - // valid credential refused by our own policy. The alternative is - // a client that authenticates fine and mysteriously cannot work. - tracing::warn!( - caller = %id, - "authenticated client matches no policy rule; denying" - ); - } - // The client id is an identifier, not a credential, and it is the - // only thing tying a connection in this log to a hive. + // The caller is an identifier, not a credential: a client id, or + // `agent:` for an agent token. One line per auth request, a connect + // or a reconnect, so `hive--agent` lines count requests made on a + // hive's shared credential, not agents. tracing::info!( user_nkey = %req.user_nkey, server_id = %req.server_id.id, @@ -281,6 +286,50 @@ async fn main() -> anyhow::Result<()> { anyhow::bail!("subscription to {AUTH_SUBJECT} ended") } +/// The OIDC path: who the `IdP` says presented `token`, and what the policy +/// grants that client. `None` permissions is a denial. +/// +/// The caller is an identity or nothing — see `introspect`'s module docs. +/// There is no "admitted, identity unknown" branch to write here because +/// there is no such value to receive. +async fn introspected( + policy: &policy::Policy, + http: &reqwest::Client, + args: &Args, + client_secret: &str, + token: &str, +) -> (Option, Option) { + let caller = introspect::identify_caller( + http, + &args.introspection_url, + &args.client_id, + client_secret, + token, + ) + .await + // An introspection that could not be *made* is a denial too. The + // failure modes of an HTTP call are exactly the conditions under + // which an attacker would most like this to fall open. + .unwrap_or_else(|e| { + tracing::warn!(error = ?e, "introspection failed; denying"); + None + }); + // Admission said who; the policy says what. A caller the `IdP` + // vouches for but no rule matches is denied — see `policy`'s module + // docs for why that is deny and not "connect with nothing". + let permissions = caller.as_deref().and_then(|id| policy.permissions(id)); + if let (Some(id), None) = (caller.as_deref(), permissions.as_ref()) { + // Loud, and the one case an operator has to be able to find: a + // valid credential refused by our own policy. The alternative is + // a client that authenticates fine and mysteriously cannot work. + tracing::warn!( + caller = %id, + "authenticated client matches no policy rule; denying" + ); + } + (caller, permissions) +} + #[cfg(test)] mod tests { use super::*; @@ -306,9 +355,16 @@ mod tests { "hive-", "--agent-client-suffix", "-agent", + "--agent-token-publish-subject", + "$SWARM.term.{agent}", + "--agent-token-publish-subject", + "$SWARM.agent-state.{agent}", + "--store-cert-role", + "swarm-nats-auth", ]) .expect("the unit's own argument vector must parse"); assert_eq!(args.agent_client_suffix, "-agent"); + assert_eq!(args.agent_token_publish_subjects.len(), 2); } /// The control for the case above: an ordinary value parses through the diff --git a/swarm-nats-auth/src/policy.rs b/swarm-nats-auth/src/policy.rs index 0edd6536..3da625eb 100644 --- a/swarm-nats-auth/src/policy.rs +++ b/swarm-nats-auth/src/policy.rs @@ -45,12 +45,16 @@ pub struct Policy { readers: Vec, extra_hive_subjects: Vec, extra_agent_subjects: Vec, + agent_token_subjects: Vec, } /// Placeholder replaced with the hive's own name in `extra_hive_subjects` and /// `extra_agent_subjects`. const HIVE_PLACEHOLDER: &str = "{hive}"; +/// Placeholder replaced with the agent's own name in `agent_token_subjects`. +const AGENT_PLACEHOLDER: &str = "{agent}"; + impl Policy { /// `hive_prefix` is the client-id prefix that marks a hive and /// `agent_suffix` what a hive's agent containers carry **on top of** it — @@ -130,9 +134,43 @@ impl Policy { readers, extra_hive_subjects, extra_agent_subjects, + agent_token_subjects: Vec::new(), }) } + /// Set the subjects an agent that proved its own credential may publish + /// to, with `{agent}` standing for its name. + /// + /// # Errors + /// + /// A template with no `{agent}` in it is refused: it would be one subject + /// shared by every agent in the swarm rather than the agent's own. + pub fn with_agent_token_subjects(mut self, subjects: Vec) -> anyhow::Result { + 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 { + let publish: Vec = self + .agent_token_subjects + .iter() + .map(|s| s.replace(AGENT_PLACEHOLDER, agent)) + .collect(); + (!publish.is_empty()).then_some(Permissions { publish }) + } + /// The permissions for `client_id`, or `None` when no rule matches. /// /// `None` is a denial. It is not "grant nothing and let them connect": @@ -1175,4 +1213,35 @@ mod tests { "an empty hive name must never be expanded into a subject: {g:?}" ); } + + #[test] + fn an_agent_token_grant_is_the_agents_own_subjects_and_nothing_else() { + let p = policy_with_agent_subject() + .with_agent_token_subjects(vec![ + "$SWARM.term.{agent}".to_owned(), + "$SWARM.agent-state.{agent}".to_owned(), + ]) + .expect("per-agent templates are valid"); + let g = p.agent_token_permissions("atlas").expect("configured"); + assert_eq!( + g.publish, + vec![ + "$SWARM.term.atlas".to_owned(), + "$SWARM.agent-state.atlas".to_owned(), + ] + ); + } + + #[test] + fn an_agent_token_subject_without_the_placeholder_is_refused() { + let err = policy() + .with_agent_token_subjects(vec!["$SWARM.term.all".to_owned()]) + .expect_err("a subject shared by every agent is not the agent's own"); + assert!(format!("{err}").contains("$SWARM.term.all"), "{err}"); + } + + #[test] + fn with_no_agent_token_subject_configured_an_agent_token_is_refused() { + assert!(policy().agent_token_permissions("atlas").is_none()); + } } diff --git a/swarm-queue-client/src/agent_token.rs b/swarm-queue-client/src/agent_token.rs new file mode 100644 index 00000000..01a6a18b --- /dev/null +++ b/swarm-queue-client/src/agent_token.rs @@ -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//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 `.` +/// 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 `.` 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 { + 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 `.` after it. +#[must_use] +pub fn parse_agent_token(token: &str) -> Option, 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 `.`. + #[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}"); + } +} diff --git a/swarm-queue-client/src/lib.rs b/swarm-queue-client/src/lib.rs index 077eca8b..8676a15c 100644 --- a/swarm-queue-client/src/lib.rs +++ b/swarm-queue-client/src/lib.rs @@ -178,6 +178,10 @@ pub mod wanted; /// inside it. See the module doc for why the two must not merge. pub mod agent_status; +/// The token an agent presents its own queue secret in. Store-free, so an agent +/// formats it without linking the secret-store client. +pub mod agent_token; + /// The subject the swarm controller publishes on when the hive-wide knowledge /// repository has changed. One writer, many readers — every hive subscribes. /// @@ -316,6 +320,15 @@ const TOKEN_REFRESH_SKEW: std::time::Duration = std::time::Duration::from_mins(2 /// degrades every other client of it. const MAX_RECONNECT_DELAY: std::time::Duration = std::time::Duration::from_mins(1); +/// Exponential from 500ms, capped at [`MAX_RECONNECT_DELAY`]. +fn reconnect_delay(attempts: usize) -> std::time::Duration { + let exp = u32::try_from(attempts.saturating_sub(1)).unwrap_or(u32::MAX); + std::cmp::min( + std::time::Duration::from_millis(500u64.saturating_mul(2u64.saturating_pow(exp.min(8)))), + MAX_RECONNECT_DELAY, + ) +} + /// Where the controller finds the queue and what it authenticates with. /// /// Every field comes from an environment variable the NixOS module sets, the @@ -738,15 +751,7 @@ pub async fn connect(cfg: QueueConfig) -> Result { // which turns an unreachable queue into a permanent 4s poll — and, before // the cache above, a permanent 4s token-request loop against authelia. // Exponential from 500ms so a momentary blip still reconnects promptly. - .reconnect_delay_callback(|attempts| { - let exp = u32::try_from(attempts.saturating_sub(1)).unwrap_or(u32::MAX); - std::cmp::min( - std::time::Duration::from_millis( - 500u64.saturating_mul(2u64.saturating_pow(exp.min(8))), - ), - MAX_RECONNECT_DELAY, - ) - }) + .reconnect_delay_callback(reconnect_delay) // The controller and the queue are separate units on (possibly) // separate hosts, and nothing orders them. Without this, a queue that // comes up one second later leaves the controller permanently @@ -780,6 +785,29 @@ pub async fn connect(cfg: QueueConfig) -> Result { 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 { + let mut options = + async_nats::ConnectOptions::with_token(token).reconnect_delay_callback(reconnect_delay); + if let Some(path) = ca_file { + options = options.add_root_certificates(path.to_path_buf()); + } + options.connect(url).await.map_err(|source| Error::Connect { + url: url.to_owned(), + source, + }) +} + #[cfg(test)] mod tests { use super::*; diff --git a/swarm-secret-client/src/queue.rs b/swarm-secret-client/src/queue.rs index 3bd5a781..852906e0 100644 --- a/swarm-secret-client/src/queue.rs +++ b/swarm-secret-client/src/queue.rs @@ -22,6 +22,9 @@ //! reaching the store and nothing else; deriving a queue identity from it would //! couple the two credentials' lifetimes, so that renewing one would mean //! renewing the other. +//! +//! The token an agent presents that secret in at the queue is spelled by +//! `swarm_queue_client::agent_token`, which needs no store client. use serde::{Deserialize, Serialize}; @@ -57,14 +60,11 @@ pub fn agent_queue_path(agent: &str) -> Result { Ok(format!("{prefix}/queue")) } -/// What [`agent_queue_path`] holds: the secret, and the principal it proves. +/// What [`agent_queue_path`] holds: the secret, and the agent it proves. /// -/// Both names ride **in the object** rather than being parsed back out of a -/// composite principal string. Hive and agent names draw from the same -/// alphabet (`hive_types::Ident`, `[a-z0-9-]`), so a principal spelled -/// `hive--agent-` parses two ways for a name containing `-agent-` -/// — and an ambiguous principal parse in an authorisation path is a caller that -/// authenticates fine and is handed somebody else's grant. +/// No hive: an agent's identity is not tied to one, and the subjects the +/// verifier grants are keyed on the agent alone. Objects written with a `hive` +/// field still decode, because unknown fields are ignored. #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] pub struct AgentCredential { /// The secret itself. Named to match [`Credential::value`] and @@ -74,13 +74,6 @@ pub struct AgentCredential { /// The agent this secret authenticates. pub agent: String, - - /// The hive that agent belongs to. - /// - /// Here because the verifying end has no roster to look it up in, and - /// because the subjects an agent is granted are hive-templated — without - /// this field the verifier would know *who* is connecting and not *where*. - pub hive: String, } /// What the path holds: the client secret, plus the client id it belongs to. @@ -198,7 +191,6 @@ mod tests { let c = AgentCredential { value: "s3cr3t".to_owned(), agent: "atlas".to_owned(), - hive: "alpha".to_owned(), }; let json = serde_json::to_string(&c).expect("serialises"); assert_eq!( @@ -212,29 +204,39 @@ mod tests { let json = serde_json::to_value(AgentCredential { value: "s3cr3t".to_owned(), agent: "atlas".to_owned(), - hive: "alpha".to_owned(), }) .expect("serialises"); assert_eq!(json["value"], "s3cr3t"); assert_eq!(json["agent"], "atlas"); - assert_eq!(json["hive"], "alpha"); + assert!(json.get("hive").is_none(), "{json}"); } - /// Neither name is optional. An object missing one is not a usable - /// credential — a verifier holding `None` for the hive can only guess at - /// the subjects to grant, and guessing is the failure this shape exists to - /// prevent. + /// A stored object may carry a `hive` field. It is ignored, and the object + /// decodes. #[test] - fn an_agent_object_missing_a_principal_does_not_decode() { - assert!( - serde_json::from_str::(r#"{"value":"s","agent":"atlas"}"#).is_err() + fn a_stored_agent_object_that_still_names_a_hive_decodes() { + let c: AgentCredential = + serde_json::from_str(r#"{"value":"s3cr3t","agent":"atlas","hive":"alpha"}"#) + .expect("a stored object carrying `hive` decodes"); + assert_eq!( + c, + AgentCredential { + value: "s3cr3t".to_owned(), + agent: "atlas".to_owned(), + } ); + } + + /// The agent is required: it is what the verifier checks the presented + /// name against. + #[test] + fn an_agent_object_missing_its_agent_does_not_decode() { + assert!(serde_json::from_str::(r#"{"value":"s"}"#).is_err()); assert!( serde_json::from_str::(r#"{"value":"s","hive":"alpha"}"#).is_err() ); } - /// The two kinds are different objects at different paths, and neither /// decodes as the other — the property that keeps a reader from picking up /// a hive-shared credential where a per-agent one was meant. #[test] @@ -249,7 +251,6 @@ mod tests { let agent_json = serde_json::to_string(&AgentCredential { value: "s".to_owned(), agent: "atlas".to_owned(), - hive: "alpha".to_owned(), }) .expect("serialises"); assert!(serde_json::from_str::(&agent_json).is_err()); diff --git a/swarmctl/src/agent.rs b/swarmctl/src/agent.rs index d22a4aab..5cada5d1 100644 --- a/swarmctl/src/agent.rs +++ b/swarmctl/src/agent.rs @@ -51,13 +51,6 @@ struct CreateAgentResponse { warnings: Vec, } -/// Body of `POST /api/agents/{name}/identity`. Mirrors the controller's own -/// `MintAgentIdentityRequest` — see this module's doc comment. -#[derive(Serialize)] -struct MintIdentityRequest<'a> { - hive: &'a str, -} - /// Success body of `POST /api/agents/{name}/identity`. #[derive(Deserialize)] struct MintIdentityResponse { @@ -117,11 +110,10 @@ pub(crate) fn parse_ident(value: &str, what: &str) -> Result { /// /// Synchronous for the same reason [`create`] is, and built on the same /// round trip. -pub(crate) fn mint_identity(socket: &Path, name: &str, hive: &str) -> Result<()> { +pub(crate) fn mint_identity(socket: &Path, name: &str) -> Result<()> { // Client-side first, so a typo is a local error rather than a 400 the - // operator waits for. The controller validates both again. + // operator waits for. The controller validates it again. let name = parse_ident(name, "agent name")?; - let hive = parse_ident(hive, "hive")?; let rt = tokio::runtime::Builder::new_current_thread() .enable_io() @@ -130,15 +122,15 @@ pub(crate) fn mint_identity(socket: &Path, name: &str, hive: &str) -> Result<()> let resp: MintIdentityResponse = rt.block_on(post( socket, &format!("/api/agents/{name}/identity"), - &MintIdentityRequest { hive: &hive }, + &serde_json::json!({}), "mint-identity", ))?; println!("queued: job node {}", resp.node_id); println!( - "agent {name:?} will have its identity re-minted on hive {hive:?} once the job graph \ - runs; `swarmctl` does not wait for it. An existing queue secret is kept as it is; the \ - store certificate is re-minted and the agent picks the new one up on its next boot" + "agent {name:?} will have its identity re-minted once the job graph runs; `swarmctl` \ + does not wait for it. An existing queue secret is kept as it is; the store \ + certificate is re-minted and the agent picks the new one up on its next boot" ); Ok(()) } diff --git a/swarmctl/src/main.rs b/swarmctl/src/main.rs index d1c032eb..eb041618 100644 --- a/swarmctl/src/main.rs +++ b/swarmctl/src/main.rs @@ -242,16 +242,6 @@ struct AgentMintForgeTokenArgs { struct AgentMintIdentityArgs { /// Name of an agent that already exists. name: String, - /// The hive that agent runs on. - /// - /// Required, and deliberately not defaulted: the credentials this mints - /// name a hive, and neither this CLI nor the controller keeps a roster of - /// which agent is on which hive. Naming the wrong one gives the agent an - /// identity scoped to a hive it doesn't run on. The controller checks - /// the value against the swarm's hive roster and names the known hives if - /// it misses. - #[arg(long, value_name = "HIVE")] - hive: String, /// swarm-controller's unix socket. /// /// Supplied by the nix module that installs this binary, from the same @@ -357,7 +347,7 @@ fn main() -> Result<()> { command: AgentVerb::MintIdentity(args), } => { let socket = path_from(args.controller_socket, "SWARM_CONTROLLER_SOCKET")?; - agent::mint_identity(&socket, &args.name, &args.hive) + agent::mint_identity(&socket, &args.name) } // Same socket-resolution reasoning as `Create` above. Verb::Agent { @@ -732,20 +722,10 @@ mod tests { assert!(args.controller_socket.is_none()); } - /// The backfill verb takes the same two names as `create`, and `--hive` - /// is required on it for the same reason: it is an address nobody can - /// infer. #[test] - fn the_backfill_verb_takes_an_agent_and_a_hive() { - let cli = Cli::try_parse_from([ - "swarmctl", - "agent", - "mint-identity", - "scribe", - "--hive", - "alpha", - ]) - .expect("the minimal form parses"); + fn the_backfill_verb_takes_an_agent_and_no_hive() { + let cli = Cli::try_parse_from(["swarmctl", "agent", "mint-identity", "scribe"]) + .expect("the minimal form parses"); let Verb::Agent { command: AgentVerb::MintIdentity(args), } = cli.command @@ -753,12 +733,18 @@ mod tests { panic!("expected `agent mint-identity`"); }; assert_eq!(args.name, "scribe"); - assert_eq!(args.hive, "alpha"); assert!(args.controller_socket.is_none()); - assert!( - Cli::try_parse_from(["swarmctl", "agent", "mint-identity", "scribe"]).is_err(), - "an omitted hive must not be defaulted" + Cli::try_parse_from([ + "swarmctl", + "agent", + "mint-identity", + "scribe", + "--hive", + "a" + ]) + .is_err(), + "the identity has no hive, so the verb must not take one" ); }