hive-agent: present this agent's own queue credential, then fall back
When the per-agent secret `queue-identity.nix` fetched is present, the harness connects with `swarm-agent.<agent>.<secret>` as a static token and publishes on `$SWARM.term.<agent>` and `$SWARM.agent-state.<agent>`. When it is absent, or that first connect fails for any reason, a refusal from a responder that does not verify agent tokens included, it connects with the hive's shared OIDC client and publishes on the hive-scoped subjects as before. Which one it took is logged once per connect. `swarm_queue_client::connect_with_token` is the static-token connect: no retry on the initial attempt, so the caller sees the refusal and can fall back. Reconnects share the existing backoff, now a named function. Closes #4630
This commit is contained in:
parent
0c1fb44a4f
commit
727065c960
5 changed files with 287 additions and 76 deletions
|
|
@ -410,10 +410,13 @@ agents sets none of the four and each agent logs that it has none; a half-set
|
||||||
environment logs an error and the harness keeps serving.
|
environment logs an error and the harness keeps serving.
|
||||||
|
|
||||||
What an agent does with that connection is publish its terminal. Every row its
|
What an agent does with that connection is publish its terminal. Every row its
|
||||||
own web UI renders also goes to `$SWARM.term.<hive>.<agent>`, one subject per
|
own web UI renders also goes to `$SWARM.term.<agent>`, one subject per agent, so
|
||||||
agent, so a swarm-level terminal can follow one agent without subscribing to
|
a swarm-level terminal can follow one agent without subscribing to the swarm's
|
||||||
the swarm's whole traffic. The `<hive>` is the one the agent's client id names,
|
whole traffic. That is the subject an agent connected with its own queue
|
||||||
which is the same string the broker builds its grant from. Publishing only: an
|
credential is granted. An agent without one, or whose own credential the queue
|
||||||
|
refused, connects with its hive's shared client and publishes to
|
||||||
|
`$SWARM.term.<hive>.<agent>` instead, the `<hive>` being the one that client id
|
||||||
|
names. The swarm controller relays both. Publishing only: an
|
||||||
agent talks about itself here and reads nothing. Rows aren't retained — a
|
agent talks about itself here and reads nothing. Rows aren't retained — a
|
||||||
subscriber that wasn't listening missed them, the same as on the agent's own
|
subscriber that wasn't listening missed them, the same as on the agent's own
|
||||||
live stream.
|
live stream.
|
||||||
|
|
@ -424,7 +427,7 @@ sending and leaves a marker in its place; the summary, level and icon still
|
||||||
arrive. The harness logs and skips a row that's too large even without its body.
|
arrive. The harness logs and skips a row that's too large even without its body.
|
||||||
|
|
||||||
The second thing an agent publishes is its **turn-state header**, on
|
The second thing an agent publishes is its **turn-state header**, on
|
||||||
`$SWARM.agent-state.<hive>.<agent>` — same shape of subject, same grant
|
`$SWARM.agent-state.<agent>` (or `$SWARM.agent-state.<hive>.<agent>`) — same shape of subject, same grant
|
||||||
mechanics, same lack of retention. It carries what a header bar wants: what the
|
mechanics, same lack of retention. It carries what a header bar wants: what the
|
||||||
turn loop is doing (`turn_state`, plus `turn_state_since` as an ISO 8601 UTC
|
turn loop is doing (`turn_state`, plus `turn_state_since` as an ISO 8601 UTC
|
||||||
stamp), which model (`model` and the resolved id the last turn actually ran on),
|
stamp), which model (`model` and the resolved id the last turn actually ran on),
|
||||||
|
|
|
||||||
|
|
@ -35,6 +35,7 @@ use tokio::sync::broadcast;
|
||||||
use swarm_queue_client::wanted::AgentState;
|
use swarm_queue_client::wanted::AgentState;
|
||||||
|
|
||||||
use crate::events::{Bus, BusEvent, LiveEvent, TurnState};
|
use crate::events::{Bus, BusEvent, LiveEvent, TurnState};
|
||||||
|
use crate::swarm_queue::{Connection, Presented};
|
||||||
use crate::term_msg::iso8601_utc;
|
use crate::term_msg::iso8601_utc;
|
||||||
|
|
||||||
/// Subject family carrying agent turn-state headers, the swarm-wide
|
/// Subject family carrying agent turn-state headers, the swarm-wide
|
||||||
|
|
@ -171,33 +172,43 @@ fn snapshot(bus: &Bus) -> AgentStateMsg {
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Start the publish task, if this agent has both queue coordinates and an
|
/// The subject this agent's header goes to under the credential it connected
|
||||||
/// identity the subject can be derived from.
|
/// with, or `None` when a hive client id names no hive.
|
||||||
///
|
fn subject(presented: &Presented, agent: &str) -> Option<String> {
|
||||||
/// Returns without spawning in every other case — no queue, an unparseable
|
match presented {
|
||||||
/// client id, no label — each of which is a legal state for an agent rather
|
Presented::Agent => Some(format!("{SUBJECT_PREFIX}.{agent}")),
|
||||||
/// than an error, and each logged once here rather than per transition.
|
Presented::Hive { client_id } => {
|
||||||
pub fn spawn(bus: &Bus) {
|
let Some(hive) = hive_from_client_id(client_id) else {
|
||||||
let Some(cfg) = crate::swarm_queue::config() else {
|
|
||||||
// `swarm_queue::init` already said why at boot; repeating it here
|
|
||||||
// would be the same fact logged twice.
|
|
||||||
return;
|
|
||||||
};
|
|
||||||
let Some(hive) = hive_from_client_id(&cfg.client_id) else {
|
|
||||||
tracing::warn!(
|
tracing::warn!(
|
||||||
client_id = %cfg.client_id,
|
%client_id,
|
||||||
expected = format!("{CLIENT_ID_PREFIX}<hive>{CLIENT_ID_SUFFIX}"),
|
expected = format!("{CLIENT_ID_PREFIX}<hive>{CLIENT_ID_SUFFIX}"),
|
||||||
"queue client id does not name a hive; not publishing turn state upward"
|
"queue client id does not name a hive; not publishing turn state upward"
|
||||||
);
|
);
|
||||||
return;
|
return None;
|
||||||
};
|
};
|
||||||
|
Some(format!("{SUBJECT_PREFIX}.{hive}.{agent}"))
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Start the publish task, if this agent has a queue credential and a label
|
||||||
|
/// to name its subject with.
|
||||||
|
///
|
||||||
|
/// Returns without spawning in every other case — no queue or no label — each
|
||||||
|
/// of which is a legal state for an agent rather than an error, and each
|
||||||
|
/// logged once here rather than per transition.
|
||||||
|
pub fn spawn(bus: &Bus) {
|
||||||
|
if !crate::swarm_queue::configured() {
|
||||||
|
// `swarm_queue::init` already said why at boot; repeating it here
|
||||||
|
// would be the same fact logged twice.
|
||||||
|
return;
|
||||||
|
}
|
||||||
let agent = crate::identity::label();
|
let agent = crate::identity::label();
|
||||||
if agent.is_empty() {
|
if agent.is_empty() {
|
||||||
tracing::warn!("this agent has no label; not publishing turn state upward");
|
tracing::warn!("this agent has no label; not publishing turn state upward");
|
||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
let subject = format!("{SUBJECT_PREFIX}.{hive}.{agent}");
|
tokio::spawn(run(bus.subscribe(), bus.clone(), agent));
|
||||||
tokio::spawn(run(bus.subscribe(), bus.clone(), subject));
|
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Watch the bus and publish whenever the header actually changed.
|
/// Watch the bus and publish whenever the header actually changed.
|
||||||
|
|
@ -218,8 +229,11 @@ pub fn spawn(bus: &Bus) {
|
||||||
/// and it carries nothing this header reads. Every other variant, including
|
/// and it carries nothing this header reads. Every other variant, including
|
||||||
/// any added later, funnels into the comparison and costs nothing when it
|
/// any added later, funnels into the comparison and costs nothing when it
|
||||||
/// changes nothing.
|
/// changes nothing.
|
||||||
async fn run(mut rx: broadcast::Receiver<BusEvent>, bus: Bus, subject: String) {
|
async fn run(mut rx: broadcast::Receiver<BusEvent>, bus: Bus, agent: String) {
|
||||||
let Some(client) = crate::swarm_queue::client().await else {
|
let Some(Connection { client, presented }) = crate::swarm_queue::client().await else {
|
||||||
|
return;
|
||||||
|
};
|
||||||
|
let Some(subject) = subject(&presented, &agent) else {
|
||||||
return;
|
return;
|
||||||
};
|
};
|
||||||
tracing::info!(subject, "publishing agent turn state to the swarm queue");
|
tracing::info!(subject, "publishing agent turn state to the swarm queue");
|
||||||
|
|
@ -307,7 +321,7 @@ async fn publish(
|
||||||
|
|
||||||
#[cfg(test)]
|
#[cfg(test)]
|
||||||
mod tests {
|
mod tests {
|
||||||
use super::{AgentStateMsg, SUBJECT_PREFIX, hive_from_client_id};
|
use super::{AgentStateMsg, Presented, hive_from_client_id, subject};
|
||||||
use crate::events::TurnState;
|
use crate::events::TurnState;
|
||||||
use swarm_queue_client::wanted::AgentState;
|
use swarm_queue_client::wanted::AgentState;
|
||||||
|
|
||||||
|
|
@ -393,10 +407,22 @@ mod tests {
|
||||||
/// this pins the half that lives here.
|
/// this pins the half that lives here.
|
||||||
#[test]
|
#[test]
|
||||||
fn the_subject_is_the_prefix_then_the_hive_then_the_agent() {
|
fn the_subject_is_the_prefix_then_the_hive_then_the_agent() {
|
||||||
let hive = hive_from_client_id("hive-alpha-agent").expect("names a hive");
|
let presented = Presented::Hive {
|
||||||
|
client_id: "hive-alpha-agent".to_owned(),
|
||||||
|
};
|
||||||
assert_eq!(
|
assert_eq!(
|
||||||
format!("{SUBJECT_PREFIX}.{hive}.mara"),
|
subject(&presented, "mara").as_deref(),
|
||||||
"$SWARM.agent-state.alpha.mara"
|
Some("$SWARM.agent-state.alpha.mara")
|
||||||
|
);
|
||||||
|
}
|
||||||
|
|
||||||
|
/// An agent that connected with its own credential publishes on the
|
||||||
|
/// subject that credential is granted, which names no hive.
|
||||||
|
#[test]
|
||||||
|
fn an_agent_on_its_own_credential_publishes_on_its_hive_free_subject() {
|
||||||
|
assert_eq!(
|
||||||
|
subject(&Presented::Agent, "mara").as_deref(),
|
||||||
|
Some("$SWARM.agent-state.mara")
|
||||||
);
|
);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -21,15 +21,16 @@
|
||||||
//! agent. The fifth is this agent's own, minted per agent at swarm level and
|
//! 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`).
|
//! fetched by the container itself (`nix/agent-modules/queue-identity.nix`).
|
||||||
//!
|
//!
|
||||||
//! It is reported here but not yet *presented*: the queue's auth-callout
|
//! When it is present, [`client`] connects with it first, and the agent
|
||||||
//! responder (`swarm-nats-auth`) validates only the hive-scoped token, and an
|
//! publishes on its own hive-free subjects. When it is absent, or the queue
|
||||||
//! agent offering a credential nothing on the other end reads back would be
|
//! refuses it, the connect falls back to the hive's shared client and the
|
||||||
//! refused. Until that responder learns the same path, the connect path below
|
//! hive-scoped subjects. [`Connection::presented`] says which, so a publisher
|
||||||
//! is unchanged and this is the fetching half.
|
//! builds the subject that credential is granted.
|
||||||
|
|
||||||
use std::path::{Path, PathBuf};
|
use std::path::{Path, PathBuf};
|
||||||
use std::sync::OnceLock;
|
use std::sync::OnceLock;
|
||||||
|
|
||||||
|
use anyhow::Context as _;
|
||||||
use tokio::sync::OnceCell;
|
use tokio::sync::OnceCell;
|
||||||
|
|
||||||
use swarm_queue_client::QueueConfig;
|
use swarm_queue_client::QueueConfig;
|
||||||
|
|
@ -42,10 +43,39 @@ const ENV_PREFIX: &str = "HIVE_AGENT";
|
||||||
/// one answer rather than re-deriving it per call.
|
/// one answer rather than re-deriving it per call.
|
||||||
static CONFIG: OnceLock<Option<QueueConfig>> = OnceLock::new();
|
static CONFIG: OnceLock<Option<QueueConfig>> = OnceLock::new();
|
||||||
|
|
||||||
|
/// Where to present this agent's own credential, resolved once at boot.
|
||||||
|
static AGENT: OnceLock<Option<AgentPath>> = OnceLock::new();
|
||||||
|
|
||||||
/// The one connection every publisher in this process shares. Separate from
|
/// The one connection every publisher in this process shares. Separate from
|
||||||
/// [`CONFIG`] because resolving the coordinates is synchronous boot work and
|
/// [`CONFIG`] because resolving the coordinates is synchronous boot work and
|
||||||
/// connecting is not — see [`client`].
|
/// connecting is not — see [`client`].
|
||||||
static CLIENT: OnceCell<Option<async_nats::Client>> = OnceCell::const_new();
|
static CLIENT: OnceCell<Option<Connection>> = OnceCell::const_new();
|
||||||
|
|
||||||
|
/// What connecting with this agent's own credential needs.
|
||||||
|
#[derive(Debug, PartialEq, Eq)]
|
||||||
|
struct AgentPath {
|
||||||
|
url: String,
|
||||||
|
ca_file: Option<PathBuf>,
|
||||||
|
/// The fetched secret. A path; the bytes are read at connect.
|
||||||
|
secret_file: PathBuf,
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Which credential the shared connection presented.
|
||||||
|
#[derive(Debug, Clone, PartialEq, Eq)]
|
||||||
|
pub enum Presented {
|
||||||
|
/// This agent's own, granted `<prefix>.<agent>`.
|
||||||
|
Agent,
|
||||||
|
/// The hive's shared client, granted `<prefix>.<hive>.>` for the hive the
|
||||||
|
/// client id names.
|
||||||
|
Hive { client_id: String },
|
||||||
|
}
|
||||||
|
|
||||||
|
/// The shared queue connection and the credential it was made with.
|
||||||
|
#[derive(Clone)]
|
||||||
|
pub struct Connection {
|
||||||
|
pub client: async_nats::Client,
|
||||||
|
pub presented: Presented,
|
||||||
|
}
|
||||||
|
|
||||||
/// The four variables the harness unit sets, before the client-id file is
|
/// The four variables the harness unit sets, before the client-id file is
|
||||||
/// read. Collected into a struct so [`decide`] is pure over them and the
|
/// read. Collected into a struct so [`decide`] is pure over them and the
|
||||||
|
|
@ -147,6 +177,16 @@ fn decide_agent_secret(path: Option<&str>) -> Option<PathBuf> {
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// Decide whether this agent can connect with its own credential: it needs the
|
||||||
|
/// queue's address and a fetched secret, and nothing of the hive's client.
|
||||||
|
fn decide_agent_path(env: &QueueEnv, secret_file: Option<PathBuf>) -> Option<AgentPath> {
|
||||||
|
Some(AgentPath {
|
||||||
|
url: env.nats_url.clone()?,
|
||||||
|
ca_file: env.ca_file.as_ref().map(Into::into),
|
||||||
|
secret_file: secret_file?,
|
||||||
|
})
|
||||||
|
}
|
||||||
|
|
||||||
/// Decide what this agent's queue configuration is, given the environment and
|
/// Decide what this agent's queue configuration is, given the environment and
|
||||||
/// whatever the client-id file held.
|
/// whatever the client-id file held.
|
||||||
///
|
///
|
||||||
|
|
@ -216,15 +256,10 @@ pub fn init() {
|
||||||
let _ = CONFIG.set(resolved);
|
let _ = CONFIG.set(resolved);
|
||||||
|
|
||||||
// Independent of everything above: this agent may hold its own secret on
|
// Independent of everything above: this agent may hold its own secret on
|
||||||
// a hive with no queue coordinates, or hold the coordinates and no secret
|
// a hive whose shared client is not published yet, or the shared client
|
||||||
// of its own yet. Reported either way, because "which credential is this
|
// and no secret of its own.
|
||||||
// agent able to present" is a question only this process can answer, and
|
let secret = decide_agent_secret(env.agent_secret_file.as_deref());
|
||||||
// it is the one the next slice's rollout will be asked repeatedly.
|
if let Some(path) = &secret {
|
||||||
//
|
|
||||||
// The answer is only logged here. Presenting it needs the queue's
|
|
||||||
// auth-callout responder to verify it, which is the next slice — see this
|
|
||||||
// module's header.
|
|
||||||
if let Some(path) = decide_agent_secret(env.agent_secret_file.as_deref()) {
|
|
||||||
// The path, never the bytes: the file holds the secret itself.
|
// The path, never the bytes: the file holds the secret itself.
|
||||||
tracing::info!(
|
tracing::info!(
|
||||||
path = %path.display(),
|
path = %path.display(),
|
||||||
|
|
@ -236,6 +271,13 @@ pub fn init() {
|
||||||
by its hive's shared client"
|
by its hive's shared client"
|
||||||
);
|
);
|
||||||
}
|
}
|
||||||
|
let _ = AGENT.set(decide_agent_path(&env, secret));
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Whether this agent has any credential to reach the queue with. `false`
|
||||||
|
/// before [`init`] has run.
|
||||||
|
pub fn configured() -> bool {
|
||||||
|
config().is_some() || AGENT.get().is_some_and(Option::is_some)
|
||||||
}
|
}
|
||||||
|
|
||||||
/// What [`init`] resolved, or `None` when this agent has no queue.
|
/// What [`init`] resolved, or `None` when this agent has no queue.
|
||||||
|
|
@ -263,14 +305,48 @@ pub fn config() -> Option<&'static QueueConfig> {
|
||||||
/// failed" — a caller does nothing differently between them, since either way
|
/// failed" — a caller does nothing differently between them, since either way
|
||||||
/// there is nothing to publish onto. Connecting is lazy so that an agent on a
|
/// there is nothing to publish onto. Connecting is lazy so that an agent on a
|
||||||
/// hive with no queue pays nothing at boot.
|
/// hive with no queue pays nothing at boot.
|
||||||
pub async fn client() -> Option<async_nats::Client> {
|
pub async fn client() -> Option<Connection> {
|
||||||
CLIENT.get_or_init(connect_once).await.clone()
|
CLIENT.get_or_init(connect_once).await.clone()
|
||||||
}
|
}
|
||||||
|
|
||||||
async fn connect_once() -> Option<async_nats::Client> {
|
/// This agent's own credential first, then the hive's. Each path taken is
|
||||||
|
/// logged once, here.
|
||||||
|
async fn connect_once() -> Option<Connection> {
|
||||||
|
if let Some(agent) = AGENT.get().and_then(Option::as_ref) {
|
||||||
|
match connect_as_agent(agent).await {
|
||||||
|
Ok(client) => {
|
||||||
|
tracing::info!(
|
||||||
|
url = %agent.url,
|
||||||
|
"connected to the swarm queue with this agent's own credential"
|
||||||
|
);
|
||||||
|
return Some(Connection {
|
||||||
|
client,
|
||||||
|
presented: Presented::Agent,
|
||||||
|
});
|
||||||
|
}
|
||||||
|
Err(e) => tracing::warn!(
|
||||||
|
error = format!("{e:#}"),
|
||||||
|
"connecting with this agent's own credential failed; falling back to \
|
||||||
|
its hive's shared client"
|
||||||
|
),
|
||||||
|
}
|
||||||
|
}
|
||||||
let cfg = config()?;
|
let cfg = config()?;
|
||||||
match swarm_queue_client::connect(cfg.clone()).await {
|
match swarm_queue_client::connect(cfg.clone()).await {
|
||||||
Ok(client) => Some(client),
|
Ok(client) => {
|
||||||
|
tracing::info!(
|
||||||
|
url = %cfg.url,
|
||||||
|
client_id = %cfg.client_id,
|
||||||
|
"connecting to the swarm queue with the hive's shared client; the \
|
||||||
|
client retries in the background until the queue accepts it"
|
||||||
|
);
|
||||||
|
Some(Connection {
|
||||||
|
client,
|
||||||
|
presented: Presented::Hive {
|
||||||
|
client_id: cfg.client_id.clone(),
|
||||||
|
},
|
||||||
|
})
|
||||||
|
}
|
||||||
Err(e) => {
|
Err(e) => {
|
||||||
// `chain`, not `{:#}`: this is `swarm_queue_client::Error`, whose
|
// `chain`, not `{:#}`: this is `swarm_queue_client::Error`, whose
|
||||||
// `Display` ignores the alternate flag, so `{:#}` renders the
|
// `Display` ignores the alternate flag, so `{:#}` renders the
|
||||||
|
|
@ -284,9 +360,26 @@ async fn connect_once() -> Option<async_nats::Client> {
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// Connect presenting this agent's own token. Any failure, a refusal
|
||||||
|
/// included, is returned for [`connect_once`] to fall back on.
|
||||||
|
async fn connect_as_agent(agent: &AgentPath) -> anyhow::Result<async_nats::Client> {
|
||||||
|
let name = crate::identity::label();
|
||||||
|
let secret = std::fs::read_to_string(&agent.secret_file)
|
||||||
|
.with_context(|| format!("reading {}", agent.secret_file.display()))?;
|
||||||
|
let token = swarm_queue_client::agent_token::format_agent_token(&name, secret.trim())?;
|
||||||
|
swarm_queue_client::connect_with_token(&agent.url, agent.ca_file.as_deref(), token)
|
||||||
|
.await
|
||||||
|
.map_err(|e| anyhow::anyhow!(swarm_queue_client::chain(&e)))
|
||||||
|
}
|
||||||
|
|
||||||
#[cfg(test)]
|
#[cfg(test)]
|
||||||
mod tests {
|
mod tests {
|
||||||
use super::{QueueEnv, Resolution, decide, decide_agent_secret, read_client_id};
|
use std::path::PathBuf;
|
||||||
|
|
||||||
|
use super::{
|
||||||
|
AgentPath, QueueEnv, Resolution, decide, decide_agent_path, decide_agent_secret,
|
||||||
|
read_client_id,
|
||||||
|
};
|
||||||
|
|
||||||
fn env(parts: [Option<&str>; 4]) -> QueueEnv {
|
fn env(parts: [Option<&str>; 4]) -> QueueEnv {
|
||||||
let [nats_url, token_endpoint, client_id_file, client_secret_file] = parts;
|
let [nats_url, token_endpoint, client_id_file, client_secret_file] = parts;
|
||||||
|
|
@ -435,4 +528,34 @@ mod tests {
|
||||||
assert!(matches!(decide(&e, None), Resolution::Absent(_)));
|
assert!(matches!(decide(&e, None), Resolution::Absent(_)));
|
||||||
assert!(decide_agent_secret(path.to_str()).is_some());
|
assert!(decide_agent_secret(path.to_str()).is_some());
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// The agent's own path needs the queue's address and its own secret, and
|
||||||
|
/// not the hive's client id: that is what lets it connect on a hive whose
|
||||||
|
/// shared client is not published.
|
||||||
|
#[test]
|
||||||
|
fn an_agent_with_its_own_secret_and_the_queue_address_connects_as_itself() {
|
||||||
|
let secret = PathBuf::from("/run/queue-identity/secret");
|
||||||
|
assert_eq!(
|
||||||
|
decide_agent_path(&full(), Some(secret.clone())),
|
||||||
|
Some(AgentPath {
|
||||||
|
url: "nats://10.42.0.1:4222".to_owned(),
|
||||||
|
ca_file: None,
|
||||||
|
secret_file: secret.clone(),
|
||||||
|
})
|
||||||
|
);
|
||||||
|
let url_only = env([Some("nats://10.42.0.1:4222"), None, None, None]);
|
||||||
|
assert!(decide_agent_path(&url_only, Some(secret)).is_some());
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn without_its_own_secret_or_the_queue_address_an_agent_does_not_connect_as_itself() {
|
||||||
|
assert_eq!(decide_agent_path(&full(), None), None);
|
||||||
|
assert_eq!(
|
||||||
|
decide_agent_path(
|
||||||
|
&env([None, None, None, None]),
|
||||||
|
Some(PathBuf::from("/run/queue-identity/secret"))
|
||||||
|
),
|
||||||
|
None
|
||||||
|
);
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -27,6 +27,7 @@
|
||||||
use tokio::sync::broadcast;
|
use tokio::sync::broadcast;
|
||||||
|
|
||||||
use crate::events::BusEvent;
|
use crate::events::BusEvent;
|
||||||
|
use crate::swarm_queue::{Connection, Presented};
|
||||||
use crate::term_msg::{ClassifyCtx, TermMsg, classify};
|
use crate::term_msg::{ClassifyCtx, TermMsg, classify};
|
||||||
|
|
||||||
/// Subject family carrying agent terminal rows, the swarm-wide agreement this
|
/// Subject family carrying agent terminal rows, the swarm-wide agreement this
|
||||||
|
|
@ -112,37 +113,50 @@ fn serialized_len(msg: &TermMsg) -> Option<usize> {
|
||||||
serde_json::to_vec(msg).ok().map(|v| v.len())
|
serde_json::to_vec(msg).ok().map(|v| v.len())
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Start the publish task, if this agent has both queue coordinates and an
|
/// The subject this agent's rows go to under the credential it connected
|
||||||
/// identity the subject can be derived from.
|
/// with, or `None` when a hive client id names no hive.
|
||||||
///
|
fn subject(presented: &Presented, agent: &str) -> Option<String> {
|
||||||
/// Returns without spawning in every other case — no queue, an unparseable
|
match presented {
|
||||||
/// client id, no label — each of which is a legal state for an agent rather
|
Presented::Agent => Some(format!("{SUBJECT_PREFIX}.{agent}")),
|
||||||
/// than an error, and each logged once here rather than per row.
|
Presented::Hive { client_id } => {
|
||||||
pub fn spawn(rx: broadcast::Receiver<BusEvent>) {
|
let Some(hive) = hive_from_client_id(client_id) else {
|
||||||
let Some(cfg) = crate::swarm_queue::config() else {
|
|
||||||
// `swarm_queue::init` already said why at boot; repeating it here
|
|
||||||
// would be the same fact logged twice.
|
|
||||||
return;
|
|
||||||
};
|
|
||||||
let Some(hive) = hive_from_client_id(&cfg.client_id) else {
|
|
||||||
tracing::warn!(
|
tracing::warn!(
|
||||||
client_id = %cfg.client_id,
|
%client_id,
|
||||||
expected = format!("{CLIENT_ID_PREFIX}<hive>{CLIENT_ID_SUFFIX}"),
|
expected = format!("{CLIENT_ID_PREFIX}<hive>{CLIENT_ID_SUFFIX}"),
|
||||||
"queue client id does not name a hive; not publishing the terminal upward"
|
"queue client id does not name a hive; not publishing the terminal upward"
|
||||||
);
|
);
|
||||||
return;
|
return None;
|
||||||
};
|
};
|
||||||
|
Some(format!("{SUBJECT_PREFIX}.{hive}.{agent}"))
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Start the publish task, if this agent has a queue credential and a label
|
||||||
|
/// to name its subject with.
|
||||||
|
///
|
||||||
|
/// Returns without spawning in every other case — no queue or no label — each
|
||||||
|
/// of which is a legal state for an agent rather than an error, and each
|
||||||
|
/// logged once here rather than per row.
|
||||||
|
pub fn spawn(rx: broadcast::Receiver<BusEvent>) {
|
||||||
|
if !crate::swarm_queue::configured() {
|
||||||
|
// `swarm_queue::init` already said why at boot; repeating it here
|
||||||
|
// would be the same fact logged twice.
|
||||||
|
return;
|
||||||
|
}
|
||||||
let agent = crate::identity::label();
|
let agent = crate::identity::label();
|
||||||
if agent.is_empty() {
|
if agent.is_empty() {
|
||||||
tracing::warn!("this agent has no label; not publishing the terminal upward");
|
tracing::warn!("this agent has no label; not publishing the terminal upward");
|
||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
let subject = format!("{SUBJECT_PREFIX}.{hive}.{agent}");
|
tokio::spawn(run(rx, agent));
|
||||||
tokio::spawn(run(rx, subject));
|
|
||||||
}
|
}
|
||||||
|
|
||||||
async fn run(mut rx: broadcast::Receiver<BusEvent>, subject: String) {
|
async fn run(mut rx: broadcast::Receiver<BusEvent>, agent: String) {
|
||||||
let Some(client) = crate::swarm_queue::client().await else {
|
let Some(Connection { client, presented }) = crate::swarm_queue::client().await else {
|
||||||
|
return;
|
||||||
|
};
|
||||||
|
let Some(subject) = subject(&presented, &agent) else {
|
||||||
return;
|
return;
|
||||||
};
|
};
|
||||||
tracing::info!(subject, "publishing the agent terminal to the swarm queue");
|
tracing::info!(subject, "publishing the agent terminal to the swarm queue");
|
||||||
|
|
@ -203,7 +217,7 @@ async fn publish(client: &async_nats::Client, subject: &str, msg: TermMsg) {
|
||||||
|
|
||||||
#[cfg(test)]
|
#[cfg(test)]
|
||||||
mod tests {
|
mod tests {
|
||||||
use super::{DROPPED_BODY, fit, hive_from_client_id};
|
use super::{DROPPED_BODY, Presented, fit, hive_from_client_id, subject};
|
||||||
use crate::events::LiveEvent;
|
use crate::events::LiveEvent;
|
||||||
use crate::term_msg::{BodyFormat, ClassifyCtx, Level, TermMsg, classify};
|
use crate::term_msg::{BodyFormat, ClassifyCtx, Level, TermMsg, classify};
|
||||||
|
|
||||||
|
|
@ -243,6 +257,27 @@ mod tests {
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// Each credential publishes on the subject it is granted: its own under
|
||||||
|
/// the agent's credential, its hive's under the shared one.
|
||||||
|
#[test]
|
||||||
|
fn the_subject_follows_the_credential_the_agent_connected_with() {
|
||||||
|
assert_eq!(
|
||||||
|
subject(&Presented::Agent, "mara").as_deref(),
|
||||||
|
Some("$SWARM.term.mara")
|
||||||
|
);
|
||||||
|
let hive = Presented::Hive {
|
||||||
|
client_id: "hive-alpha-agent".to_owned(),
|
||||||
|
};
|
||||||
|
assert_eq!(
|
||||||
|
subject(&hive, "mara").as_deref(),
|
||||||
|
Some("$SWARM.term.alpha.mara")
|
||||||
|
);
|
||||||
|
let unparseable = Presented::Hive {
|
||||||
|
client_id: "hive-alpha".to_owned(),
|
||||||
|
};
|
||||||
|
assert_eq!(subject(&unparseable, "mara"), None);
|
||||||
|
}
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
fn a_row_that_already_fits_is_published_unchanged() {
|
fn a_row_that_already_fits_is_published_unchanged() {
|
||||||
let msg = TermMsg::new(Level::Info, "turn ok")
|
let msg = TermMsg::new(Level::Info, "turn ok")
|
||||||
|
|
|
||||||
|
|
@ -320,6 +320,15 @@ const TOKEN_REFRESH_SKEW: std::time::Duration = std::time::Duration::from_mins(2
|
||||||
/// degrades every other client of it.
|
/// degrades every other client of it.
|
||||||
const MAX_RECONNECT_DELAY: std::time::Duration = std::time::Duration::from_mins(1);
|
const MAX_RECONNECT_DELAY: std::time::Duration = std::time::Duration::from_mins(1);
|
||||||
|
|
||||||
|
/// Exponential from 500ms, capped at [`MAX_RECONNECT_DELAY`].
|
||||||
|
fn reconnect_delay(attempts: usize) -> std::time::Duration {
|
||||||
|
let exp = u32::try_from(attempts.saturating_sub(1)).unwrap_or(u32::MAX);
|
||||||
|
std::cmp::min(
|
||||||
|
std::time::Duration::from_millis(500u64.saturating_mul(2u64.saturating_pow(exp.min(8)))),
|
||||||
|
MAX_RECONNECT_DELAY,
|
||||||
|
)
|
||||||
|
}
|
||||||
|
|
||||||
/// Where the controller finds the queue and what it authenticates with.
|
/// Where the controller finds the queue and what it authenticates with.
|
||||||
///
|
///
|
||||||
/// Every field comes from an environment variable the NixOS module sets, the
|
/// Every field comes from an environment variable the NixOS module sets, the
|
||||||
|
|
@ -742,15 +751,7 @@ pub async fn connect(cfg: QueueConfig) -> Result<async_nats::Client, Error> {
|
||||||
// which turns an unreachable queue into a permanent 4s poll — and, before
|
// which turns an unreachable queue into a permanent 4s poll — and, before
|
||||||
// the cache above, a permanent 4s token-request loop against authelia.
|
// the cache above, a permanent 4s token-request loop against authelia.
|
||||||
// Exponential from 500ms so a momentary blip still reconnects promptly.
|
// Exponential from 500ms so a momentary blip still reconnects promptly.
|
||||||
.reconnect_delay_callback(|attempts| {
|
.reconnect_delay_callback(reconnect_delay)
|
||||||
let exp = u32::try_from(attempts.saturating_sub(1)).unwrap_or(u32::MAX);
|
|
||||||
std::cmp::min(
|
|
||||||
std::time::Duration::from_millis(
|
|
||||||
500u64.saturating_mul(2u64.saturating_pow(exp.min(8))),
|
|
||||||
),
|
|
||||||
MAX_RECONNECT_DELAY,
|
|
||||||
)
|
|
||||||
})
|
|
||||||
// The controller and the queue are separate units on (possibly)
|
// The controller and the queue are separate units on (possibly)
|
||||||
// separate hosts, and nothing orders them. Without this, a queue that
|
// separate hosts, and nothing orders them. Without this, a queue that
|
||||||
// comes up one second later leaves the controller permanently
|
// comes up one second later leaves the controller permanently
|
||||||
|
|
@ -784,6 +785,29 @@ pub async fn connect(cfg: QueueConfig) -> Result<async_nats::Client, Error> {
|
||||||
Ok(client)
|
Ok(client)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// Connect presenting a fixed `token`, as an agent presents its own queue
|
||||||
|
/// credential. Reconnects present the same token.
|
||||||
|
///
|
||||||
|
/// Unlike [`connect`], the first attempt must succeed: a refused or failed
|
||||||
|
/// connect is returned rather than retried in the background, so the caller
|
||||||
|
/// can fall back to another credential. `ca_file` is the queue's trust anchor,
|
||||||
|
/// as [`QueueConfig::ca_file`].
|
||||||
|
pub async fn connect_with_token(
|
||||||
|
url: &str,
|
||||||
|
ca_file: Option<&std::path::Path>,
|
||||||
|
token: String,
|
||||||
|
) -> Result<async_nats::Client, Error> {
|
||||||
|
let mut options =
|
||||||
|
async_nats::ConnectOptions::with_token(token).reconnect_delay_callback(reconnect_delay);
|
||||||
|
if let Some(path) = ca_file {
|
||||||
|
options = options.add_root_certificates(path.to_path_buf());
|
||||||
|
}
|
||||||
|
options.connect(url).await.map_err(|source| Error::Connect {
|
||||||
|
url: url.to_owned(),
|
||||||
|
source,
|
||||||
|
})
|
||||||
|
}
|
||||||
|
|
||||||
#[cfg(test)]
|
#[cfg(test)]
|
||||||
mod tests {
|
mod tests {
|
||||||
use super::*;
|
use super::*;
|
||||||
|
|
|
||||||
Loading…
Reference in a new issue