diff --git a/docs/swarm/README.md b/docs/swarm/README.md index 88ec8360..224e0e27 100644 --- a/docs/swarm/README.md +++ b/docs/swarm/README.md @@ -174,10 +174,13 @@ subagent daemon: each subagent's rows go to `$SWARM.term..sub.` 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 identity, exactly as the harness does, and publishes nothing without it. The -queue keeps these rows: the daemon creates the stream `term-sub-` on -first use, holding rows for 24 hours, and the swarm lists an agent's subagents -from that stream's subjects. The grant covers `CREATE` and `INFO` on that one stream name and no -other `JetStream` subject. Subagents publish output only and read nothing. +queue keeps these rows in the stream `term-sub-` for 24 hours, and the +swarm lists an agent's subagents from that stream's subjects. The swarm +controller creates that stream for every agent a hive's wanted state names, +within a minute of the agent appearing there. The agent's grant is publish on its own +`$SWARM.term..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 take the connection down with it, so the harness drops such a row's body before diff --git a/docs/swarm/ui.md b/docs/swarm/ui.md index 43dc0744..9e396c82 100644 --- a/docs/swarm/ui.md +++ b/docs/swarm/ui.md @@ -37,8 +37,8 @@ terminal has no turn-state badges. The agent's subagent daemon publishes each subagent's terminal rows on `$SWARM.term..sub.` with the agent's own queue credential, -into a stream named `term-sub-` that it creates on first use. The -stream keeps rows for 24 hours. swarm-controller serves the list as +into a stream named `term-sub-` that swarm-controller creates for every +declared agent. The stream keeps rows for 24 hours. swarm-controller serves the list as `GET /api/agents//subagents`, the subagents named by the subjects in that stream, and relays one subagent's rows as SSE on `GET /api/agents//subagents//term/stream`. An agent with no diff --git a/hive-subagent-mcp/Cargo.toml b/hive-subagent-mcp/Cargo.toml index fcc059d2..cb4d2340 100644 --- a/hive-subagent-mcp/Cargo.toml +++ b/hive-subagent-mcp/Cargo.toml @@ -28,8 +28,8 @@ rmcp.workspace = true schemars.workspace = true serde.workspace = true serde_json.workspace = true -# `subagent_term`: the subject, the stream, and opening it as the agent. -swarm-queue-client = { workspace = true, features = ["subagent-term"] } +# `subagent_term`: the subject each subagent's rows are published on. +swarm-queue-client.workspace = true # The agent's own queue credential, read under its store identity. swarm-secret-client.workspace = true tokio.workspace = true diff --git a/hive-subagent-mcp/src/swarm_term.rs b/hive-subagent-mcp/src/swarm_term.rs index 9934e0e5..46f69284 100644 --- a/hive-subagent-mcp/src/swarm_term.rs +++ b/hive-subagent-mcp/src/swarm_term.rs @@ -17,14 +17,14 @@ //! 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. //! -//! **The agent creates its stream.** Before publishing, the task opens -//! `term-sub-`, creating it when it is missing. The stream's subjects -//! are how the swarm lists this agent's subagents; a row published while the -//! stream cannot be opened still reaches a live subscriber. +//! **Publish only.** The rows land in `term-sub-`, the stream +//! `swarm-controller` creates for this agent; this daemon never opens it. A +//! row published while that stream does not exist yet still reaches a live +//! 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 -//! refused credential, a full queue, a failed stream create and a failed -//! publish are each a log line, and the run goes on. +//! refused credential, a full queue and a failed publish are each a log +//! line, and the run goes on. use std::collections::HashMap; use std::sync::Arc; @@ -56,8 +56,8 @@ const QUEUE_DEPTH: usize = 4096; /// connection attempt. const STORE_READ_TIMEOUT: Duration = Duration::from_secs(10); -/// How long after a failed connect, or a failed stream open, the task tries -/// again. Rows that arrive before a first connect succeeds are dropped. +/// How long after a failed connect the task tries again. Rows that arrive +/// before a first connect succeeds are dropped. const RETRY_AFTER: Duration = Duration::from_mins(1); /// Room left under the server's payload limit for subject, headers and @@ -243,19 +243,15 @@ async fn connect(target: &Arc) -> anyhow::Result { .map_err(|e| anyhow!(swarm_queue_client::chain(&e))) } -/// The connection and the stream, each retried at most once per -/// [`RETRY_AFTER`]. +/// The connection, retried at most once per [`RETRY_AFTER`]. #[derive(Default)] struct Queue { client: Option, connect_after: Option, - stream_open: bool, - stream_after: Option, } impl Queue { - /// A client to publish on, connecting and opening the stream first when - /// they are due. + /// A client to publish on, connecting first when a connect is due. async fn ready(&mut self, target: &Arc) -> Option { let now = Instant::now(); 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()?; - 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) + self.client.clone() } } diff --git a/nix/host-modules/swarm-nats.nix b/nix/host-modules/swarm-nats.nix index f5cb2fb9..493848b6 100644 --- a/nix/host-modules/swarm-nats.nix +++ b/nix/host-modules/swarm-nats.nix @@ -844,16 +844,12 @@ in # 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}"}" - # Its subagents' terminal rows, and the one stream that keeps - # them (`swarm_queue_client::subagent_term`). The agent's - # subagent daemon creates the stream on first use, so it gets - # `CREATE` and `INFO` on that stream name alone: no `UPDATE`, - # no `DELETE`, no consumer, no other stream. A `CREATE`'s - # config travels in its payload, which no subject grant can - # narrow. + # Its subagents' terminal rows + # (`swarm_queue_client::subagent_term`), publish only and no + # `$JS.API.*` subject: swarm-controller creates the + # `term-sub-{agent}` stream that keeps them, so its subjects + # and limits are never the agent's to choose. "--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}" ]; # Every credential arrives by `LoadCredential` and is named diff --git a/nix/module-eval/nats-authelia.nix b/nix/module-eval/nats-authelia.nix index 142807ea..92f46889 100644 --- a/nix/module-eval/nats-authelia.nix +++ b/nix/module-eval/nats-authelia.nix @@ -165,7 +165,7 @@ let # 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) # 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 = let args = lib.splitString " " (responderOf natsOldPath).serviceConfig.ExecStart; @@ -183,8 +183,6 @@ let "'$$SWARM.agent-state.{agent}'" "'$$KV.agent-icons.{agent}'" "'$$SWARM.term.{agent}.sub.>'" - "'$$JS.API.STREAM.CREATE.term-sub-{agent}'" - "'$$JS.API.STREAM.INFO.term-sub-{agent}'" ]; } { diff --git a/swarm-controller/Cargo.toml b/swarm-controller/Cargo.toml index d5929e7b..9a876b5a 100644 --- a/swarm-controller/Cargo.toml +++ b/swarm-controller/Cargo.toml @@ -90,7 +90,9 @@ swarm-authelia-bridge-sock.workspace = true # creation config are shared with the hive that writes it, so this end does # 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 # 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 diff --git a/swarm-controller/src/main.rs b/swarm-controller/src/main.rs index 2a89d872..047175ce 100644 --- a/swarm-controller/src/main.rs +++ b/swarm-controller/src/main.rs @@ -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>, 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. /// /// 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()); 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, // 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 diff --git a/swarm-controller/src/subagent_term.rs b/swarm-controller/src/subagent_term.rs index 40ac5e9b..d5a4e58c 100644 --- a/swarm-controller/src/subagent_term.rs +++ b/swarm-controller/src/subagent_term.rs @@ -1,14 +1,13 @@ -//! An agent's subagents: which ones the swarm has rows for, and a live SSE -//! relay of one subagent's terminal. +//! An agent's subagents: the stream that keeps their rows, which ones the +//! 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 //! `$SWARM.term.{agent}.sub.{subagent}` under the agent's own queue credential, -//! into the stream `term-sub-{agent}` it creates itself -//! (`swarm_queue_client::subagent_term`). Subagents are spawned on demand +//! which may publish there and nothing else of these streams. [`spawn`] +//! 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 -//! holds rows for, read with `STREAM.INFO`. This daemon never creates the -//! stream: a missing one is an agent with no subagent output yet, and lists -//! nothing. +//! holds rows for, read with `STREAM.INFO`. A missing stream lists nothing. //! //! **No hive lookup.** Unlike [`crate::term_stream`], there is no //! 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 //! input from the swarm. +use std::collections::BTreeSet; use std::convert::Infallible; +use std::sync::Arc; +use std::time::Duration; use axum::Json; use axum::extract::{Path, State}; @@ -29,6 +31,58 @@ use futures_util::{Stream, StreamExt as _, TryStreamExt as _}; use swarm_queue_client::subagent_term; 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, hives: Vec) { + 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 { + 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( get, diff --git a/swarm-nats-auth/src/main.rs b/swarm-nats-auth/src/main.rs index a1607d6c..7a2ddb99 100644 --- a/swarm-nats-auth/src/main.rs +++ b/swarm-nats-auth/src/main.rs @@ -363,16 +363,12 @@ mod tests { "$KV.agent-icons.{agent}", "--agent-token-publish-subject", "$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", "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(), 6); + assert_eq!(args.agent_token_publish_subjects.len(), 4); } /// The control for the case above: an ordinary value parses through the diff --git a/swarm-nats-auth/src/policy.rs b/swarm-nats-auth/src/policy.rs index 6bbe9716..74f20eb8 100644 --- a/swarm-nats-auth/src/policy.rs +++ b/swarm-nats-auth/src/policy.rs @@ -1219,34 +1219,32 @@ mod tests { ); } - /// The subagent templates as `swarm-nats.nix` spells them expand to the - /// subjects and stream name the subagent daemon uses, and to `CREATE` and - /// `INFO` on that one stream only. + /// The subagent template as `swarm-nats.nix` spells it expands to the + /// subjects the subagent daemon publishes on, publish only. The stream is + /// the controller's to create, which its own grant covers. #[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}; let p = policy_with_agent_subject() - .with_agent_token_subjects(vec![ - "$SWARM.term.{agent}.sub.>".to_owned(), - "$JS.API.STREAM.CREATE.term-sub-{agent}".to_owned(), - "$JS.API.STREAM.INFO.term-sub-{agent}".to_owned(), - ]) - .expect("per-agent templates are valid"); + .with_agent_token_subjects(vec!["$SWARM.term.{agent}.sub.>".to_owned()]) + .expect("a per-agent template is valid"); let g = p.agent_token_permissions("atlas").expect("configured"); - assert_eq!( - g.publish, - vec![ - stream_subjects("atlas"), - format!("$JS.API.STREAM.CREATE.{}", stream_name("atlas")), - format!("$JS.API.STREAM.INFO.{}", stream_name("atlas")), - ] - ); + assert_eq!(g.publish, vec![stream_subjects("atlas")]); assert!(subject("atlas", "scout").starts_with(g.publish[0].trim_end_matches('>'))); 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!( - g.publish.iter().all(|s| !argus.publish.contains(s)), - "atlas and argus share a subject: {g:?} {argus:?}" + !stream.contains('.'), + "{stream} is one token, so `*` covers it" + ); + assert!( + policy() + .reader_subjects() + .contains(&"$JS.API.STREAM.CREATE.*".to_owned()) ); } diff --git a/swarm-queue-client/Cargo.toml b/swarm-queue-client/Cargo.toml index aa62f848..7caec0bb 100644 --- a/swarm-queue-client/Cargo.toml +++ b/swarm-queue-client/Cargo.toml @@ -30,8 +30,8 @@ kv = ["async-nats/kv"] # `cargo check -p swarm-nats-auth` — no `kv` anywhere in that build — # surfaced it as `cannot find jetstream in async_nats`). notices = ["async-nats/jetstream"] -# `subagent_term::open_or_create`, for the one publisher that creates its own -# stream. Explicit for the same reason as `notices`. +# `subagent_term::open_or_create`, for swarm-controller, which creates every +# agent's subagent stream. Explicit for the same reason as `notices`. subagent-term = ["async-nats/jetstream"] [dependencies] diff --git a/swarm-queue-client/src/subagent_term.rs b/swarm-queue-client/src/subagent_term.rs index 69d25915..53914d3b 100644 --- a/swarm-queue-client/src/subagent_term.rs +++ b/swarm-queue-client/src/subagent_term.rs @@ -3,21 +3,20 @@ //! //! An agent's subagent daemon publishes each subagent's classified terminal //! rows on `$SWARM.term..sub.`, under the agent's own queue -//! credential. The queue grants that credential this subject family and -//! `CREATE`/`INFO` on the one stream [`crate::subagent_term::stream_name`] -//! names, nothing else of `JetStream`, so the agent creates its own stream on -//! first use. +//! credential. The queue grants that credential publish on this subject +//! family and no `JetStream` subject. //! //! 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 -//! the stream holds is the list. A reader takes it from the stream's subject -//! counts ([`crate::subagent_term::subagent_names`]) and never creates the -//! stream. +//! the stream holds is the list. `swarm-controller` creates one stream per +//! agent with the fixed config [`crate::subagent_term::open_or_create`] +//! 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 -//! them: the publisher in `hive-subagent-mcp`, the reader in -//! `swarm-controller`, and the queue grant in `swarm-nats.nix`, whose -//! `{agent}` templates spell the same strings. +//! them: the publisher in `hive-subagent-mcp`, `swarm-controller`, and the +//! queue grant in `swarm-nats.nix`, whose `{agent}` template spells the same +//! subjects. /// The subject family every agent terminal shares. const TERM_PREFIX: &str = "$SWARM.term"; @@ -79,10 +78,8 @@ pub fn subagent_names<'a>(agent: &str, subjects: impl IntoIterator