diff --git a/docs/swarm/README.md b/docs/swarm/README.md index 224e0e27..eb2fbd3a 100644 --- a/docs/swarm/README.md +++ b/docs/swarm/README.md @@ -174,10 +174,12 @@ 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 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 +queue keeps these rows in the stream `term-sub-` for 24 hours or 64 MiB, +whichever comes first, dropping the oldest rows once either limit kicks in so +a publish never fails because of it, 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. diff --git a/swarm-queue-client/src/subagent_term.rs b/swarm-queue-client/src/subagent_term.rs index 53914d3b..965c9480 100644 --- a/swarm-queue-client/src/subagent_term.rs +++ b/swarm-queue-client/src/subagent_term.rs @@ -31,6 +31,13 @@ const SUBAGENT_TOKEN: &str = "sub"; /// its last output. pub const MAX_AGE: std::time::Duration = std::time::Duration::from_hours(24); +/// How large one agent's stream may grow before the discard policy reclaims +/// space. 64 MiB is enough rows from a day of ordinary subagent chatter to +/// stay well inside it while still bounding the disk one compromised or +/// runaway agent can claim for 24h. +#[cfg(feature = "subagent-term")] +const MAX_BYTES: i64 = 64 * 1024 * 1024; + /// The stream holding `agent`'s subagent rows: `term-sub-`. #[must_use] pub fn stream_name(agent: &str) -> String { @@ -78,6 +85,27 @@ pub fn subagent_names<'a>(agent: &str, subjects: impl IntoIterator async_nats::jetstream::stream::Config { + async_nats::jetstream::stream::Config { + name: stream_name(agent), + description: Some(format!("Terminal rows of {agent}'s subagents")), + subjects: vec![stream_subjects(agent)], + max_age: MAX_AGE, + max_bytes: MAX_BYTES, + discard: async_nats::jetstream::stream::DiscardPolicy::Old, + max_message_size, + ..Default::default() + } +} + /// Open `agent`'s subagent stream, creating it if it does not exist yet. An /// existing stream is opened as it is. #[cfg(feature = "subagent-term")] @@ -91,18 +119,13 @@ pub async fn open_or_create( Ok(stream) => Ok(stream), Err(e) => { tracing::info!(stream = %name, reason = %e, "subagent terminal stream not available, creating it"); - js.create_stream(async_nats::jetstream::stream::Config { - name: name.clone(), - description: Some(format!("Terminal rows of {agent}'s subagents")), - subjects: vec![stream_subjects(agent)], - max_age: MAX_AGE, - ..Default::default() - }) - .await - .map_err(|source| crate::Error::CreateAgentStream { - stream: name, - source, - }) + let max_message_size = i32::try_from(crate::max_payload(client)).unwrap_or(i32::MAX); + js.create_stream(stream_config(agent, max_message_size)) + .await + .map_err(|source| crate::Error::CreateAgentStream { + stream: name, + source, + }) } } } @@ -155,4 +178,19 @@ mod tests { ]; assert_eq!(subagent_names("iris", subjects), ["alpha", "zed"]); } + + /// The config the controller creates caps disk per agent and never + /// fails a publish because of that cap — `discard: Old` reclaims the + /// oldest rows instead of rejecting the newest. + #[cfg(feature = "subagent-term")] + #[test] + fn stream_config_caps_bytes_and_discards_the_oldest() { + let cfg = stream_config("iris", 8_388_608); + assert_eq!(cfg.max_bytes, MAX_BYTES); + assert_eq!( + cfg.discard, + async_nats::jetstream::stream::DiscardPolicy::Old + ); + assert_eq!(cfg.max_message_size, 8_388_608); + } }