Compare commits
30 changed files with 353 additions and 1435 deletions
2
Cargo.lock
generated
2
Cargo.lock
generated
|
|
@ -4828,9 +4828,7 @@ dependencies = [
|
|||
"serde",
|
||||
"serde_json",
|
||||
"sha2 0.11.0",
|
||||
"subtle",
|
||||
"swarm-queue-client",
|
||||
"swarm-secret-client",
|
||||
"tokio",
|
||||
"tracing",
|
||||
"tracing-subscriber",
|
||||
|
|
|
|||
|
|
@ -230,9 +230,6 @@ 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
|
||||
|
|
|
|||
|
|
@ -33,8 +33,9 @@ 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.
|
||||
swarmctl agent mint-identity ruth
|
||||
# On the swarm-controller host: give ruth her store identity. <hive> is the
|
||||
# name of the hive she runs on.
|
||||
swarmctl agent mint-identity ruth --hive <hive>
|
||||
|
||||
# On ruth's hive: re-apply her container config, which is when hive-c0re
|
||||
# hands the new identity to the container.
|
||||
|
|
|
|||
|
|
@ -410,13 +410,10 @@ 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.<agent>`, one subject per agent, so
|
||||
a swarm-level terminal can follow one agent without subscribing to the swarm's
|
||||
whole traffic. An agent connected with its own queue credential gets that
|
||||
subject. An agent without one, or whose own credential the queue
|
||||
refused, connects with its hive's shared client and publishes to
|
||||
`$SWARM.term.<hive>.<agent>` instead, the `<hive>` being the one that client id
|
||||
names. The swarm controller relays both. Publishing only: an
|
||||
own web UI renders also goes to `$SWARM.term.<hive>.<agent>`, one subject per
|
||||
agent, so a swarm-level terminal can follow one agent without subscribing to
|
||||
the swarm's whole traffic. The `<hive>` is the one the agent's client id names,
|
||||
which is the same string the broker builds its grant from. Publishing only: an
|
||||
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.
|
||||
|
|
@ -427,7 +424,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.<agent>` (or `$SWARM.agent-state.<hive>.<agent>`) — same shape of subject, same grant
|
||||
`$SWARM.agent-state.<hive>.<agent>` — same shape of subject, same grant
|
||||
mechanics, same lack of retention. It carries what a header bar wants: what the
|
||||
turn loop is doing (`turn_state`, plus `turn_state_since` as an ISO 8601 UTC
|
||||
stamp), which model (`model` and the resolved id the last turn actually ran on),
|
||||
|
|
|
|||
|
|
@ -101,10 +101,12 @@ mint therefore never receives one, and nothing will ever come back around to
|
|||
it. Re-run the mint for one agent with:
|
||||
|
||||
```sh
|
||||
swarmctl agent mint-identity <agent>
|
||||
swarmctl agent mint-identity <agent> --hive <hive>
|
||||
```
|
||||
|
||||
The queue secret half is idempotent — an agent that already has one keeps exactly the
|
||||
`--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
|
||||
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.
|
||||
|
|
|
|||
|
|
@ -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] <NAME>`
|
||||
**Usage:** `swarmctl agent mint-identity [OPTIONS] --hive <HIVE> <NAME>`
|
||||
|
||||
###### **Arguments:**
|
||||
|
||||
|
|
@ -100,6 +100,9 @@ Queues and returns, the same way `agent create` does — watch the swarm UI's jo
|
|||
|
||||
###### **Options:**
|
||||
|
||||
* `--hive <HIVE>` — The hive that agent runs on.
|
||||
|
||||
Required, and deliberately not defaulted: the credentials this mints name a hive, and neither this CLI nor the controller keeps a roster of which agent is on which hive. Naming the wrong one gives the agent an identity scoped to a hive it doesn't run on. The controller checks the value against the swarm's hive roster and names the known hives if it misses.
|
||||
* `--controller-socket <PATH>` — swarm-controller's unix socket.
|
||||
|
||||
Supplied by the nix module that installs this binary, from the same `socketPath` option the daemon binds; falls back to `SWARM_CONTROLLER_SOCKET`.
|
||||
|
|
|
|||
|
|
@ -35,7 +35,6 @@ 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
|
||||
|
|
@ -172,43 +171,33 @@ fn snapshot(bus: &Bus) -> AgentStateMsg {
|
|||
}
|
||||
}
|
||||
|
||||
/// The subject this agent's header goes to under the credential it connected
|
||||
/// with, or `None` when a hive client id names no hive.
|
||||
fn subject(presented: &Presented, agent: &str) -> Option<String> {
|
||||
match presented {
|
||||
Presented::Agent => Some(format!("{SUBJECT_PREFIX}.{agent}")),
|
||||
Presented::Hive { client_id } => {
|
||||
let Some(hive) = hive_from_client_id(client_id) else {
|
||||
tracing::warn!(
|
||||
%client_id,
|
||||
expected = format!("{CLIENT_ID_PREFIX}<hive>{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.
|
||||
/// Start the publish task, if this agent has both queue coordinates and an
|
||||
/// identity the subject can be derived from.
|
||||
///
|
||||
/// Returns without spawning in every other case — no queue 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.
|
||||
/// Returns without spawning in every other case — no queue, an unparseable
|
||||
/// client id, no label — each of which is a legal state for an agent rather
|
||||
/// than an error, and each logged once here rather than per transition.
|
||||
pub fn spawn(bus: &Bus) {
|
||||
if !crate::swarm_queue::configured() {
|
||||
let Some(cfg) = crate::swarm_queue::config() else {
|
||||
// `swarm_queue::init` already said why at boot; repeating it here
|
||||
// would be the same fact logged twice.
|
||||
return;
|
||||
}
|
||||
};
|
||||
let Some(hive) = hive_from_client_id(&cfg.client_id) else {
|
||||
tracing::warn!(
|
||||
client_id = %cfg.client_id,
|
||||
expected = format!("{CLIENT_ID_PREFIX}<hive>{CLIENT_ID_SUFFIX}"),
|
||||
"queue client id does not name a hive; not publishing turn state upward"
|
||||
);
|
||||
return;
|
||||
};
|
||||
let agent = crate::identity::label();
|
||||
if agent.is_empty() {
|
||||
tracing::warn!("this agent has no label; not publishing turn state upward");
|
||||
return;
|
||||
}
|
||||
tokio::spawn(run(bus.subscribe(), bus.clone(), agent));
|
||||
let subject = format!("{SUBJECT_PREFIX}.{hive}.{agent}");
|
||||
tokio::spawn(run(bus.subscribe(), bus.clone(), subject));
|
||||
}
|
||||
|
||||
/// Watch the bus and publish whenever the header actually changed.
|
||||
|
|
@ -229,11 +218,8 @@ pub fn spawn(bus: &Bus) {
|
|||
/// and it carries nothing this header reads. Every other variant, including
|
||||
/// any added later, funnels into the comparison and costs nothing when it
|
||||
/// changes nothing.
|
||||
async fn run(mut rx: broadcast::Receiver<BusEvent>, bus: Bus, agent: String) {
|
||||
let Some(Connection { client, presented }) = crate::swarm_queue::client().await else {
|
||||
return;
|
||||
};
|
||||
let Some(subject) = subject(&presented, &agent) else {
|
||||
async fn run(mut rx: broadcast::Receiver<BusEvent>, bus: Bus, subject: String) {
|
||||
let Some(client) = crate::swarm_queue::client().await else {
|
||||
return;
|
||||
};
|
||||
tracing::info!(subject, "publishing agent turn state to the swarm queue");
|
||||
|
|
@ -321,7 +307,7 @@ async fn publish(
|
|||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::{AgentStateMsg, Presented, hive_from_client_id, subject};
|
||||
use super::{AgentStateMsg, SUBJECT_PREFIX, hive_from_client_id};
|
||||
use crate::events::TurnState;
|
||||
use swarm_queue_client::wanted::AgentState;
|
||||
|
||||
|
|
@ -407,22 +393,10 @@ mod tests {
|
|||
/// this pins the half that lives here.
|
||||
#[test]
|
||||
fn the_subject_is_the_prefix_then_the_hive_then_the_agent() {
|
||||
let presented = Presented::Hive {
|
||||
client_id: "hive-alpha-agent".to_owned(),
|
||||
};
|
||||
let hive = hive_from_client_id("hive-alpha-agent").expect("names a hive");
|
||||
assert_eq!(
|
||||
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")
|
||||
format!("{SUBJECT_PREFIX}.{hive}.mara"),
|
||||
"$SWARM.agent-state.alpha.mara"
|
||||
);
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -21,16 +21,15 @@
|
|||
//! 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`).
|
||||
//!
|
||||
//! 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.
|
||||
//! 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.
|
||||
|
||||
use std::path::{Path, PathBuf};
|
||||
use std::sync::OnceLock;
|
||||
|
||||
use anyhow::Context as _;
|
||||
use tokio::sync::OnceCell;
|
||||
|
||||
use swarm_queue_client::QueueConfig;
|
||||
|
|
@ -43,39 +42,10 @@ const ENV_PREFIX: &str = "HIVE_AGENT";
|
|||
/// one answer rather than re-deriving it per call.
|
||||
static CONFIG: OnceLock<Option<QueueConfig>> = OnceLock::new();
|
||||
|
||||
/// Where to present this agent's own credential, resolved once at boot.
|
||||
static AGENT: OnceLock<Option<AgentPath>> = OnceLock::new();
|
||||
|
||||
/// The one connection every publisher in this process shares. Separate from
|
||||
/// [`CONFIG`] because resolving the coordinates is synchronous boot work and
|
||||
/// connecting is not — see [`client`].
|
||||
static CLIENT: OnceCell<Option<Connection>> = OnceCell::const_new();
|
||||
|
||||
/// What connecting with this agent's own credential needs.
|
||||
#[derive(Debug, PartialEq, Eq)]
|
||||
struct AgentPath {
|
||||
url: String,
|
||||
ca_file: Option<PathBuf>,
|
||||
/// The fetched secret. A path; the bytes are read at connect.
|
||||
secret_file: PathBuf,
|
||||
}
|
||||
|
||||
/// Which credential the shared connection presented.
|
||||
#[derive(Debug, Clone, PartialEq, Eq)]
|
||||
pub enum Presented {
|
||||
/// This agent's own, granted `<prefix>.<agent>`.
|
||||
Agent,
|
||||
/// The hive's shared client, granted `<prefix>.<hive>.>` for the hive the
|
||||
/// client id names.
|
||||
Hive { client_id: String },
|
||||
}
|
||||
|
||||
/// The shared queue connection and the credential it was made with.
|
||||
#[derive(Clone)]
|
||||
pub struct Connection {
|
||||
pub client: async_nats::Client,
|
||||
pub presented: Presented,
|
||||
}
|
||||
static CLIENT: OnceCell<Option<async_nats::Client>> = OnceCell::const_new();
|
||||
|
||||
/// 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
|
||||
|
|
@ -177,16 +147,6 @@ fn decide_agent_secret(path: Option<&str>) -> Option<PathBuf> {
|
|||
}
|
||||
}
|
||||
|
||||
/// Decide whether this agent can connect with its own credential: it needs the
|
||||
/// queue's address and a fetched secret, and nothing of the hive's client.
|
||||
fn decide_agent_path(env: &QueueEnv, secret_file: Option<PathBuf>) -> Option<AgentPath> {
|
||||
Some(AgentPath {
|
||||
url: env.nats_url.clone()?,
|
||||
ca_file: env.ca_file.as_ref().map(Into::into),
|
||||
secret_file: secret_file?,
|
||||
})
|
||||
}
|
||||
|
||||
/// Decide what this agent's queue configuration is, given the environment and
|
||||
/// whatever the client-id file held.
|
||||
///
|
||||
|
|
@ -256,10 +216,15 @@ pub fn init() {
|
|||
let _ = CONFIG.set(resolved);
|
||||
|
||||
// Independent of everything above: this agent may hold its own secret on
|
||||
// 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 {
|
||||
// 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()) {
|
||||
// The path, never the bytes: the file holds the secret itself.
|
||||
tracing::info!(
|
||||
path = %path.display(),
|
||||
|
|
@ -271,13 +236,6 @@ 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.
|
||||
|
|
@ -305,48 +263,14 @@ 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<Connection> {
|
||||
pub async fn client() -> Option<async_nats::Client> {
|
||||
CLIENT.get_or_init(connect_once).await.clone()
|
||||
}
|
||||
|
||||
/// This agent's own credential first, then the hive's. Each path taken is
|
||||
/// logged once, here.
|
||||
async fn connect_once() -> Option<Connection> {
|
||||
if let Some(agent) = AGENT.get().and_then(Option::as_ref) {
|
||||
match connect_as_agent(agent).await {
|
||||
Ok(client) => {
|
||||
tracing::info!(
|
||||
url = %agent.url,
|
||||
"connected to the swarm queue with this agent's own credential"
|
||||
);
|
||||
return Some(Connection {
|
||||
client,
|
||||
presented: Presented::Agent,
|
||||
});
|
||||
}
|
||||
Err(e) => tracing::warn!(
|
||||
error = format!("{e:#}"),
|
||||
"connecting with this agent's own credential failed; falling back to \
|
||||
its hive's shared client"
|
||||
),
|
||||
}
|
||||
}
|
||||
async fn connect_once() -> Option<async_nats::Client> {
|
||||
let cfg = config()?;
|
||||
match swarm_queue_client::connect(cfg.clone()).await {
|
||||
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(),
|
||||
},
|
||||
})
|
||||
}
|
||||
Ok(client) => Some(client),
|
||||
Err(e) => {
|
||||
// `chain`, not `{:#}`: this is `swarm_queue_client::Error`, whose
|
||||
// `Display` ignores the alternate flag, so `{:#}` renders the
|
||||
|
|
@ -360,26 +284,9 @@ async fn connect_once() -> Option<Connection> {
|
|||
}
|
||||
}
|
||||
|
||||
/// Connect presenting this agent's own token. Any failure, a refusal
|
||||
/// included, is returned for [`connect_once`] to fall back on.
|
||||
async fn connect_as_agent(agent: &AgentPath) -> anyhow::Result<async_nats::Client> {
|
||||
let name = crate::identity::label();
|
||||
let secret = std::fs::read_to_string(&agent.secret_file)
|
||||
.with_context(|| format!("reading {}", agent.secret_file.display()))?;
|
||||
let token = swarm_queue_client::agent_token::format_agent_token(&name, secret.trim())?;
|
||||
swarm_queue_client::connect_with_token(&agent.url, agent.ca_file.as_deref(), token)
|
||||
.await
|
||||
.map_err(|e| anyhow::anyhow!(swarm_queue_client::chain(&e)))
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use std::path::PathBuf;
|
||||
|
||||
use super::{
|
||||
AgentPath, QueueEnv, Resolution, decide, decide_agent_path, decide_agent_secret,
|
||||
read_client_id,
|
||||
};
|
||||
use super::{QueueEnv, Resolution, decide, 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;
|
||||
|
|
@ -528,34 +435,4 @@ 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
|
||||
);
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -27,7 +27,6 @@
|
|||
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
|
||||
|
|
@ -113,50 +112,37 @@ fn serialized_len(msg: &TermMsg) -> Option<usize> {
|
|||
serde_json::to_vec(msg).ok().map(|v| v.len())
|
||||
}
|
||||
|
||||
/// The subject this agent's rows go to under the credential it connected
|
||||
/// with, or `None` when a hive client id names no hive.
|
||||
fn subject(presented: &Presented, agent: &str) -> Option<String> {
|
||||
match presented {
|
||||
Presented::Agent => Some(format!("{SUBJECT_PREFIX}.{agent}")),
|
||||
Presented::Hive { client_id } => {
|
||||
let Some(hive) = hive_from_client_id(client_id) else {
|
||||
tracing::warn!(
|
||||
%client_id,
|
||||
expected = format!("{CLIENT_ID_PREFIX}<hive>{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.
|
||||
/// Start the publish task, if this agent has both queue coordinates and an
|
||||
/// identity the subject can be derived from.
|
||||
///
|
||||
/// Returns without spawning in every other case — no queue 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.
|
||||
/// Returns without spawning in every other case — no queue, an unparseable
|
||||
/// client id, no label — each of which is a legal state for an agent rather
|
||||
/// than an error, and each logged once here rather than per row.
|
||||
pub fn spawn(rx: broadcast::Receiver<BusEvent>) {
|
||||
if !crate::swarm_queue::configured() {
|
||||
let Some(cfg) = crate::swarm_queue::config() else {
|
||||
// `swarm_queue::init` already said why at boot; repeating it here
|
||||
// would be the same fact logged twice.
|
||||
return;
|
||||
}
|
||||
};
|
||||
let Some(hive) = hive_from_client_id(&cfg.client_id) else {
|
||||
tracing::warn!(
|
||||
client_id = %cfg.client_id,
|
||||
expected = format!("{CLIENT_ID_PREFIX}<hive>{CLIENT_ID_SUFFIX}"),
|
||||
"queue client id does not name a hive; not publishing the terminal upward"
|
||||
);
|
||||
return;
|
||||
};
|
||||
let agent = crate::identity::label();
|
||||
if agent.is_empty() {
|
||||
tracing::warn!("this agent has no label; not publishing the terminal upward");
|
||||
return;
|
||||
}
|
||||
tokio::spawn(run(rx, agent));
|
||||
let subject = format!("{SUBJECT_PREFIX}.{hive}.{agent}");
|
||||
tokio::spawn(run(rx, subject));
|
||||
}
|
||||
|
||||
async fn run(mut rx: broadcast::Receiver<BusEvent>, agent: String) {
|
||||
let Some(Connection { client, presented }) = crate::swarm_queue::client().await else {
|
||||
return;
|
||||
};
|
||||
let Some(subject) = subject(&presented, &agent) else {
|
||||
async fn run(mut rx: broadcast::Receiver<BusEvent>, subject: String) {
|
||||
let Some(client) = crate::swarm_queue::client().await else {
|
||||
return;
|
||||
};
|
||||
tracing::info!(subject, "publishing the agent terminal to the swarm queue");
|
||||
|
|
@ -217,7 +203,7 @@ async fn publish(client: &async_nats::Client, subject: &str, msg: TermMsg) {
|
|||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::{DROPPED_BODY, Presented, fit, hive_from_client_id, subject};
|
||||
use super::{DROPPED_BODY, fit, hive_from_client_id};
|
||||
use crate::events::LiveEvent;
|
||||
use crate::term_msg::{BodyFormat, ClassifyCtx, Level, TermMsg, classify};
|
||||
|
||||
|
|
@ -257,27 +243,6 @@ 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")
|
||||
|
|
|
|||
|
|
@ -32,7 +32,6 @@
|
|||
./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
|
||||
|
|
|
|||
|
|
@ -260,12 +260,6 @@ 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
|
||||
'';
|
||||
};
|
||||
};
|
||||
|
|
|
|||
|
|
@ -1,37 +0,0 @@
|
|||
# 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";
|
||||
};
|
||||
};
|
||||
}
|
||||
|
|
@ -754,18 +754,6 @@ 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.
|
||||
|
|
@ -1450,24 +1438,6 @@ in
|
|||
'';
|
||||
};
|
||||
|
||||
natsAuthCommonName = lib.mkOption {
|
||||
type = lib.types.str;
|
||||
default = "swarm-nats-auth";
|
||||
description = ''
|
||||
Subject the store's `swarm-nats-auth` cert-auth role accepts: the
|
||||
identity the queue's auth-callout responder presents to read agent
|
||||
queue credentials, `swarm/agents/<agent>/queue` and nothing else under
|
||||
an agent.
|
||||
|
||||
Its own principal rather than
|
||||
{option}`services.hyperhive.deploy.bao.natsCommonName`: that one
|
||||
issues the queue's TLS leaf from a host unit, this one reads secrets
|
||||
from inside the queue's container.
|
||||
|
||||
⚠️ Reserved as a hive name by ./swarm.nix, like its siblings.
|
||||
'';
|
||||
};
|
||||
|
||||
matrixCtlHiveName = lib.mkOption {
|
||||
type = lib.types.str;
|
||||
default = toString hyperhiveCfg.hiveName;
|
||||
|
|
@ -2051,7 +2021,6 @@ 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
|
||||
|
|
@ -2754,8 +2723,6 @@ 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
|
||||
|
|
|
|||
|
|
@ -219,25 +219,6 @@ 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.
|
||||
|
|
@ -517,34 +498,6 @@ 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 {
|
||||
|
|
@ -952,11 +905,6 @@ 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
|
||||
|
|
@ -966,25 +914,12 @@ 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
|
||||
|
|
@ -1022,10 +957,7 @@ 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"
|
||||
# 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";
|
||||
++ lib.optional deployCfg.authelia.enable "container@${autheliaCfg.machine}.service";
|
||||
requires = lib.optional deployCfg.nats.autoGenerateCallout "swarm-nats-callout-keys.service";
|
||||
serviceConfig = {
|
||||
Type = "oneshot";
|
||||
|
|
@ -1074,22 +1006,7 @@ 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
|
||||
|
|
|
|||
|
|
@ -53,7 +53,6 @@ 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
|
||||
|
|
|
|||
|
|
@ -151,7 +151,7 @@ let
|
|||
_: u: (u.environment.BAO_CLIENT_CERT or null) == granterCertFile
|
||||
) baoGrantWithConsumers.systemd.services;
|
||||
|
||||
# The twelve units that write a `swarm-*` grant, by name, for the discovery
|
||||
# The eleven units that write a `swarm-*` grant, by name, for the discovery
|
||||
# control below.
|
||||
grantingUnitNames = [
|
||||
"swarm-bao-controller-policy"
|
||||
|
|
@ -165,7 +165,6 @@ 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
|
||||
|
|
@ -562,33 +561,6 @@ let
|
|||
&& !(lib.hasInfix "swarm-grafana" s)
|
||||
&& !(lib.hasInfix "sys/policies/acl" s);
|
||||
}
|
||||
{
|
||||
# `+` is one path segment, so this reaches `swarm/agents/<agent>/queue`
|
||||
# and no other credential an agent holds; `swarm/agents/*` would reach
|
||||
# all of them.
|
||||
name = "the queue responder's grant is every agent's queue credential and nothing else";
|
||||
ok =
|
||||
let
|
||||
s = baoGrantHere.systemd.services.swarm-bao-nats-auth-policy.script;
|
||||
in
|
||||
lib.hasInfix "path \"secret/data/swarm/agents/+/queue\" {" s
|
||||
&& lib.hasInfix "capabilities = [\"read\"]" s
|
||||
&& lib.length (lib.filter lib.isList (builtins.split "path \"" s)) == 1
|
||||
&& !(lib.hasInfix "secret/data/swarm/agents/*" s)
|
||||
&& !(lib.hasInfix "secret/data/swarm/hives" s)
|
||||
&& !(lib.hasInfix "secret/data/swarm/services" s)
|
||||
&& lib.hasInfix "auth/cert/certs/swarm-nats-auth" s
|
||||
&& lib.hasInfix "allowed_common_names=swarm-nats-auth" s
|
||||
&& lib.hasInfix "token_policies=swarm-nats-auth" s;
|
||||
}
|
||||
{
|
||||
name = "the PKI unit signs the queue responder's leaf under its own subject";
|
||||
ok =
|
||||
let
|
||||
s = baoGrantHere.systemd.services.swarm-bao-pki.script;
|
||||
in
|
||||
lib.hasInfix "/nats-auth.pem ]" s && lib.hasInfix "swarm-nats-auth \"\" clientAuth" s;
|
||||
}
|
||||
{
|
||||
# 🩸 The half that makes the policies above bind: a policy grants only
|
||||
# through a token that carries it, and a token is minted by a cert-auth
|
||||
|
|
@ -776,24 +748,24 @@ let
|
|||
}
|
||||
{
|
||||
# A store host without the granter's pair writes its grants some other
|
||||
# way, so none of the twelve units may exist. Without this arm
|
||||
# way, so none of the eleven 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 twelve granting units render";
|
||||
name = "without the granter's pair none of the eleven 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 twelve.
|
||||
# The control: the same store with the pair renders all eleven.
|
||||
&& lib.all (unit: baoGrantHere.systemd.services ? ${unit}) grantingUnitNames;
|
||||
}
|
||||
{
|
||||
# 🩸 What replaced the silent skip. With no bootstrap token the twelve still
|
||||
# 🩸 What replaced the silent skip. With no bootstrap token the eleven 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 twelve, each failing loudly with the one-time step";
|
||||
name = "a store host without a bootstrap token renders the eleven, each failing loudly with the one-time step";
|
||||
ok =
|
||||
let
|
||||
s = baoGranterNoToken.systemd.services;
|
||||
|
|
@ -1157,8 +1129,8 @@ let
|
|||
}
|
||||
{
|
||||
# What makes the case above mean something: discovery by the granter's
|
||||
# certificate reaches all twelve units, and each yields calls.
|
||||
name = "the granter-policy check sees all twelve granting units, and parses calls from each";
|
||||
# certificate reaches all eleven units, and each yields calls.
|
||||
name = "the granter-policy check sees all eleven 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)
|
||||
|
|
|
|||
|
|
@ -60,28 +60,6 @@ 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
|
||||
|
|
@ -139,73 +117,6 @@ 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
|
||||
|
|
|
|||
|
|
@ -136,7 +136,7 @@ fn generate_queue_secret() -> Result<String> {
|
|||
/// Anything that stops one of those five steps, with the step named. A
|
||||
/// failure here fails the job node and nothing else — the agent is still
|
||||
/// created, without a store identity.
|
||||
pub async fn mint_and_verify(agent: &str) -> Result<()> {
|
||||
pub async fn mint_and_verify(agent: &str, hive: &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,15 +161,18 @@ pub async fn mint_and_verify(agent: &str) -> Result<()> {
|
|||
.read_optional(&queue_path)
|
||||
.await
|
||||
.with_context(|| format!("checking whether {queue_path} already holds a credential"))?;
|
||||
// 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.
|
||||
// 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.
|
||||
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!(
|
||||
|
|
@ -184,8 +187,9 @@ pub async fn mint_and_verify(agent: &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 name
|
||||
// Never "rotated": the secret is the same one, only the names
|
||||
// around it moved.
|
||||
corrected = existing.is_some(),
|
||||
"agent queue credential published"
|
||||
|
|
@ -276,9 +280,9 @@ async fn read_back_as_agent(
|
|||
.read(queue_path)
|
||||
.await
|
||||
.with_context(|| format!("reading {queue_path} back under {role}'s own token"))?;
|
||||
// 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.
|
||||
// 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.
|
||||
if read_back != *queue_credential {
|
||||
bail!("the store returned a different object at {queue_path} than the one just published");
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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. `crate::term_stream::agent_subjects`
|
||||
/// spells the full subjects; as with the terminal relay's own copy, there is no
|
||||
/// ends never see the constant together. `$SWARM.agent-state.<hive>.<agent>`
|
||||
/// is the full subject; as with the terminal relay's own copy, there is no
|
||||
/// 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,12 +117,16 @@ pub(crate) async fn stream_agent_state(
|
|||
)
|
||||
})?;
|
||||
|
||||
let subscriber = crate::term_stream::subscribe_all(
|
||||
&client,
|
||||
&crate::term_stream::agent_subjects(SUBJECT_PREFIX, &hive, &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}"),
|
||||
)
|
||||
.await?;
|
||||
})?;
|
||||
|
||||
tracing::info!(%subject, "state stream: client attached");
|
||||
let stream =
|
||||
subscriber.map(|msg| Ok(Event::default().data(String::from_utf8_lossy(&msg.payload))));
|
||||
Ok(Sse::new(stream).keep_alive(KeepAlive::default()))
|
||||
|
|
|
|||
|
|
@ -100,10 +100,11 @@ enum SwarmNodeKind {
|
|||
/// `agent_identity::mint_and_verify` — including why this node does not
|
||||
/// report success on a write.
|
||||
///
|
||||
/// 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 },
|
||||
/// 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 },
|
||||
/// 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
|
||||
|
|
@ -123,7 +124,8 @@ 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 because the wanted-state bucket is keyed per hive.
|
||||
/// Carries the hive for the same reason `TriggerDeploy`/`MintAgentIdentity`
|
||||
/// do: 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.
|
||||
|
|
@ -164,12 +166,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 })
|
||||
}
|
||||
|
|
@ -312,7 +314,7 @@ async fn run_swarm_node(
|
|||
Err(e) => Outcome::Failed(format!("{e:#}")),
|
||||
},
|
||||
},
|
||||
SwarmNodeKind::MintAgentIdentity { agent } => mint_identity(&agent).await,
|
||||
SwarmNodeKind::MintAgentIdentity { hive, agent } => mint_identity(&agent, &hive).await,
|
||||
SwarmNodeKind::MintAgentForgeToken { agent } => mint_forge_token(deps.forge, &agent).await,
|
||||
SwarmNodeKind::MintAgentMatrixAccount { agent } => {
|
||||
mint_matrix_account(deps.matrix_homeserver.as_deref(), &agent).await
|
||||
|
|
@ -338,10 +340,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_jobq::scheduler::Outcome {
|
||||
async fn mint_identity(agent: &str, hive: &str) -> hive_jobq::scheduler::Outcome {
|
||||
use hive_jobq::scheduler::Outcome;
|
||||
|
||||
match agent_identity::mint_and_verify(agent).await {
|
||||
match agent_identity::mint_and_verify(agent, hive).await {
|
||||
Ok(()) => Outcome::Done,
|
||||
Err(e) => Outcome::Failed(format!("{e:#}")),
|
||||
}
|
||||
|
|
@ -1510,6 +1512,7 @@ fn declare_agent_job(
|
|||
// authelia.
|
||||
let mint_identity = b
|
||||
.node(SwarmNodeKind::MintAgentIdentity {
|
||||
hive: hive.to_owned(),
|
||||
agent: agent.to_owned(),
|
||||
})
|
||||
.after_ok(create_identity);
|
||||
|
|
@ -1667,6 +1670,20 @@ 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 {
|
||||
|
|
@ -1695,9 +1712,10 @@ 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` is not a valid identifier (problem+json)", body = String),
|
||||
(status = 400, description = "`name` or `hive` is not a valid identifier, or `hive` is not in this swarm (problem+json)", body = String),
|
||||
(status = 500, description = "the job could not be queued (problem+json)", body = String),
|
||||
),
|
||||
tag = "agents"
|
||||
|
|
@ -1705,12 +1723,28 @@ struct MintAgentIdentityResponse {
|
|||
async fn mint_agent_identity(
|
||||
State(state): State<AppState>,
|
||||
Path(name): Path<String>,
|
||||
Json(req): Json<MintAgentIdentityRequest>,
|
||||
) -> Result<Json<MintAgentIdentityResponse>, problem_details::ProblemDetails> {
|
||||
// The name is interpolated into store paths and policy documents
|
||||
// downstream, so it is validated here as well as there.
|
||||
// Both names are interpolated into store paths and policy documents
|
||||
// downstream, so both are 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
|
||||
|
|
@ -1723,6 +1757,7 @@ async fn mint_agent_identity(
|
|||
.insert_job(None, |b| {
|
||||
vec![
|
||||
b.node(SwarmNodeKind::MintAgentIdentity {
|
||||
hive: hive.clone(),
|
||||
agent: agent.clone(),
|
||||
})
|
||||
.guid(),
|
||||
|
|
@ -2768,6 +2803,39 @@ 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
|
||||
|
|
@ -2780,6 +2848,9 @@ 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");
|
||||
|
|
@ -2805,9 +2876,12 @@ 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 valid agent name must be accepted");
|
||||
.expect("a hive in the roster must be accepted");
|
||||
|
||||
let guard = sched
|
||||
.lock()
|
||||
|
|
@ -2819,9 +2893,10 @@ mod tests {
|
|||
assert!(
|
||||
matches!(
|
||||
&node.payload,
|
||||
SwarmNodeKind::MintAgentIdentity { agent } if agent == "atlas"
|
||||
SwarmNodeKind::MintAgentIdentity { hive, agent }
|
||||
if hive == "pr1ma" && agent == "atlas"
|
||||
),
|
||||
"the one node must be the mint, for the named agent: {:?}",
|
||||
"the one node must be the mint, carrying both names: {:?}",
|
||||
node.payload
|
||||
);
|
||||
// The id the operator is told to watch has to be the node that was
|
||||
|
|
@ -2950,6 +3025,7 @@ mod tests {
|
|||
let id = sched
|
||||
.append(
|
||||
SwarmNodeKind::MintAgentIdentity {
|
||||
hive: "pr1ma".to_owned(),
|
||||
agent: "atlas".to_owned(),
|
||||
},
|
||||
Vec::new(),
|
||||
|
|
@ -3116,6 +3192,25 @@ 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.
|
||||
|
|
|
|||
|
|
@ -9,8 +9,7 @@
|
|||
//!
|
||||
//! **Live tail only, on purpose.** `hive-agent` publishes each
|
||||
//! already-classified `TermMsg` row to the core subject
|
||||
//! `$SWARM.term.{agent}` (`$SWARM.term.{hive}.{agent}` from an agent on its
|
||||
//! hive's shared credential; see `hive-agent::swarm_term`'s module
|
||||
//! `$SWARM.term.{hive}.{agent}` (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
|
||||
|
|
@ -35,8 +34,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. [`agent_subjects`] spells the full
|
||||
/// subjects; there is no way to check the two sides agree short of
|
||||
/// never see the constant together. `$SWARM.term.<hive>.<agent>` is the
|
||||
/// full subject; there is no way to check the two sides agree short of
|
||||
/// this comment and the module docs on both ends staying honest about it.
|
||||
const SUBJECT_PREFIX: &str = "$SWARM.term";
|
||||
|
||||
|
|
@ -117,63 +116,17 @@ pub(crate) async fn stream_agent_term(
|
|||
)
|
||||
})?;
|
||||
|
||||
let subscriber = subscribe_all(
|
||||
&client,
|
||||
&agent_subjects(SUBJECT_PREFIX, &hive, &agent),
|
||||
"term",
|
||||
)
|
||||
.await?;
|
||||
let stream =
|
||||
subscriber.map(|msg| Ok(Event::default().data(String::from_utf8_lossy(&msg.payload))));
|
||||
Ok(Sse::new(stream).keep_alive(KeepAlive::default()))
|
||||
}
|
||||
|
||||
/// The subjects one agent publishes on under `prefix`: its own
|
||||
/// `<prefix>.<agent>`, and `<prefix>.<hive>.<agent>`, which an agent still on
|
||||
/// its hive's shared queue credential publishes to.
|
||||
pub(crate) fn agent_subjects(prefix: &str, hive: &str, agent: &str) -> [String; 2] {
|
||||
[
|
||||
format!("{prefix}.{agent}"),
|
||||
format!("{prefix}.{hive}.{agent}"),
|
||||
]
|
||||
}
|
||||
|
||||
/// Subscribe to every one of `subjects`, merged into one stream. `what` names
|
||||
/// the relay in logs and errors.
|
||||
pub(crate) async fn subscribe_all(
|
||||
client: &async_nats::Client,
|
||||
subjects: &[String],
|
||||
what: &str,
|
||||
) -> Result<futures_util::stream::SelectAll<async_nats::Subscriber>, problem_details::ProblemDetails>
|
||||
{
|
||||
let mut subscribers = Vec::with_capacity(subjects.len());
|
||||
for subject in subjects {
|
||||
let subject = format!("{SUBJECT_PREFIX}.{hive}.{agent}");
|
||||
let subscriber = client.subscribe(subject.clone()).await.map_err(|e| {
|
||||
tracing::warn!(%subject, error = %e, "{what} stream: subscribe failed");
|
||||
tracing::warn!(%subject, error = %e, "term 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(),
|
||||
]
|
||||
);
|
||||
}
|
||||
tracing::info!(%subject, "term stream: client attached");
|
||||
let stream =
|
||||
subscriber.map(|msg| Ok(Event::default().data(String::from_utf8_lossy(&msg.payload))));
|
||||
Ok(Sse::new(stream).keep_alive(KeepAlive::default()))
|
||||
}
|
||||
|
|
|
|||
|
|
@ -22,17 +22,12 @@ 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
|
||||
|
|
|
|||
|
|
@ -1,382 +0,0 @@
|
|||
//! Verifying an agent's own credential, presented as `auth_token`.
|
||||
//!
|
||||
//! The spelling is `swarm_queue_client::agent_token`'s, shared with the agent
|
||||
//! that presents it: a prefix, the agent's name, and its secret. The name is only a
|
||||
//! claim. It becomes an identity when the secret equals the one stored at
|
||||
//! `swarm/agents/<agent>/queue`, and nothing in this module grants on the name
|
||||
//! alone: every path that does not reach that comparison, and every lookup
|
||||
//! that fails, is a denial.
|
||||
//!
|
||||
//! A token without the prefix is not this module's: it goes to introspection
|
||||
//! unchanged.
|
||||
|
||||
use std::future::Future;
|
||||
|
||||
use anyhow::Context;
|
||||
use subtle::ConstantTimeEq;
|
||||
use swarm_queue_client::agent_token::{AgentToken, Malformed, parse_agent_token};
|
||||
use swarm_secret_client::client::{DEFAULT_CERT_MOUNT, Settings};
|
||||
use swarm_secret_client::queue::{self, AgentCredential};
|
||||
|
||||
use crate::policy::{Permissions, Policy};
|
||||
|
||||
/// What a presented `auth_token` is.
|
||||
pub enum Presented<'a> {
|
||||
/// An agent's own credential, to be checked against the store.
|
||||
Agent(AgentToken<'a>),
|
||||
/// Carries the agent-token prefix but is not well-formed. Denied without a
|
||||
/// lookup, and never handed to introspection: it holds a secret meant for
|
||||
/// the store, not for the `IdP`.
|
||||
Malformed(Malformed),
|
||||
/// Anything else, which is an OIDC access token for introspection.
|
||||
Bearer(&'a str),
|
||||
}
|
||||
|
||||
/// Sort `token` onto the agent path or the introspection path.
|
||||
pub fn classify(token: &str) -> Presented<'_> {
|
||||
match parse_agent_token(token) {
|
||||
Some(Ok(agent)) => Presented::Agent(agent),
|
||||
Some(Err(e)) => Presented::Malformed(e),
|
||||
None => Presented::Bearer(token),
|
||||
}
|
||||
}
|
||||
|
||||
/// Where the stored credential for an agent comes from.
|
||||
pub trait CredentialSource {
|
||||
/// The object at `agent`'s queue path, `None` when nothing is stored
|
||||
/// there, or an error when the store could not answer.
|
||||
fn lookup(
|
||||
&self,
|
||||
agent: &str,
|
||||
) -> impl Future<Output = anyhow::Result<Option<AgentCredential>>> + Send;
|
||||
}
|
||||
|
||||
/// The swarm secret store, reached with this responder's own certificate.
|
||||
pub struct Store {
|
||||
settings: Settings,
|
||||
cert_role: String,
|
||||
}
|
||||
|
||||
impl Store {
|
||||
/// The store named by the `BAO_*` environment, or `None` when that is
|
||||
/// unset: a responder with no store identity denies every agent token and
|
||||
/// serves the OIDC path unchanged.
|
||||
pub fn from_env(cert_role: String) -> Option<Self> {
|
||||
match Settings::from_env() {
|
||||
Ok(settings) => Some(Self {
|
||||
settings,
|
||||
cert_role,
|
||||
}),
|
||||
Err(e) => {
|
||||
tracing::info!(reason = %e, "no secret store configured; agent tokens are denied");
|
||||
None
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl CredentialSource for Store {
|
||||
/// Logs in per lookup. An agent connects once per boot and per reconnect,
|
||||
/// so a login each time costs little, and it leaves no token to expire
|
||||
/// inside a long-lived process.
|
||||
async fn lookup(&self, agent: &str) -> anyhow::Result<Option<AgentCredential>> {
|
||||
let path = queue::agent_queue_path(agent)?;
|
||||
let store = swarm_secret_client::SecretStore::connect(
|
||||
&self.settings,
|
||||
&self.cert_role,
|
||||
DEFAULT_CERT_MOUNT,
|
||||
)
|
||||
.await
|
||||
.context("logging in to the secret store")?;
|
||||
store
|
||||
.read_optional(&path)
|
||||
.await
|
||||
.with_context(|| format!("reading {path}"))
|
||||
}
|
||||
}
|
||||
|
||||
/// Whether `token`'s secret is the one stored for the agent it names.
|
||||
///
|
||||
/// The lookup is bounded by introspection's budget, under the server's
|
||||
/// `authorization.timeout`. The callout loop answers one request at a time, so
|
||||
/// a store that does not answer holds every request behind it, OIDC ones
|
||||
/// included, for at most that long, and this then denies.
|
||||
async fn verify(source: &impl CredentialSource, token: &AgentToken<'_>) -> bool {
|
||||
let agent = token.agent;
|
||||
let stored = match tokio::time::timeout(
|
||||
crate::introspect::INTROSPECTION_TIMEOUT,
|
||||
source.lookup(agent),
|
||||
)
|
||||
.await
|
||||
{
|
||||
Ok(Ok(Some(stored))) => stored,
|
||||
Ok(Ok(None)) => {
|
||||
tracing::warn!(
|
||||
agent,
|
||||
"no queue credential is stored for this agent; denying"
|
||||
);
|
||||
return false;
|
||||
}
|
||||
Ok(Err(e)) => {
|
||||
tracing::warn!(
|
||||
agent,
|
||||
error = format!("{e:#}"),
|
||||
"credential lookup failed; denying"
|
||||
);
|
||||
return false;
|
||||
}
|
||||
Err(_) => {
|
||||
tracing::warn!(agent, "credential lookup timed out; denying");
|
||||
return false;
|
||||
}
|
||||
};
|
||||
if stored.agent != agent {
|
||||
tracing::warn!(
|
||||
agent,
|
||||
stored_agent = %stored.agent,
|
||||
"the stored credential names a different agent than its path; denying"
|
||||
);
|
||||
return false;
|
||||
}
|
||||
// Constant time, so the time to refuse says nothing about how much of the
|
||||
// secret was right. A length mismatch returns early; the length is not
|
||||
// secret, every minted secret has the same one.
|
||||
stored
|
||||
.value
|
||||
.as_bytes()
|
||||
.ct_eq(token.secret.as_bytes())
|
||||
.into()
|
||||
}
|
||||
|
||||
/// The grant for an agent token, or `None` for a denial.
|
||||
///
|
||||
/// `source` is `None` when this responder has no store identity, which denies.
|
||||
pub async fn authorize<S: CredentialSource>(
|
||||
policy: &Policy,
|
||||
source: Option<&S>,
|
||||
token: &AgentToken<'_>,
|
||||
) -> Option<Permissions> {
|
||||
let Some(source) = source else {
|
||||
tracing::warn!(
|
||||
agent = token.agent,
|
||||
"agent token presented, but this responder has no secret store; denying"
|
||||
);
|
||||
return None;
|
||||
};
|
||||
if !verify(source, token).await {
|
||||
return None;
|
||||
}
|
||||
let permissions = policy.agent_token_permissions(token.agent);
|
||||
if permissions.is_none() {
|
||||
tracing::warn!(
|
||||
agent = token.agent,
|
||||
"verified agent token, but no --agent-token-publish-subject is configured; denying"
|
||||
);
|
||||
}
|
||||
permissions
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
|
||||
/// A store holding at most one credential, or failing outright.
|
||||
enum Fake {
|
||||
/// `credential` is stored at `path_agent`'s queue path.
|
||||
Holds {
|
||||
path_agent: &'static str,
|
||||
credential: AgentCredential,
|
||||
},
|
||||
Fails,
|
||||
/// A store that accepts the request and never answers.
|
||||
Hangs,
|
||||
}
|
||||
|
||||
impl CredentialSource for Fake {
|
||||
async fn lookup(&self, agent: &str) -> anyhow::Result<Option<AgentCredential>> {
|
||||
match self {
|
||||
Self::Holds {
|
||||
path_agent,
|
||||
credential,
|
||||
} => Ok((agent == *path_agent).then(|| credential.clone())),
|
||||
Self::Fails => anyhow::bail!("store unreachable"),
|
||||
Self::Hangs => std::future::pending().await,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
fn stored(path_agent: &'static str, named: &str) -> Fake {
|
||||
Fake::Holds {
|
||||
path_agent,
|
||||
credential: AgentCredential {
|
||||
value: "Ab9_-zSECRET".to_owned(),
|
||||
agent: named.to_owned(),
|
||||
},
|
||||
}
|
||||
}
|
||||
|
||||
fn atlas() -> Fake {
|
||||
stored("atlas", "atlas")
|
||||
}
|
||||
|
||||
fn policy() -> Policy {
|
||||
Policy::new(
|
||||
"hive-".to_owned(),
|
||||
"-agent".to_owned(),
|
||||
"hive-status".to_owned(),
|
||||
vec!["swarm-controller".to_owned()],
|
||||
vec![],
|
||||
vec!["$SWARM.term.{hive}.>".to_owned()],
|
||||
)
|
||||
.expect("valid")
|
||||
.with_agent_token_subjects(vec![
|
||||
"$SWARM.term.{agent}".to_owned(),
|
||||
"$SWARM.agent-state.{agent}".to_owned(),
|
||||
])
|
||||
.expect("valid")
|
||||
}
|
||||
|
||||
async fn grant(source: Option<&Fake>, token: &str) -> Option<Permissions> {
|
||||
let Presented::Agent(token) = classify(token) else {
|
||||
panic!("{token:?} must classify as an agent token");
|
||||
};
|
||||
authorize(&policy(), source, &token).await
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn a_valid_agent_token_is_granted_exactly_its_own_subjects() {
|
||||
let g = grant(Some(&atlas()), "swarm-agent.atlas.Ab9_-zSECRET")
|
||||
.await
|
||||
.expect("granted");
|
||||
assert_eq!(
|
||||
g.publish,
|
||||
vec![
|
||||
"$SWARM.term.atlas".to_owned(),
|
||||
"$SWARM.agent-state.atlas".to_owned(),
|
||||
]
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn a_wrong_secret_is_denied() {
|
||||
assert!(
|
||||
grant(Some(&atlas()), "swarm-agent.atlas.Ab9_-zSECREX")
|
||||
.await
|
||||
.is_none()
|
||||
);
|
||||
assert!(
|
||||
grant(Some(&atlas()), "swarm-agent.atlas.Ab9_-zSECRE")
|
||||
.await
|
||||
.is_none()
|
||||
);
|
||||
}
|
||||
|
||||
/// The secret check is what stops one agent claiming another's name: the
|
||||
/// right secret under someone else's name finds that agent's credential,
|
||||
/// or none.
|
||||
#[tokio::test]
|
||||
async fn another_agents_name_with_this_agents_secret_is_denied() {
|
||||
assert!(
|
||||
grant(Some(&atlas()), "swarm-agent.argus.Ab9_-zSECRET")
|
||||
.await
|
||||
.is_none()
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn an_agent_with_nothing_stored_is_denied() {
|
||||
assert!(
|
||||
grant(
|
||||
Some(&stored("argus", "argus")),
|
||||
"swarm-agent.atlas.Ab9_-zSECRET"
|
||||
)
|
||||
.await
|
||||
.is_none()
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn a_failed_lookup_is_denied() {
|
||||
assert!(
|
||||
grant(Some(&Fake::Fails), "swarm-agent.atlas.Ab9_-zSECRET")
|
||||
.await
|
||||
.is_none()
|
||||
);
|
||||
}
|
||||
|
||||
/// The callout loop answers one request at a time, so a store that never
|
||||
/// answers must cost at most the bound, and then deny.
|
||||
#[tokio::test]
|
||||
async fn a_store_that_never_answers_is_denied_within_the_bound() {
|
||||
let bound = crate::introspect::INTROSPECTION_TIMEOUT;
|
||||
let started = std::time::Instant::now();
|
||||
assert!(
|
||||
grant(Some(&Fake::Hangs), "swarm-agent.atlas.Ab9_-zSECRET")
|
||||
.await
|
||||
.is_none()
|
||||
);
|
||||
let took = started.elapsed();
|
||||
assert!(took >= bound, "denied before the bound: {took:?}");
|
||||
assert!(
|
||||
took < bound + std::time::Duration::from_millis(250),
|
||||
"denied well after the bound: {took:?}"
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn no_store_is_denied() {
|
||||
assert!(
|
||||
grant(None, "swarm-agent.atlas.Ab9_-zSECRET")
|
||||
.await
|
||||
.is_none()
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn a_stored_object_naming_another_agent_is_denied() {
|
||||
assert!(
|
||||
grant(
|
||||
Some(&stored("argus", "atlas")),
|
||||
"swarm-agent.argus.Ab9_-zSECRET"
|
||||
)
|
||||
.await
|
||||
.is_none()
|
||||
);
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn a_verified_agent_with_no_subjects_configured_is_denied() {
|
||||
let bare = Policy::new(
|
||||
"hive-".to_owned(),
|
||||
"-agent".to_owned(),
|
||||
"hive-status".to_owned(),
|
||||
vec![],
|
||||
vec![],
|
||||
vec![],
|
||||
)
|
||||
.expect("valid");
|
||||
let Presented::Agent(token) = classify("swarm-agent.atlas.Ab9_-zSECRET") else {
|
||||
panic!("an agent token");
|
||||
};
|
||||
assert!(authorize(&bare, Some(&atlas()), &token).await.is_none());
|
||||
}
|
||||
|
||||
/// Everything without the prefix reaches introspection as it arrived.
|
||||
#[test]
|
||||
fn a_token_without_the_prefix_goes_to_introspection_unchanged() {
|
||||
for token in ["authelia_at_abc.def", "atlas.Ab9_-zSECRET"] {
|
||||
let Presented::Bearer(t) = classify(token) else {
|
||||
panic!("{token:?} is not an agent token");
|
||||
};
|
||||
assert_eq!(t, token);
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn a_malformed_agent_token_goes_nowhere() {
|
||||
assert!(matches!(
|
||||
classify("swarm-agent.atlas"),
|
||||
Presented::Malformed(_)
|
||||
));
|
||||
}
|
||||
}
|
||||
|
|
@ -8,8 +8,7 @@
|
|||
//! 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 (or, for an agent's own token, against the
|
||||
//! secret store: see `agent_token`), and answers with a NATS user JWT signed
|
||||
//! authelia's introspection endpoint, 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.
|
||||
|
|
@ -28,7 +27,6 @@ use anyhow::Context;
|
|||
use clap::Parser;
|
||||
use futures_util::StreamExt;
|
||||
|
||||
mod agent_token;
|
||||
mod introspect;
|
||||
mod policy;
|
||||
mod request;
|
||||
|
|
@ -91,7 +89,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 under it.
|
||||
/// agent: two agents on one hive are indistinguishable to this responder.
|
||||
///
|
||||
/// A suffix on the hive's id rather than a prefix of its own, because
|
||||
/// `agent-<name>` reads as *the agent called `<name>`* — the one thing
|
||||
|
|
@ -129,18 +127,6 @@ struct Args {
|
|||
/// is refused, which is loud rather than silently over-broad.
|
||||
#[arg(long = "agent-publish-subject")]
|
||||
agent_publish_subjects: Vec<String>,
|
||||
|
||||
/// Subjects an agent that presented its own credential may publish to,
|
||||
/// with `{agent}` standing for its name. Repeatable, empty by default,
|
||||
/// which denies every agent token.
|
||||
#[arg(long = "agent-token-publish-subject")]
|
||||
agent_token_publish_subjects: Vec<String>,
|
||||
|
||||
/// Role on the secret store's `cert` auth mount this responder logs in
|
||||
/// as to read agent credentials. The store is found through the `BAO_*`
|
||||
/// environment; with none, every agent token is denied.
|
||||
#[arg(long, default_value = "swarm-nats-auth")]
|
||||
store_cert_role: String,
|
||||
}
|
||||
|
||||
/// Read a secret file and strip surrounding whitespace.
|
||||
|
|
@ -191,9 +177,7 @@ 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")?;
|
||||
|
|
@ -226,33 +210,44 @@ 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. 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),
|
||||
// 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,
|
||||
};
|
||||
// The caller is an identifier, not a credential: a client id, or
|
||||
// `agent:<name>` for an agent token. One line per auth request, a connect
|
||||
// or a reconnect, so `hive-<h>-agent` lines count requests made on a
|
||||
// hive's shared credential, not agents.
|
||||
// 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.
|
||||
tracing::info!(
|
||||
user_nkey = %req.user_nkey,
|
||||
server_id = %req.server_id.id,
|
||||
|
|
@ -286,50 +281,6 @@ async fn main() -> anyhow::Result<()> {
|
|||
anyhow::bail!("subscription to {AUTH_SUBJECT} ended")
|
||||
}
|
||||
|
||||
/// The OIDC path: who the `IdP` says presented `token`, and what the policy
|
||||
/// grants that client. `None` permissions is a denial.
|
||||
///
|
||||
/// The caller is an identity or nothing — see `introspect`'s module docs.
|
||||
/// There is no "admitted, identity unknown" branch to write here because
|
||||
/// there is no such value to receive.
|
||||
async fn introspected(
|
||||
policy: &policy::Policy,
|
||||
http: &reqwest::Client,
|
||||
args: &Args,
|
||||
client_secret: &str,
|
||||
token: &str,
|
||||
) -> (Option<String>, Option<policy::Permissions>) {
|
||||
let caller = introspect::identify_caller(
|
||||
http,
|
||||
&args.introspection_url,
|
||||
&args.client_id,
|
||||
client_secret,
|
||||
token,
|
||||
)
|
||||
.await
|
||||
// An introspection that could not be *made* is a denial too. The
|
||||
// failure modes of an HTTP call are exactly the conditions under
|
||||
// which an attacker would most like this to fall open.
|
||||
.unwrap_or_else(|e| {
|
||||
tracing::warn!(error = ?e, "introspection failed; denying");
|
||||
None
|
||||
});
|
||||
// Admission said who; the policy says what. A caller the `IdP`
|
||||
// vouches for but no rule matches is denied — see `policy`'s module
|
||||
// docs for why that is deny and not "connect with nothing".
|
||||
let permissions = caller.as_deref().and_then(|id| policy.permissions(id));
|
||||
if let (Some(id), None) = (caller.as_deref(), permissions.as_ref()) {
|
||||
// Loud, and the one case an operator has to be able to find: a
|
||||
// valid credential refused by our own policy. The alternative is
|
||||
// a client that authenticates fine and mysteriously cannot work.
|
||||
tracing::warn!(
|
||||
caller = %id,
|
||||
"authenticated client matches no policy rule; denying"
|
||||
);
|
||||
}
|
||||
(caller, permissions)
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
|
|
@ -355,16 +306,9 @@ 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
|
||||
|
|
|
|||
|
|
@ -45,16 +45,12 @@ pub struct Policy {
|
|||
readers: Vec<String>,
|
||||
extra_hive_subjects: Vec<String>,
|
||||
extra_agent_subjects: Vec<String>,
|
||||
agent_token_subjects: Vec<String>,
|
||||
}
|
||||
|
||||
/// Placeholder replaced with the hive's own name in `extra_hive_subjects` and
|
||||
/// `extra_agent_subjects`.
|
||||
const HIVE_PLACEHOLDER: &str = "{hive}";
|
||||
|
||||
/// Placeholder replaced with the agent's own name in `agent_token_subjects`.
|
||||
const AGENT_PLACEHOLDER: &str = "{agent}";
|
||||
|
||||
impl Policy {
|
||||
/// `hive_prefix` is the client-id prefix that marks a hive and
|
||||
/// `agent_suffix` what a hive's agent containers carry **on top of** it —
|
||||
|
|
@ -134,43 +130,9 @@ impl Policy {
|
|||
readers,
|
||||
extra_hive_subjects,
|
||||
extra_agent_subjects,
|
||||
agent_token_subjects: Vec::new(),
|
||||
})
|
||||
}
|
||||
|
||||
/// Set the subjects an agent that proved its own credential may publish
|
||||
/// to, with `{agent}` standing for its name.
|
||||
///
|
||||
/// # Errors
|
||||
///
|
||||
/// A template with no `{agent}` in it is refused: it would be one subject
|
||||
/// shared by every agent in the swarm rather than the agent's own.
|
||||
pub fn with_agent_token_subjects(mut self, subjects: Vec<String>) -> anyhow::Result<Self> {
|
||||
if let Some(bad) = subjects.iter().find(|s| !s.contains(AGENT_PLACEHOLDER)) {
|
||||
anyhow::bail!(
|
||||
"--agent-token-publish-subject {bad:?} contains no {AGENT_PLACEHOLDER}: every \
|
||||
agent in the swarm would be granted that exact subject"
|
||||
);
|
||||
}
|
||||
self.agent_token_subjects = subjects;
|
||||
Ok(self)
|
||||
}
|
||||
|
||||
/// The permissions for an agent whose own credential was verified, or
|
||||
/// `None` when none are configured.
|
||||
///
|
||||
/// Keyed on the agent alone: its identity is not tied to a hive. `agent`
|
||||
/// must already be a single `[A-Za-z0-9_-]` segment, which the token parse
|
||||
/// guarantees, so it cannot widen a subject with `.`, `*` or `>`.
|
||||
pub fn agent_token_permissions(&self, agent: &str) -> Option<Permissions> {
|
||||
let publish: Vec<String> = self
|
||||
.agent_token_subjects
|
||||
.iter()
|
||||
.map(|s| s.replace(AGENT_PLACEHOLDER, agent))
|
||||
.collect();
|
||||
(!publish.is_empty()).then_some(Permissions { publish })
|
||||
}
|
||||
|
||||
/// The permissions for `client_id`, or `None` when no rule matches.
|
||||
///
|
||||
/// `None` is a denial. It is not "grant nothing and let them connect":
|
||||
|
|
@ -1213,35 +1175,4 @@ 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());
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -1,155 +0,0 @@
|
|||
//! The one spelling of the token an agent presents its own queue secret in,
|
||||
//! shared by the agent that formats it and the auth-callout responder that
|
||||
//! parses it back:
|
||||
//! [`AGENT_TOKEN_PREFIX`](crate::agent_token::AGENT_TOKEN_PREFIX), the agent's
|
||||
//! name, `.`, the secret.
|
||||
//!
|
||||
//! The secret itself lives at `swarm/agents/<agent>/queue` in the swarm's
|
||||
//! secret store (`swarm_secret_client::queue`). This module knows nothing of
|
||||
//! the store, so an agent formats its token without linking a store client.
|
||||
|
||||
/// Marks an `auth_token` as an agent's own credential.
|
||||
///
|
||||
/// An OIDC access token may itself contain `.`, so the `<agent>.<secret>`
|
||||
/// shape alone does not tell the two apart. Authelia's access tokens start
|
||||
/// `authelia_at_`.
|
||||
pub const AGENT_TOKEN_PREFIX: &str = "swarm-agent.";
|
||||
|
||||
/// An agent's own credential as presented at the queue.
|
||||
///
|
||||
/// No `Debug`: `secret` is the credential itself.
|
||||
pub struct AgentToken<'a> {
|
||||
/// The agent the presenter claims to be. Unproven until `secret` is
|
||||
/// checked against that agent's stored credential.
|
||||
pub agent: &'a str,
|
||||
/// The presented secret.
|
||||
pub secret: &'a str,
|
||||
}
|
||||
|
||||
/// A token that is not `<agent>.<secret>` after its prefix, or parts that
|
||||
/// would not make one. Carries the reason only: a token holds a secret.
|
||||
#[derive(Debug, thiserror::Error)]
|
||||
#[error("malformed agent token: {0}")]
|
||||
pub struct Malformed(pub &'static str);
|
||||
|
||||
/// Spell `agent`'s credential as the token it presents at the queue.
|
||||
///
|
||||
/// # Errors
|
||||
/// [`Malformed`] when `agent` or `secret` is empty or holds anything outside
|
||||
/// `[A-Za-z0-9_-]`. Either would make [`parse_agent_token`] split the token
|
||||
/// differently than it was joined.
|
||||
pub fn format_agent_token(agent: &str, secret: &str) -> Result<String, Malformed> {
|
||||
check_agent(agent)?;
|
||||
check_secret(secret)?;
|
||||
Ok(format!("{AGENT_TOKEN_PREFIX}{agent}.{secret}"))
|
||||
}
|
||||
|
||||
/// Read a presented token back into an [`AgentToken`].
|
||||
///
|
||||
/// `None` when `token` does not start with [`AGENT_TOKEN_PREFIX`]: it is some
|
||||
/// other kind of token, not a malformed one of these. `Some(Err(_))` when it
|
||||
/// does and is not exactly `<agent>.<secret>` after it.
|
||||
#[must_use]
|
||||
pub fn parse_agent_token(token: &str) -> Option<Result<AgentToken<'_>, Malformed>> {
|
||||
let rest = token.strip_prefix(AGENT_TOKEN_PREFIX)?;
|
||||
Some(
|
||||
rest.split_once('.')
|
||||
.ok_or(Malformed("no `.` between the agent and the secret"))
|
||||
.and_then(|(agent, secret)| {
|
||||
check_agent(agent)?;
|
||||
check_secret(secret)?;
|
||||
Ok(AgentToken { agent, secret })
|
||||
}),
|
||||
)
|
||||
}
|
||||
|
||||
fn is_segment(s: &str) -> bool {
|
||||
!s.is_empty()
|
||||
&& s.bytes()
|
||||
.all(|b| b.is_ascii_alphanumeric() || b == b'-' || b == b'_')
|
||||
}
|
||||
|
||||
/// The alphabet a store path segment allows, so the name cannot widen a
|
||||
/// subject with `.`, `*` or `>`, or address another agent's path.
|
||||
fn check_agent(agent: &str) -> Result<(), Malformed> {
|
||||
if is_segment(agent) {
|
||||
Ok(())
|
||||
} else {
|
||||
Err(Malformed("the agent is not a single [A-Za-z0-9_-] segment"))
|
||||
}
|
||||
}
|
||||
|
||||
/// The alphabet the controller mints secrets in, base64url without padding.
|
||||
/// The error never carries the value.
|
||||
fn check_secret(secret: &str) -> Result<(), Malformed> {
|
||||
if is_segment(secret) {
|
||||
Ok(())
|
||||
} else {
|
||||
Err(Malformed(
|
||||
"the secret is empty or holds a byte outside [A-Za-z0-9_-]",
|
||||
))
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
|
||||
#[test]
|
||||
fn an_agent_token_round_trips() {
|
||||
let token = format_agent_token("atlas", "Ab9_-z").expect("legal");
|
||||
assert_eq!(token, "swarm-agent.atlas.Ab9_-z");
|
||||
let parsed = parse_agent_token(&token)
|
||||
.expect("carries the prefix")
|
||||
.expect("well-formed");
|
||||
assert_eq!(parsed.agent, "atlas");
|
||||
assert_eq!(parsed.secret, "Ab9_-z");
|
||||
}
|
||||
|
||||
/// Anything without the prefix belongs to the OIDC path, including a
|
||||
/// token that happens to look like `<agent>.<secret>`.
|
||||
#[test]
|
||||
fn a_token_without_the_prefix_is_not_an_agent_token() {
|
||||
for token in ["authelia_at_abc.def", "atlas.s3cr3t", "", "swarm-agent"] {
|
||||
assert!(parse_agent_token(token).is_none(), "{token:?}");
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn a_malformed_agent_token_is_refused() {
|
||||
for token in [
|
||||
"swarm-agent.",
|
||||
"swarm-agent.atlas",
|
||||
"swarm-agent.atlas.",
|
||||
"swarm-agent..s3cr3t",
|
||||
"swarm-agent.atlas.s3.cr3t",
|
||||
"swarm-agent.at*las.s3cr3t",
|
||||
"swarm-agent.at>las.s3cr3t",
|
||||
"swarm-agent.at/las.s3cr3t",
|
||||
"swarm-agent.atlas.s3cr3t=",
|
||||
"swarm-agent.atlas.s3 cr3t",
|
||||
] {
|
||||
assert!(
|
||||
matches!(parse_agent_token(token), Some(Err(_))),
|
||||
"{token:?} must be refused"
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
/// What formats always parses back into the same two parts.
|
||||
#[test]
|
||||
fn formatting_refuses_what_parsing_would_split_differently() {
|
||||
assert!(format_agent_token("at.las", "s3cr3t").is_err());
|
||||
assert!(format_agent_token("atlas", "s3.cr3t").is_err());
|
||||
assert!(format_agent_token("atlas", "").is_err());
|
||||
assert!(format_agent_token("", "s3cr3t").is_err());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn a_malformed_secret_is_not_echoed_in_the_error() {
|
||||
let Some(Err(e)) = parse_agent_token("swarm-agent.atlas.hunter2!") else {
|
||||
panic!("must be refused");
|
||||
};
|
||||
assert!(!e.to_string().contains("hunter2"), "{e}");
|
||||
}
|
||||
}
|
||||
|
|
@ -178,10 +178,6 @@ 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.
|
||||
///
|
||||
|
|
@ -320,15 +316,6 @@ 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
|
||||
|
|
@ -751,7 +738,15 @@ pub async fn connect(cfg: QueueConfig) -> Result<async_nats::Client, Error> {
|
|||
// which turns an unreachable queue into a permanent 4s poll — and, before
|
||||
// the cache above, a permanent 4s token-request loop against authelia.
|
||||
// Exponential from 500ms so a momentary blip still reconnects promptly.
|
||||
.reconnect_delay_callback(reconnect_delay)
|
||||
.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,
|
||||
)
|
||||
})
|
||||
// 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
|
||||
|
|
@ -785,29 +780,6 @@ pub async fn connect(cfg: QueueConfig) -> Result<async_nats::Client, Error> {
|
|||
Ok(client)
|
||||
}
|
||||
|
||||
/// Connect presenting a fixed `token`, as an agent presents its own queue
|
||||
/// credential. Reconnects present the same token.
|
||||
///
|
||||
/// Unlike [`connect`], the first attempt must succeed: a refused or failed
|
||||
/// connect is returned rather than retried in the background, so the caller
|
||||
/// can fall back to another credential. `ca_file` is the queue's trust anchor,
|
||||
/// as [`QueueConfig::ca_file`].
|
||||
pub async fn connect_with_token(
|
||||
url: &str,
|
||||
ca_file: Option<&std::path::Path>,
|
||||
token: String,
|
||||
) -> Result<async_nats::Client, Error> {
|
||||
let mut options =
|
||||
async_nats::ConnectOptions::with_token(token).reconnect_delay_callback(reconnect_delay);
|
||||
if let Some(path) = ca_file {
|
||||
options = options.add_root_certificates(path.to_path_buf());
|
||||
}
|
||||
options.connect(url).await.map_err(|source| Error::Connect {
|
||||
url: url.to_owned(),
|
||||
source,
|
||||
})
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
|
|
|
|||
|
|
@ -22,9 +22,6 @@
|
|||
//! 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};
|
||||
|
||||
|
|
@ -60,11 +57,14 @@ pub fn agent_queue_path(agent: &str) -> Result<String, Error> {
|
|||
Ok(format!("{prefix}/queue"))
|
||||
}
|
||||
|
||||
/// What [`agent_queue_path`] holds: the secret, and the agent it proves.
|
||||
/// What [`agent_queue_path`] holds: the secret, and the principal it proves.
|
||||
///
|
||||
/// 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.
|
||||
/// Both names ride **in the object** rather than being parsed back out of a
|
||||
/// composite principal string. Hive and agent names draw from the same
|
||||
/// alphabet (`hive_types::Ident`, `[a-z0-9-]`), so a principal spelled
|
||||
/// `hive-<hive>-agent-<agent>` parses two ways for a name containing `-agent-`
|
||||
/// — and an ambiguous principal parse in an authorisation path is a caller that
|
||||
/// authenticates fine and is handed somebody else's grant.
|
||||
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
|
||||
pub struct AgentCredential {
|
||||
/// The secret itself. Named to match [`Credential::value`] and
|
||||
|
|
@ -74,6 +74,13 @@ 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.
|
||||
|
|
@ -191,6 +198,7 @@ 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!(
|
||||
|
|
@ -204,39 +212,29 @@ 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!(json.get("hive").is_none(), "{json}");
|
||||
assert_eq!(json["hive"], "alpha");
|
||||
}
|
||||
|
||||
/// A stored object may carry a `hive` field. It is ignored, and the object
|
||||
/// decodes.
|
||||
/// 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.
|
||||
#[test]
|
||||
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(),
|
||||
}
|
||||
fn an_agent_object_missing_a_principal_does_not_decode() {
|
||||
assert!(
|
||||
serde_json::from_str::<AgentCredential>(r#"{"value":"s","agent":"atlas"}"#).is_err()
|
||||
);
|
||||
}
|
||||
|
||||
/// The agent is required: it is what the verifier checks the presented
|
||||
/// name against.
|
||||
#[test]
|
||||
fn an_agent_object_missing_its_agent_does_not_decode() {
|
||||
assert!(serde_json::from_str::<AgentCredential>(r#"{"value":"s"}"#).is_err());
|
||||
assert!(
|
||||
serde_json::from_str::<AgentCredential>(r#"{"value":"s","hive":"alpha"}"#).is_err()
|
||||
);
|
||||
}
|
||||
|
||||
/// The two kinds are different objects at different paths, and neither
|
||||
/// decodes as the other — the property that keeps a reader from picking up
|
||||
/// a hive-shared credential where a per-agent one was meant.
|
||||
#[test]
|
||||
|
|
@ -251,6 +249,7 @@ mod tests {
|
|||
let agent_json = serde_json::to_string(&AgentCredential {
|
||||
value: "s".to_owned(),
|
||||
agent: "atlas".to_owned(),
|
||||
hive: "alpha".to_owned(),
|
||||
})
|
||||
.expect("serialises");
|
||||
assert!(serde_json::from_str::<Credential>(&agent_json).is_err());
|
||||
|
|
|
|||
|
|
@ -51,6 +51,13 @@ struct CreateAgentResponse {
|
|||
warnings: Vec<String>,
|
||||
}
|
||||
|
||||
/// Body of `POST /api/agents/{name}/identity`. Mirrors the controller's own
|
||||
/// `MintAgentIdentityRequest` — see this module's doc comment.
|
||||
#[derive(Serialize)]
|
||||
struct MintIdentityRequest<'a> {
|
||||
hive: &'a str,
|
||||
}
|
||||
|
||||
/// Success body of `POST /api/agents/{name}/identity`.
|
||||
#[derive(Deserialize)]
|
||||
struct MintIdentityResponse {
|
||||
|
|
@ -110,10 +117,11 @@ pub(crate) fn parse_ident(value: &str, what: &str) -> Result<String> {
|
|||
///
|
||||
/// Synchronous for the same reason [`create`] is, and built on the same
|
||||
/// round trip.
|
||||
pub(crate) fn mint_identity(socket: &Path, name: &str) -> Result<()> {
|
||||
pub(crate) fn mint_identity(socket: &Path, name: &str, hive: &str) -> Result<()> {
|
||||
// Client-side first, so a typo is a local error rather than a 400 the
|
||||
// operator waits for. The controller validates it again.
|
||||
// operator waits for. The controller validates both again.
|
||||
let name = parse_ident(name, "agent name")?;
|
||||
let hive = parse_ident(hive, "hive")?;
|
||||
|
||||
let rt = tokio::runtime::Builder::new_current_thread()
|
||||
.enable_io()
|
||||
|
|
@ -122,15 +130,15 @@ pub(crate) fn mint_identity(socket: &Path, name: &str) -> Result<()> {
|
|||
let resp: MintIdentityResponse = rt.block_on(post(
|
||||
socket,
|
||||
&format!("/api/agents/{name}/identity"),
|
||||
&serde_json::json!({}),
|
||||
&MintIdentityRequest { hive: &hive },
|
||||
"mint-identity",
|
||||
))?;
|
||||
|
||||
println!("queued: job node {}", resp.node_id);
|
||||
println!(
|
||||
"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"
|
||||
"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"
|
||||
);
|
||||
Ok(())
|
||||
}
|
||||
|
|
|
|||
|
|
@ -242,6 +242,16 @@ 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
|
||||
|
|
@ -347,7 +357,7 @@ fn main() -> Result<()> {
|
|||
command: AgentVerb::MintIdentity(args),
|
||||
} => {
|
||||
let socket = path_from(args.controller_socket, "SWARM_CONTROLLER_SOCKET")?;
|
||||
agent::mint_identity(&socket, &args.name)
|
||||
agent::mint_identity(&socket, &args.name, &args.hive)
|
||||
}
|
||||
// Same socket-resolution reasoning as `Create` above.
|
||||
Verb::Agent {
|
||||
|
|
@ -722,9 +732,19 @@ 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_no_hive() {
|
||||
let cli = Cli::try_parse_from(["swarmctl", "agent", "mint-identity", "scribe"])
|
||||
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");
|
||||
let Verb::Agent {
|
||||
command: AgentVerb::MintIdentity(args),
|
||||
|
|
@ -733,18 +753,12 @@ 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",
|
||||
"--hive",
|
||||
"a"
|
||||
])
|
||||
.is_err(),
|
||||
"the identity has no hive, so the verb must not take one"
|
||||
Cli::try_parse_from(["swarmctl", "agent", "mint-identity", "scribe"]).is_err(),
|
||||
"an omitted hive must not be defaulted"
|
||||
);
|
||||
}
|
||||
|
||||
|
|
|
|||
Loading…
Reference in a new issue