Watch
0
0
Fork
You've already forked hyperhive
0

swarm: let an agent publish its own icon

The auth callout grants an agent that presents its own queue credential
one more subject, `$KV.agent-icons.<agent>`: its own key in the
agent-icons bucket and no other. The hive's shared agent client is
granted none of the bucket, since every agent on a hive presents it.

hive-agent writes `/etc/hyperhive/icon.svg`, the file its `GET /icon`
serves, to that key once per start, as a JetStream publish straight to
the subject (what `kv::Store::put` sends, minus the bucket lookup), so
the one subject is the whole grant. No icon deletes the key. A failed
write, including one that arrives before the bucket exists, is retried
with backoff until acked. An agent connected with the hive's shared
client publishes nothing.

swarm-controller creates the bucket as soon as its queue connection is
up, instead of on the first icon read, so an agent's write does not
wait for someone to look.

Measured against a local nats-server with a user allowed publish on
`$KV.agent-icons.atlas` only: the write to its own key is stored and
readable, a write to `$KV.agent-icons.argus` is refused (the ack times
out), the DEL marker makes the key read as absent, and a write before
the bucket exists fails with "no responders".
This commit is contained in:
atlas 2026-09-28 10:55:45 +02:00 • committed by mara
commit e974194e3a
13 changed files with 348 additions and 27 deletions

View file

@ -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

View file

@ -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<S: Surface>(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

View file

@ -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<u8>),
/// 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<Icon> {
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<String> {
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"<svg/>\n").expect("write");
assert_eq!(
read_icon(&path).expect("read"),
Icon::Set(b"<svg/>\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());
}
}