From 727065c9601aa2bab639f77daa258dfbc2bba84f Mon Sep 17 00:00:00 2001 From: atlas Date: Sat, 26 Sep 2026 01:51:13 +0200 Subject: [PATCH] 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..` as a static token and publishes on `$SWARM.term.` and `$SWARM.agent-state.`. 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 --- docs/swarm/README.md | 13 ++- hive-agent/src/swarm_agent_state.rs | 72 +++++++++---- hive-agent/src/swarm_queue.rs | 161 ++++++++++++++++++++++++---- hive-agent/src/swarm_term.rs | 75 +++++++++---- swarm-queue-client/src/lib.rs | 42 ++++++-- 5 files changed, 287 insertions(+), 76 deletions(-) diff --git a/docs/swarm/README.md b/docs/swarm/README.md index 37808df0..615d5281 100644 --- a/docs/swarm/README.md +++ b/docs/swarm/README.md @@ -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. What an agent does with that connection is publish its terminal. Every row its -own web UI renders also goes to `$SWARM.term..`, one subject per -agent, so a swarm-level terminal can follow one agent without subscribing to -the swarm's whole traffic. The `` is the one the agent's client id names, -which is the same string the broker builds its grant from. Publishing only: an +own web UI renders also goes to `$SWARM.term.`, one subject per agent, so +a swarm-level terminal can follow one agent without subscribing to the swarm's +whole traffic. That is the subject an agent connected with its own queue +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..` instead, the `` 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 subscriber that wasn't listening missed them, the same as on the agent's own 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. The second thing an agent publishes is its **turn-state header**, on -`$SWARM.agent-state..` — same shape of subject, same grant +`$SWARM.agent-state.` (or `$SWARM.agent-state..`) — same shape of subject, same grant mechanics, same lack of retention. It carries what a header bar wants: what the turn loop is doing (`turn_state`, plus `turn_state_since` as an ISO 8601 UTC stamp), which model (`model` and the resolved id the last turn actually ran on), diff --git a/hive-agent/src/swarm_agent_state.rs b/hive-agent/src/swarm_agent_state.rs index 4c7a988e..f4d48d77 100644 --- a/hive-agent/src/swarm_agent_state.rs +++ b/hive-agent/src/swarm_agent_state.rs @@ -35,6 +35,7 @@ use tokio::sync::broadcast; use swarm_queue_client::wanted::AgentState; use crate::events::{Bus, BusEvent, LiveEvent, TurnState}; +use crate::swarm_queue::{Connection, Presented}; use crate::term_msg::iso8601_utc; /// Subject family carrying agent turn-state headers, the swarm-wide @@ -171,33 +172,43 @@ fn snapshot(bus: &Bus) -> AgentStateMsg { } } -/// Start the publish task, if this agent has both queue coordinates and an -/// identity the subject can be derived from. +/// The subject this agent's header goes to under the credential it connected +/// with, or `None` when a hive client id names no hive. +fn subject(presented: &Presented, agent: &str) -> Option { + match presented { + Presented::Agent => Some(format!("{SUBJECT_PREFIX}.{agent}")), + Presented::Hive { client_id } => { + let Some(hive) = hive_from_client_id(client_id) else { + tracing::warn!( + %client_id, + expected = format!("{CLIENT_ID_PREFIX}{CLIENT_ID_SUFFIX}"), + "queue client id does not name a hive; not publishing turn state upward" + ); + return None; + }; + Some(format!("{SUBJECT_PREFIX}.{hive}.{agent}")) + } + } +} + +/// Start the publish task, if this agent has a queue credential and a label +/// to name its subject with. /// -/// Returns without spawning in every other case — no queue, an unparseable -/// client id, no label — each of which is a legal state for an agent rather -/// than an error, and each logged once here rather than per transition. +/// 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) { - let Some(cfg) = crate::swarm_queue::config() else { + 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 Some(hive) = hive_from_client_id(&cfg.client_id) else { - tracing::warn!( - client_id = %cfg.client_id, - expected = format!("{CLIENT_ID_PREFIX}{CLIENT_ID_SUFFIX}"), - "queue client id does not name a hive; not publishing turn state upward" - ); - return; - }; + } let agent = crate::identity::label(); if agent.is_empty() { tracing::warn!("this agent has no label; not publishing turn state upward"); return; } - let subject = format!("{SUBJECT_PREFIX}.{hive}.{agent}"); - tokio::spawn(run(bus.subscribe(), bus.clone(), subject)); + tokio::spawn(run(bus.subscribe(), bus.clone(), agent)); } /// 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 /// any added later, funnels into the comparison and costs nothing when it /// changes nothing. -async fn run(mut rx: broadcast::Receiver, bus: Bus, subject: String) { - let Some(client) = crate::swarm_queue::client().await else { +async fn run(mut rx: broadcast::Receiver, bus: Bus, agent: String) { + let Some(Connection { client, presented }) = crate::swarm_queue::client().await else { + return; + }; + let Some(subject) = subject(&presented, &agent) else { return; }; tracing::info!(subject, "publishing agent turn state to the swarm queue"); @@ -307,7 +321,7 @@ async fn publish( #[cfg(test)] 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 swarm_queue_client::wanted::AgentState; @@ -393,10 +407,22 @@ mod tests { /// this pins the half that lives here. #[test] 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!( - format!("{SUBJECT_PREFIX}.{hive}.mara"), - "$SWARM.agent-state.alpha.mara" + subject(&presented, "mara").as_deref(), + Some("$SWARM.agent-state.alpha.mara") + ); + } + + /// An agent that connected with its own credential publishes on the + /// subject that credential is granted, which names no hive. + #[test] + fn an_agent_on_its_own_credential_publishes_on_its_hive_free_subject() { + assert_eq!( + subject(&Presented::Agent, "mara").as_deref(), + Some("$SWARM.agent-state.mara") ); } diff --git a/hive-agent/src/swarm_queue.rs b/hive-agent/src/swarm_queue.rs index b8d267e9..3613e02a 100644 --- a/hive-agent/src/swarm_queue.rs +++ b/hive-agent/src/swarm_queue.rs @@ -21,15 +21,16 @@ //! 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`). //! -//! It is reported here but not yet *presented*: the queue's auth-callout -//! responder (`swarm-nats-auth`) validates only the hive-scoped token, and an -//! agent offering a credential nothing on the other end reads back would be -//! refused. Until that responder learns the same path, the connect path below -//! is unchanged and this is the fetching half. +//! 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. use std::path::{Path, PathBuf}; use std::sync::OnceLock; +use anyhow::Context as _; use tokio::sync::OnceCell; use swarm_queue_client::QueueConfig; @@ -42,10 +43,39 @@ const ENV_PREFIX: &str = "HIVE_AGENT"; /// one answer rather than re-deriving it per call. static CONFIG: OnceLock> = OnceLock::new(); +/// Where to present this agent's own credential, resolved once at boot. +static AGENT: OnceLock> = OnceLock::new(); + /// The one connection every publisher in this process shares. Separate from /// [`CONFIG`] because resolving the coordinates is synchronous boot work and /// connecting is not — see [`client`]. -static CLIENT: OnceCell> = OnceCell::const_new(); +static CLIENT: OnceCell> = OnceCell::const_new(); + +/// What connecting with this agent's own credential needs. +#[derive(Debug, PartialEq, Eq)] +struct AgentPath { + url: String, + ca_file: Option, + /// 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 `.`. + Agent, + /// The hive's shared client, granted `..>` 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 /// 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 { } } +/// Decide whether this agent can connect with its own credential: it needs the +/// queue's address and a fetched secret, and nothing of the hive's client. +fn decide_agent_path(env: &QueueEnv, secret_file: Option) -> Option { + Some(AgentPath { + url: env.nats_url.clone()?, + ca_file: env.ca_file.as_ref().map(Into::into), + secret_file: secret_file?, + }) +} + /// Decide what this agent's queue configuration is, given the environment and /// whatever the client-id file held. /// @@ -216,15 +256,10 @@ pub fn init() { let _ = CONFIG.set(resolved); // 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 - // of its own yet. Reported either way, because "which credential is this - // agent able to present" is a question only this process can answer, and - // it is the one the next slice's rollout will be asked repeatedly. - // - // The answer is only logged here. Presenting it needs the queue's - // auth-callout responder to verify it, which is the next slice — see this - // module's header. - if let Some(path) = decide_agent_secret(env.agent_secret_file.as_deref()) { + // 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. tracing::info!( path = %path.display(), @@ -236,6 +271,13 @@ pub fn init() { by its hive's shared client" ); } + let _ = AGENT.set(decide_agent_path(&env, secret)); +} + +/// Whether this agent has any credential to reach the queue with. `false` +/// before [`init`] has run. +pub fn configured() -> bool { + config().is_some() || AGENT.get().is_some_and(Option::is_some) } /// What [`init`] resolved, or `None` when this agent has no queue. @@ -263,14 +305,48 @@ pub fn config() -> Option<&'static QueueConfig> { /// failed" — a caller does nothing differently between them, since either way /// there is nothing to publish onto. Connecting is lazy so that an agent on a /// hive with no queue pays nothing at boot. -pub async fn client() -> Option { +pub async fn client() -> Option { CLIENT.get_or_init(connect_once).await.clone() } -async fn connect_once() -> Option { +/// This agent's own credential first, then the hive's. Each path taken is +/// logged once, here. +async fn connect_once() -> Option { + if let Some(agent) = AGENT.get().and_then(Option::as_ref) { + match connect_as_agent(agent).await { + 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()?; 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) => { // `chain`, not `{:#}`: this is `swarm_queue_client::Error`, whose // `Display` ignores the alternate flag, so `{:#}` renders the @@ -284,9 +360,26 @@ async fn connect_once() -> Option { } } +/// Connect presenting this agent's own token. Any failure, a refusal +/// included, is returned for [`connect_once`] to fall back on. +async fn connect_as_agent(agent: &AgentPath) -> anyhow::Result { + let name = crate::identity::label(); + let secret = std::fs::read_to_string(&agent.secret_file) + .with_context(|| format!("reading {}", agent.secret_file.display()))?; + let token = swarm_queue_client::agent_token::format_agent_token(&name, secret.trim())?; + swarm_queue_client::connect_with_token(&agent.url, agent.ca_file.as_deref(), token) + .await + .map_err(|e| anyhow::anyhow!(swarm_queue_client::chain(&e))) +} + #[cfg(test)] 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 { 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!(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 + ); + } } diff --git a/hive-agent/src/swarm_term.rs b/hive-agent/src/swarm_term.rs index 249e2b57..620b81d4 100644 --- a/hive-agent/src/swarm_term.rs +++ b/hive-agent/src/swarm_term.rs @@ -27,6 +27,7 @@ use tokio::sync::broadcast; use crate::events::BusEvent; +use crate::swarm_queue::{Connection, Presented}; use crate::term_msg::{ClassifyCtx, TermMsg, classify}; /// Subject family carrying agent terminal rows, the swarm-wide agreement this @@ -112,37 +113,50 @@ fn serialized_len(msg: &TermMsg) -> Option { serde_json::to_vec(msg).ok().map(|v| v.len()) } -/// Start the publish task, if this agent has both queue coordinates and an -/// identity the subject can be derived from. +/// The subject this agent's rows go to under the credential it connected +/// with, or `None` when a hive client id names no hive. +fn subject(presented: &Presented, agent: &str) -> Option { + match presented { + Presented::Agent => Some(format!("{SUBJECT_PREFIX}.{agent}")), + Presented::Hive { client_id } => { + let Some(hive) = hive_from_client_id(client_id) else { + tracing::warn!( + %client_id, + expected = format!("{CLIENT_ID_PREFIX}{CLIENT_ID_SUFFIX}"), + "queue client id does not name a hive; not publishing the terminal upward" + ); + return None; + }; + Some(format!("{SUBJECT_PREFIX}.{hive}.{agent}")) + } + } +} + +/// Start the publish task, if this agent has a queue credential and a label +/// to name its subject with. /// -/// Returns without spawning in every other case — no queue, an unparseable -/// client id, no label — each of which is a legal state for an agent rather -/// than an error, and each logged once here rather than per row. +/// 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) { - let Some(cfg) = crate::swarm_queue::config() else { + 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 Some(hive) = hive_from_client_id(&cfg.client_id) else { - tracing::warn!( - client_id = %cfg.client_id, - expected = format!("{CLIENT_ID_PREFIX}{CLIENT_ID_SUFFIX}"), - "queue client id does not name a hive; not publishing the terminal upward" - ); - return; - }; + } let agent = crate::identity::label(); if agent.is_empty() { tracing::warn!("this agent has no label; not publishing the terminal upward"); return; } - let subject = format!("{SUBJECT_PREFIX}.{hive}.{agent}"); - tokio::spawn(run(rx, subject)); + tokio::spawn(run(rx, agent)); } -async fn run(mut rx: broadcast::Receiver, subject: String) { - let Some(client) = crate::swarm_queue::client().await else { +async fn run(mut rx: broadcast::Receiver, agent: String) { + let Some(Connection { client, presented }) = crate::swarm_queue::client().await else { + return; + }; + let Some(subject) = subject(&presented, &agent) else { return; }; 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)] 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::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] fn a_row_that_already_fits_is_published_unchanged() { let msg = TermMsg::new(Level::Info, "turn ok") diff --git a/swarm-queue-client/src/lib.rs b/swarm-queue-client/src/lib.rs index 20c5a666..8676a15c 100644 --- a/swarm-queue-client/src/lib.rs +++ b/swarm-queue-client/src/lib.rs @@ -320,6 +320,15 @@ const TOKEN_REFRESH_SKEW: std::time::Duration = std::time::Duration::from_mins(2 /// degrades every other client of it. const MAX_RECONNECT_DELAY: std::time::Duration = std::time::Duration::from_mins(1); +/// Exponential from 500ms, capped at [`MAX_RECONNECT_DELAY`]. +fn reconnect_delay(attempts: usize) -> std::time::Duration { + let exp = u32::try_from(attempts.saturating_sub(1)).unwrap_or(u32::MAX); + std::cmp::min( + std::time::Duration::from_millis(500u64.saturating_mul(2u64.saturating_pow(exp.min(8)))), + MAX_RECONNECT_DELAY, + ) +} + /// Where the controller finds the queue and what it authenticates with. /// /// Every field comes from an environment variable the NixOS module sets, the @@ -742,15 +751,7 @@ pub async fn connect(cfg: QueueConfig) -> Result { // which turns an unreachable queue into a permanent 4s poll — and, before // the cache above, a permanent 4s token-request loop against authelia. // Exponential from 500ms so a momentary blip still reconnects promptly. - .reconnect_delay_callback(|attempts| { - let exp = u32::try_from(attempts.saturating_sub(1)).unwrap_or(u32::MAX); - std::cmp::min( - std::time::Duration::from_millis( - 500u64.saturating_mul(2u64.saturating_pow(exp.min(8))), - ), - MAX_RECONNECT_DELAY, - ) - }) + .reconnect_delay_callback(reconnect_delay) // The controller and the queue are separate units on (possibly) // separate hosts, and nothing orders them. Without this, a queue that // comes up one second later leaves the controller permanently @@ -784,6 +785,29 @@ pub async fn connect(cfg: QueueConfig) -> Result { Ok(client) } +/// Connect presenting a fixed `token`, as an agent presents its own queue +/// credential. Reconnects present the same token. +/// +/// 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 { + let mut options = + async_nats::ConnectOptions::with_token(token).reconnect_delay_callback(reconnect_delay); + if let Some(path) = ca_file { + options = options.add_root_certificates(path.to_path_buf()); + } + options.connect(url).await.map_err(|source| Error::Connect { + url: url.to_owned(), + source, + }) +} + #[cfg(test)] mod tests { use super::*;