diff --git a/swarm-controller/src/agent_state_stream.rs b/swarm-controller/src/agent_state_stream.rs index a9adc67c..7eca9f93 100644 --- a/swarm-controller/src/agent_state_stream.rs +++ b/swarm-controller/src/agent_state_stream.rs @@ -33,8 +33,8 @@ use crate::AppState; /// Subject family carrying agent turn-state headers — must match /// `hive-agent::swarm_agent_state::SUBJECT_PREFIX` exactly, since the two -/// ends never see the constant together. `$SWARM.agent-state..` -/// is the full subject; as with the terminal relay's own copy, there is no +/// ends never see the constant together. `crate::term_stream::agent_subjects` +/// spells the full subjects; as with the terminal relay's own copy, there is no /// way to check the two sides agree short of this comment and the module /// docs on both ends staying honest about it. const SUBJECT_PREFIX: &str = "$SWARM.agent-state"; @@ -117,16 +117,12 @@ pub(crate) async fn stream_agent_state( ) })?; - let subject = format!("{SUBJECT_PREFIX}.{hive}.{agent}"); - let subscriber = client.subscribe(subject.clone()).await.map_err(|e| { - tracing::warn!(%subject, error = %e, "state stream: subscribe failed"); - crate::error_problem( - axum::http::StatusCode::INTERNAL_SERVER_ERROR, - &format!("subscribing to {subject} failed: {e}"), - ) - })?; - - tracing::info!(%subject, "state stream: client attached"); + let subscriber = crate::term_stream::subscribe_all( + &client, + &crate::term_stream::agent_subjects(SUBJECT_PREFIX, &hive, &agent), + "state", + ) + .await?; let stream = subscriber.map(|msg| Ok(Event::default().data(String::from_utf8_lossy(&msg.payload)))); Ok(Sse::new(stream).keep_alive(KeepAlive::default())) diff --git a/swarm-controller/src/term_stream.rs b/swarm-controller/src/term_stream.rs index 628f364f..3cd6a878 100644 --- a/swarm-controller/src/term_stream.rs +++ b/swarm-controller/src/term_stream.rs @@ -9,7 +9,8 @@ //! //! **Live tail only, on purpose.** `hive-agent` publishes each //! already-classified `TermMsg` row to the core subject -//! `$SWARM.term.{hive}.{agent}` (see `hive-agent::swarm_term`'s module +//! `$SWARM.term.{agent}` (`$SWARM.term.{hive}.{agent}` from an agent on its +//! hive's shared credential; see `hive-agent::swarm_term`'s module //! doc) — a core subject, not `JetStream`, so a subscriber who was not //! listening missed the row, same as on the agent's own local SSE //! stream. This handler relays exactly that: no replay, no last-N @@ -34,8 +35,8 @@ use crate::AppState; /// Subject family carrying agent terminal rows — must match /// `hive-agent::swarm_term::SUBJECT_PREFIX` exactly, since the two ends -/// never see the constant together. `$SWARM.term..` is the -/// full subject; there is no way to check the two sides agree short of +/// never see the constant together. [`agent_subjects`] spells the full +/// subjects; there is no way to check the two sides agree short of /// this comment and the module docs on both ends staying honest about it. const SUBJECT_PREFIX: &str = "$SWARM.term"; @@ -116,17 +117,63 @@ pub(crate) async fn stream_agent_term( ) })?; - let subject = format!("{SUBJECT_PREFIX}.{hive}.{agent}"); - let subscriber = client.subscribe(subject.clone()).await.map_err(|e| { - tracing::warn!(%subject, error = %e, "term stream: subscribe failed"); - crate::error_problem( - axum::http::StatusCode::INTERNAL_SERVER_ERROR, - &format!("subscribing to {subject} failed: {e}"), - ) - })?; - - tracing::info!(%subject, "term stream: client attached"); + let subscriber = subscribe_all( + &client, + &agent_subjects(SUBJECT_PREFIX, &hive, &agent), + "term", + ) + .await?; let stream = subscriber.map(|msg| Ok(Event::default().data(String::from_utf8_lossy(&msg.payload)))); Ok(Sse::new(stream).keep_alive(KeepAlive::default())) } + +/// The subjects one agent publishes on under `prefix`: its own +/// `.`, and `..`, which an agent still on +/// its hive's shared queue credential publishes to. +pub(crate) fn agent_subjects(prefix: &str, hive: &str, agent: &str) -> [String; 2] { + [ + format!("{prefix}.{agent}"), + format!("{prefix}.{hive}.{agent}"), + ] +} + +/// Subscribe to every one of `subjects`, merged into one stream. `what` names +/// the relay in logs and errors. +pub(crate) async fn subscribe_all( + client: &async_nats::Client, + subjects: &[String], + what: &str, +) -> Result, problem_details::ProblemDetails> +{ + let mut subscribers = Vec::with_capacity(subjects.len()); + for subject in subjects { + let subscriber = client.subscribe(subject.clone()).await.map_err(|e| { + tracing::warn!(%subject, error = %e, "{what} stream: subscribe failed"); + crate::error_problem( + axum::http::StatusCode::INTERNAL_SERVER_ERROR, + &format!("subscribing to {subject} failed: {e}"), + ) + })?; + subscribers.push(subscriber); + } + tracing::info!(?subjects, "{what} stream: client attached"); + Ok(futures_util::stream::select_all(subscribers)) +} + +#[cfg(test)] +mod tests { + use super::*; + + /// Both shapes, so an agent on either credential is streamed. + #[test] + fn a_relay_subscribes_to_the_agents_own_subject_and_its_hive_scoped_one() { + assert_eq!( + agent_subjects(SUBJECT_PREFIX, "h1", "atlas"), + [ + "$SWARM.term.atlas".to_owned(), + "$SWARM.term.h1.atlas".to_owned(), + ] + ); + } +}