Watch
0
0
Fork
You've already forked hyperhive
0

swarm-controller: cap a subagent terminal stream's disk use

The controller-created term-sub-<agent> stream had max_age only, so a
publishing agent could grow it without bound for 24h. Add a 64 MiB
max_bytes cap with discard: Old (oldest rows drop first, publish never
fails on the cap), and size max_message_size off the queue's live
max_payload rather than a hardcoded guess.
This commit is contained in:
atlas 2026-10-03 01:13:26 +02:00 • committed by mara
commit e21546d08d
2 changed files with 56 additions and 16 deletions

View file

@ -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-<agent>`.
#[must_use]
pub fn stream_name(agent: &str) -> String {
@ -78,6 +85,27 @@ pub fn subagent_names<'a>(agent: &str, subjects: impl IntoIterator<Item = &'a st
names
}
/// The config [`open_or_create`] creates `agent`'s stream with.
///
/// `max_message_size` is taken as a parameter rather than read from a
/// constant here: it tracks the queue's live `max_payload`
/// ([`crate::max_payload`]), the same value the publisher already sizes its
/// rows against, so the stream never rejects a row the queue itself
/// accepted.
#[cfg(feature = "subagent-term")]
fn stream_config(agent: &str, max_message_size: i32) -> 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);
}
}