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
|
|
@ -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);
|
||||
}
|
||||
}
|
||||
|
|
|
|||
Loading…
Reference in a new issue