Compare commits
16 changed files with 43 additions and 752 deletions
|
|
@ -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
|
||||||
|
|
|
||||||
|
|
@ -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;
|
|
||||||
}
|
|
||||||
|
|
|
||||||
|
|
@ -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,84 +46,48 @@ 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">
|
<div class="ui-agent-card-line1">
|
||||||
<AgentIcon name={row.name} />
|
<span class="ui-agent-card-name">{row.name}</span>
|
||||||
<div class="ui-agent-card-body">
|
<Badge
|
||||||
<div class="ui-agent-card-line1">
|
tone={tone}
|
||||||
<span class="ui-agent-card-name">{row.name}</span>
|
title={
|
||||||
<Badge
|
row.freshness === "unknown"
|
||||||
tone={tone}
|
? "reported by its hive but not registered in swarm-level identity — needs migration"
|
||||||
title={
|
: undefined
|
||||||
row.freshness === "unknown"
|
}
|
||||||
? "reported by its hive but not registered in swarm-level identity — needs migration"
|
value={
|
||||||
: undefined
|
<>
|
||||||
}
|
{label}
|
||||||
value={
|
{row.last_seen_unix !== null ? (
|
||||||
<>
|
<>
|
||||||
{label}
|
{" "}
|
||||||
{row.last_seen_unix !== null ? (
|
(<RelativeTime epochMs={row.last_seen_unix * 1000} />)
|
||||||
<>
|
|
||||||
{" "}
|
|
||||||
(<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 ? (
|
|
||||||
<Badge
|
|
||||||
tone="negative"
|
|
||||||
value="failed"
|
|
||||||
title={error.detail ?? "the declaration failed"}
|
|
||||||
/>
|
|
||||||
) : null}
|
) : null}
|
||||||
</span>
|
</>
|
||||||
</div>
|
}
|
||||||
<div class="ui-agent-card-message">
|
/>
|
||||||
{row.snapshot?.status_text ?? "—"}
|
<span class="ui-agent-card-wanted" onClick={(e) => e.stopPropagation()}>
|
||||||
</div>
|
<WantedMenu
|
||||||
</div>
|
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>
|
||||||
</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>
|
|
||||||
);
|
|
||||||
}
|
|
||||||
|
|
|
||||||
|
|
@ -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
|
||||||
|
|
|
||||||
|
|
@ -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
|
||||||
|
|
|
||||||
|
|
@ -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());
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
@ -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
|
||||||
|
|
|
||||||
|
|
@ -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.
|
||||||
|
|
|
||||||
|
|
@ -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))
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
@ -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)]
|
||||||
|
|
|
||||||
|
|
@ -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(),
|
||||||
|
|
|
||||||
|
|
@ -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(),
|
|
||||||
]
|
]
|
||||||
);
|
);
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -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
|
||||||
|
|
|
||||||
|
|
@ -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()
|
||||||
|
|
|
||||||
|
|
@ -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");
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
@ -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.
|
||||||
///
|
///
|
||||||
|
|
|
||||||
Loading…
Reference in a new issue