Watch
0
0
Fork
You've already forked hyperhive
0

swarm-controller: create every agent's subagent stream

swarm-controller now creates `term-sub-<agent>` for every agent a hive is
declared to run, at start and every minute after, with the config
`swarm_queue_client::subagent_term::open_or_create` spells (subjects
`$SWARM.term.<agent>.sub.>`, max_age 24h). An existing stream is opened as
it is, as the controller does for its other streams and buckets, under the
`$JS.API.STREAM.CREATE.*` grant it already holds.

The agent no longer creates the stream: its token is granted publish on
`$SWARM.term.<agent>.sub.>` and no `$JS.API.STREAM.CREATE|INFO` subject,
and the subagent daemon only publishes. A `CREATE` carries the stream's
config in its payload, which no subject grant narrows, so the agent could
otherwise pick the stream's subjects and limits.
This commit is contained in:
atlas 2026-10-02 23:34:27 +02:00 • committed by mara
commit 318f67cda9
13 changed files with 135 additions and 102 deletions

View file

@ -174,10 +174,13 @@ subagent daemon: each subagent's rows go to `$SWARM.term.<agent>.sub.<subagent>`
classified the same way. The queue grants that family to the agent's own queue classified the same way. The queue grants that family to the agent's own queue
credential alone, so the daemon reads it from the store under the agent's store credential alone, so the daemon reads it from the store under the agent's store
identity, exactly as the harness does, and publishes nothing without it. The identity, exactly as the harness does, and publishes nothing without it. The
queue keeps these rows: the daemon creates the stream `term-sub-<agent>` on queue keeps these rows in the stream `term-sub-<agent>` for 24 hours, and the
first use, holding rows for 24 hours, and the swarm lists an agent's subagents swarm lists an agent's subagents from that stream's subjects. The swarm
from that stream's subjects. The grant covers `CREATE` and `INFO` on that one stream name and no controller creates that stream for every agent a hive's wanted state names,
other `JetStream` subject. Subagents publish output only and read nothing. within a minute of the agent appearing there. The agent's grant is publish on its own
`$SWARM.term.<agent>.sub.>` and no `JetStream` subject, so the stream's
subjects and limits are the controller's and never the agent's. Subagents
publish output only and read nothing.
The queue would refuse a row too large for its `max_payload` outright and The queue would refuse a row too large for its `max_payload` outright and
take the connection down with it, so the harness drops such a row's body before take the connection down with it, so the harness drops such a row's body before

View file

@ -37,8 +37,8 @@ terminal has no turn-state badges.
The agent's subagent daemon publishes each subagent's terminal rows on The agent's subagent daemon publishes each subagent's terminal rows on
`$SWARM.term.<agent>.sub.<subagent>` with the agent's own queue credential, `$SWARM.term.<agent>.sub.<subagent>` with the agent's own queue credential,
into a stream named `term-sub-<agent>` that it creates on first use. The into a stream named `term-sub-<agent>` that swarm-controller creates for every
stream keeps rows for 24 hours. swarm-controller serves the list as declared agent. The stream keeps rows for 24 hours. swarm-controller serves the list as
`GET /api/agents/<name>/subagents`, the subagents named by the subjects in that `GET /api/agents/<name>/subagents`, the subagents named by the subjects in that
stream, and relays one subagent's rows as SSE on stream, and relays one subagent's rows as SSE on
`GET /api/agents/<name>/subagents/<subagent>/term/stream`. An agent with no `GET /api/agents/<name>/subagents/<subagent>/term/stream`. An agent with no

View file

@ -28,8 +28,8 @@ rmcp.workspace = true
schemars.workspace = true schemars.workspace = true
serde.workspace = true serde.workspace = true
serde_json.workspace = true serde_json.workspace = true
# `subagent_term`: the subject, the stream, and opening it as the agent. # `subagent_term`: the subject each subagent's rows are published on.
swarm-queue-client = { workspace = true, features = ["subagent-term"] } swarm-queue-client.workspace = true
# The agent's own queue credential, read under its store identity. # The agent's own queue credential, read under its store identity.
swarm-secret-client.workspace = true swarm-secret-client.workspace = true
tokio.workspace = true tokio.workspace = true

View file

@ -17,14 +17,14 @@
//! both daemons run as the agent's user. The hive's shared client is never //! both daemons run as the agent's user. The hive's shared client is never
//! tried, because the queue grants it nothing under these subjects. //! tried, because the queue grants it nothing under these subjects.
//! //!
//! **The agent creates its stream.** Before publishing, the task opens //! **Publish only.** The rows land in `term-sub-<agent>`, the stream
//! `term-sub-<agent>`, creating it when it is missing. The stream's subjects //! `swarm-controller` creates for this agent; this daemon never opens it. A
//! are how the swarm lists this agent's subagents; a row published while the //! row published while that stream does not exist yet still reaches a live
//! stream cannot be opened still reaches a live subscriber. //! subscriber, but the swarm does not list the subagent from it.
//! //!
//! **Best-effort.** A subagent's run never waits on the queue: no store, a //! **Best-effort.** A subagent's run never waits on the queue: no store, a
//! refused credential, a full queue, a failed stream create and a failed //! refused credential, a full queue and a failed publish are each a log
//! publish are each a log line, and the run goes on. //! line, and the run goes on.
use std::collections::HashMap; use std::collections::HashMap;
use std::sync::Arc; use std::sync::Arc;
@ -56,8 +56,8 @@ const QUEUE_DEPTH: usize = 4096;
/// connection attempt. /// connection attempt.
const STORE_READ_TIMEOUT: Duration = Duration::from_secs(10); const STORE_READ_TIMEOUT: Duration = Duration::from_secs(10);
/// How long after a failed connect, or a failed stream open, the task tries /// How long after a failed connect the task tries again. Rows that arrive
/// again. Rows that arrive before a first connect succeeds are dropped. /// before a first connect succeeds are dropped.
const RETRY_AFTER: Duration = Duration::from_mins(1); const RETRY_AFTER: Duration = Duration::from_mins(1);
/// Room left under the server's payload limit for subject, headers and /// Room left under the server's payload limit for subject, headers and
@ -243,19 +243,15 @@ async fn connect(target: &Arc<Target>) -> anyhow::Result<async_nats::Client> {
.map_err(|e| anyhow!(swarm_queue_client::chain(&e))) .map_err(|e| anyhow!(swarm_queue_client::chain(&e)))
} }
/// The connection and the stream, each retried at most once per /// The connection, retried at most once per [`RETRY_AFTER`].
/// [`RETRY_AFTER`].
#[derive(Default)] #[derive(Default)]
struct Queue { struct Queue {
client: Option<async_nats::Client>, client: Option<async_nats::Client>,
connect_after: Option<Instant>, connect_after: Option<Instant>,
stream_open: bool,
stream_after: Option<Instant>,
} }
impl Queue { impl Queue {
/// A client to publish on, connecting and opening the stream first when /// A client to publish on, connecting first when a connect is due.
/// they are due.
async fn ready(&mut self, target: &Arc<Target>) -> Option<async_nats::Client> { async fn ready(&mut self, target: &Arc<Target>) -> Option<async_nats::Client> {
let now = Instant::now(); let now = Instant::now();
if self.client.is_none() && self.connect_after.is_none_or(|t| now >= t) { if self.client.is_none() && self.connect_after.is_none_or(|t| now >= t) {
@ -275,25 +271,7 @@ impl Queue {
} }
} }
} }
let client = self.client.clone()?; self.client.clone()
if !self.stream_open
&& self.stream_after.is_none_or(|t| now >= t)
&& swarm_queue_client::ensure_connected(&client).is_ok()
{
match subagent_term::open_or_create(&client, &target.agent).await {
Ok(_) => self.stream_open = true,
Err(e) => {
tracing::warn!(
error = %swarm_queue_client::chain(&e),
retry_in = ?RETRY_AFTER,
"subagent terminal: opening this agent's stream failed; rows still \
reach live subscribers but the swarm cannot list the subagent"
);
self.stream_after = Some(now + RETRY_AFTER);
}
}
}
Some(client)
} }
} }

View file

@ -844,16 +844,12 @@ in
# on a hive presents that one, so it cannot be scoped to one # on a hive presents that one, so it cannot be scoped to one
# agent's key. # agent's key.
"--agent-token-publish-subject ${lib.escapeShellArg "\$\$KV.agent-icons.{agent}"}" "--agent-token-publish-subject ${lib.escapeShellArg "\$\$KV.agent-icons.{agent}"}"
# Its subagents' terminal rows, and the one stream that keeps # Its subagents' terminal rows
# them (`swarm_queue_client::subagent_term`). The agent's # (`swarm_queue_client::subagent_term`), publish only and no
# subagent daemon creates the stream on first use, so it gets # `$JS.API.*` subject: swarm-controller creates the
# `CREATE` and `INFO` on that stream name alone: no `UPDATE`, # `term-sub-{agent}` stream that keeps them, so its subjects
# no `DELETE`, no consumer, no other stream. A `CREATE`'s # and limits are never the agent's to choose.
# config travels in its payload, which no subject grant can
# narrow.
"--agent-token-publish-subject ${lib.escapeShellArg "\$\$SWARM.term.{agent}.sub.>"}" "--agent-token-publish-subject ${lib.escapeShellArg "\$\$SWARM.term.{agent}.sub.>"}"
"--agent-token-publish-subject ${lib.escapeShellArg "\$\$JS.API.STREAM.CREATE.term-sub-{agent}"}"
"--agent-token-publish-subject ${lib.escapeShellArg "\$\$JS.API.STREAM.INFO.term-sub-{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

@ -165,7 +165,7 @@ let
# The whole per-agent grant, compared as a list rather than by infix: a # The whole per-agent grant, compared as a list rather than by infix: a
# wider subject added beside these (`$$JS.API.>`, another stream's name) # wider subject added beside these (`$$JS.API.>`, another stream's name)
# would pass every presence check above. # would pass every presence check above.
name = "a verified agent's grant is exactly its own subjects and its subagent stream"; name = "a verified agent's grant is exactly its own subjects, with no JetStream API subject";
ok = ok =
let let
args = lib.splitString " " (responderOf natsOldPath).serviceConfig.ExecStart; args = lib.splitString " " (responderOf natsOldPath).serviceConfig.ExecStart;
@ -183,8 +183,6 @@ let
"'$$SWARM.agent-state.{agent}'" "'$$SWARM.agent-state.{agent}'"
"'$$KV.agent-icons.{agent}'" "'$$KV.agent-icons.{agent}'"
"'$$SWARM.term.{agent}.sub.>'" "'$$SWARM.term.{agent}.sub.>'"
"'$$JS.API.STREAM.CREATE.term-sub-{agent}'"
"'$$JS.API.STREAM.INFO.term-sub-{agent}'"
]; ];
} }
{ {

View file

@ -90,7 +90,9 @@ swarm-authelia-bridge-sock.workspace = true
# creation config are shared with the hive that writes it, so this end does # creation config are shared with the hive that writes it, so this end does
# not get to declare them privately. # not get to declare them privately.
# #
swarm-queue-client = { workspace = true, features = ["kv"] } # `subagent-term`: this daemon creates every agent's subagent stream, with the
# config the crate spells.
swarm-queue-client = { workspace = true, features = ["kv", "subagent-term"] }
# `matrix_account.rs` writes the credential this daemon's route accepts. Same # `matrix_account.rs` writes the credential this daemon's route accepts. Same
# crate the hive reads it back with, which is the point: the path, the field # crate the hive reads it back with, which is the point: the path, the field
# name and the object's shape are agreements between the two ends, and a # name and the object's shape are agreements between the two ends, and a

View file

@ -2392,6 +2392,16 @@ fn spawn_agent_renewal(
}); });
} }
/// Start creating every live agent's subagent stream when a swarm queue is
/// wired up. Lifted out of `main` for `clippy::too_many_lines`.
fn spawn_subagent_streams(status: Option<&Arc<status::StatusReader>>, hives: &[HiveEntry]) {
let (Some(status), Some(wanted)) = (status, wanted_writer(status)) else {
return;
};
let hives = hives.iter().map(|h| h.name.clone()).collect();
subagent_term::spawn(status.queue_client(), wanted, hives);
}
/// Check an agent's forge token now, and mint one if it is missing or stale. /// Check an agent's forge token now, and mint one if it is missing or stale.
/// ///
/// The periodic pass (`forge::agent_token::spawn`) does the same every five /// The periodic pass (`forge::agent_token::spawn`) does the same every five
@ -2920,6 +2930,7 @@ async fn main() -> Result<()> {
let state_forge = keep_forge_for_state(forge_client, webhook_secret.clone()); let state_forge = keep_forge_for_state(forge_client, webhook_secret.clone());
spawn_agent_renewal(&jobq, wanted_writer(status.as_ref()), &hives); spawn_agent_renewal(&jobq, wanted_writer(status.as_ref()), &hives);
spawn_subagent_streams(status.as_ref(), &hives);
// Before serving, because a hive whose role does not exist cannot log in, // Before serving, because a hive whose role does not exist cannot log in,
// and one whose policy does not exist logs in able to read nothing — // and one whose policy does not exist logs in able to read nothing —
// either way it cannot collect what this daemon writes for it. A store // either way it cannot collect what this daemon writes for it. A store

View file

@ -1,14 +1,13 @@
//! An agent's subagents: which ones the swarm has rows for, and a live SSE //! An agent's subagents: the stream that keeps their rows, which ones the
//! relay of one subagent's terminal. //! swarm has rows for, and a live SSE relay of one subagent's terminal.
//! //!
//! The agent's subagent daemon publishes each subagent's classified rows on //! The agent's subagent daemon publishes each subagent's classified rows on
//! `$SWARM.term.{agent}.sub.{subagent}` under the agent's own queue credential, //! `$SWARM.term.{agent}.sub.{subagent}` under the agent's own queue credential,
//! into the stream `term-sub-{agent}` it creates itself //! which may publish there and nothing else of these streams. [`spawn`]
//! (`swarm_queue_client::subagent_term`). Subagents are spawned on demand //! creates the stream `term-sub-{agent}` for every live agent, with the config
//! `swarm_queue_client::subagent_term` spells. Subagents are spawned on demand
//! under caller-chosen names, so the list is the set of subjects that stream //! under caller-chosen names, so the list is the set of subjects that stream
//! holds rows for, read with `STREAM.INFO`. This daemon never creates the //! holds rows for, read with `STREAM.INFO`. A missing stream lists nothing.
//! stream: a missing one is an agent with no subagent output yet, and lists
//! nothing.
//! //!
//! **No hive lookup.** Unlike [`crate::term_stream`], there is no //! **No hive lookup.** Unlike [`crate::term_stream`], there is no
//! hive-scoped subject to follow: only the agent's own credential is granted //! hive-scoped subject to follow: only the agent's own credential is granted
@ -19,7 +18,10 @@
//! while it is attached. Nothing is sent toward a subagent: subagents take no //! while it is attached. Nothing is sent toward a subagent: subagents take no
//! input from the swarm. //! input from the swarm.
use std::collections::BTreeSet;
use std::convert::Infallible; use std::convert::Infallible;
use std::sync::Arc;
use std::time::Duration;
use axum::Json; use axum::Json;
use axum::extract::{Path, State}; use axum::extract::{Path, State};
@ -29,6 +31,58 @@ use futures_util::{Stream, StreamExt as _, TryStreamExt as _};
use swarm_queue_client::subagent_term; use swarm_queue_client::subagent_term;
use crate::AppState; use crate::AppState;
use crate::wanted::WantedWriter;
/// How often [`spawn`] checks every live agent's stream. Also how long a newly
/// declared agent's subagents can go unlisted.
const ENSURE_INTERVAL: Duration = Duration::from_mins(1);
/// Create the subagent stream of every live agent that has none, now and every
/// [`ENSURE_INTERVAL`] after. An existing stream is left as it is.
///
/// The live agents are the ones some hive's wanted-state declaration names in
/// a state other than `Destroyed` (`agent_renewal::live_agents`). A pass skips
/// a hive whose declaration cannot be read and an agent whose stream cannot be
/// created, logging each; a pass while the queue is down does nothing. It
/// never stops the daemon.
pub(crate) fn spawn(client: async_nats::Client, wanted: Arc<WantedWriter>, hives: Vec<String>) {
tokio::spawn(async move {
let mut ticker = tokio::time::interval(ENSURE_INTERVAL);
loop {
ticker.tick().await;
if swarm_queue_client::ensure_connected(&client).is_err() {
tracing::debug!("subagent streams: queue not connected; next pass");
continue;
}
for agent in live_agents(&wanted, &hives).await {
if let Err(e) = subagent_term::open_or_create(&client, &agent).await {
tracing::warn!(
%agent,
error = %swarm_queue_client::chain(&e),
"subagent streams: creating the agent's stream failed; next pass"
);
}
}
}
});
}
/// Every agent the readable declarations of `hives` keep alive.
async fn live_agents(wanted: &WantedWriter, hives: &[String]) -> BTreeSet<String> {
let mut declarations = Vec::with_capacity(hives.len());
for hive in hives {
match wanted.view(hive).await {
Ok(Some(d)) => declarations.push(d),
Ok(None) => {}
Err(e) => tracing::warn!(
%hive,
error = %format!("{e:#}"),
"subagent streams: reading the hive's wanted-state declaration failed"
),
}
}
crate::agent_renewal::live_agents(&declarations)
}
#[utoipa::path( #[utoipa::path(
get, get,

View file

@ -363,16 +363,12 @@ mod tests {
"$KV.agent-icons.{agent}", "$KV.agent-icons.{agent}",
"--agent-token-publish-subject", "--agent-token-publish-subject",
"$SWARM.term.{agent}.sub.>", "$SWARM.term.{agent}.sub.>",
"--agent-token-publish-subject",
"$JS.API.STREAM.CREATE.term-sub-{agent}",
"--agent-token-publish-subject",
"$JS.API.STREAM.INFO.term-sub-{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(), 6); assert_eq!(args.agent_token_publish_subjects.len(), 4);
} }
/// 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

@ -1219,34 +1219,32 @@ mod tests {
); );
} }
/// The subagent templates as `swarm-nats.nix` spells them expand to the /// The subagent template as `swarm-nats.nix` spells it expands to the
/// subjects and stream name the subagent daemon uses, and to `CREATE` and /// subjects the subagent daemon publishes on, publish only. The stream is
/// `INFO` on that one stream only. /// the controller's to create, which its own grant covers.
#[test] #[test]
fn an_agents_subagent_grant_is_its_own_stream_and_nothing_wider() { fn an_agents_subagent_grant_is_publish_on_its_own_family_only() {
use swarm_queue_client::subagent_term::{stream_name, stream_subjects, subject}; use swarm_queue_client::subagent_term::{stream_name, stream_subjects, subject};
let p = policy_with_agent_subject() let p = policy_with_agent_subject()
.with_agent_token_subjects(vec![ .with_agent_token_subjects(vec!["$SWARM.term.{agent}.sub.>".to_owned()])
"$SWARM.term.{agent}.sub.>".to_owned(), .expect("a per-agent template is valid");
"$JS.API.STREAM.CREATE.term-sub-{agent}".to_owned(),
"$JS.API.STREAM.INFO.term-sub-{agent}".to_owned(),
])
.expect("per-agent templates are valid");
let g = p.agent_token_permissions("atlas").expect("configured"); let g = p.agent_token_permissions("atlas").expect("configured");
assert_eq!( assert_eq!(g.publish, vec![stream_subjects("atlas")]);
g.publish,
vec![
stream_subjects("atlas"),
format!("$JS.API.STREAM.CREATE.{}", stream_name("atlas")),
format!("$JS.API.STREAM.INFO.{}", stream_name("atlas")),
]
);
assert!(subject("atlas", "scout").starts_with(g.publish[0].trim_end_matches('>'))); assert!(subject("atlas", "scout").starts_with(g.publish[0].trim_end_matches('>')));
let argus = p.agent_token_permissions("argus").expect("configured"); let argus = p.agent_token_permissions("argus").expect("configured");
assert_eq!(argus.publish, vec![stream_subjects("argus")]);
assert_ne!(g.publish, argus.publish);
let stream = stream_name("atlas");
assert!( assert!(
g.publish.iter().all(|s| !argus.publish.contains(s)), !stream.contains('.'),
"atlas and argus share a subject: {g:?} {argus:?}" "{stream} is one token, so `*` covers it"
);
assert!(
policy()
.reader_subjects()
.contains(&"$JS.API.STREAM.CREATE.*".to_owned())
); );
} }

View file

@ -30,8 +30,8 @@ kv = ["async-nats/kv"]
# `cargo check -p swarm-nats-auth` — no `kv` anywhere in that build — # `cargo check -p swarm-nats-auth` — no `kv` anywhere in that build —
# surfaced it as `cannot find jetstream in async_nats`). # surfaced it as `cannot find jetstream in async_nats`).
notices = ["async-nats/jetstream"] notices = ["async-nats/jetstream"]
# `subagent_term::open_or_create`, for the one publisher that creates its own # `subagent_term::open_or_create`, for swarm-controller, which creates every
# stream. Explicit for the same reason as `notices`. # agent's subagent stream. Explicit for the same reason as `notices`.
subagent-term = ["async-nats/jetstream"] subagent-term = ["async-nats/jetstream"]
[dependencies] [dependencies]

View file

@ -3,21 +3,20 @@
//! //!
//! An agent's subagent daemon publishes each subagent's classified terminal //! An agent's subagent daemon publishes each subagent's classified terminal
//! rows on `$SWARM.term.<agent>.sub.<subagent>`, under the agent's own queue //! rows on `$SWARM.term.<agent>.sub.<subagent>`, under the agent's own queue
//! credential. The queue grants that credential this subject family and //! credential. The queue grants that credential publish on this subject
//! `CREATE`/`INFO` on the one stream [`crate::subagent_term::stream_name`] //! family and no `JetStream` subject.
//! names, nothing else of `JetStream`, so the agent creates its own stream on
//! first use.
//! //!
//! The stream exists so the swarm can list an agent's subagents: subagents //! The stream exists so the swarm can list an agent's subagents: subagents
//! are spawned on demand under caller-chosen names, and the set of subjects //! are spawned on demand under caller-chosen names, and the set of subjects
//! the stream holds is the list. A reader takes it from the stream's subject //! the stream holds is the list. `swarm-controller` creates one stream per
//! counts ([`crate::subagent_term::subagent_names`]) and never creates the //! agent with the fixed config [`crate::subagent_term::open_or_create`]
//! stream. //! spells, and takes the list from its subject counts
//! ([`crate::subagent_term::subagent_names`]).
//! //!
//! The name and the subject live here because three crates must agree on //! The name and the subject live here because three crates must agree on
//! them: the publisher in `hive-subagent-mcp`, the reader in //! them: the publisher in `hive-subagent-mcp`, `swarm-controller`, and the
//! `swarm-controller`, and the queue grant in `swarm-nats.nix`, whose //! queue grant in `swarm-nats.nix`, whose `{agent}` template spells the same
//! `{agent}` templates spell the same strings. //! subjects.
/// The subject family every agent terminal shares. /// The subject family every agent terminal shares.
const TERM_PREFIX: &str = "$SWARM.term"; const TERM_PREFIX: &str = "$SWARM.term";
@ -79,10 +78,8 @@ pub fn subagent_names<'a>(agent: &str, subjects: impl IntoIterator<Item = &'a st
names names
} }
/// Open `agent`'s subagent stream, creating it if it does not exist yet. /// Open `agent`'s subagent stream, creating it if it does not exist yet. An
/// /// existing stream is opened as it is.
/// Called by the publisher only. The reader opens the stream with `get_stream`
/// and reads a missing one as "no subagents".
#[cfg(feature = "subagent-term")] #[cfg(feature = "subagent-term")]
pub async fn open_or_create( pub async fn open_or_create(
client: &async_nats::Client, client: &async_nats::Client,