hive-agent: read the per-agent queue secret from bao in process
The harness now reads swarm/agents/<agent>/queue from the store itself, under the agent's own store certificate, and holds it in memory only. It reads once before the first connect and again on every reconnect attempt (async-nats `ConnectOptions::with_auth_callback`), so an agent whose secret was re-minted reconnects with the new value instead of being refused until the container restarts. hive-agent-queue-credential.service, the /run file it wrote, and HIVE_AGENT_QUEUE_AGENT_SECRET_FILE are gone; queue-identity.nix now hands hive-agent.service the store address, its certificate paths and the agent name. A failed or empty read before the first connect still falls back to the hive's shared client. Each read is bounded by a 10s timeout, and retries wait out the existing reconnect backoff (500ms doubling, capped at 60s). Closes #4783
This commit is contained in:
parent
f7437a4773
commit
ccb5bd3b38
10 changed files with 770 additions and 529 deletions
1
Cargo.lock
generated
1
Cargo.lock
generated
|
|
@ -1690,6 +1690,7 @@ dependencies = [
|
|||
"serde",
|
||||
"serde_json",
|
||||
"swarm-queue-client",
|
||||
"swarm-secret-client",
|
||||
"tempfile",
|
||||
"tokio",
|
||||
"tokio-stream",
|
||||
|
|
|
|||
|
|
@ -59,20 +59,20 @@ of the cell says how.
|
|||
|
||||
<!-- vale write-good.Passive = NO -->
|
||||
|
||||
| store path | minter | reader — pulls at runtime, holds in memory | automatic re-mint | automatic re-pull |
|
||||
| ----------------------------------------------------------- | ---------------------------------------------------------------------------------------------------------------------------------------- | ------------------------------------------------------------------------------------------------------------------------------------------------------------------ | -------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- | --------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------- |
|
||||
| `swarm/agents/<agent>/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/<agent>/matrix/<account>` | `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/<agent>/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/<agent>/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/<agent>/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/<hive>/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/<hive>/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/<hive>/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/<clientId>/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/<agent>/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/<agent>/matrix/<account>` | `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/<agent>/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/<agent>/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/<agent>/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/<hive>/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/<hive>/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/<hive>/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/<clientId>/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 |
|
||||
|
||||
<!-- vale write-good.Passive = YES -->
|
||||
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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<Option<QueueConfig>> = 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<Option<AgentPath>> = OnceLock::new();
|
||||
|
||||
|
|
@ -56,8 +72,148 @@ static CLIENT: OnceCell<Option<Connection>> = OnceCell::const_new();
|
|||
struct AgentPath {
|
||||
url: String,
|
||||
ca_file: Option<PathBuf>,
|
||||
/// The fetched secret. A path; the bytes are read at connect.
|
||||
secret_file: PathBuf,
|
||||
store: Arc<StoreSource>,
|
||||
}
|
||||
|
||||
/// 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<Output = anyhow::Result<Option<String>>> + Send;
|
||||
}
|
||||
|
||||
impl SecretSource for Arc<StoreSource> {
|
||||
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<Option<String>> {
|
||||
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<queue::AgentCredential> = 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<S> {
|
||||
source: S,
|
||||
timeout: Duration,
|
||||
held: Mutex<Held>,
|
||||
}
|
||||
|
||||
#[derive(Default)]
|
||||
struct Held {
|
||||
value: Option<String>,
|
||||
/// `value` was read by [`AgentSecret::prime`] and no attempt has
|
||||
/// presented it yet.
|
||||
unpresented: bool,
|
||||
}
|
||||
|
||||
impl<S: SecretSource> AgentSecret<S> {
|
||||
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<String> {
|
||||
{
|
||||
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<Option<String>> {
|
||||
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<String>,
|
||||
/// 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<String>,
|
||||
}
|
||||
|
||||
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<String> {
|
|||
}
|
||||
}
|
||||
|
||||
/// 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<PathBuf> {
|
||||
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<String>) -> anyhow::Result<Option<StoreSource>> {
|
||||
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<PathBuf>) -> Option<AgentPath> {
|
||||
/// 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<StoreSource>) -> Option<AgentPath> {
|
||||
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<Connection> {
|
|||
/// 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 {
|
||||
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<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)))
|
||||
/// 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<S: SecretSource>(
|
||||
url: &str,
|
||||
ca_file: Option<&Path>,
|
||||
name: String,
|
||||
secret: AgentSecret<S>,
|
||||
) -> anyhow::Result<async_nats::Client> {
|
||||
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<String> {
|
||||
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<VecDeque<Answer>>,
|
||||
reads: AtomicUsize,
|
||||
}
|
||||
|
||||
fn scripted(answers: impl IntoIterator<Item = Answer>) -> Arc<Scripted> {
|
||||
Arc::new(Scripted {
|
||||
answers: Mutex::new(answers.into_iter().collect()),
|
||||
reads: AtomicUsize::new(0),
|
||||
})
|
||||
}
|
||||
|
||||
impl SecretSource for Arc<Scripted> {
|
||||
fn origin(&self) -> &'static str {
|
||||
"swarm/agents/a1/queue"
|
||||
}
|
||||
|
||||
async fn fetch(&self) -> anyhow::Result<Option<String>> {
|
||||
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<Scripted>) -> 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);
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -5,15 +5,14 @@
|
|||
# to `swarm/agents/<agent>/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 `<state>/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.
|
||||
|
|
|
|||
|
|
@ -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/<pid>/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: <status>` 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/<agent>/*`, 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;
|
||||
};
|
||||
};
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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/<unit>/` 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";
|
||||
}
|
||||
|
|
|
|||
|
|
@ -22,7 +22,7 @@ if (: >"$probe") 2>/dev/null; then
|
|||
exit 1
|
||||
fi
|
||||
|
||||
# setup FORGE|QUEUE <description>
|
||||
# setup FORGE <description>
|
||||
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"
|
||||
|
|
|
|||
|
|
@ -827,20 +827,36 @@ 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.
|
||||
/// 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<F, Fut>(
|
||||
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);
|
||||
token: F,
|
||||
) -> Result<async_nats::Client, Error>
|
||||
where
|
||||
F: Fn() -> Fut + Send + Sync + 'static,
|
||||
Fut: Future<Output = Result<String, String>> + 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());
|
||||
}
|
||||
|
|
|
|||
Loading…
Reference in a new issue