swarm-controller: relay an agent's hive-free subjects beside the old ones
The terminal and turn-state relays now subscribe to `$SWARM.term.<agent>` and `$SWARM.agent-state.<agent>`, the subjects a verified agent token is granted, as well as the hive-scoped `<prefix>.<hive>.<agent>` an agent still on its hive's shared credential publishes to. Both at once, so a swarm-ui terminal keeps working whichever credential an agent connected with and in whatever order hosts deploy.
This commit is contained in:
parent
9bad58d86d
commit
0c1fb44a4f
2 changed files with 68 additions and 25 deletions
|
|
@ -33,8 +33,8 @@ use crate::AppState;
|
||||||
|
|
||||||
/// Subject family carrying agent turn-state headers — must match
|
/// Subject family carrying agent turn-state headers — must match
|
||||||
/// `hive-agent::swarm_agent_state::SUBJECT_PREFIX` exactly, since the two
|
/// `hive-agent::swarm_agent_state::SUBJECT_PREFIX` exactly, since the two
|
||||||
/// ends never see the constant together. `$SWARM.agent-state.<hive>.<agent>`
|
/// ends never see the constant together. `crate::term_stream::agent_subjects`
|
||||||
/// is the full subject; as with the terminal relay's own copy, there is no
|
/// 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
|
/// way to check the two sides agree short of this comment and the module
|
||||||
/// docs on both ends staying honest about it.
|
/// docs on both ends staying honest about it.
|
||||||
const SUBJECT_PREFIX: &str = "$SWARM.agent-state";
|
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 = crate::term_stream::subscribe_all(
|
||||||
let subscriber = client.subscribe(subject.clone()).await.map_err(|e| {
|
&client,
|
||||||
tracing::warn!(%subject, error = %e, "state stream: subscribe failed");
|
&crate::term_stream::agent_subjects(SUBJECT_PREFIX, &hive, &agent),
|
||||||
crate::error_problem(
|
"state",
|
||||||
axum::http::StatusCode::INTERNAL_SERVER_ERROR,
|
)
|
||||||
&format!("subscribing to {subject} failed: {e}"),
|
.await?;
|
||||||
)
|
|
||||||
})?;
|
|
||||||
|
|
||||||
tracing::info!(%subject, "state stream: client attached");
|
|
||||||
let stream =
|
let stream =
|
||||||
subscriber.map(|msg| Ok(Event::default().data(String::from_utf8_lossy(&msg.payload))));
|
subscriber.map(|msg| Ok(Event::default().data(String::from_utf8_lossy(&msg.payload))));
|
||||||
Ok(Sse::new(stream).keep_alive(KeepAlive::default()))
|
Ok(Sse::new(stream).keep_alive(KeepAlive::default()))
|
||||||
|
|
|
||||||
|
|
@ -9,7 +9,8 @@
|
||||||
//!
|
//!
|
||||||
//! **Live tail only, on purpose.** `hive-agent` publishes each
|
//! **Live tail only, on purpose.** `hive-agent` publishes each
|
||||||
//! already-classified `TermMsg` row to the core subject
|
//! 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
|
//! 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
|
//! listening missed the row, same as on the agent's own local SSE
|
||||||
//! stream. This handler relays exactly that: no replay, no last-N
|
//! 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
|
/// Subject family carrying agent terminal rows — must match
|
||||||
/// `hive-agent::swarm_term::SUBJECT_PREFIX` exactly, since the two ends
|
/// `hive-agent::swarm_term::SUBJECT_PREFIX` exactly, since the two ends
|
||||||
/// never see the constant together. `$SWARM.term.<hive>.<agent>` is the
|
/// never see the constant together. [`agent_subjects`] spells the full
|
||||||
/// full subject; there is no way to check the two sides agree short of
|
/// 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.
|
/// this comment and the module docs on both ends staying honest about it.
|
||||||
const SUBJECT_PREFIX: &str = "$SWARM.term";
|
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 = subscribe_all(
|
||||||
let subscriber = client.subscribe(subject.clone()).await.map_err(|e| {
|
&client,
|
||||||
tracing::warn!(%subject, error = %e, "term stream: subscribe failed");
|
&agent_subjects(SUBJECT_PREFIX, &hive, &agent),
|
||||||
crate::error_problem(
|
"term",
|
||||||
axum::http::StatusCode::INTERNAL_SERVER_ERROR,
|
)
|
||||||
&format!("subscribing to {subject} failed: {e}"),
|
.await?;
|
||||||
)
|
|
||||||
})?;
|
|
||||||
|
|
||||||
tracing::info!(%subject, "term stream: client attached");
|
|
||||||
let stream =
|
let stream =
|
||||||
subscriber.map(|msg| Ok(Event::default().data(String::from_utf8_lossy(&msg.payload))));
|
subscriber.map(|msg| Ok(Event::default().data(String::from_utf8_lossy(&msg.payload))));
|
||||||
Ok(Sse::new(stream).keep_alive(KeepAlive::default()))
|
Ok(Sse::new(stream).keep_alive(KeepAlive::default()))
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// The subjects one agent publishes on under `prefix`: its own
|
||||||
|
/// `<prefix>.<agent>`, and `<prefix>.<hive>.<agent>`, 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<futures_util::stream::SelectAll<async_nats::Subscriber>, 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(),
|
||||||
|
]
|
||||||
|
);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
|
||||||
Loading…
Reference in a new issue