diff --git a/docs/swarm/README.md b/docs/swarm/README.md index 8dbd3f25..d498fddf 100644 --- a/docs/swarm/README.md +++ b/docs/swarm/README.md @@ -456,6 +456,17 @@ Swarm-side, `GET /api/agents//state/stream` relays that subject as SSE, resolving the agent's hive at request time exactly as the terminal stream does. The payload passes through opaquely — the controller never parses a header. +The third thing an agent publishes is its **icon**, the same SVG its own +`GET /icon` serves. It goes into the `agent-icons` KV bucket under the key +``, with no hive in it, so the swarm can show the icon of an agent that's +stopped or has moved hives. Only an agent connected with its own queue +credential publishes it: the queue grants that credential +`$KV.agent-icons.` and no other key, and grants the hive's shared client +none of the bucket. The harness writes once per start, because a config change +reaches an agent by restarting its container. An agent with no icon deletes its +key. Swarm-side, `GET /api/agents//icon` serves the stored bytes, and 404 +means the agent has no icon. + ### Swarm-wide forge objects The controller also keeps the forge objects that are one per swarm, not diff --git a/hive-agent/Cargo.toml b/hive-agent/Cargo.toml index b928b7eb..6571aff1 100644 --- a/hive-agent/Cargo.toml +++ b/hive-agent/Cargo.toml @@ -9,9 +9,10 @@ workspace = true [dependencies] anyhow.workspace = true -# Named directly only for the client type the terminal publisher holds; the -# connect and the credential handling live in `swarm-queue-client` below. -async-nats.workspace = true +# Named directly for the client type the publishers hold, and for `jetstream`, +# which the icon publisher's acked write needs. The connect and the credential +# handling live in `swarm-queue-client` below. +async-nats = { workspace = true, features = ["jetstream"] } axum.workspace = true chrono.workspace = true reqwest.workspace = true diff --git a/hive-agent/src/main.rs b/hive-agent/src/main.rs index ac1e0ca2..e1614007 100644 --- a/hive-agent/src/main.rs +++ b/hive-agent/src/main.rs @@ -28,6 +28,7 @@ mod serve_common; mod state_entry_watch; mod stats; mod stream_enrich; +mod swarm_agent_icon; mod swarm_agent_state; mod swarm_queue; mod swarm_term; @@ -521,6 +522,8 @@ async fn serve_main(socket: &Path, poll_ms: u64) -> Result<()> { // state from this agent's first transition, not from whenever a reader // first asked. swarm_agent_state::spawn(&bus); + // Once per start: a new icon only arrives with a new harness process. + swarm_agent_icon::spawn(); // Set by the web UI's `/api/cancel` on a successful SIGINT, read-and- // cleared by `handle_turn` before building the next wake prompt — see // `hive_sh4re::inbox::INTERRUPTED_HINT`. Shared between the web server task and diff --git a/hive-agent/src/swarm_agent_icon.rs b/hive-agent/src/swarm_agent_icon.rs new file mode 100644 index 00000000..3a421ba2 --- /dev/null +++ b/hive-agent/src/swarm_agent_icon.rs @@ -0,0 +1,220 @@ +//! Publishing this agent's icon into the swarm's `agent-icons` bucket. +//! +//! The bytes are the file the harness's own `GET /icon` serves +//! ([`ICON_PATH`]), unmodified. The key is the agent name, so the swarm can +//! show the icon while the agent is stopped or after it has moved hives. +//! +//! **Once per harness start.** The icon comes from the agent's NixOS config, +//! and a config change is applied by stopping the container, switching it and +//! starting it again, so a new icon always arrives with a new harness process. +//! An agent with no icon deletes its key, so an icon removed from the config +//! also disappears from the swarm. +//! +//! **Only under the agent's own credential.** That credential is granted +//! exactly this agent's key. The hive's shared client is granted none of the +//! bucket, so under it this publishes nothing. +//! +//! The write is a `JetStream` publish straight to the key's subject, the same +//! message `kv::Store::put` sends, without opening the bucket first. That +//! keeps the grant to the one subject. swarm-controller creates the bucket, +//! so a write can arrive before it exists; that write and any other failed +//! one are retried with backoff until one is acked. + +use std::path::Path; +use std::time::Duration; + +use async_nats::jetstream::context::PublishErrorKind; + +use crate::swarm_queue::{Connection, Presented}; + +/// Where the agent's NixOS config puts the icon +/// (`services.hyperhive.agent.icon`), and the file `GET /icon` serves. Absent +/// when none is configured. +const ICON_PATH: &str = "/etc/hyperhive/icon.svg"; + +/// First retry delay after a failed write, doubled per failure up to +/// [`MAX_RETRY`]. +const FIRST_RETRY: Duration = Duration::from_secs(5); + +/// Longest wait between retries. +const MAX_RETRY: Duration = Duration::from_mins(5); + +/// The KV header and value that mark an entry deleted. `kv::Store::get` +/// answers `None` for an entry carrying them. +const KV_OPERATION: &str = "KV-Operation"; +const KV_OPERATION_DELETE: &str = "DEL"; + +/// What to write under this agent's key. +#[derive(Debug, PartialEq, Eq)] +enum Icon { + /// The SVG's bytes. + Set(Vec), + /// No icon is configured, so the key is deleted. + Unset, +} + +/// Read the icon at `path`. A missing file is [`Icon::Unset`]; any other read +/// error is returned, and the key is left as it is. +fn read_icon(path: &Path) -> std::io::Result { + match std::fs::read(path) { + Ok(bytes) => Ok(Icon::Set(bytes)), + Err(e) if e.kind() == std::io::ErrorKind::NotFound => Ok(Icon::Unset), + Err(e) => Err(e), + } +} + +/// The subject this agent's icon is written to, or `None` under a credential +/// that is not granted it. +fn subject(presented: &Presented, agent: &str) -> Option { + match presented { + Presented::Agent => Some(swarm_queue_client::agent_icon::subject(agent)), + Presented::Hive { .. } => None, + } +} + +/// Start the publish task, if this agent has a queue credential and a label. +pub fn spawn() { + if !crate::swarm_queue::configured() { + return; + } + let agent = crate::identity::label(); + if agent.is_empty() { + tracing::warn!("this agent has no label; not publishing its icon upward"); + return; + } + tokio::spawn(run(agent)); +} + +async fn run(agent: String) { + let Some(Connection { client, presented }) = crate::swarm_queue::client().await else { + return; + }; + let Some(subject) = subject(&presented, &agent) else { + tracing::info!( + "connected with the hive's shared client, which may not write an agent's \ + icon; not publishing it upward" + ); + return; + }; + let icon = match read_icon(Path::new(ICON_PATH)) { + Ok(icon) => icon, + Err(e) => { + tracing::warn!(path = ICON_PATH, error = %e, "reading the agent icon failed; not publishing it"); + return; + } + }; + let js = async_nats::jetstream::new(client.clone()); + let mut delay = FIRST_RETRY; + loop { + match write(&client, &js, &subject, &icon).await { + Ok(()) => { + tracing::info!( + subject, + set = matches!(icon, Icon::Set(_)), + "published this agent's icon" + ); + return; + } + Err(Failure::TooLarge(e)) => { + tracing::warn!(error = %e, "the agent icon is over the queue's payload limit; not publishing it"); + return; + } + Err(Failure::Retry(e)) => { + tracing::warn!(error = e, retry_in = ?delay, "publishing the agent icon failed"); + } + } + tokio::time::sleep(delay).await; + delay = (delay * 2).min(MAX_RETRY); + } +} + +/// Why a write did not land. +enum Failure { + /// The payload can never fit; retrying cannot help. + TooLarge(async_nats::jetstream::context::PublishError), + /// Anything else, rendered for the log. + Retry(String), +} + +/// One acked write of `icon` to `subject`. +async fn write( + client: &async_nats::Client, + js: &async_nats::jetstream::Context, + subject: &str, + icon: &Icon, +) -> Result<(), Failure> { + // An unconnected client does not fail a publish, it buffers it, and the + // payload limit it checks against is the library default until the + // server has announced its own. + swarm_queue_client::ensure_connected(client) + .map_err(|e| Failure::Retry(swarm_queue_client::chain(&e)))?; + let sent = match icon { + Icon::Set(bytes) => js.publish(subject.to_owned(), bytes.clone().into()).await, + Icon::Unset => { + let mut headers = async_nats::HeaderMap::new(); + headers.insert(KV_OPERATION, KV_OPERATION_DELETE); + js.publish_with_headers(subject.to_owned(), headers, Vec::new().into()) + .await + } + }; + let ack = sent.map_err(|e| match e.kind() { + PublishErrorKind::MaxPayloadExceeded => Failure::TooLarge(e), + _ => Failure::Retry(e.to_string()), + })?; + // `StreamNotFound` here is the bucket not existing yet. + ack.await.map_err(|e| Failure::Retry(e.to_string()))?; + Ok(()) +} + +#[cfg(test)] +mod tests { + use super::{Icon, Presented, read_icon, subject}; + + #[test] + fn under_its_own_credential_an_agent_writes_its_own_key() { + assert_eq!( + subject(&Presented::Agent, "atlas").as_deref(), + Some("$KV.agent-icons.atlas") + ); + } + + /// The hive's shared client is granted no icon key, so a publish under it + /// would be refused. + #[test] + fn under_the_hives_shared_client_nothing_is_written() { + let hive = Presented::Hive { + client_id: "hive-alpha-agent".to_owned(), + }; + assert_eq!(subject(&hive, "atlas"), None); + } + + #[test] + fn a_configured_icon_is_published_byte_for_byte() { + let dir = tempfile::tempdir().expect("tempdir"); + let path = dir.path().join("icon.svg"); + std::fs::write(&path, b"\n").expect("write"); + assert_eq!( + read_icon(&path).expect("read"), + Icon::Set(b"\n".to_vec()) + ); + } + + /// No file is the unconfigured state, which clears the key rather than + /// leaving a removed icon in the swarm. + #[test] + fn no_icon_file_clears_the_key() { + let dir = tempfile::tempdir().expect("tempdir"); + assert_eq!( + read_icon(&dir.path().join("icon.svg")).expect("read"), + Icon::Unset + ); + } + + /// Any other read failure is an error, so a transient one does not delete + /// an icon the agent still has. + #[test] + fn an_unreadable_icon_is_an_error_not_an_unset_icon() { + let dir = tempfile::tempdir().expect("tempdir"); + assert!(read_icon(dir.path()).is_err()); + } +} diff --git a/nix/host-modules/swarm-nats.nix b/nix/host-modules/swarm-nats.nix index fe1bff9b..c8b30c39 100644 --- a/nix/host-modules/swarm-nats.nix +++ b/nix/host-modules/swarm-nats.nix @@ -956,6 +956,16 @@ in # the same two streams, keyed on the agent alone. "--agent-token-publish-subject ${lib.escapeShellArg "\$\$SWARM.term.{agent}"}" "--agent-token-publish-subject ${lib.escapeShellArg "\$\$SWARM.agent-state.{agent}"}" + # Its own key in the `agent-icons` KV bucket, and nothing else + # of JetStream. A KV put is a JetStream publish to + # `$KV..`, acked on the client's inbox, so this + # one subject is the whole write. The harness publishes to the + # subject directly instead of opening the bucket, so it needs + # no `$JS.API.*` subject; swarm-controller creates the bucket. + # Not granted to the hive's shared agent client: every agent + # on a hive presents that one, so it cannot be scoped to one + # agent's key. + "--agent-token-publish-subject ${lib.escapeShellArg "\$\$KV.agent-icons.{agent}"}" "--store-cert-role ${lib.escapeShellArg authCertRole}" ]; # Every credential arrives by `LoadCredential` and is named diff --git a/nix/module-eval/nats-authelia.nix b/nix/module-eval/nats-authelia.nix index cce9c1aa..0ed2f2d0 100644 --- a/nix/module-eval/nats-authelia.nix +++ b/nix/module-eval/nats-authelia.nix @@ -151,6 +151,16 @@ let && lib.hasInfix "--agent-token-publish-subject '$$SWARM.agent-state.{agent}'" exec && lib.hasInfix "--store-cert-role swarm-nats-auth" exec; } + { + # Its own arm, like the two above: the icon write is a separate feature + # and dropping it must not hide behind the streams still being granted. + name = "the responder grants a verified agent its own icon key"; + ok = + let + exec = (responderOf natsOldPath).serviceConfig.ExecStart; + in + lib.hasInfix "--agent-token-publish-subject '$$KV.agent-icons.{agent}'" exec; + } { # Each credential by `LoadCredential`, and the environment naming where # the unit sees it: a DynamicUser cannot read the copies directly. diff --git a/swarm-controller/src/agent_icon.rs b/swarm-controller/src/agent_icon.rs index c735af16..bf4f019d 100644 --- a/swarm-controller/src/agent_icon.rs +++ b/swarm-controller/src/agent_icon.rs @@ -7,8 +7,10 @@ //! property of it, so there is nothing for a timestamp to mean here. //! //! **No hive is named on this path**, because the key does not carry one — -//! see `swarm_queue_client::agent_icon` for why, and for the grant work -//! the publish side is still waiting on. +//! see `swarm_queue_client::agent_icon` for why, and for who writes it. +//! +//! This reader is also what creates the bucket: the agents that write it are +//! granted their own key and nothing that could create a stream. use anyhow::{Context, Result}; @@ -37,6 +39,26 @@ impl AgentIconReader { .await } + /// Open the bucket, creating it, as soon as the queue is connected. + /// + /// Agents write their keys without opening the bucket, so until this or + /// a first [`Self::get`] has run, every write is refused and retried. + pub async fn create_when_connected(&self) { + const POLL: std::time::Duration = std::time::Duration::from_secs(5); + loop { + if swarm_queue_client::ensure_connected(&self.client).is_ok() { + match self.store().await { + Ok(_) => return, + Err(e) => tracing::warn!( + error = swarm_queue_client::chain(&e), + "creating the agent-icon bucket failed; retrying" + ), + } + } + tokio::time::sleep(POLL).await; + } + } + /// `agent`'s icon, or `None` when it has published none. /// /// The bytes are returned unopened: this daemon does not parse, rewrite diff --git a/swarm-controller/src/main.rs b/swarm-controller/src/main.rs index 6fc30988..63a2f6b7 100644 --- a/swarm-controller/src/main.rs +++ b/swarm-controller/src/main.rs @@ -946,10 +946,15 @@ fn agent_status_reader( /// same rationale as [`wanted_writer`]. No staleness threshold, unlike /// [`agent_status_reader`]: an icon is a property of the agent rather than /// a report about it, so there is no cadence it could be late against. +/// +/// Also starts creating the bucket, which the agents writing it cannot do. fn agent_icon_reader( status: Option<&Arc>, ) -> Option> { - status.map(|s| Arc::new(agent_icon::AgentIconReader::new(s.queue_client()))) + let reader = Arc::new(agent_icon::AgentIconReader::new(status?.queue_client())); + let creating = Arc::clone(&reader); + tokio::spawn(async move { creating.create_when_connected().await }); + Some(reader) } /// A hive name that is shaped like one and names a hive this swarm has. @@ -2034,9 +2039,6 @@ const ICON_MAX_AGE: u32 = 300; /// harness's own `GET /icon` already has, so a caller falls back /// client-side on a failed load rather than probing first. /// -/// ⚠️ Until the agent-side publisher lands (see -/// `swarm_queue_client::agent_icon`), that 404 is every agent's answer. -/// /// 🩸 **The body is an untrusted document served from this daemon's own /// origin**, and an SVG can carry script. The response headers say it may /// not run any: `Content-Security-Policy: sandbox` with no `allow-scripts` diff --git a/swarm-nats-auth/src/agent_token.rs b/swarm-nats-auth/src/agent_token.rs index 4a929df4..c38898f2 100644 --- a/swarm-nats-auth/src/agent_token.rs +++ b/swarm-nats-auth/src/agent_token.rs @@ -232,6 +232,7 @@ mod tests { .with_agent_token_subjects(vec![ "$SWARM.term.{agent}".to_owned(), "$SWARM.agent-state.{agent}".to_owned(), + "$KV.agent-icons.{agent}".to_owned(), ]) .expect("valid") } @@ -253,6 +254,7 @@ mod tests { vec![ "$SWARM.term.atlas".to_owned(), "$SWARM.agent-state.atlas".to_owned(), + "$KV.agent-icons.atlas".to_owned(), ] ); } diff --git a/swarm-nats-auth/src/main.rs b/swarm-nats-auth/src/main.rs index d88a88ed..65673f73 100644 --- a/swarm-nats-auth/src/main.rs +++ b/swarm-nats-auth/src/main.rs @@ -359,12 +359,14 @@ mod tests { "$SWARM.term.{agent}", "--agent-token-publish-subject", "$SWARM.agent-state.{agent}", + "--agent-token-publish-subject", + "$KV.agent-icons.{agent}", "--store-cert-role", "swarm-nats-auth", ]) .expect("the unit's own argument vector must parse"); assert_eq!(args.agent_client_suffix, "-agent"); - assert_eq!(args.agent_token_publish_subjects.len(), 2); + assert_eq!(args.agent_token_publish_subjects.len(), 3); } /// The control for the case above: an ordinary value parses through the diff --git a/swarm-nats-auth/src/policy.rs b/swarm-nats-auth/src/policy.rs index 3da625eb..a88bf116 100644 --- a/swarm-nats-auth/src/policy.rs +++ b/swarm-nats-auth/src/policy.rs @@ -1232,6 +1232,34 @@ mod tests { ); } + /// A verified agent may write its own icon key and no other agent's. The + /// hive's shared agent client, which every agent on the hive presents, may + /// write none of the bucket. + #[test] + fn an_agent_may_write_only_its_own_icon_key() { + use swarm_queue_client::agent_icon::{BUCKET, subject}; + + let p = policy_with_agent_subject() + .with_agent_token_subjects(vec!["$KV.agent-icons.{agent}".to_owned()]) + .expect("a per-agent template is valid"); + let atlas = p.agent_token_permissions("atlas").expect("configured"); + assert_eq!(atlas.publish, vec![subject("atlas")]); + assert!(!atlas.publish.contains(&subject("argus"))); + // Control: the same template does grant argus its own key, so atlas + // lacking it is atlas's grant being scoped. + let argus = p.agent_token_permissions("argus").expect("configured"); + assert_eq!(argus.publish, vec![subject("argus")]); + + let shared = p.permissions("hive-alpha-agent").expect("alpha's agents"); + assert!( + !shared + .publish + .iter() + .any(|s| s.starts_with(&format!("$KV.{BUCKET}."))), + "the hive's shared agent client got an icon key: {shared:?}" + ); + } + #[test] fn an_agent_token_subject_without_the_placeholder_is_refused() { let err = policy() diff --git a/swarm-queue-client/src/agent_icon.rs b/swarm-queue-client/src/agent_icon.rs index 88187aa5..b7c08966 100644 --- a/swarm-queue-client/src/agent_icon.rs +++ b/swarm-queue-client/src/agent_icon.rs @@ -14,14 +14,17 @@ //! container is stopped: the value's lifetime is the agent's, not its //! placement's. //! -//! ⚠️ **Nothing writes this bucket yet.** The publisher runs inside the -//! agent's own container, and an agent's NATS grants are hive-scoped -//! (`$KV...*`), which cannot authorise a write to a -//! single-token agent key — so the publish side is blocked on a grant -//! shape the policy layer does not have today. Until it lands, every read -//! here answers "no icon", which is the same answer an agent that never -//! set one gets, and the same 404 the per-agent harness's own `GET /icon` -//! has always returned for an unconfigured agent. +//! **Each agent writes its own key** with its own queue credential. The auth +//! callout grants that credential `$KV.agent-icons.` and no other key. +//! The hive's shared agent client gets nothing in this bucket: every agent on +//! a hive presents that same client, so its grant cannot be narrowed to one +//! agent's key. +//! +//! The writer publishes to [`crate::agent_icon::subject`] directly instead of +//! opening the bucket, so that one subject is its whole grant. Creating the +//! bucket is left to the reader, with `open_or_create`. A write that arrives +//! before the bucket exists is refused with "no stream", and the writer +//! retries. #[cfg(feature = "kv")] use crate::Error; @@ -42,6 +45,16 @@ pub const BUCKET: &str = "agent-icons"; /// than an SVG does not have an agent icon. pub const MEDIA_TYPE: &str = "image/svg+xml"; +/// The subject `agent`'s entry is written on: `$KV..`, where the +/// key is the agent name and nothing else. +/// +/// A KV put is a `JetStream` publish to this subject, so a writer that +/// publishes directly has to use the same subject the bucket's `get` reads. +#[must_use] +pub fn subject(agent: &str) -> String { + format!("$KV.{BUCKET}.{agent}") +} + /// Open the agent-icon bucket, creating it if nothing has yet. /// /// `history: 1`, same rationale as [`crate::agent_status::open_or_create`]: @@ -77,15 +90,13 @@ pub async fn open_or_create( #[cfg(test)] mod tests { - use super::BUCKET; + use super::subject; - /// The published subject is `$KV..`, and this key is the - /// agent name alone — so the subject carries **one** token after the - /// bucket. Pinned here because that is precisely what a hive-scoped - /// grant (`$KV...*`, two tokens) cannot match, and the - /// reason the publish side needs a grant shape of its own. + /// One token after the bucket, the agent name alone. That is what the + /// per-agent grant `$KV.agent-icons.{agent}` expands to, and what a + /// hive-scoped grant (`$KV...*`, two tokens) cannot match. #[test] fn the_published_subject_carries_the_agent_and_no_placement() { - assert_eq!(format!("$KV.{BUCKET}.iris"), "$KV.agent-icons.iris"); + assert_eq!(subject("iris"), "$KV.agent-icons.iris"); } } diff --git a/swarm-queue-client/src/lib.rs b/swarm-queue-client/src/lib.rs index d7e673a6..05fed75e 100644 --- a/swarm-queue-client/src/lib.rs +++ b/swarm-queue-client/src/lib.rs @@ -184,8 +184,7 @@ pub mod agent_token; /// The bucket agent icons are published into — one key per agent, with no /// hive in it, unlike [`agent_status`]. See the module doc for why the -/// placement stays out of the key, and for what still has to land before -/// anything can write it. +/// placement stays out of the key, and for who may write it. pub mod agent_icon; /// The subject the swarm controller publishes on when the hive-wide knowledge