Watch
0
0
Fork
You've already forked hyperhive
0

Compare commits

..
16 changed files with 43 additions and 752 deletions

View file

@ -456,18 +456,6 @@ Swarm-side, `GET /api/agents/<name>/state/stream` relays that subject as SSE,
resolving the agent's hive at request time exactly as the terminal stream does. 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 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
`<agent>`, 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.<agent>` 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/<name>/icon` serves the stored bytes, and 404
means the agent has no icon. swarm-ui's agent cards load it as an `<img>` and show the
dimmed hyperhive mark for an agent without one, as the hive dashboard does.
### Swarm-wide forge objects ### Swarm-wide forge objects
The controller also keeps the forge objects that are one per swarm, not The controller also keeps the forge objects that are one per swarm, not

View file

@ -22,37 +22,3 @@
white-space: normal; white-space: normal;
word-break: break-word; word-break: break-word;
} }
/* Icon beside the name/message, sized and dimmed like the hive
dashboard's `.container-icon` (dashboard.css). */
.ui-agent-card-row {
display: flex;
align-items: flex-start;
gap: 0.7em;
}
.ui-agent-card-body {
display: flex;
flex-direction: column;
gap: 0.35em;
flex: 1;
min-width: 0;
}
.ui-agent-card-icon {
position: relative;
overflow: hidden;
flex: none;
width: 5em;
aspect-ratio: 1;
border-radius: 6px;
background-color: color-mix(in srgb, var(--crust) 60%, transparent);
}
.ui-agent-card-icon > img {
position: absolute;
inset: 0;
width: 100%;
height: 100%;
object-fit: contain;
}
.ui-agent-card-icon-none {
filter: grayscale(1);
opacity: 0.4;
}

View file

@ -12,7 +12,6 @@
// `WantedMenu` below keeps its own clicks (mouse or keyboard-synthesized, // `WantedMenu` below keeps its own clicks (mouse or keyboard-synthesized,
// `WantedMenu` hosts real `<button>`s) from also firing `Card`'s // `WantedMenu` hosts real `<button>`s) from also firing `Card`'s
// `onClick` and opening the detail panel. // `onClick` and opening the detail panel.
import { useState } from "preact/hooks";
import { Badge } from "@hive/shared/badge.js"; import { Badge } from "@hive/shared/badge.js";
import type { ProblemDetails } from "@hive/shared/api-error.js"; import type { ProblemDetails } from "@hive/shared/api-error.js";
import { Card } from "../../ui/card/Card.js"; import { Card } from "../../ui/card/Card.js";
@ -47,9 +46,6 @@ export function AgentCard({
const { tone, label } = FRESHNESS[row.freshness]; const { tone, label } = FRESHNESS[row.freshness];
return ( return (
<Card selected={selected} onClick={() => onOpenDetail(row)}> <Card selected={selected} onClick={() => onOpenDetail(row)}>
<div class="ui-agent-card-row">
<AgentIcon name={row.name} />
<div class="ui-agent-card-body">
<div class="ui-agent-card-line1"> <div class="ui-agent-card-line1">
<span class="ui-agent-card-name">{row.name}</span> <span class="ui-agent-card-name">{row.name}</span>
<Badge <Badge
@ -71,10 +67,7 @@ export function AgentCard({
</> </>
} }
/> />
<span <span class="ui-agent-card-wanted" onClick={(e) => e.stopPropagation()}>
class="ui-agent-card-wanted"
onClick={(e) => e.stopPropagation()}
>
<WantedMenu <WantedMenu
row={row} row={row}
pending={pending} pending={pending}
@ -95,36 +88,6 @@ export function AgentCard({
<div class="ui-agent-card-message"> <div class="ui-agent-card-message">
{row.snapshot?.status_text ?? "—"} {row.snapshot?.status_text ?? "—"}
</div> </div>
</div>
</div>
</Card> </Card>
); );
} }
// The agent's icon from swarm-controller's `GET /api/agents/{name}/icon`,
// the same square and fallback as the hive dashboard's container row: an
// agent with no icon (404) or an unreachable controller shows the dimmed
// hyperhive mark instead of a broken image. Always an `<img>`: the body is
// an agent-authored SVG, and only an image load keeps its script from running.
function AgentIcon({ name }: { name: string }) {
const [fallback, setFallback] = useState(false);
return (
<div
class={
fallback
? "ui-agent-card-icon ui-agent-card-icon-none"
: "ui-agent-card-icon"
}
>
<img
alt=""
src={
fallback
? "/favicon.svg"
: `/api/agents/${encodeURIComponent(name)}/icon`
}
onError={() => setFallback(true)}
/>
</div>
);
}

View file

@ -9,10 +9,9 @@ workspace = true
[dependencies] [dependencies]
anyhow.workspace = true anyhow.workspace = true
# Named directly for the client type the publishers hold, and for `jetstream`, # Named directly only for the client type the terminal publisher holds; the
# which the icon publisher's acked write needs. The connect and the credential # connect and the credential handling live in `swarm-queue-client` below.
# handling live in `swarm-queue-client` below. async-nats.workspace = true
async-nats = { workspace = true, features = ["jetstream"] }
axum.workspace = true axum.workspace = true
chrono.workspace = true chrono.workspace = true
reqwest.workspace = true reqwest.workspace = true

View file

@ -28,7 +28,6 @@ mod serve_common;
mod state_entry_watch; mod state_entry_watch;
mod stats; mod stats;
mod stream_enrich; mod stream_enrich;
mod swarm_agent_icon;
mod swarm_agent_state; mod swarm_agent_state;
mod swarm_queue; mod swarm_queue;
mod swarm_term; mod swarm_term;
@ -522,8 +521,6 @@ async fn serve_main<S: Surface>(socket: &Path, poll_ms: u64) -> Result<()> {
// state from this agent's first transition, not from whenever a reader // state from this agent's first transition, not from whenever a reader
// first asked. // first asked.
swarm_agent_state::spawn(&bus); 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- // 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 // cleared by `handle_turn` before building the next wake prompt — see
// `hive_sh4re::inbox::INTERRUPTED_HINT`. Shared between the web server task and // `hive_sh4re::inbox::INTERRUPTED_HINT`. Shared between the web server task and

View file

@ -1,220 +0,0 @@
//! 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());
}
}

View file

@ -956,16 +956,6 @@ in
# the same two streams, keyed on the agent alone. # 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.term.{agent}"}"
"--agent-token-publish-subject ${lib.escapeShellArg "\$\$SWARM.agent-state.{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.<bucket>.<key>`, 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}" "--store-cert-role ${lib.escapeShellArg authCertRole}"
]; ];
# Every credential arrives by `LoadCredential` and is named # Every credential arrives by `LoadCredential` and is named

View file

@ -151,16 +151,6 @@ let
&& lib.hasInfix "--agent-token-publish-subject '$$SWARM.agent-state.{agent}'" exec && lib.hasInfix "--agent-token-publish-subject '$$SWARM.agent-state.{agent}'" exec
&& lib.hasInfix "--store-cert-role swarm-nats-auth" 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 # Each credential by `LoadCredential`, and the environment naming where
# the unit sees it: a DynamicUser cannot read the copies directly. # the unit sees it: a DynamicUser cannot read the copies directly.

View file

@ -1,79 +0,0 @@
//! Reads an agent's icon out of the swarm's agent-icon KV bucket.
//!
//! Sibling of [`crate::agent_status`] and built the same way — the bucket
//! is resolved on first use and cached, a resolution failure is not — but
//! a much smaller read: one key, no roster to be complete against and no
//! freshness to derive. An icon is not a report about the agent, it is a
//! 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 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};
/// Reads the agent-icon bucket. Holds a NATS client rather than a bucket
/// handle, so a controller that starts before the bucket exists picks it
/// up without a restart.
pub struct AgentIconReader {
client: async_nats::Client,
store: tokio::sync::OnceCell<async_nats::jetstream::kv::Store>,
}
impl AgentIconReader {
#[must_use]
pub fn new(client: async_nats::Client) -> Self {
Self {
client,
store: tokio::sync::OnceCell::new(),
}
}
async fn store(
&self,
) -> std::result::Result<&async_nats::jetstream::kv::Store, swarm_queue_client::Error> {
self.store
.get_or_try_init(|| swarm_queue_client::agent_icon::open_or_create(&self.client))
.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
/// or validate the SVG, and a consumer that renders it must treat it as
/// the untrusted document it is — see `get_agent_icon`'s response
/// headers.
pub async fn get(&self, agent: &str) -> Result<Option<Vec<u8>>> {
// An unconnected client does not fail a JetStream request, it hangs
// on it — see `swarm_queue_client::ensure_connected`.
swarm_queue_client::ensure_connected(&self.client)?;
let store = self.store().await?;
let icon = store
.get(agent)
.await
.with_context(|| format!("reading the agent-icon entry for {agent}"))?;
Ok(icon.map(Into::into))
}
}

View file

@ -40,7 +40,6 @@ use swarm_authelia_bridge_sock::BridgeResponse;
use utoipa::{OpenApi, ToSchema}; use utoipa::{OpenApi, ToSchema};
use utoipa_axum::{router::OpenApiRouter, routes}; use utoipa_axum::{router::OpenApiRouter, routes};
mod agent_icon;
mod agent_identity; mod agent_identity;
mod agent_state_stream; mod agent_state_stream;
mod agent_status; mod agent_status;
@ -608,10 +607,6 @@ struct AppState {
/// connection, several consumers" rationale as `wanted` above. `None` /// connection, several consumers" rationale as `wanted` above. `None`
/// in exactly the state `status` is. /// in exactly the state `status` is.
agent_status: Option<Arc<agent_status::AgentStatusReader>>, agent_status: Option<Arc<agent_status::AgentStatusReader>>,
/// Per-agent icons, sharing `status`'s connection — same "one queue
/// connection, several consumers" rationale as `agent_status` above,
/// and `None` in exactly the state `status` is.
agent_icons: Option<Arc<agent_icon::AgentIconReader>>,
/// The swarm-level job graph, wrapped in its /// The swarm-level job graph, wrapped in its
/// [`hive_jobq::scheduler::Scheduler`] now that something drives it /// [`hive_jobq::scheduler::Scheduler`] now that something drives it
/// (`spawn_jobq_worker`) — the graph alone was enough for the /// (`spawn_jobq_worker`) — the graph alone was enough for the
@ -867,8 +862,7 @@ fn error_problem(status: axum::http::StatusCode, detail: &str) -> problem_detail
problem_details::ProblemDetails::from_status_code(status).with_detail(detail) problem_details::ProblemDetails::from_status_code(status).with_detail(detail)
} }
/// The status a [`wanted::WantedWriter`] or [`agent_icon::AgentIconReader`] /// The status a [`wanted::WantedWriter`] failure should answer with.
/// failure should answer with.
/// ///
/// `WantedWriter::view`/`set` call `ensure_connected` internally and /// `WantedWriter::view`/`set` call `ensure_connected` internally and
/// propagate through `anyhow`, so the concrete /// propagate through `anyhow`, so the concrete
@ -942,21 +936,6 @@ fn agent_status_reader(
}) })
} }
/// The per-agent icon reader, sharing the status reader's connection —
/// 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<status::StatusReader>>,
) -> Option<Arc<agent_icon::AgentIconReader>> {
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. /// A hive name that is shaped like one and names a hive this swarm has.
/// ///
/// Reports the status and the detail rather than a rendered /// Reports the status and the detail rather than a rendered
@ -2017,92 +1996,6 @@ where
Ok(Json(MakeForgeAdminResponse { change })) Ok(Json(MakeForgeAdminResponse { change }))
} }
/// How long a served icon may be reused without asking again.
///
/// An agent roster asks once per row, and an icon changes about as often
/// as an agent is reconfigured — so a few minutes of browser caching is
/// what keeps a page render from being one bucket read per agent, and
/// costs nothing anyone will notice.
const ICON_MAX_AGE: u32 = 300;
/// `agent`'s icon — the swarm-level answer to "what does this agent look
/// like".
///
/// **Nothing on this path names a hive.** The icon is a property of the
/// agent, held at swarm scope under the agent's own key, so it is
/// answerable for an agent that has moved hives or that no hive is
/// currently running. That is the whole point of it living here rather
/// than being fetched from wherever the agent happens to run.
///
/// **404 means this agent has no icon**, which is the expected
/// unconfigured state and not an error — the same contract the per-agent
/// harness's own `GET /icon` already has, so a caller falls back
/// client-side on a failed load rather than probing first.
///
/// 🩸 **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`
/// keeps a direct navigation to this URL from executing it, and `nosniff`
/// keeps a browser from deciding the bytes are something else. An `<img>`
/// load — the one consumer — never runs script regardless; the headers are
/// for every other way a URL gets opened.
#[utoipa::path(
get,
path = "/api/agents/{name}/icon",
params(("name" = String, Path, description = "agent name")),
responses(
(status = 200, description = "the agent's icon, an SVG", content_type = "image/svg+xml"),
(status = 400, description = "not an identifier (problem+json)", body = String),
(status = 404, description = "this agent has no icon (problem+json)", body = String),
(status = 500, description = "the icon bucket could not be read (problem+json)", body = String),
(status = 503, description = "no swarm queue is wired up, or it is not connected (problem+json)", body = String),
),
tag = "agents"
)]
async fn get_agent_icon(
State(state): State<AppState>,
Path(name): Path<String>,
) -> Result<axum::response::Response, problem_details::ProblemDetails> {
use axum::http::{StatusCode, header};
use axum::response::IntoResponse as _;
let Some(reader) = state.agent_icons.as_ref() else {
return Err(error_problem(
StatusCode::SERVICE_UNAVAILABLE,
"no swarm queue is wired up on this host",
));
};
let agent = hive_types::Ident::parse(&name)
.map_err(|reason| error_problem(StatusCode::BAD_REQUEST, reason))?
.into_string();
let icon = reader.get(&agent).await.map_err(|e| {
tracing::warn!(agent = %agent, error = %format!("{e:#}"), "reading the agent icon failed");
error_problem(wanted_error_status(&e), &format!("{e:#}"))
})?;
let Some(icon) = icon else {
return Err(error_problem(
StatusCode::NOT_FOUND,
"this agent has no icon",
));
};
Ok((
[
(
header::CONTENT_TYPE,
swarm_queue_client::agent_icon::MEDIA_TYPE.to_owned(),
),
(
header::CACHE_CONTROL,
format!("public, max-age={ICON_MAX_AGE}"),
),
(header::CONTENT_SECURITY_POLICY, "sandbox".to_owned()),
(header::X_CONTENT_TYPE_OPTIONS, "nosniff".to_owned()),
],
icon,
)
.into_response())
}
/// Every agent with an open config PR, in one response — the bulk /// Every agent with an open config PR, in one response — the bulk
/// counterpart to [`get_agent_config_pr`]. swarm-ui's config-PR table needs /// counterpart to [`get_agent_config_pr`]. swarm-ui's config-PR table needs
/// every agent's status to render, and fetching them one at a time doesn't /// every agent's status to render, and fetching them one at a time doesn't
@ -2410,7 +2303,6 @@ async fn main() -> Result<()> {
links: Arc::new(load_links()), links: Arc::new(load_links()),
wanted: wanted_writer(status.as_ref()), wanted: wanted_writer(status.as_ref()),
agent_status: agent_status_reader(status.as_ref()), agent_status: agent_status_reader(status.as_ref()),
agent_icons: agent_icon_reader(status.as_ref()),
status, status,
jobq, jobq,
webhook_secret, webhook_secret,
@ -2442,7 +2334,6 @@ fn build_app(state: AppState) -> axum::Router {
.routes(routes!(get_jobq_graph)) .routes(routes!(get_jobq_graph))
.routes(routes!(get_jobq_rollup)) .routes(routes!(get_jobq_rollup))
.routes(routes!(get_agent_config_pr)) .routes(routes!(get_agent_config_pr))
.routes(routes!(get_agent_icon))
.routes(routes!(get_config_prs)) .routes(routes!(get_config_prs))
.routes(routes!(create_agent)) .routes(routes!(create_agent))
.routes(routes!(mint_agent_identity)) .routes(routes!(mint_agent_identity))
@ -2564,7 +2455,6 @@ mod tests {
// that would write a declaration. // that would write a declaration.
wanted: None, wanted: None,
agent_status: None, agent_status: None,
agent_icons: None,
jobq: std::sync::Arc::clone(&sched), jobq: std::sync::Arc::clone(&sched),
webhook_secret: None, webhook_secret: None,
config_prs: None, config_prs: None,
@ -2616,33 +2506,6 @@ mod tests {
); );
} }
/// An icon route with no bucket behind it must say it could not ask,
/// not that the agent has no icon.
///
/// The two are one status code apart and read identically to an
/// `<img>` — both end in the fallback glyph — which is exactly why the
/// distinction has to be asserted here: nothing downstream would ever
/// notice the 503 silently becoming a 404.
#[tokio::test]
async fn an_icon_with_no_queue_refuses_rather_than_answering_no_icon() {
use axum::response::IntoResponse as _;
let (state, _sched) = state_with_roster();
assert!(
state.agent_icons.is_none(),
"the fixture must have no queue"
);
let err = super::get_agent_icon(
axum::extract::State(state),
axum::extract::Path("iris".to_owned()),
)
.await
.expect_err("a queue-less controller cannot answer an icon");
let resp = err.into_response();
assert_eq!(resp.status(), axum::http::StatusCode::SERVICE_UNAVAILABLE);
}
/// The roster check is the half that makes the recorded hive worth /// The roster check is the half that makes the recorded hive worth
/// having, so assert it by EFFECT rather than by the message: a hive /// having, so assert it by EFFECT rather than by the message: a hive
/// that is not in this swarm must be refused **before anything is /// that is not in this swarm must be refused **before anything is
@ -3825,32 +3688,6 @@ mod tests {
"{problem:?}" "{problem:?}"
); );
} }
/// `GET /api/agents/{name}/icon` against a disconnected queue answers
/// 503 like its siblings, not the 404 an `<img>` would show identically
/// or the 500 of a failed read.
#[tokio::test]
async fn get_agent_icon_answers_503_on_a_disconnected_queue() {
let (state, _sched) = state_with_roster();
let state = super::AppState {
agent_icons: Some(std::sync::Arc::new(
super::agent_icon::AgentIconReader::new(disconnected_client().await),
)),
..state
};
let problem = super::get_agent_icon(
axum::extract::State(state),
axum::extract::Path("iris".to_owned()),
)
.await
.expect_err("a disconnected queue must not read as an icon");
assert_eq!(
problem.status,
Some(axum::http::StatusCode::SERVICE_UNAVAILABLE),
"{problem:?}"
);
}
} }
#[cfg(test)] #[cfg(test)]

View file

@ -573,7 +573,6 @@ mod tests {
status: Some(std::sync::Arc::new(status)), status: Some(std::sync::Arc::new(status)),
wanted: None, wanted: None,
agent_status: None, agent_status: None,
agent_icons: None,
jobq: std::sync::Arc::new(std::sync::Mutex::new(hive_jobq::scheduler::Scheduler::new( jobq: std::sync::Arc::new(std::sync::Mutex::new(hive_jobq::scheduler::Scheduler::new(
hive_jobq::Graph::new(), hive_jobq::Graph::new(),
hive_jobq::resources::ResourceTable::new(), hive_jobq::resources::ResourceTable::new(),

View file

@ -232,7 +232,6 @@ mod tests {
.with_agent_token_subjects(vec![ .with_agent_token_subjects(vec![
"$SWARM.term.{agent}".to_owned(), "$SWARM.term.{agent}".to_owned(),
"$SWARM.agent-state.{agent}".to_owned(), "$SWARM.agent-state.{agent}".to_owned(),
"$KV.agent-icons.{agent}".to_owned(),
]) ])
.expect("valid") .expect("valid")
} }
@ -254,7 +253,6 @@ mod tests {
vec![ vec![
"$SWARM.term.atlas".to_owned(), "$SWARM.term.atlas".to_owned(),
"$SWARM.agent-state.atlas".to_owned(), "$SWARM.agent-state.atlas".to_owned(),
"$KV.agent-icons.atlas".to_owned(),
] ]
); );
} }

View file

@ -359,14 +359,12 @@ mod tests {
"$SWARM.term.{agent}", "$SWARM.term.{agent}",
"--agent-token-publish-subject", "--agent-token-publish-subject",
"$SWARM.agent-state.{agent}", "$SWARM.agent-state.{agent}",
"--agent-token-publish-subject",
"$KV.agent-icons.{agent}",
"--store-cert-role", "--store-cert-role",
"swarm-nats-auth", "swarm-nats-auth",
]) ])
.expect("the unit's own argument vector must parse"); .expect("the unit's own argument vector must parse");
assert_eq!(args.agent_client_suffix, "-agent"); assert_eq!(args.agent_client_suffix, "-agent");
assert_eq!(args.agent_token_publish_subjects.len(), 3); assert_eq!(args.agent_token_publish_subjects.len(), 2);
} }
/// The control for the case above: an ordinary value parses through the /// The control for the case above: an ordinary value parses through the

View file

@ -1232,34 +1232,6 @@ 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] #[test]
fn an_agent_token_subject_without_the_placeholder_is_refused() { fn an_agent_token_subject_without_the_placeholder_is_refused() {
let err = policy() let err = policy()

View file

@ -1,102 +0,0 @@
//! The per-agent icon KV bucket: one key per agent, holding the icon the
//! swarm shows for it.
//!
//! **Keyed by the agent alone — no hive token**, which is the one way this
//! bucket deliberately differs from its otherwise-identical sibling
//! [`crate::agent_status`] (`{hive}.{agent}`). An agent is not fixed to a
//! hive and can move between them; an icon keyed by where the agent
//! currently runs would be stranded under the old key by a migration, and
//! a reader would have to know the placement to ask the question at all.
//! Operator ruling: *"the bucket is per agent, not per hive. agents can
//! move hives."*
//!
//! That is also what makes the icon answerable for an agent whose
//! container is stopped: the value's lifetime is the agent's, not its
//! placement's.
//!
//! **Each agent writes its own key** with its own queue credential. The auth
//! callout grants that credential `$KV.agent-icons.<agent>` 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;
/// The KV bucket agent icons are published into, keyed by agent name.
///
/// A constant and not an option, for the reason
/// [`crate::agent_status::BUCKET`] gives: writer and reader must name the
/// same bucket, and an option is a way for the two to disagree.
pub const BUCKET: &str = "agent-icons";
/// The media type of every value in this bucket.
///
/// The value is the icon's bytes **verbatim, not a JSON envelope**: an SVG
/// carries no metadata this bucket would have to describe, and the one
/// consumer serves the bytes straight back out. So the type is fixed here
/// rather than stored per entry — a publisher that has something other
/// 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.<bucket>.<key>`, 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`]:
/// a consumer wants each agent's current icon, not every icon it has ever
/// had.
#[cfg(feature = "kv")]
pub async fn open_or_create(
client: &async_nats::Client,
) -> Result<async_nats::jetstream::kv::Store, Error> {
let js = async_nats::jetstream::new(client.clone());
match js.get_key_value(BUCKET).await {
Ok(store) => Ok(store),
Err(e) => {
tracing::info!(
bucket = BUCKET,
reason = %e,
"agent-icon bucket not available, creating it"
);
js.create_key_value(async_nats::jetstream::kv::Config {
bucket: BUCKET.to_owned(),
description: "Current icon published by each agent".to_owned(),
history: 1,
..Default::default()
})
.await
.map_err(|source| Error::CreateBucket {
bucket: BUCKET.to_owned(),
source,
})
}
}
}
#[cfg(test)]
mod tests {
use super::subject;
/// 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.<bucket>.<hive>.*`, two tokens) cannot match.
#[test]
fn the_published_subject_carries_the_agent_and_no_placement() {
assert_eq!(subject("iris"), "$KV.agent-icons.iris");
}
}

View file

@ -182,11 +182,6 @@ pub mod agent_status;
/// formats it without linking the secret-store client. /// formats it without linking the secret-store client.
pub mod agent_token; 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 who may write it.
pub mod agent_icon;
/// The subject the swarm controller publishes on when the hive-wide knowledge /// The subject the swarm controller publishes on when the hive-wide knowledge
/// repository has changed. One writer, many readers — every hive subscribes. /// repository has changed. One writer, many readers — every hive subscribes.
/// ///