diff --git a/Cargo.lock b/Cargo.lock index ed646445..194abcef 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1690,6 +1690,7 @@ dependencies = [ "serde", "serde_json", "swarm-queue-client", + "swarm-secret-client", "tempfile", "tokio", "tokio-stream", diff --git a/docs/swarm/credentials.md b/docs/swarm/credentials.md index bc1bc7fa..8723a49f 100644 --- a/docs/swarm/credentials.md +++ b/docs/swarm/credentials.md @@ -59,20 +59,20 @@ of the cell says how. -| store path | minter | reader — pulls at runtime, holds in memory | automatic re-mint | automatic re-pull | -| ----------------------------------------------------------- | ---------------------------------------------------------------------------------------------------------------------------------------- | ------------------------------------------------------------------------------------------------------------------------------------------------------------------ | -------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- | --------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- | -| `swarm/agents//matrix/main` | `swarm-controller`, with the swarm's appservice token, at agent creation and in a five-minute pass | the agent container itself, under the certificate its hive passed in | ✅ the pass re-mints when the stored token is missing, unknown to the homeserver, or someone else's | ✅ `hive-matrix-daemon` exits when the homeserver rejects its token, and a five-minute timer restarts it, which reads the store again | -| `swarm/agents//matrix/` | `swarm-controller` | the agent container itself, under the certificate its hive passed in | must be stated | must be stated | -| `swarm/controller/swarm-controller/matrix/appservice-token` | `swarm-matrix-ctl`, inside the `hive-matrix` container, once | `swarm-controller`, under its own certificate | ❌ `swarm-matrix-ctl` mints it once; the container keeps its copy and republishes it when the store's differs | ✅ the controller reads it on every five-minute matrix pass | -| `swarm/controller/swarm-controller/oidc/client` | authelia, at its first boot, where the controller registers its client; `swarm-secret-publish` copies it in | `swarm-controller`, under its own certificate, once at start | ❌ authelia mints it once. A re-mint is republished by `swarm-secret-publish`'s path unit | ❌ read once at start; the controller holds the old value until it restarts | -| `swarm/agents//bao-mtls` | the store's agent PKI mount (`deploy.bao.agentPkiMountPath`), which generates the key, at `swarm-controller`'s request at agent creation | `hive-c0re`, under the hive's own certificate, when it writes the agent's container config | ✅ `swarm-controller`'s five-minute pass re-issues a live agent's leaf once it's past half its validity (45 of 90 days, read from the certificate itself) | ❌ `hive-c0re` reads it when it writes the container config, so the agent presents a new leaf from its next start; the old leaf stays valid until it expires | -| `swarm/agents//queue` | `swarm-controller`, at agent creation | the agent container itself, under its own certificate — the identity it presents to the swarm queue, naming that one agent rather than its hive | ✅ `swarm-controller`'s five-minute pass re-mints a live agent's secret once it's 45 days old by `minted_at` on the stored object; a secret with no `minted_at` gets one stamped, value unchanged. The pass skips agents declared `Destroyed` — declaring an agent destroyed deletes every version of the path instead, the undo of the mint rather than another one | ❌ fetched when the container starts. The queue checks the secret only at connect, so an open connection survives a re-mint, but a reconnect before the next restart is denied — including a reconnect after the credential was revoked | -| `swarm/agents//forge-token` | `swarm-controller`, at agent creation and in a pass every 5 minutes over every agent with a store identity | the agent container itself, under its own certificate, fetched to `/run/hive-agent-forge-token/token` | ✅ the controller re-mints when the stored token is missing or no longer matches the forge (last eight characters and scopes) | ✅ the agent re-fetches on a 10-minute timer | -| `swarm/hives//matrix/appservice-token` | one minter, on the authelia host | the hive process that presents the token to its homeserver, under the hive's own certificate | must be stated | must be stated | -| `swarm/hives//matrix/sender-token` | `swarm-matrix-ctl`, in the `hive-matrix` container | `swarm-matrix-ctl` itself, under its own certificate, before it decides whether to mint, and hive-c0re's `stored_sender_token()`, under the hive's own certificate | must be stated | must be stated | -| `swarm/hives//queue/agent` | authelia | `swarm-bao-queue-agent` on the hive's host, under its own per-hive certificate; no agent's policy reaches it | must be stated | must be stated | -| `swarm/services//oidc/client` | authelia | the service process that presents the client secret, under the certificate of the host it runs on | must be stated | must be stated | -| _(not in the store)_ a hive's mTLS leaf | the store's own PKI, or an operator placing it by hand | its own client, off disk — the exception above, because it's what makes every other row's pull possible | must be stated | must be stated | +| store path | minter | reader — pulls at runtime, holds in memory | automatic re-mint | automatic re-pull | +| ----------------------------------------------------------- | ---------------------------------------------------------------------------------------------------------------------------------------- | -------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- | -------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- | ----------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- | +| `swarm/agents//matrix/main` | `swarm-controller`, with the swarm's appservice token, at agent creation and in a five-minute pass | the agent container itself, under the certificate its hive passed in | ✅ the pass re-mints when the stored token is missing, unknown to the homeserver, or someone else's | ✅ `hive-matrix-daemon` exits when the homeserver rejects its token, and a five-minute timer restarts it, which reads the store again | +| `swarm/agents//matrix/` | `swarm-controller` | the agent container itself, under the certificate its hive passed in | must be stated | must be stated | +| `swarm/controller/swarm-controller/matrix/appservice-token` | `swarm-matrix-ctl`, inside the `hive-matrix` container, once | `swarm-controller`, under its own certificate | ❌ `swarm-matrix-ctl` mints it once; the container keeps its copy and republishes it when the store's differs | ✅ the controller reads it on every five-minute matrix pass | +| `swarm/controller/swarm-controller/oidc/client` | authelia, at its first boot, where the controller registers its client; `swarm-secret-publish` copies it in | `swarm-controller`, under its own certificate, once at start | ❌ authelia mints it once. A re-mint is republished by `swarm-secret-publish`'s path unit | ❌ read once at start; the controller holds the old value until it restarts | +| `swarm/agents//bao-mtls` | the store's agent PKI mount (`deploy.bao.agentPkiMountPath`), which generates the key, at `swarm-controller`'s request at agent creation | `hive-c0re`, under the hive's own certificate, when it writes the agent's container config | ✅ `swarm-controller`'s five-minute pass re-issues a live agent's leaf once it's past half its validity (45 of 90 days, read from the certificate itself) | ❌ `hive-c0re` reads it when it writes the container config, so the agent presents a new leaf from its next start; the old leaf stays valid until it expires | +| `swarm/agents//queue` | `swarm-controller`, at agent creation | `hive-agent` in the agent container, under the agent's own certificate, held in memory — the identity it presents to the swarm queue, naming that one agent rather than its hive | ✅ `swarm-controller`'s five-minute pass re-mints a live agent's secret once it's 45 days old by `minted_at` on the stored object; a secret with no `minted_at` gets one stamped, value unchanged. The pass skips agents declared `Destroyed` — declaring an agent destroyed deletes every version of the path instead, the undo of the mint rather than another one | ✅ `hive-agent` reads the path before its first connect and again on every reconnect attempt, so a reconnect after a re-mint presents the new secret. An open connection keeps the secret it connected with; after a revocation the agent keeps retrying under the queue client's backoff | +| `swarm/agents//forge-token` | `swarm-controller`, at agent creation and in a pass every 5 minutes over every agent with a store identity | the agent container itself, under its own certificate, fetched to `/run/hive-agent-forge-token/token` | ✅ the controller re-mints when the stored token is missing or no longer matches the forge (last eight characters and scopes) | ✅ the agent re-fetches on a 10-minute timer | +| `swarm/hives//matrix/appservice-token` | one minter, on the authelia host | the hive process that presents the token to its homeserver, under the hive's own certificate | must be stated | must be stated | +| `swarm/hives//matrix/sender-token` | `swarm-matrix-ctl`, in the `hive-matrix` container | `swarm-matrix-ctl` itself, under its own certificate, before it decides whether to mint, and hive-c0re's `stored_sender_token()`, under the hive's own certificate | must be stated | must be stated | +| `swarm/hives//queue/agent` | authelia | `swarm-bao-queue-agent` on the hive's host, under its own per-hive certificate; no agent's policy reaches it | must be stated | must be stated | +| `swarm/services//oidc/client` | authelia | the service process that presents the client secret, under the certificate of the host it runs on | must be stated | must be stated | +| _(not in the store)_ a hive's mTLS leaf | the store's own PKI, or an operator placing it by hand | its own client, off disk — the exception above, because it's what makes every other row's pull possible | must be stated | must be stated | diff --git a/hive-agent/Cargo.toml b/hive-agent/Cargo.toml index 6571aff1..a34b24df 100644 --- a/hive-agent/Cargo.toml +++ b/hive-agent/Cargo.toml @@ -40,6 +40,8 @@ serde_json.workspace = true # of. The terminal publisher is a plain core-subject publish, so it needs the # connect and the payload limit and nothing from JetStream. swarm-queue-client.workspace = true +# This agent's own queue secret, read from the store under its own certificate. +swarm-secret-client.workspace = true tokio.workspace = true tokio-stream.workspace = true tower-http.workspace = true diff --git a/hive-agent/src/swarm_queue.rs b/hive-agent/src/swarm_queue.rs index b7aa4bc1..23cb2017 100644 --- a/hive-agent/src/swarm_queue.rs +++ b/hive-agent/src/swarm_queue.rs @@ -14,26 +14,33 @@ //! here over the inputs this consumer actually has, and [`decide`] is the //! single place it lives. //! -//! Beside all four sits a fifth coordinate, resolved by -//! [`decide_agent_secret`] and not part of their group. The four are the -//! *hive's* — one OIDC client shared by every container on it — so at the -//! queue's auth callout they say which hive is connecting and never which -//! 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`). +//! Beside all four sits a fifth coordinate, resolved by [`store_source`] and +//! not part of their group. The four are the *hive's* — one OIDC client +//! shared by every container on it — so at the queue's auth callout they say +//! which hive is connecting and never which agent. The fifth is this agent's +//! own, minted per agent at swarm level. This process reads it from the swarm +//! secret store under the agent's own certificate and holds it in memory only +//! ([`AgentSecret`]). //! -//! 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. +//! When the store holds one, [`client`] connects with it first, and the agent +//! publishes on its own hive-free subjects. When it holds none, or the queue +//! refuses it, the first connect falls back to the hive's shared client and +//! the hive-scoped subjects. [`Connection::presented`] says which, so a +//! publisher builds the subject that credential is granted. use std::path::{Path, PathBuf}; -use std::sync::OnceLock; +use std::sync::{Arc, Mutex, OnceLock}; +use std::time::Duration; -use anyhow::Context as _; +use anyhow::{Context as _, anyhow}; use tokio::sync::OnceCell; use swarm_queue_client::{ClientSecret, QueueConfig}; +use swarm_secret_client::{ + SecretStore, + client::{DEFAULT_CERT_MOUNT, ENV_ADDR, ENV_CACERT, Settings}, + policy, queue, +}; /// Variable prefix for this agent's coordinates. Distinct from `HIVE_C0RE`'s /// on purpose: an agent authenticates as its own client, not as its hive. @@ -43,6 +50,15 @@ const ENV_PREFIX: &str = "HIVE_AGENT"; /// one answer rather than re-deriving it per call. static CONFIG: OnceLock> = OnceLock::new(); +/// Names the agent whose queue secret this process reads. The cert-auth role +/// and the store path are both built from it by `swarm_secret_client`. +const ENV_AGENT_NAME: &str = "HIVE_AGENT_NAME"; + +/// How long one read of the store may take. The read runs inside a queue +/// connection attempt, and a store that never answers would otherwise stall +/// every reconnect behind it. +const STORE_READ_TIMEOUT: Duration = Duration::from_secs(10); + /// Where to present this agent's own credential, resolved once at boot. static AGENT: OnceLock> = OnceLock::new(); @@ -56,8 +72,148 @@ static CLIENT: OnceCell> = OnceCell::const_new(); struct AgentPath { url: String, ca_file: Option, - /// The fetched secret. A path; the bytes are read at connect. - secret_file: PathBuf, + store: Arc, +} + +/// This agent's queue secret as the store holds it: where the store is, the +/// role this agent logs in as, and the path it reads. +#[derive(Debug, PartialEq, Eq)] +struct StoreSource { + settings: Settings, + role: String, + path: String, +} + +/// Where [`AgentSecret`] reads the secret from. A trait so the retry rules can +/// be exercised without a store. +trait SecretSource: Send + Sync + 'static { + /// The store path, for log lines and errors. Never the value. + fn origin(&self) -> &str; + + /// The value stored now, or `None` when nothing is stored. + fn fetch(&self) -> impl Future>> + Send; +} + +impl SecretSource for Arc { + fn origin(&self) -> &str { + &self.path + } + + /// A fresh login per read: reads happen once per connection attempt, and a + /// token held between them would be one more secret in memory to expire. + async fn fetch(&self) -> anyhow::Result> { + let store = SecretStore::connect(&self.settings, &self.role, DEFAULT_CERT_MOUNT) + .await + .context("logging in to the swarm secret store as this agent")?; + // `read_optional`: only a 404 is "nothing stored". A 403 stays an + // error, so a policy that stopped covering the path is not reported + // as a secret that was deleted. + let credential: Option = store + .read_optional(&self.path) + .await + .with_context(|| format!("reading {} from the store", self.path))?; + Ok(credential.map(|c| c.value)) + } +} + +/// This agent's queue secret, held in memory between connection attempts. +/// +/// No `Debug`: [`Held`] carries the secret. +struct AgentSecret { + source: S, + timeout: Duration, + held: Mutex, +} + +#[derive(Default)] +struct Held { + value: Option, + /// `value` was read by [`AgentSecret::prime`] and no attempt has + /// presented it yet. + unpresented: bool, +} + +impl AgentSecret { + fn new(source: S) -> Self { + Self { + source, + timeout: STORE_READ_TIMEOUT, + held: Mutex::default(), + } + } + + fn held(&self) -> std::sync::MutexGuard<'_, Held> { + // A poisoned lock only means another attempt panicked mid-update; the + // worst it left is a stale value, which the next read replaces. + self.held + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner) + } + + /// Read the secret before the first connection attempt. An `Err` here is + /// what sends [`connect_once`] to the hive's shared client. + async fn prime(&self) -> anyhow::Result<()> { + let value = self + .read() + .await? + .with_context(|| format!("nothing is stored at {}", self.source.origin()))?; + *self.held() = Held { + value: Some(value), + unpresented: true, + }; + Ok(()) + } + + /// The secret to present on one connection attempt. + /// + /// The value [`prime`](Self::prime) read serves the first attempt. Every + /// later attempt reads the store again, so after the queue refuses a + /// re-minted secret the next attempt presents the current one. A store + /// that cannot be read leaves the held value in use; a store that holds + /// nothing drops it. Either way the attempt is one read, and the next one + /// waits out the connection's reconnect backoff. + async fn for_attempt(&self) -> anyhow::Result { + { + let mut held = self.held(); + if std::mem::take(&mut held.unpresented) + && let Some(value) = &held.value + { + return Ok(value.clone()); + } + } + match self.read().await { + Ok(Some(value)) => { + self.held().value = Some(value.clone()); + Ok(value) + } + Ok(None) => { + self.held().value = None; + tracing::warn!( + path = %self.source.origin(), + "nothing is stored at this agent's queue credential path; its queue \ + connection retries under backoff" + ); + Err(anyhow!("nothing is stored at {}", self.source.origin())) + } + Err(e) => { + let held = self.held().value.clone(); + tracing::warn!( + error = format!("{e:#}"), + presenting_held = held.is_some(), + "re-reading this agent's queue credential from the store failed" + ); + held.ok_or(e) + } + } + } + + /// One bounded read. An empty stored value is no value. + async fn read(&self) -> anyhow::Result> { + let value = tokio::time::timeout(self.timeout, self.source.fetch()) + .await + .map_err(|_| anyhow!("the store did not answer within {:?}", self.timeout))??; + Ok(value.map(|v| v.trim().to_owned()).filter(|v| !v.is_empty())) + } } /// Which credential the shared connection presented. @@ -92,13 +248,6 @@ struct QueueEnv { /// `hive_c0re::meta` embeds at build time — so it is here for a /// deployment that needs a different one, not for ours. ca_file: Option, - /// Where `nix/agent-modules/queue-identity.nix` fetched this agent's own - /// per-agent secret to. Outside the all-or-none group above because it is - /// governed by a different switch entirely — that unit is generated by - /// the agent having a *store* address, not by its hive having queue - /// coordinates — so an agent can legally have this and none of the four, - /// or the four and not this. - agent_secret_file: Option, } impl QueueEnv { @@ -110,7 +259,6 @@ impl QueueEnv { client_id_file: var("OIDC_CLIENT_ID_FILE"), client_secret_file: var("OIDC_CLIENT_SECRET_FILE"), ca_file: var("OIDC_CA_FILE"), - agent_secret_file: var("QUEUE_AGENT_SECRET_FILE"), } } } @@ -147,43 +295,47 @@ fn read_client_id(path: &Path) -> Option { } } -/// Decide whether this agent has its own per-agent queue secret, given the -/// path the fetch unit was told to write it to. +/// Where this agent reads its own queue secret from, or `None` when this +/// container was given no store. /// -/// Three states collapse to two answers. No variable means no store address -/// for this container, so no fetch unit was generated at all. A variable -/// naming a file that is missing or empty means the unit ran and found -/// nothing minted — the ordinary state of an agent created before its swarm -/// knew to mint one, which that unit reports and survives. Only a non-empty -/// file is a credential. +/// Taken through a lookup so it can be asserted without mutating the process +/// environment. The shape follows `hive_matrix_mcp::credential`'s +/// `store_settings`, which reads the same variables for the same identity. /// -/// The file is not read. Its *contents* are the secret and belong nowhere but -/// the moment of use; what a caller needs from here is whether there is one -/// and where, which `metadata` answers without opening it. -fn decide_agent_secret(path: Option<&str>) -> Option { - let path = PathBuf::from(path?); - match std::fs::metadata(&path) { - Ok(m) if m.len() > 0 => Some(path), - Ok(_) => None, - Err(e) if e.kind() == std::io::ErrorKind::NotFound => None, - Err(e) => { - tracing::warn!( - path = %path.display(), - error = %e, - "checking for this agent's own queue credential failed" - ); - None - } +/// # Errors +/// A store address without the agent name or the identity beside it, which +/// is a half-delivered container rather than one without a store. +fn store_source(get: impl Fn(&str) -> Option) -> anyhow::Result> { + if get(ENV_ADDR).is_none_or(|v| v.is_empty()) { + return Ok(None); } + let agent = get(ENV_AGENT_NAME) + .filter(|v| !v.is_empty()) + .with_context(|| format!("{ENV_AGENT_NAME} is unset or empty"))?; + let settings = Settings::from_lookup(|k| get(k).filter(|v| k != ENV_CACERT || ca_is_usable(v))) + .context("reading the swarm secret store's coordinates from the environment")?; + Ok(Some(StoreSource { + settings, + role: policy::agent_object_name(&agent)?, + path: queue::agent_queue_path(&agent)?, + })) +} + +/// Whether the CA bundle at `path` is a file with bytes in it. `BAO_CACERT` +/// names a systemd credential that may not have been delivered; absent means +/// the container's own trust store. +fn ca_is_usable(path: &str) -> bool { + std::fs::metadata(path).is_ok_and(|m| m.len() > 0) } /// Decide whether this agent can connect with its own credential: it needs the -/// queue's address and a fetched secret, and nothing of the hive's client. -fn decide_agent_path(env: &QueueEnv, secret_file: Option) -> Option { +/// queue's address and a store to read its secret from, and nothing of the +/// hive's client. +fn decide_agent_path(env: &QueueEnv, store: Option) -> Option { Some(AgentPath { url: env.nats_url.clone()?, ca_file: env.ca_file.as_ref().map(Into::into), - secret_file: secret_file?, + store: Arc::new(store?), }) } @@ -258,20 +410,26 @@ pub fn init() { // 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 { - // The path, never the bytes: the file holds the secret itself. + let store = store_source(|k| std::env::var(k).ok()).unwrap_or_else(|e| { + tracing::warn!( + error = format!("{e:#}"), + "this agent's secret store coordinates are incomplete, so it cannot read \ + its own swarm queue credential" + ); + None + }); + if let Some(store) = &store { tracing::info!( - path = %path.display(), - "this agent has its own swarm queue credential" + path = %store.path, + "this agent reads its own swarm queue credential from the secret store" ); } else { tracing::info!( - "no per-agent swarm queue credential; this agent is known to the queue \ - by its hive's shared client" + "no secret store for this agent; it is known to the queue by its hive's \ + shared client" ); } - let _ = AGENT.set(decide_agent_path(&env, secret)); + let _ = AGENT.set(decide_agent_path(&env, store)); } /// Whether this agent has any credential to reach the queue with. `false` @@ -313,7 +471,9 @@ pub async fn client() -> Option { /// logged once, here. async fn connect_once() -> Option { if let Some(agent) = AGENT.get().and_then(Option::as_ref) { - match connect_as_agent(agent).await { + let secret = AgentSecret::new(Arc::clone(&agent.store)); + let name = crate::identity::label(); + match connect_as_agent(&agent.url, agent.ca_file.as_deref(), name, secret).await { Ok(client) => { tracing::info!( url = %agent.url, @@ -360,25 +520,47 @@ async fn connect_once() -> Option { } } -/// Connect presenting this agent's own token. Any failure, a refusal -/// included, is returned for [`connect_once`] to fall back on. -async fn connect_as_agent(agent: &AgentPath) -> anyhow::Result { - let name = crate::identity::label(); - let secret = std::fs::read_to_string(&agent.secret_file) - .with_context(|| format!("reading {}", agent.secret_file.display()))?; - let token = swarm_queue_client::agent_token::format_agent_token(&name, secret.trim())?; - swarm_queue_client::connect_with_token(&agent.url, agent.ca_file.as_deref(), token) - .await - .map_err(|e| anyhow::anyhow!(swarm_queue_client::chain(&e))) +/// Connect presenting this agent's own token as `name`. +/// +/// The secret is read before the first attempt, and any failure of that read +/// or of the first attempt, a refusal included, is returned for +/// [`connect_once`] to fall back on. Once connected, each reconnect attempt +/// takes its secret from [`AgentSecret::for_attempt`]. +async fn connect_as_agent( + url: &str, + ca_file: Option<&Path>, + name: String, + secret: AgentSecret, +) -> anyhow::Result { + secret.prime().await?; + let secret = Arc::new(secret); + swarm_queue_client::connect_with_token(url, ca_file, move || { + let secret = Arc::clone(&secret); + let name = name.clone(); + // Spawned so the future handed back is only a `JoinHandle`, which is + // `Sync` as async-nats's auth callback requires; the store read is not. + let attempt = tokio::spawn(async move { + let value = secret.for_attempt().await.map_err(|e| format!("{e:#}"))?; + swarm_queue_client::agent_token::format_agent_token(&name, &value) + .map_err(|e| e.to_string()) + }); + async move { attempt.await.map_err(|e| e.to_string())? } + }) + .await + .map_err(|e| anyhow!(swarm_queue_client::chain(&e))) } - #[cfg(test)] mod tests { - use std::path::PathBuf; + use std::collections::VecDeque; + use std::sync::atomic::{AtomicUsize, Ordering}; + use std::sync::{Arc, Mutex}; + use std::time::Duration; + + use anyhow::anyhow; use super::{ - AgentPath, ClientSecret, QueueEnv, Resolution, decide, decide_agent_path, - decide_agent_secret, read_client_id, + AgentPath, AgentSecret, ClientSecret, QueueEnv, Resolution, SecretSource, StoreSource, + connect_as_agent, decide, decide_agent_path, read_client_id, store_source, }; fn env(parts: [Option<&str>; 4]) -> QueueEnv { @@ -389,7 +571,6 @@ mod tests { client_id_file: client_id_file.map(str::to_owned), client_secret_file: client_secret_file.map(str::to_owned), ca_file: None, - agent_secret_file: None, } } @@ -484,81 +665,313 @@ mod tests { } } - /// The whole point of the fetch: a non-empty file is this agent's own - /// credential, and the answer is the path rather than what is in it. + /// A lookup standing in for a container that was handed a store. + fn with_store(k: &str) -> Option { + match k { + "BAO_ADDR" => Some("https://bao.t.local:8200".to_owned()), + "BAO_CLIENT_CERT" => Some("/run/credentials/x/cert".to_owned()), + "BAO_CLIENT_KEY" => Some("/run/credentials/x/key".to_owned()), + "HIVE_AGENT_NAME" => Some("a1".to_owned()), + _ => None, + } + } + + fn source() -> StoreSource { + store_source(with_store) + .expect("a complete store environment") + .expect("a store address was given") + } + #[test] - fn a_fetched_secret_resolves_to_its_path() { - let dir = tempfile::tempdir().expect("tempdir"); - let path = dir.path().join("secret"); - std::fs::write(&path, "s3cr3t").expect("write"); + fn a_store_address_reads_the_agents_own_path_as_its_own_role() { + let s = source(); + assert_eq!(s.path, "swarm/agents/a1/queue"); + assert_eq!(s.role, "hive-agent-a1"); + } + + /// No address: this container was given no store, so there is nothing to + /// read from and nothing wrong. + #[test] + fn no_store_address_is_no_store() { + let none = store_source(|k| if k == "BAO_ADDR" { None } else { with_store(k) }); + assert!(none.expect("absent is not an error").is_none()); + let empty = store_source(|k| { + if k == "BAO_ADDR" { + Some(String::new()) + } else { + with_store(k) + } + }); + assert!(empty.expect("empty is absent").is_none()); + } + + #[test] + fn a_store_address_without_the_agent_or_its_identity_is_an_error() { + for var in ["HIVE_AGENT_NAME", "BAO_CLIENT_CERT", "BAO_CLIENT_KEY"] { + let e = store_source(|k| if k == var { None } else { with_store(k) }) + .expect_err("a half-delivered store"); + assert!(format!("{e:#}").contains(var), "dropping {var}: {e:#}"); + } + } + + /// `BAO_CACERT` is always named and may point at nothing; that is the + /// container's own trust store, not a CA path that fails the handshake. + #[test] + fn an_undelivered_ca_is_dropped() { + let with_missing_ca = store_source(|k| { + if k == "BAO_CACERT" { + Some("/nonexistent/ca".to_owned()) + } else { + with_store(k) + } + }) + .expect("complete") + .expect("given a store"); + assert_eq!(with_missing_ca, source()); + } + + /// The per-agent secret is resolved by a separate switch from the hive's + /// client: an agent whose hive has no queue client can still read its own + /// swarm-minted secret, and needs only the queue's address for it. + #[test] + fn an_agent_with_a_store_and_the_queue_address_connects_as_itself() { + let url_only = env([Some("nats://10.42.0.1:4222"), None, None, None]); + assert!(!matches!( + decide(&url_only, None), + Resolution::Configured(_) + )); assert_eq!( - decide_agent_secret(path.to_str()).as_deref(), - Some(path.as_path()) - ); - } - - /// The rollout state, and the one this must not confuse with a - /// credential: the fetch unit ran, found nothing minted for this agent, - /// and left no file. Treating that as a secret would have the harness - /// present zero bytes to the queue. - #[test] - fn a_missing_or_empty_fetched_secret_is_no_credential() { - let dir = tempfile::tempdir().expect("tempdir"); - let path = dir.path().join("secret"); - assert_eq!(decide_agent_secret(path.to_str()), None); - std::fs::write(&path, "").expect("write"); - assert_eq!(decide_agent_secret(path.to_str()), None); - } - - /// No variable at all: this container was given no store address, so no - /// fetch unit exists to have written anything. - #[test] - fn no_fetch_path_is_no_credential() { - assert_eq!(decide_agent_secret(None), None); - } - - /// The two credentials are resolved by separate switches, and this is the - /// asymmetry that makes keeping them apart worth it: an agent whose hive - /// has no queue can still hold its own swarm-minted secret, because that - /// one is minted with no hive in the chain. - #[test] - fn the_per_agent_secret_is_independent_of_the_hive_coordinates() { - let dir = tempfile::tempdir().expect("tempdir"); - let path = dir.path().join("secret"); - std::fs::write(&path, "s3cr3t").expect("write"); - - let e = env([None, None, None, None]); - 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())), + decide_agent_path(&url_only, Some(source())), Some(AgentPath { url: "nats://10.42.0.1:4222".to_owned(), ca_file: None, - secret_file: secret.clone(), + store: Arc::new(source()), }) ); - 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() { + fn without_a_store_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")) - ), + decide_agent_path(&env([None, None, None, None]), Some(source())), None ); } + + enum Answer { + Value(&'static str), + Absent, + Fails, + Hangs, + } + + /// A store that answers each read with the next scripted answer and + /// counts the reads. + struct Scripted { + answers: Mutex>, + reads: AtomicUsize, + } + + fn scripted(answers: impl IntoIterator) -> Arc { + Arc::new(Scripted { + answers: Mutex::new(answers.into_iter().collect()), + reads: AtomicUsize::new(0), + }) + } + + impl SecretSource for Arc { + fn origin(&self) -> &'static str { + "swarm/agents/a1/queue" + } + + async fn fetch(&self) -> anyhow::Result> { + self.reads.fetch_add(1, Ordering::SeqCst); + let next = self.answers.lock().expect("lock").pop_front(); + match next.expect("a read nobody scripted") { + Answer::Value(v) => Ok(Some(v.to_owned())), + Answer::Absent => Ok(None), + Answer::Fails => Err(anyhow!("the store is unreachable")), + Answer::Hangs => std::future::pending().await, + } + } + } + + fn reads(s: &Arc) -> usize { + s.reads.load(Ordering::SeqCst) + } + + /// (a) The first attempt presents what was just read, and does not read + /// the store a second time for it. + #[tokio::test] + async fn the_first_attempt_presents_the_value_read_before_it() { + let store = scripted([Answer::Value("v1\n")]); + let secret = AgentSecret::new(Arc::clone(&store)); + secret.prime().await.expect("the store holds a secret"); + assert_eq!(secret.for_attempt().await.expect("primed"), "v1"); + assert_eq!(reads(&store), 1); + } + + /// (b) After the queue refuses a re-minted secret, the next attempt reads + /// the store and presents what it holds now: one read per attempt. + #[tokio::test] + async fn each_later_attempt_presents_what_the_store_holds_now() { + let store = scripted([ + Answer::Value("v1"), + Answer::Value("v1"), + Answer::Value("v2"), + ]); + let secret = AgentSecret::new(Arc::clone(&store)); + secret.prime().await.expect("the store holds a secret"); + assert_eq!(secret.for_attempt().await.expect("primed"), "v1"); + assert_eq!(secret.for_attempt().await.expect("re-read"), "v1"); + assert_eq!(secret.for_attempt().await.expect("re-read"), "v2"); + assert_eq!(reads(&store), 3); + } + + /// (b) Bounded: a store that never answers costs one attempt the read + /// timeout, and the attempt goes ahead with the value it already holds. + #[tokio::test] + async fn a_store_that_does_not_answer_is_given_up_on() { + let store = scripted([Answer::Value("v1"), Answer::Hangs, Answer::Hangs]); + let mut secret = AgentSecret::new(Arc::clone(&store)); + secret.timeout = Duration::from_millis(20); + secret.prime().await.expect("the store holds a secret"); + assert_eq!(secret.for_attempt().await.expect("primed"), "v1"); + assert_eq!( + secret + .for_attempt() + .await + .expect("the held value stands in"), + "v1" + ); + assert_eq!(reads(&store), 2); + + let cold = AgentSecret::new(scripted([Answer::Hangs])); + let mut cold = cold; + cold.timeout = Duration::from_millis(20); + let e = cold.prime().await.expect_err("nothing to stand in"); + assert!(format!("{e:#}").contains("did not answer"), "{e:#}"); + } + + /// (b) A secret the store no longer holds (the agent was destroyed) is not + /// presented again, even when a later read fails. + #[tokio::test] + async fn a_secret_removed_from_the_store_is_dropped() { + let store = scripted([Answer::Value("v1"), Answer::Absent, Answer::Fails]); + let secret = AgentSecret::new(Arc::clone(&store)); + secret.prime().await.expect("the store holds a secret"); + assert_eq!(secret.for_attempt().await.expect("primed"), "v1"); + let gone = secret.for_attempt().await.expect_err("nothing stored"); + assert!( + format!("{gone:#}").contains("swarm/agents/a1/queue"), + "{gone:#}" + ); + secret + .for_attempt() + .await + .expect_err("the dropped value does not come back"); + } + + /// (c) Nothing to connect with before the first attempt is an error, which + /// is what sends `connect_once` to the hive's shared client. The queue is + /// never dialled: nothing listens at the address given. + #[tokio::test] + async fn no_secret_before_the_first_connect_falls_back() { + for (answer, says) in [ + (Answer::Fails, "unreachable"), + (Answer::Absent, "nothing is stored"), + (Answer::Value(" \n"), "nothing is stored"), + ] { + let secret = AgentSecret::new(scripted([answer])); + let e = connect_as_agent("nats://127.0.0.1:9", None, "a1".to_owned(), secret) + .await + .expect_err("no secret to present"); + assert!(format!("{e:#}").contains(says), "{e:#}"); + } + } + + /// Serve one NATS handshake on `sock`: send INFO, read CONNECT up to the + /// PING, and return the token presented. + async fn handshake(sock: &mut tokio::net::TcpStream, port: u16) -> String { + use tokio::io::{AsyncBufReadExt as _, AsyncWriteExt as _}; + let info = format!( + "INFO {{\"server_id\":\"fake\",\"version\":\"2.10.0\",\"go\":\"go\",\ + \"host\":\"127.0.0.1\",\"port\":{port},\"headers\":true,\ + \"max_payload\":1048576,\"proto\":1,\"auth_required\":true,\ + \"nonce\":\"n0nce\"}}\r\n" + ); + sock.write_all(info.as_bytes()).await.expect("send INFO"); + let mut lines = tokio::io::BufReader::new(sock).lines(); + let mut token = String::new(); + while let Some(line) = lines.next_line().await.expect("read") { + if let Some(json) = line.strip_prefix("CONNECT ") { + let v: serde_json::Value = serde_json::from_str(json).expect("CONNECT is JSON"); + token = v["auth_token"].as_str().unwrap_or_default().to_owned(); + } else if line.starts_with("PING") { + break; + } + } + token + } + + /// (b) end to end through async-nats: connected with `v1`, the connection + /// drops, the reconnect presenting `v1` is refused, and the next attempt + /// reads `v2` and is accepted. + #[tokio::test] + async fn a_refused_reconnect_reconnects_with_the_secret_read_after_it() { + use tokio::io::AsyncWriteExt as _; + + let listener = tokio::net::TcpListener::bind("127.0.0.1:0") + .await + .expect("bind"); + let port = listener.local_addr().expect("addr").port(); + let (seen_tx, mut seen_rx) = tokio::sync::mpsc::unbounded_channel(); + tokio::spawn(async move { + let mut kept = Vec::new(); + for n in 0.. { + let (mut sock, _) = listener.accept().await.expect("accept"); + let token = handshake(&mut sock, port).await; + let accept = n == 0 || token.rsplit_once('.').is_some_and(|(_, s)| s == "v2"); + let reply: &[u8] = if accept { + b"PONG\r\n" + } else { + b"-ERR 'Authorization Violation'\r\n" + }; + let _ = sock.write_all(reply).await; + let _ = seen_tx.send(token); + // The first connection is dropped to force a reconnect; the + // one accepted after the refusal stays open. + if accept && n > 0 { + kept.push(sock); + } + } + }); + + let store = scripted([ + Answer::Value("v1"), + Answer::Value("v1"), + Answer::Value("v2"), + ]); + let secret = AgentSecret::new(Arc::clone(&store)); + let url = format!("nats://127.0.0.1:{port}"); + let _client = connect_as_agent(&url, None, "a1".to_owned(), secret) + .await + .expect("the first connect is accepted"); + + let mut seen = Vec::new(); + while seen.len() < 3 { + let token = tokio::time::timeout(Duration::from_secs(20), seen_rx.recv()) + .await + .expect("the client keeps reconnecting") + .expect("the server is running"); + seen.push(token); + } + let secrets: Vec<&str> = seen + .iter() + .map(|t| t.rsplit_once('.').map_or("", |(_, s)| s)) + .collect(); + assert_eq!(secrets, ["v1", "v1", "v2"]); + assert_eq!(reads(&store), 3); + } } diff --git a/nix/agent-modules/forge-token.nix b/nix/agent-modules/forge-token.nix index 5f081f91..6268d0e5 100644 --- a/nix/agent-modules/forge-token.nix +++ b/nix/agent-modules/forge-token.nix @@ -5,15 +5,14 @@ # to `swarm/agents//forge-token` # (`swarm_secret_client::forge::agent_token_path`); no hive is in that chain. # This unit logs in to the store with the certificate ./bao.nix already proves -# it can log in with, and reads its own path. The shape is ./queue-identity.nix's, -# for the same reason: a hive handing the token over would be the hive reading -# a secret on the agent's behalf. +# it can log in with, and reads its own path. A hive handing the token over +# would be the hive reading a secret on the agent's behalf. # -# Unlike the queue credential this one rotates: the controller replaces the -# token when the forge's copy stops matching the stored one. So the unit is -# re-run by a timer, and it swaps the file in by rename, only when the value -# changed, so a reader never sees half a token and a watcher on the file -# (forge-avatar-sync.path) fires only on a real change. +# The token rotates: the controller replaces it when the forge's copy stops +# matching the stored one. So the unit is re-run by a timer, and it swaps the +# file in by rename, only when the value changed, so a reader never sees half +# a token and a watcher on the file (forge-avatar-sync.path) fires only on a +# real change. # # Consumers read `services.hyperhive.agent.forge.tokenFile` first and fall back # to `/forge-token`, the file the hive wrote before this existed. @@ -79,7 +78,9 @@ in description = "fetch this agent's own forge token from the secret store"; after = [ "network.target" - # Ordering only, for the reason ./queue-identity.nix gives. + # Ordering only: that unit reports a store this agent cannot reach as + # itself, and should get to say so before this one reports a path it + # could not read. "hive-agent-bao-identity.service" ]; before = [ "hive-forge-notify.service" ]; @@ -93,10 +94,10 @@ in startLimitIntervalSec = 300; serviceConfig = { Type = "oneshot"; - # Not `RemainAfterExit`, unlike the queue fetch: the timer below has to - # be able to start this unit again, and an active unit cannot be - # started. `RuntimeDirectoryPreserve` is what keeps the directory, and - # the token in it, alive between runs instead. + # Not `RemainAfterExit`: the timer below has to be able to start this + # unit again, and an active unit cannot be started. + # `RuntimeDirectoryPreserve` is what keeps the directory, and the token + # in it, alive between runs instead. RemainAfterExit = false; TimeoutStartSec = 30; Restart = "on-failure"; @@ -167,8 +168,8 @@ in fi export BAO_TOKEN - # 🩸 Degrades rather than fails, for the reason ./queue-identity.nix - # gives: the policy stanza that let the login read `bao-mtls` covers + # 🩸 Degrades where ./bao.nix's check fails, because one policy stanza + # governs both reads: the one that let the login read `bao-mtls` covers # this path too, so a refusal here is a token not minted yet. The # file already in place, if any, is kept: a store that is briefly # unreachable must not take a working token away. diff --git a/nix/agent-modules/queue-identity.nix b/nix/agent-modules/queue-identity.nix index 01f03812..148ed360 100644 --- a/nix/agent-modules/queue-identity.nix +++ b/nix/agent-modules/queue-identity.nix @@ -1,5 +1,5 @@ -# This agent's own credential for the swarm queue, fetched from the store by -# the agent itself. +# This agent's own credential for the swarm queue, read from the store by the +# harness itself. # # The credential beside it in ./queue.nix is keyed per **hive**: one OIDC # client minted at deploy time and handed to every agent container on the @@ -9,17 +9,21 @@ # (`swarm_secret_client::queue::agent_queue_path`), at swarm level, with no # hive anywhere in the chain. # -# That is also why this unit *fetches* rather than being handed a credential: -# a hive courier in the path would be the hive vouching for which agent this -# is, which is the property the per-agent credential exists to remove. The -# agent authenticates to the store as itself, with the certificate ./bao.nix -# already proves it can log in with, and reads its own path. +# The agent authenticates to the store as itself, with the certificate +# ./bao.nix already proves it can log in with, and reads its own path. A hive +# courier in the path would be the hive vouching for which agent this is, +# which is the property the per-agent credential exists to remove. +# +# Nothing here writes the secret anywhere. The harness reads it into memory +# before its first queue connection and again on every reconnect attempt +# (`hive-agent`'s `swarm_queue`), so a re-minted secret is picked up without a +# restart. What this module hands the harness is the store's address and the +# paths of its certificate. # # ⛔ The store certificate is for reaching the store and nothing else. It is # never presented to the queue: what goes to the queue is the secret read # back from this path. { - pkgs, lib, config, ... @@ -29,32 +33,17 @@ let # This container's agent name — the same string `swarm-controller` minted # the credential under, because the agent's unix user is named for the - # agent (see ./user.nix). + # agent (see ./user.nix). The harness builds the path and the cert-auth + # role from it. agentName = config.services.hyperhive.agent.user.name; - # The same three ids ./bao.nix loads. Both units present the same - # certificate because there is one store identity per agent; the ids are - # `hive_c0re::lifecycle::agent_identity`'s and a rename is a rename there - # too. + # The same three ids ./bao.nix loads, because there is one store identity + # per agent; the ids are `hive_c0re::lifecycle::agent_identity`'s and a + # rename is a rename there too. certCredential = "hive-agent-bao-cert"; keyCredential = "hive-agent-bao-key"; serverCaCredential = "hive-agent-bao-server-ca"; - unitName = "hive-agent-queue-credential"; - - # The nix half of `swarm_secret_client::queue::agent_queue_path` plus - # `path::MOUNT`, spelled exactly as ./bao.nix spells its own sibling path. - queuePath = "secret/swarm/agents/${agentName}/queue"; - - # `RuntimeDirectory=` under the unit's own `User=`, so the file is owned by - # the agent and readable by the harness without a mode change. /run and not - # the state dir on purpose: a secret fetched at boot has no business - # surviving one. - runtimeDir = "${unitName}"; - secretFile = "/run/${runtimeDir}/secret"; - # bao's stderr, in the unit's own `0700` directory rather than `/tmp`. - errFile = "/run/${runtimeDir}/bao.err"; - # The store's address is the whole switch, exactly as in ./bao.nix — and # deliberately *not* the queue coordinates in ./queue.nix. Those are the # hive's, and gating a swarm-minted per-agent credential on a hive-level @@ -63,167 +52,27 @@ let configured = cfg.addr != null; in { - options.services.hyperhive.agent.queue.agentSecretFile = lib.mkOption { - type = lib.types.str; - readOnly = true; - default = secretFile; - description = '' - Path this agent's own swarm-queue secret is fetched to, for a consumer - outside the harness unit. Read-only for the same reason - {option}`services.hyperhive.agent.queue.clientSecretFile` is: it is a - fact about where the fetch writes, not a knob. - - 🩸 A PATH and never a value. The file is `0400` to the agent user and - is read at the moment it is needed; nothing in this tree puts its - contents in an environment variable, where `/proc//environ` would - publish them to every process in the container. - - The file exists only once the swarm has minted a credential for this - agent. An agent created before its swarm did so has none, and - `${unitName}.service` says so in the journal rather than failing — - see that unit. - ''; - }; - config = lib.mkIf configured { - systemd.services.${unitName} = { - description = "fetch this agent's own swarm-queue credential from the secret store"; - after = [ - "network.target" - # Ordering only, not a requirement: that unit is the one that reports - # a store this agent cannot reach as itself, and it should get to say - # so before this one reports a path it could not read. - "hive-agent-bao-identity.service" + systemd.services.hive-agent = { + # Bare ids, no paths: the terse `LoadCredential=` form that inherits a + # credential the service *manager* received, which is what the container + # manager passed in. ./bao.nix and ./queue.nix state the same shape. + serviceConfig.LoadCredential = [ + certCredential + keyCredential + serverCaCredential ]; - before = [ "hive-agent.service" ]; - wantedBy = [ "multi-user.target" ]; - path = [ - pkgs.openbao - pkgs.coreutils - ]; - # Same sizing and the same `[Unit]`-not-`[Service]` placement as - # ./bao.nix's check: a few short attempts cover a store that comes up - # alongside this container, and a longer window only delays the report. - startLimitBurst = 4; - startLimitIntervalSec = 300; - serviceConfig = { - Type = "oneshot"; - # Keeps the unit active, which is what keeps `RuntimeDirectory=` - # from being removed out from under the harness. - RemainAfterExit = true; - TimeoutStartSec = 30; - Restart = "on-failure"; - RestartSec = 15; - User = agentName; - Group = agentName; - RuntimeDirectory = runtimeDir; - RuntimeDirectoryMode = "0700"; - UMask = "0377"; - # Bare ids, no paths: the terse `LoadCredential=` form that inherits - # a credential the service *manager* received, which is what the - # container manager passed in. ./bao.nix and ./queue.nix state the - # same shape. - LoadCredential = [ - certCredential - keyCredential - serverCaCredential - ]; - }; + # Paths and a name only. `%d` is this unit's own credentials directory. environment = { + HIVE_AGENT_NAME = agentName; BAO_ADDR = cfg.addr; - # `%d` is `$CREDENTIALS_DIRECTORY`, per-unit and owned by `User=`. BAO_CLIENT_CERT = "%d/${certCredential}"; BAO_CLIENT_KEY = "%d/${keyCredential}"; + # Named even when no CA was delivered: a bare `LoadCredential=` is + # non-fatal when absent, and the harness treats a missing or empty file + # as "use the container's own trust store". + BAO_CACERT = "%d/${serverCaCredential}"; }; - script = '' - set -euo pipefail - - # No identity delivered at all. ./bao.nix's check reports this as the - # failure it is; there is nothing for this unit to add, and failing - # here too would be the same cause stated twice. - for id in ${lib.escapeShellArg certCredential} ${lib.escapeShellArg keyCredential}; do - if [ ! -s "$CREDENTIALS_DIRECTORY/$id" ]; then - echo "this agent has no store identity, so it cannot fetch its own queue credential." >&2 - exit 0 - fi - if [ ! -r "$CREDENTIALS_DIRECTORY/$id" ]; then - echo "cannot read $CREDENTIALS_DIRECTORY/$id, so this agent cannot present its store identity." >&2 - exit 1 - fi - done - - # Only when one was delivered — absent means verify the store's - # listener against the container's own trust store, which is what a - # deployment with a real CA wants. ./bao.nix says the same. - if [ -s "$CREDENTIALS_DIRECTORY/${serverCaCredential}" ]; then - export BAO_CACERT="$CREDENTIALS_DIRECTORY/${serverCaCredential}" - fi - - # `UMask=0377` makes every file this script creates `0400`, so only - # the redirect that creates a file can write to it: `$err` is removed - # before each redirect into it. - err=${lib.escapeShellArg errFile} - trap 'rm -f "$err"' EXIT - - # Cert auth is a login, not a transport setting: the `BAO_CLIENT_*` - # variables above only pick the certificate the handshake presents. - # `-token-only` answers on stdout and skips the token helper, which - # is a `sh` this unit's `path` does not carry. - rm -f "$err" - if ! BAO_TOKEN="$(bao login -method=cert -token-only 2>"$err")"; then - # The redirect creates `$err` before bao starts, so no file means - # bao never ran. bao prints `Code: ` only for an HTTP - # answer, and `remote error: tls:` only for an alert the store sent. - re='Code: ([0-9]{3})' - if [ ! -e "$err" ]; then - echo "could not create $err, so bao never ran and the store at $BAO_ADDR was not asked." >&2 - elif [[ "$(<"$err")" =~ $re ]]; then - case "''${BASH_REMATCH[1]}" in - 4*) echo "the swarm secret store at $BAO_ADDR refused this agent's certificate login with HTTP ''${BASH_REMATCH[1]}:" >&2 ;; - *) echo "the swarm secret store at $BAO_ADDR failed this agent's certificate login with HTTP ''${BASH_REMATCH[1]}:" >&2 ;; - esac - elif [[ "$(<"$err")" == *"remote error: tls:"* ]]; then - echo "the swarm secret store at $BAO_ADDR refused this agent's certificate in the TLS handshake:" >&2 - else - echo "bao got no answer from the swarm secret store at $BAO_ADDR (network, DNS, or TLS on this side):" >&2 - fi - if [ -s "$err" ]; then cat "$err" >&2; fi - exit 1 - fi - export BAO_TOKEN - - # 🩸 Degrades where ./bao.nix's check fails, and the reason is that - # the two reads are governed by the *same* policy stanza: - # `swarm_secret_client::policy::render_agent` grants read on - # `secret/data/swarm/agents//*`, which covers this path and - # the `bao-mtls` one beside it alike. So a refusal this unit sees and - # that check did not cannot be a policy that drifted — it is an - # object that has not been minted, which is the ordinary state of - # every agent created before its swarm knew to mint one. A unit that - # failed at every boot over that would be loud about a deployment - # doing nothing wrong. - # - # ⚠️ Written by redirect into the runtime directory, never echoed: - # the field is the secret itself. - rm -f "$err" - if ! bao kv get -field=value ${lib.escapeShellArg queuePath} > ${lib.escapeShellArg secretFile} 2>"$err"; then - rm -f ${lib.escapeShellArg secretFile} - echo "no per-agent queue credential at ${queuePath} yet; this agent falls back to its hive's shared one." >&2 - if [ -s "$err" ]; then cat "$err" >&2; fi - exit 0 - fi - - echo "fetched this agent's own queue credential from ${queuePath}." - ''; - }; - - # The harness reads a path and never a value, the same shape ./queue.nix - # hands it the hive-scoped secret in. Not `%d` here: this credential is - # not a systemd credential at all — it is a file this container fetched - # for itself, which is the whole point. - systemd.services.hive-agent = { - after = [ "${unitName}.service" ]; - environment.HIVE_AGENT_QUEUE_AGENT_SECRET_FILE = secretFile; }; }; } diff --git a/nix/module-eval/agent-queue-bao.nix b/nix/module-eval/agent-queue-bao.nix index 4768d71a..3e117387 100644 --- a/nix/module-eval/agent-queue-bao.nix +++ b/nix/module-eval/agent-queue-bao.nix @@ -40,8 +40,15 @@ let agentNoBao = agentWith { }; + # Both at once, for the property that only exists when both are configured: + # the store certificate never reaches a queue variable. + agentQueueBao = agentWith { + hyperhive.queue.natsUrl = "nats://10.42.0.1:4222"; + hyperhive.queue.tokenEndpoint = "https://auth.t.local/api/oidc/token"; + services.hyperhive.agent.bao.addr = "https://bao.t.local:8200"; + }; + agentBaoIdentity = machine: machine.systemd.services.hive-agent-bao-identity; - agentQueueCredential = machine: machine.systemd.services.hive-agent-queue-credential; cases = [ { # Both ids or neither: the secret authenticates nobody without the id it @@ -166,91 +173,62 @@ let ok = !(agentNoBao.systemd.services ? hive-agent-bao-identity); } { - # The per-agent credential's fetch rides on the same store identity the - # check above proves, because there is one identity per agent. A fetch - # unit loading different ids would be a second certificate nothing - # mints. - name = "the queue-credential fetch presents the agent's own store identity"; + # The harness reads its per-agent queue secret from the store itself, + # under the same store identity the check above proves, because there is + # one identity per agent. + name = "the harness presents the agent's own store identity"; ok = let - u = agentQueueCredential agentBao; + u = agentHarness agentBao; in builtins.elem "hive-agent-bao-cert" u.serviceConfig.LoadCredential + && builtins.elem "hive-agent-bao-key" u.serviceConfig.LoadCredential && u.environment.BAO_CLIENT_CERT == "%d/hive-agent-bao-cert" - && u.environment.BAO_CLIENT_KEY == "%d/hive-agent-bao-key"; + && u.environment.BAO_CLIENT_KEY == "%d/hive-agent-bao-key" + && u.environment.BAO_ADDR == agentBao.services.hyperhive.agent.bao.addr; } { - # Same 403-not-a-miss reason as every other reader here: the path - # `swarm_secret_client::queue::agent_queue_path` builds is the one this - # agent's own policy stanza covers. Built from the agent's own name, - # because the name is what makes it this agent's credential and not a - # neighbour's — which is the entire point of minting one per agent. - name = "the queue-credential fetch reads the agent's own per-agent path"; + # `swarm_secret_client` builds both the store path and the cert-auth role + # from this name, so it has to be the name the credential was minted + # under. + name = "the harness is told the agent name its queue secret is minted under"; ok = - let - m = agentBao; - name = m.services.hyperhive.agent.user.name; - in - lib.hasInfix "secret/swarm/agents/${name}/queue" (agentQueueCredential m).script; + (agentHarness agentBao).environment.HIVE_AGENT_NAME == agentBao.services.hyperhive.agent.user.name; + } + { + # The per-agent queue secret lives in the harness's memory only: no unit + # fetches it to a file, and the harness is handed no file to read it + # from. + name = "no unit writes the per-agent queue secret to disk"; + ok = + !(agentBao.systemd.services ? hive-agent-queue-credential) + && !((agentHarness agentBao).environment ? HIVE_AGENT_QUEUE_AGENT_SECRET_FILE); } { # ⛔ The store certificate authenticates against the store and nothing - # else. It must never reach the queue, so nothing in this unit may hand - # a `BAO_CLIENT_*` path to anything queue-shaped. - name = "the queue-credential fetch never points the queue at the store certificate"; + # else. With both the store and the queue configured, no queue variable + # may name one of its paths. + name = "the harness never points the queue at the store certificate"; ok = let - u = agentQueueCredential agentBao; - harness = agentHarness agentBao; + e = (agentHarness agentQueueBao).environment; + queueVars = lib.filterAttrs ( + n: _: lib.hasPrefix "HIVE_AGENT_OIDC_" n || n == "HIVE_AGENT_NATS_URL" + ) e; in - !(lib.hasInfix "nats" u.script) - && !(lib.any (lib.hasPrefix "BAO_") (builtins.attrNames harness.environment)); - } - { - # A secret fetched at boot has no business surviving one, and the - # directory has to be the unit's own so the file is owned by the agent - # rather than needing a mode change. `RemainAfterExit` is what keeps - # systemd from removing it out from under the harness. - name = "the fetched credential lands in a runtime directory the unit keeps alive"; - ok = - let - u = agentQueueCredential agentBao; - c = u.serviceConfig; - in - c.RuntimeDirectory == "hive-agent-queue-credential" - && c.RemainAfterExit - && lib.hasInfix "/run/hive-agent-queue-credential/secret" u.script; - } - { - # 🩸 Degrades where the identity check fails, and the reason is that - # both reads are governed by one policy stanza: a refusal this unit - # sees and that check did not is an object not yet minted, not a policy - # that drifted. An agent created before its swarm minted one is a - # deployment doing nothing wrong. - name = "a per-agent credential that was never minted does not fail the unit"; - ok = lib.hasInfix "exit 0" (agentQueueCredential agentBao).script; - } - { - # The harness reads a path and never a value — the same discipline - # ../agent-modules/queue.nix keeps for the hive-scoped secret. Not - # `%d`: this one is not a systemd credential, it is a file the - # container fetched for itself. - name = "the harness is handed the fetched credential as a path"; - ok = - let - m = agentBao; - in - (agentHarness m).environment.HIVE_AGENT_QUEUE_AGENT_SECRET_FILE - == m.services.hyperhive.agent.queue.agentSecretFile; + queueVars != { } && !(lib.any (lib.hasInfix "hive-agent-bao") (builtins.attrValues queueVars)); } { # The absence arm. An agent whose swarm gave it no store has nothing to - # log in with, so there is nothing to fetch with either — and the - # harness is then told no path rather than one that never fills. - name = "an agent told no store address fetches no queue credential"; + # log in with, so the harness is handed no store coordinates at all. + name = "an agent told no store address hands the harness no store identity"; ok = - !(agentNoBao.systemd.services ? hive-agent-queue-credential) - && !((agentHarness agentNoBao).environment ? HIVE_AGENT_QUEUE_AGENT_SECRET_FILE); + let + u = agentHarness agentNoBao; + in + !(u.environment ? BAO_ADDR) + && !(u.environment ? HIVE_AGENT_NAME) + && !(builtins.elem "hive-agent-bao-cert" (u.serviceConfig.LoadCredential or [ ])); } ]; in diff --git a/nix/script-tests/agent-bao-fetch.nix b/nix/script-tests/agent-bao-fetch.nix index 00156389..b985995b 100644 --- a/nix/script-tests/agent-bao-fetch.nix +++ b/nix/script-tests/agent-bao-fetch.nix @@ -1,9 +1,8 @@ -# `checks.script-test-agent-bao-fetch` — runs the two agent units that log in -# to the swarm secret store and fetch a secret, ../agent-modules/forge-token.nix -# and ../agent-modules/queue-identity.nix, against a stub `bao`. The cases are -# in ./agent-bao-fetch.sh. +# `checks.script-test-agent-bao-fetch` — runs the agent unit that logs in to +# the swarm secret store and fetches a secret, ../agent-modules/forge-token.nix, +# against a stub `bao`. The cases are in ./agent-bao-fetch.sh. # -# What runs is each unit's rendered `ExecStart`, on the unit's own `PATH` with +# What runs is the unit's rendered `ExecStart`, on the unit's own `PATH` with # the stub in front. The one edit is the unit's `/run//` prefix, moved # under the build directory because the sandbox has no writable `/run`. { @@ -71,7 +70,6 @@ in pkgs.runCommand "hyperhive-script-test-agent-bao-fetch" ( unitEnv "FORGE" "hive-agent-forge-token" - // unitEnv "QUEUE" "hive-agent-queue-credential" // { FAKE_BAO_BIN = "${fakeBao}/bin"; } diff --git a/nix/script-tests/agent-bao-fetch.sh b/nix/script-tests/agent-bao-fetch.sh index a69efc47..b86a81e8 100644 --- a/nix/script-tests/agent-bao-fetch.sh +++ b/nix/script-tests/agent-bao-fetch.sh @@ -22,7 +22,7 @@ if (: >"$probe") 2>/dev/null; then exit 1 fi -# setup FORGE|QUEUE +# setup FORGE setup() { prefix=$1 case="$prefix: $2" @@ -111,125 +111,108 @@ expect_no_file() { if [ -e "$1" ]; then fail "$1 left behind"; fi } -for prefix in FORGE QUEUE; do - if [ "$prefix" = FORGE ]; then - secret_name=token - path=secret/swarm/agents/a1/forge-token - missing_msg="no forge token at $path yet" - else - secret_name=secret - path=secret/swarm/agents/a1/queue - missing_msg="no per-agent queue credential at $path yet" - fi - login_call="login -method=cert -token-only" - kv_call="kv get -field=value $path" +secret_name=token +path=secret/swarm/agents/a1/forge-token +missing_msg="no forge token at $path yet" +login_call="login -method=cert -token-only" +kv_call="kv get -field=value $path" - setup "$prefix" "no store identity delivered" - : >"$creds/hive-agent-bao-cert" - go - expect_rc 0 - expect_err "has no store identity" - expect_calls +setup FORGE "no store identity delivered" +: >"$creds/hive-agent-bao-cert" +go +expect_rc 0 +expect_err "has no store identity" +expect_calls - setup "$prefix" "a credential that cannot be read" - chmod 0000 "$creds/hive-agent-bao-key" - go - expect_rc 1 - expect_err "cannot read $creds/hive-agent-bao-key" - expect_calls +setup FORGE "a credential that cannot be read" +chmod 0000 "$creds/hive-agent-bao-key" +go +expect_rc 1 +expect_err "cannot read $creds/hive-agent-bao-key" +expect_calls - setup "$prefix" "bao's stderr file cannot be created" - chmod 0500 "$rt" - go - chmod 0700 "$rt" - expect_rc 1 - expect_err "could not create $rt/bao.err, so bao never ran" - expect_calls +setup FORGE "bao's stderr file cannot be created" +chmod 0500 "$rt" +go +chmod 0700 "$rt" +expect_rc 1 +expect_err "could not create $rt/bao.err, so bao never ran" +expect_calls - setup "$prefix" "the store refuses the login with an HTTP 4xx" - login_rc=2 - login_err=$'Error authenticating: Error making API request.\n\nURL: PUT '"$bao_addr"$'/v1/auth/cert/login\nCode: 403. Errors:\n\n* permission denied\n' - go - expect_rc 1 - expect_err "refused this agent's certificate login with HTTP 403:" - expect_err "* permission denied" - expect_calls "$login_call" +setup FORGE "the store refuses the login with an HTTP 4xx" +login_rc=2 +login_err=$'Error authenticating: Error making API request.\n\nURL: PUT '"$bao_addr"$'/v1/auth/cert/login\nCode: 403. Errors:\n\n* permission denied\n' +go +expect_rc 1 +expect_err "refused this agent's certificate login with HTTP 403:" +expect_err "* permission denied" +expect_calls "$login_call" - setup "$prefix" "the store fails the login with another HTTP status" - login_rc=2 - login_err=$'Error authenticating: Error making API request.\n\nCode: 500. Errors:\n\n* internal error\n' - go - expect_rc 1 - expect_err "failed this agent's certificate login with HTTP 500:" - expect_err "* internal error" - expect_calls "$login_call" +setup FORGE "the store fails the login with another HTTP status" +login_rc=2 +login_err=$'Error authenticating: Error making API request.\n\nCode: 500. Errors:\n\n* internal error\n' +go +expect_rc 1 +expect_err "failed this agent's certificate login with HTTP 500:" +expect_err "* internal error" +expect_calls "$login_call" - setup "$prefix" "the store refuses the certificate with a TLS alert" - login_rc=2 - login_err="Error authenticating: Put \"$bao_addr/v1/auth/cert/login\": remote error: tls: unknown certificate authority" - go - expect_rc 1 - expect_err "refused this agent's certificate in the TLS handshake:" - expect_err "remote error: tls: unknown certificate authority" - expect_calls "$login_call" +setup FORGE "the store refuses the certificate with a TLS alert" +login_rc=2 +login_err="Error authenticating: Put \"$bao_addr/v1/auth/cert/login\": remote error: tls: unknown certificate authority" +go +expect_rc 1 +expect_err "refused this agent's certificate in the TLS handshake:" +expect_err "remote error: tls: unknown certificate authority" +expect_calls "$login_call" - setup "$prefix" "the store does not answer" - login_rc=2 - login_err="Error authenticating: Put \"$bao_addr/v1/auth/cert/login\": dial tcp: lookup bao.t.local: no such host" - go - expect_rc 1 - expect_err "got no answer from the swarm secret store at $bao_addr" - expect_err "no such host" - expect_calls "$login_call" +setup FORGE "the store does not answer" +login_rc=2 +login_err="Error authenticating: Put \"$bao_addr/v1/auth/cert/login\": dial tcp: lookup bao.t.local: no such host" +go +expect_rc 1 +expect_err "got no answer from the swarm secret store at $bao_addr" +expect_err "no such host" +expect_calls "$login_call" - # Only the forge unit keeps its runtime directory between runs; the queue - # unit starts each run with an empty one. - setup "$prefix" "nothing minted at the agent's path yet" - if [ "$prefix" = FORGE ]; then - printf old >"$rt/$secret_name" - chmod 0400 "$rt/$secret_name" - fi - kv_rc=2 - kv_err="No value found at secret/data/swarm/agents/a1" - go - expect_rc 0 - expect_err "$missing_msg" - expect_err "No value found" - expect_calls "$login_call" "$kv_call" - if [ "$prefix" = FORGE ]; then - expect_file "$rt/$secret_name" old - else - expect_no_file "$rt/$secret_name" - fi +# The unit keeps its runtime directory between runs. +setup FORGE "nothing minted at the agent's path yet" +printf old >"$rt/$secret_name" +chmod 0400 "$rt/$secret_name" +kv_rc=2 +kv_err="No value found at secret/data/swarm/agents/a1" +go +expect_rc 0 +expect_err "$missing_msg" +expect_err "No value found" +expect_calls "$login_call" "$kv_call" +expect_file "$rt/$secret_name" old - setup "$prefix" "a first fetch writes the secret 0400" - go - expect_rc 0 - expect_out "fetched this agent's" - expect_no_err "Permission denied" - expect_calls "$login_call" "$kv_call" - expect_file "$rt/$secret_name" the-secret - if [ "$(stat -c %a "$rt/$secret_name")" != 400 ]; then fail "$secret_name is mode $(stat -c %a "$rt/$secret_name"), expected 400"; fi - expect_no_file "$rt/bao.err" - expect_no_file "$rt/token.new" +setup FORGE "a first fetch writes the secret 0400" +go +expect_rc 0 +expect_out "fetched this agent's" +expect_no_err "Permission denied" +expect_calls "$login_call" "$kv_call" +expect_file "$rt/$secret_name" the-secret +if [ "$(stat -c %a "$rt/$secret_name")" != 400 ]; then fail "$secret_name is mode $(stat -c %a "$rt/$secret_name"), expected 400"; fi +expect_no_file "$rt/bao.err" +expect_no_file "$rt/token.new" - # `UMask=0377` makes every file the script creates 0400, the login's own - # `bao.err` included, and a redirect into a 0400 file fails before bao starts. - setup "$prefix" "0400 files already in place do not block the fetch" - printf stale >"$rt/bao.err" - chmod 0400 "$rt/bao.err" - if [ "$prefix" = FORGE ]; then - printf stale >"$rt/token.new" - chmod 0400 "$rt/token.new" - fi - go - expect_rc 0 - expect_out "fetched this agent's" - expect_no_err "Permission denied" - expect_no_err "refused" - expect_calls "$login_call" "$kv_call" - expect_file "$rt/$secret_name" the-secret -done +# `UMask=0377` makes every file the script creates 0400, the login's own +# `bao.err` included, and a redirect into a 0400 file fails before bao starts. +setup FORGE "0400 files already in place do not block the fetch" +printf stale >"$rt/bao.err" +chmod 0400 "$rt/bao.err" +printf stale >"$rt/token.new" +chmod 0400 "$rt/token.new" +go +expect_rc 0 +expect_out "fetched this agent's" +expect_no_err "Permission denied" +expect_no_err "refused" +expect_calls "$login_call" "$kv_call" +expect_file "$rt/$secret_name" the-secret setup FORGE "an unchanged token is left in place" printf the-secret >"$rt/token" diff --git a/swarm-queue-client/src/lib.rs b/swarm-queue-client/src/lib.rs index af3ce9e5..ca46594c 100644 --- a/swarm-queue-client/src/lib.rs +++ b/swarm-queue-client/src/lib.rs @@ -827,20 +827,36 @@ pub async fn connect(cfg: QueueConfig) -> Result { Ok(client) } -/// Connect presenting a fixed `token`, as an agent presents its own queue -/// credential. Reconnects present the same token. +/// Connect presenting the token `token` resolves to, as an agent presents its +/// own queue credential. +/// +/// `token` runs once per connection attempt, reconnects included, so a +/// credential replaced at its source is presented on the next attempt. An +/// `Err` from it fails that attempt, and the next one waits out the same +/// backoff as a refused connect. /// /// 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( +pub async fn connect_with_token( url: &str, ca_file: Option<&std::path::Path>, - token: String, -) -> Result { - let mut options = - async_nats::ConnectOptions::with_token(token).reconnect_delay_callback(reconnect_delay); + token: F, +) -> Result +where + F: Fn() -> Fut + Send + Sync + 'static, + Fut: Future> + Send + Sync + 'static, +{ + let mut options = async_nats::ConnectOptions::with_auth_callback(move |_nonce| { + let token = token(); + async move { + let mut auth = async_nats::Auth::new(); + auth.token = Some(token.await.map_err(async_nats::AuthError::new)?); + Ok(auth) + } + }) + .reconnect_delay_callback(reconnect_delay); if let Some(path) = ca_file { options = options.add_root_certificates(path.to_path_buf()); }