Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
3b8df27dcf | ||
|
|
e974194e3a | ||
|
|
4dbab024da | ||
|
|
9513058a71 |
16 changed files with 752 additions and 43 deletions
|
|
@ -456,6 +456,18 @@ 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.
|
||||
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
|
||||
|
||||
The controller also keeps the forge objects that are one per swarm, not
|
||||
|
|
|
|||
|
|
@ -22,3 +22,37 @@
|
|||
white-space: normal;
|
||||
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;
|
||||
}
|
||||
|
|
|
|||
|
|
@ -12,6 +12,7 @@
|
|||
// `WantedMenu` below keeps its own clicks (mouse or keyboard-synthesized,
|
||||
// `WantedMenu` hosts real `<button>`s) from also firing `Card`'s
|
||||
// `onClick` and opening the detail panel.
|
||||
import { useState } from "preact/hooks";
|
||||
import { Badge } from "@hive/shared/badge.js";
|
||||
import type { ProblemDetails } from "@hive/shared/api-error.js";
|
||||
import { Card } from "../../ui/card/Card.js";
|
||||
|
|
@ -46,48 +47,84 @@ export function AgentCard({
|
|||
const { tone, label } = FRESHNESS[row.freshness];
|
||||
return (
|
||||
<Card selected={selected} onClick={() => onOpenDetail(row)}>
|
||||
<div class="ui-agent-card-line1">
|
||||
<span class="ui-agent-card-name">{row.name}</span>
|
||||
<Badge
|
||||
tone={tone}
|
||||
title={
|
||||
row.freshness === "unknown"
|
||||
? "reported by its hive but not registered in swarm-level identity — needs migration"
|
||||
: undefined
|
||||
}
|
||||
value={
|
||||
<>
|
||||
{label}
|
||||
{row.last_seen_unix !== null ? (
|
||||
<>
|
||||
{" "}
|
||||
(<RelativeTime epochMs={row.last_seen_unix * 1000} />)
|
||||
</>
|
||||
) : null}
|
||||
</>
|
||||
}
|
||||
/>
|
||||
<span class="ui-agent-card-wanted" onClick={(e) => e.stopPropagation()}>
|
||||
<WantedMenu
|
||||
row={row}
|
||||
pending={pending}
|
||||
showDestroy={false}
|
||||
onSelectUp={onSelectUp}
|
||||
onSelectOffline={onSelectOffline}
|
||||
onSelectPaused={onSelectPaused}
|
||||
/>
|
||||
{error ? (
|
||||
<div class="ui-agent-card-row">
|
||||
<AgentIcon name={row.name} />
|
||||
<div class="ui-agent-card-body">
|
||||
<div class="ui-agent-card-line1">
|
||||
<span class="ui-agent-card-name">{row.name}</span>
|
||||
<Badge
|
||||
tone="negative"
|
||||
value="failed"
|
||||
title={error.detail ?? "the declaration failed"}
|
||||
tone={tone}
|
||||
title={
|
||||
row.freshness === "unknown"
|
||||
? "reported by its hive but not registered in swarm-level identity — needs migration"
|
||||
: undefined
|
||||
}
|
||||
value={
|
||||
<>
|
||||
{label}
|
||||
{row.last_seen_unix !== null ? (
|
||||
<>
|
||||
{" "}
|
||||
(<RelativeTime epochMs={row.last_seen_unix * 1000} />)
|
||||
</>
|
||||
) : null}
|
||||
</>
|
||||
}
|
||||
/>
|
||||
) : null}
|
||||
</span>
|
||||
</div>
|
||||
<div class="ui-agent-card-message">
|
||||
{row.snapshot?.status_text ?? "—"}
|
||||
<span
|
||||
class="ui-agent-card-wanted"
|
||||
onClick={(e) => e.stopPropagation()}
|
||||
>
|
||||
<WantedMenu
|
||||
row={row}
|
||||
pending={pending}
|
||||
showDestroy={false}
|
||||
onSelectUp={onSelectUp}
|
||||
onSelectOffline={onSelectOffline}
|
||||
onSelectPaused={onSelectPaused}
|
||||
/>
|
||||
{error ? (
|
||||
<Badge
|
||||
tone="negative"
|
||||
value="failed"
|
||||
title={error.detail ?? "the declaration failed"}
|
||||
/>
|
||||
) : null}
|
||||
</span>
|
||||
</div>
|
||||
<div class="ui-agent-card-message">
|
||||
{row.snapshot?.status_text ?? "—"}
|
||||
</div>
|
||||
</div>
|
||||
</div>
|
||||
</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>
|
||||
);
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
220
hive-agent/src/swarm_agent_icon.rs
Normal file
220
hive-agent/src/swarm_agent_icon.rs
Normal 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());
|
||||
}
|
||||
}
|
||||
|
|
@ -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.<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}"
|
||||
];
|
||||
# Every credential arrives by `LoadCredential` and is named
|
||||
|
|
|
|||
|
|
@ -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.
|
||||
|
|
|
|||
79
swarm-controller/src/agent_icon.rs
Normal file
79
swarm-controller/src/agent_icon.rs
Normal file
|
|
@ -0,0 +1,79 @@
|
|||
//! 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))
|
||||
}
|
||||
}
|
||||
|
|
@ -40,6 +40,7 @@ use swarm_authelia_bridge_sock::BridgeResponse;
|
|||
use utoipa::{OpenApi, ToSchema};
|
||||
use utoipa_axum::{router::OpenApiRouter, routes};
|
||||
|
||||
mod agent_icon;
|
||||
mod agent_identity;
|
||||
mod agent_state_stream;
|
||||
mod agent_status;
|
||||
|
|
@ -607,6 +608,10 @@ struct AppState {
|
|||
/// connection, several consumers" rationale as `wanted` above. `None`
|
||||
/// in exactly the state `status` is.
|
||||
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
|
||||
/// [`hive_jobq::scheduler::Scheduler`] now that something drives it
|
||||
/// (`spawn_jobq_worker`) — the graph alone was enough for the
|
||||
|
|
@ -862,7 +867,8 @@ fn error_problem(status: axum::http::StatusCode, detail: &str) -> problem_detail
|
|||
problem_details::ProblemDetails::from_status_code(status).with_detail(detail)
|
||||
}
|
||||
|
||||
/// The status a [`wanted::WantedWriter`] failure should answer with.
|
||||
/// The status a [`wanted::WantedWriter`] or [`agent_icon::AgentIconReader`]
|
||||
/// failure should answer with.
|
||||
///
|
||||
/// `WantedWriter::view`/`set` call `ensure_connected` internally and
|
||||
/// propagate through `anyhow`, so the concrete
|
||||
|
|
@ -936,6 +942,21 @@ 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.
|
||||
///
|
||||
/// Reports the status and the detail rather than a rendered
|
||||
|
|
@ -1996,6 +2017,92 @@ where
|
|||
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
|
||||
/// 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
|
||||
|
|
@ -2303,6 +2410,7 @@ async fn main() -> Result<()> {
|
|||
links: Arc::new(load_links()),
|
||||
wanted: wanted_writer(status.as_ref()),
|
||||
agent_status: agent_status_reader(status.as_ref()),
|
||||
agent_icons: agent_icon_reader(status.as_ref()),
|
||||
status,
|
||||
jobq,
|
||||
webhook_secret,
|
||||
|
|
@ -2334,6 +2442,7 @@ fn build_app(state: AppState) -> axum::Router {
|
|||
.routes(routes!(get_jobq_graph))
|
||||
.routes(routes!(get_jobq_rollup))
|
||||
.routes(routes!(get_agent_config_pr))
|
||||
.routes(routes!(get_agent_icon))
|
||||
.routes(routes!(get_config_prs))
|
||||
.routes(routes!(create_agent))
|
||||
.routes(routes!(mint_agent_identity))
|
||||
|
|
@ -2455,6 +2564,7 @@ mod tests {
|
|||
// that would write a declaration.
|
||||
wanted: None,
|
||||
agent_status: None,
|
||||
agent_icons: None,
|
||||
jobq: std::sync::Arc::clone(&sched),
|
||||
webhook_secret: None,
|
||||
config_prs: None,
|
||||
|
|
@ -2506,6 +2616,33 @@ 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
|
||||
/// 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
|
||||
|
|
@ -3688,6 +3825,32 @@ mod tests {
|
|||
"{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)]
|
||||
|
|
|
|||
|
|
@ -573,6 +573,7 @@ mod tests {
|
|||
status: Some(std::sync::Arc::new(status)),
|
||||
wanted: None,
|
||||
agent_status: None,
|
||||
agent_icons: None,
|
||||
jobq: std::sync::Arc::new(std::sync::Mutex::new(hive_jobq::scheduler::Scheduler::new(
|
||||
hive_jobq::Graph::new(),
|
||||
hive_jobq::resources::ResourceTable::new(),
|
||||
|
|
|
|||
|
|
@ -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(),
|
||||
]
|
||||
);
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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()
|
||||
|
|
|
|||
102
swarm-queue-client/src/agent_icon.rs
Normal file
102
swarm-queue-client/src/agent_icon.rs
Normal file
|
|
@ -0,0 +1,102 @@
|
|||
//! 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");
|
||||
}
|
||||
}
|
||||
|
|
@ -182,6 +182,11 @@ pub mod agent_status;
|
|||
/// formats it without linking the secret-store client.
|
||||
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
|
||||
/// repository has changed. One writer, many readers — every hive subscribes.
|
||||
///
|
||||
|
|
|
|||
Loading…
Reference in a new issue