swarm: show subagent terminals in the swarm UI
An agent's subagent daemon publishes each subagent's output as terminal
rows on `$SWARM.term.<agent>.sub.<subagent>`, as the agent, into a
per-agent stream it creates itself; swarm-controller lists an agent's
subagents from that stream's subjects and relays one subagent's rows as
SSE; the swarm UI lists them under the agent's terminal preview and
reuses AgentTermPreview, full-screen tab included, with no input.
- swarm-nats.nix: the agent token may also publish
`$SWARM.term.{agent}.sub.>` and `$JS.API.STREAM.CREATE|INFO` on
`term-sub-{agent}`, and nothing else of JetStream. A module-eval arm
pins the agent-token grant as an exact list.
- mcp.nix: hive-subagent-daemon loads the agent's store identity
(`hive-agent-bao-cert/-key/-server-ca`, the ones hive-agent loads)
whenever the agent has a store, not only on the opencode preset. The
agent's own queue secret lives in the store, so this is the credential
the harness connects with.
- hive-subagent-mcp: `swarm_term` reads the agent's queue secret under
that identity, connects with the agent token, opens or creates
`term-sub-<agent>` (max_age 24h), and publishes classified rows from
the sink every subagent line already passes through. The sink only
queues (bounded, drop-and-count); a missing store, refused credential,
failed stream create or failed publish is a log line.
- The stream-json classifier (`stream_enrich`) and the `TermMsg` row
types plus `fit` move from the hive-agent binary into hive-sh4re, so
the subagent daemon publishes the rows AgentTermPreview already
renders. hive-agent keeps its LiveEvent classifier on top.
- swarm-controller: `GET /api/agents/{name}/subagents` and
`GET /api/agents/{name}/subagents/{subagent}/term/stream`.
- docs/swarm: what the UI shows and what the queue carries.
Closes #4827
This commit is contained in:
parent
b90be9e65e
commit
d6f94e5247
35 changed files with 1529 additions and 340 deletions
|
|
@ -9,6 +9,8 @@ workspace = true
|
|||
|
||||
[dependencies]
|
||||
anyhow.workspace = true
|
||||
# The swarm-queue client for publishing subagent terminals (`swarm_term`).
|
||||
async-nats.workspace = true
|
||||
axum.workspace = true
|
||||
clap.workspace = true
|
||||
hive-agent-sock.workspace = true
|
||||
|
|
@ -17,6 +19,7 @@ hive-runtime.workspace = true
|
|||
# `permissions::builtin_tools_arg` — the same `--tools` resolution the parent
|
||||
# harness spawns its own claude with, so a subagent's built-in surface is its
|
||||
# parent's rather than a second list that drifts. See `session::build_config`.
|
||||
# Also `term_msg`/`stream_enrich`, the rows `swarm_term` publishes.
|
||||
hive-sh4re.workspace = true
|
||||
hive-sock-client.workspace = true
|
||||
hive-types.workspace = true
|
||||
|
|
@ -25,6 +28,10 @@ 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"] }
|
||||
# The agent's own queue credential, read under its store identity.
|
||||
swarm-secret-client.workspace = true
|
||||
tokio.workspace = true
|
||||
tracing.workspace = true
|
||||
tracing-subscriber.workspace = true
|
||||
|
|
|
|||
|
|
@ -21,3 +21,4 @@ pub mod mcp_config;
|
|||
pub mod paths;
|
||||
pub mod role;
|
||||
pub mod session;
|
||||
pub mod swarm_term;
|
||||
|
|
|
|||
|
|
@ -54,7 +54,9 @@ async fn main() -> Result<()> {
|
|||
// The parent agent's runtime, from the same variables the harness reads.
|
||||
let runtime = hive_runtime::RuntimeSpec::from_env()?;
|
||||
let state = Arc::new(
|
||||
hive_subagent_mcp::session::State::new(todo_socket, signal_base).on_runtime(runtime),
|
||||
hive_subagent_mcp::session::State::new(todo_socket, signal_base)
|
||||
.on_runtime(runtime)
|
||||
.publishing_to(hive_subagent_mcp::swarm_term::Publisher::spawn()),
|
||||
);
|
||||
|
||||
// Serve the MCP tools over streamable-http forever. No background poll
|
||||
|
|
|
|||
|
|
@ -20,7 +20,8 @@
|
|||
//! **A running turn is not necessarily a working turn.** A wedged child
|
||||
//! satisfies "running" as fully as a busy one, so every line of the child's
|
||||
//! streams bumps `last_event_at` (`LivenessSink`) and `status` reports its
|
||||
//! age — the timestamp says the child is alive, never what it said.
|
||||
//! age — the timestamp says the child is alive, never what it said. What it
|
||||
//! said goes to the swarm as the subagent's terminal (`crate::swarm_term`).
|
||||
//!
|
||||
//! **`continue` doesn't pre-check the session's existence — it waits for
|
||||
//! the answer instead.** claude's own `--resume` is the authority, so
|
||||
|
|
@ -98,6 +99,8 @@ use hive_runtime::{
|
|||
use hive_sh4re::permissions::ToolGroup;
|
||||
use tokio::sync::oneshot;
|
||||
|
||||
use crate::swarm_term::Output;
|
||||
|
||||
/// How long a `continue` holds its tool call open waiting to find out
|
||||
/// whether the resume landed. A cap, not a delay: both real outcomes settle
|
||||
/// it well inside this, and it is only ever reached by a child that neither
|
||||
|
|
@ -380,6 +383,9 @@ pub struct State {
|
|||
/// What a subagent's turns run on: claude unless [`State::on_runtime`]
|
||||
/// says otherwise.
|
||||
runtime: SubagentRuntime,
|
||||
/// Where every subagent's output lines go up to the swarm, when this
|
||||
/// agent can publish there ([`State::publishing_to`]).
|
||||
term: Option<crate::swarm_term::Publisher>,
|
||||
}
|
||||
|
||||
/// Which opaque URL segment belongs to which session — the whole of a
|
||||
|
|
@ -432,9 +438,17 @@ impl State {
|
|||
signal_base,
|
||||
signal_tokens: Mutex::new(SignalTokens::default()),
|
||||
runtime: SubagentRuntime::Claude,
|
||||
term: None,
|
||||
}
|
||||
}
|
||||
|
||||
/// Hand every subagent's output lines to `term` as they arrive.
|
||||
#[must_use]
|
||||
pub fn publishing_to(mut self, term: Option<crate::swarm_term::Publisher>) -> Self {
|
||||
self.term = term;
|
||||
self
|
||||
}
|
||||
|
||||
/// Run every subagent on `runtime` — the daemon passes the parent
|
||||
/// agent's, read from its environment at startup. On ACP this resolves
|
||||
/// the harness dir, which only exists inside a container.
|
||||
|
|
@ -1488,23 +1502,23 @@ fn continue_on_runtime(
|
|||
}
|
||||
|
||||
/// Bumps `name`'s liveness clock on every line of the turn's output, and
|
||||
/// does nothing else with it. Replaces the `NoopSink` this daemon used to
|
||||
/// run turns against, which discarded the stream wholesale and left `status`
|
||||
/// unable to tell a working child from a wedged one.
|
||||
/// offers the line to the swarm terminal publisher when this agent has one
|
||||
/// ([`crate::swarm_term`]).
|
||||
///
|
||||
/// **All three callbacks, deliberately.** A stderr line or a stdout line
|
||||
/// that didn't parse as JSON is proof the child is alive every bit as much
|
||||
/// as a stream-json event is, and the failure that matters here is reporting
|
||||
/// a live subagent as wedged — so anything the child says counts, and what
|
||||
/// it said is never read: classifying *what* the subagent is doing is a
|
||||
/// separate question from whether it's doing anything.
|
||||
/// a live subagent as wedged — so anything the child says counts. Liveness
|
||||
/// never reads what it said: classifying *what* the subagent is doing is the
|
||||
/// publisher's job and a separate question from whether it's doing anything.
|
||||
///
|
||||
/// Sink methods are called synchronously from the driver's stream readers as
|
||||
/// lines arrive, so the body has to stay cheap — one uncontended map write
|
||||
/// is, and forwarding to a channel to do the same write elsewhere would cost
|
||||
/// more than it saved. `settle` adds a second uncontended lock on a
|
||||
/// `resume`d turn only, and finds an already-emptied slot after the first
|
||||
/// event — strictly less work than the map write next to it.
|
||||
/// event — strictly less work than the map write next to it. The publisher's
|
||||
/// `offer` copies the line into a bounded channel and never waits.
|
||||
///
|
||||
/// **Liveness counts every callback; "the turn is underway" does not.** A
|
||||
/// resume that matched nothing is not silent: claude writes the reason to
|
||||
|
|
@ -1513,8 +1527,7 @@ fn continue_on_runtime(
|
|||
/// every missed resume as a successful start. The narrowest fact that
|
||||
/// separates the two is the event's own kind — a `result` is stream-json's
|
||||
/// end-of-turn marker, so an event that isn't one is a turn still in
|
||||
/// progress. That's the envelope, not the content: nothing here reads what
|
||||
/// the subagent said.
|
||||
/// progress. That's the envelope, not the content.
|
||||
struct LivenessSink {
|
||||
state: Arc<State>,
|
||||
name: String,
|
||||
|
|
@ -1530,14 +1543,27 @@ impl hive_claude::Sink for LivenessSink {
|
|||
if event.get("type").and_then(serde_json::Value::as_str) != Some("result") {
|
||||
settle(self.verdict.as_ref(), ResumeVerdict::Underway);
|
||||
}
|
||||
self.publish(|| Output::Event(event.clone()));
|
||||
}
|
||||
|
||||
fn on_stdout_line(&self, _line: &str) {
|
||||
fn on_stdout_line(&self, line: &str) {
|
||||
self.state.note_event(&self.name);
|
||||
self.publish(|| Output::Stdout(line.to_owned()));
|
||||
}
|
||||
|
||||
fn on_stderr_line(&self, _line: &str) {
|
||||
fn on_stderr_line(&self, line: &str) {
|
||||
self.state.note_event(&self.name);
|
||||
self.publish(|| Output::Stderr(line.to_owned()));
|
||||
}
|
||||
}
|
||||
|
||||
impl LivenessSink {
|
||||
/// Offer a line to the publisher, building the copy only when there is
|
||||
/// one to take it.
|
||||
fn publish(&self, line: impl FnOnce() -> Output) {
|
||||
if let Some(term) = &self.state.term {
|
||||
term.offer(&self.name, line());
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
|
|
|
|||
459
hive-subagent-mcp/src/swarm_term.rs
Normal file
459
hive-subagent-mcp/src/swarm_term.rs
Normal file
|
|
@ -0,0 +1,459 @@
|
|||
//! Publishing each subagent's terminal rows onto the swarm queue, as this
|
||||
//! agent.
|
||||
//!
|
||||
//! Every line a subagent writes passes through `session`'s sink, which hands
|
||||
//! it to [`Publisher::offer`]. The sink runs synchronously on the subagent's
|
||||
//! output reader, so `offer` only queues. One task drains the queue: it
|
||||
//! classifies each `stream-json` event into the rows the agent's own terminal
|
||||
//! carries (`hive_sh4re::stream_enrich`) and publishes them on
|
||||
//! `$SWARM.term.<agent>.sub.<subagent>` (`swarm_queue_client::subagent_term`).
|
||||
//!
|
||||
//! **Output only.** Nothing here subscribes: a subagent takes no input from
|
||||
//! the swarm.
|
||||
//!
|
||||
//! **As the agent, never as the hive.** The connection presents the agent's
|
||||
//! own queue credential, read from the swarm secret store under the agent's
|
||||
//! store identity: the same credential and identity the harness uses, since
|
||||
//! 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-<agent>`, 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.
|
||||
//!
|
||||
//! **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.
|
||||
|
||||
use std::collections::HashMap;
|
||||
use std::sync::Arc;
|
||||
use std::sync::atomic::{AtomicU64, Ordering};
|
||||
use std::time::{Duration, Instant, SystemTime, UNIX_EPOCH};
|
||||
|
||||
use anyhow::{Context as _, anyhow};
|
||||
use hive_sh4re::term_msg::{ClassifyCtx, Level, TermMsg, fit, iso8601_utc};
|
||||
use swarm_queue_client::subagent_term;
|
||||
use swarm_secret_client::{
|
||||
SecretStore,
|
||||
client::{DEFAULT_CERT_MOUNT, ENV_ADDR, ENV_CACERT, Settings},
|
||||
policy, queue,
|
||||
};
|
||||
use tokio::sync::mpsc;
|
||||
|
||||
/// Names the agent. The store role, the store path, the queue token and the
|
||||
/// subjects are all built from it.
|
||||
const ENV_AGENT_NAME: &str = "HIVE_AGENT_NAME";
|
||||
|
||||
/// Where the swarm queue listens, as this container reaches it.
|
||||
const ENV_NATS_URL: &str = "HIVE_AGENT_NATS_URL";
|
||||
|
||||
/// Lines waiting for the publish task. Past this, `offer` drops the line
|
||||
/// rather than block a subagent's output reader.
|
||||
const QUEUE_DEPTH: usize = 4096;
|
||||
|
||||
/// How long one read of the store may take; it runs inside a queue
|
||||
/// 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.
|
||||
const RETRY_AFTER: Duration = Duration::from_mins(1);
|
||||
|
||||
/// Room left under the server's payload limit for subject, headers and
|
||||
/// framing, as in `hive-agent`'s `swarm_term`.
|
||||
const HEADROOM: usize = 1024;
|
||||
|
||||
/// One line of a subagent's output, as the runtime's sink hands it over.
|
||||
pub enum Output {
|
||||
/// A `stream-json` event.
|
||||
Event(serde_json::Value),
|
||||
/// A stdout line that was not JSON.
|
||||
Stdout(String),
|
||||
/// A stderr line.
|
||||
Stderr(String),
|
||||
}
|
||||
|
||||
struct Line {
|
||||
subagent: String,
|
||||
output: Output,
|
||||
/// Unix seconds when the sink saw the line.
|
||||
ts: i64,
|
||||
}
|
||||
|
||||
/// The sink's handle on the publish task.
|
||||
pub struct Publisher {
|
||||
tx: mpsc::Sender<Line>,
|
||||
dropped: Arc<AtomicU64>,
|
||||
}
|
||||
|
||||
impl Publisher {
|
||||
/// Start the publish task, or return `None` when this daemon has no
|
||||
/// queue address or no store to read the agent's credential from. Logs
|
||||
/// which, once.
|
||||
#[must_use]
|
||||
pub fn spawn() -> Option<Self> {
|
||||
let target = match Target::from_lookup(|k| std::env::var(k).ok()) {
|
||||
Ok(Some(target)) => target,
|
||||
Ok(None) => {
|
||||
tracing::info!(
|
||||
"no swarm queue address or no secret store for this agent; subagent \
|
||||
terminals are not published"
|
||||
);
|
||||
return None;
|
||||
}
|
||||
Err(e) => {
|
||||
tracing::warn!(
|
||||
error = format!("{e:#}"),
|
||||
"this agent's secret store coordinates are incomplete; subagent \
|
||||
terminals are not published"
|
||||
);
|
||||
return None;
|
||||
}
|
||||
};
|
||||
swarm_queue_client::install_crypto_provider();
|
||||
tracing::info!(
|
||||
url = %target.url,
|
||||
stream = %subagent_term::stream_name(&target.agent),
|
||||
"publishing subagent terminals to the swarm queue as this agent"
|
||||
);
|
||||
let (publisher, rx) = Self::channel(QUEUE_DEPTH);
|
||||
tokio::spawn(run(rx, target, Arc::clone(&publisher.dropped)));
|
||||
Some(publisher)
|
||||
}
|
||||
|
||||
fn channel(depth: usize) -> (Self, mpsc::Receiver<Line>) {
|
||||
let (tx, rx) = mpsc::channel(depth);
|
||||
let publisher = Self {
|
||||
tx,
|
||||
dropped: Arc::new(AtomicU64::new(0)),
|
||||
};
|
||||
(publisher, rx)
|
||||
}
|
||||
|
||||
/// Queue `output` from `subagent` for publishing. Never blocks: a full
|
||||
/// queue drops the line and counts it.
|
||||
pub fn offer(&self, subagent: &str, output: Output) {
|
||||
let ts = SystemTime::now()
|
||||
.duration_since(UNIX_EPOCH)
|
||||
.map_or(0, |d| i64::try_from(d.as_secs()).unwrap_or(i64::MAX));
|
||||
let line = Line {
|
||||
subagent: subagent.to_owned(),
|
||||
output,
|
||||
ts,
|
||||
};
|
||||
if self.tx.try_send(line).is_err() {
|
||||
self.dropped.fetch_add(1, Ordering::Relaxed);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// Where to publish, and what to connect with.
|
||||
struct Target {
|
||||
url: String,
|
||||
agent: String,
|
||||
store: StoreSource,
|
||||
}
|
||||
|
||||
/// This agent's queue secret as the store holds it.
|
||||
struct StoreSource {
|
||||
settings: Settings,
|
||||
role: String,
|
||||
path: String,
|
||||
}
|
||||
|
||||
impl Target {
|
||||
/// `None` when the queue address or the store address is absent, which
|
||||
/// is a container without a swarm queue or without a store.
|
||||
///
|
||||
/// # Errors
|
||||
/// A store address without the agent name or the identity beside it.
|
||||
fn from_lookup(get: impl Fn(&str) -> Option<String>) -> anyhow::Result<Option<Self>> {
|
||||
let Some(url) = get(ENV_NATS_URL).filter(|v| !v.is_empty()) else {
|
||||
return Ok(None);
|
||||
};
|
||||
if get(ENV_ADDR).is_none_or(|v| v.is_empty()) {
|
||||
return Ok(None);
|
||||
}
|
||||
let agent = get(ENV_AGENT_NAME)
|
||||
.filter(|v| !v.is_empty())
|
||||
.with_context(|| format!("{ENV_AGENT_NAME} is unset or empty"))?;
|
||||
let settings =
|
||||
Settings::from_lookup(|k| get(k).filter(|v| k != ENV_CACERT || ca_is_usable(v)))
|
||||
.context("reading the swarm secret store's coordinates from the environment")?;
|
||||
let store = StoreSource {
|
||||
settings,
|
||||
role: policy::agent_object_name(&agent)?,
|
||||
path: queue::agent_queue_path(&agent)?,
|
||||
};
|
||||
Ok(Some(Self { url, agent, store }))
|
||||
}
|
||||
|
||||
fn subject(&self, subagent: &str) -> String {
|
||||
subagent_term::subject(&self.agent, subagent)
|
||||
}
|
||||
}
|
||||
|
||||
/// Whether the CA bundle at `path` is a file with bytes in it. Absent means
|
||||
/// the container's own trust store.
|
||||
fn ca_is_usable(path: &str) -> bool {
|
||||
std::fs::metadata(path).is_ok_and(|m| m.len() > 0)
|
||||
}
|
||||
|
||||
impl StoreSource {
|
||||
/// The stored secret, read with a fresh login. Errors name the path and
|
||||
/// never the value.
|
||||
async fn read(&self) -> anyhow::Result<String> {
|
||||
let fetch = async {
|
||||
let store = SecretStore::connect(&self.settings, &self.role, DEFAULT_CERT_MOUNT)
|
||||
.await
|
||||
.context("logging in to the swarm secret store as this agent")?;
|
||||
let credential: Option<queue::AgentCredential> = store
|
||||
.read_optional(&self.path)
|
||||
.await
|
||||
.with_context(|| format!("reading {} from the store", self.path))?;
|
||||
credential
|
||||
.map(|c| c.value.trim().to_owned())
|
||||
.filter(|v| !v.is_empty())
|
||||
.with_context(|| format!("nothing is stored at {}", self.path))
|
||||
};
|
||||
tokio::time::timeout(STORE_READ_TIMEOUT, fetch)
|
||||
.await
|
||||
.map_err(|_| anyhow!("the store did not answer within {STORE_READ_TIMEOUT:?}"))?
|
||||
}
|
||||
}
|
||||
|
||||
/// Connect presenting the agent's own token, read from the store on every
|
||||
/// connection attempt so a re-minted secret is picked up on reconnect.
|
||||
async fn connect(target: &Arc<Target>) -> anyhow::Result<async_nats::Client> {
|
||||
let url = target.url.clone();
|
||||
let target = Arc::clone(target);
|
||||
swarm_queue_client::connect_with_token(&url, None, move || {
|
||||
let target = Arc::clone(&target);
|
||||
// Spawned so the future handed back is only a `JoinHandle`, which is
|
||||
// `Sync` as async-nats's auth callback requires; the store read is not.
|
||||
let attempt = tokio::spawn(async move {
|
||||
let value = target.store.read().await.map_err(|e| format!("{e:#}"))?;
|
||||
swarm_queue_client::agent_token::format_agent_token(&target.agent, &value)
|
||||
.map_err(|e| e.to_string())
|
||||
});
|
||||
async move { attempt.await.map_err(|e| e.to_string())? }
|
||||
})
|
||||
.await
|
||||
.map_err(|e| anyhow!(swarm_queue_client::chain(&e)))
|
||||
}
|
||||
|
||||
/// The connection and the stream, each retried at most once per
|
||||
/// [`RETRY_AFTER`].
|
||||
#[derive(Default)]
|
||||
struct Queue {
|
||||
client: Option<async_nats::Client>,
|
||||
connect_after: Option<Instant>,
|
||||
stream_open: bool,
|
||||
stream_after: Option<Instant>,
|
||||
}
|
||||
|
||||
impl Queue {
|
||||
/// A client to publish on, connecting and opening the stream first when
|
||||
/// they are due.
|
||||
async fn ready(&mut self, target: &Arc<Target>) -> Option<async_nats::Client> {
|
||||
let now = Instant::now();
|
||||
if self.client.is_none() && self.connect_after.is_none_or(|t| now >= t) {
|
||||
match connect(target).await {
|
||||
Ok(client) => {
|
||||
tracing::info!(url = %target.url, "subagent terminal: connected as this agent");
|
||||
self.client = Some(client);
|
||||
}
|
||||
Err(e) => {
|
||||
tracing::warn!(
|
||||
error = format!("{e:#}"),
|
||||
retry_in = ?RETRY_AFTER,
|
||||
"subagent terminal: connecting with this agent's own queue credential \
|
||||
failed; rows are dropped until a connect succeeds"
|
||||
);
|
||||
self.connect_after = Some(now + RETRY_AFTER);
|
||||
}
|
||||
}
|
||||
}
|
||||
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)
|
||||
}
|
||||
}
|
||||
|
||||
async fn run(mut rx: mpsc::Receiver<Line>, target: Target, dropped: Arc<AtomicU64>) {
|
||||
let target = Arc::new(target);
|
||||
let mut ctxs: HashMap<String, ClassifyCtx> = HashMap::new();
|
||||
let mut queue = Queue::default();
|
||||
while let Some(line) = rx.recv().await {
|
||||
let missed = dropped.swap(0, Ordering::Relaxed);
|
||||
if missed > 0 {
|
||||
tracing::warn!(missed, "subagent terminal: queue full, lines dropped");
|
||||
}
|
||||
let ctx = ctxs.entry(line.subagent.clone()).or_default();
|
||||
let rows = classify(&line.output, line.ts, ctx);
|
||||
if rows.is_empty() {
|
||||
continue;
|
||||
}
|
||||
let Some(client) = queue.ready(&target).await else {
|
||||
continue;
|
||||
};
|
||||
let subject = target.subject(&line.subagent);
|
||||
for row in rows {
|
||||
publish(&client, &subject, row).await;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// The rows one line of output renders as, each stamped with `ts`.
|
||||
fn classify(output: &Output, ts: i64, ctx: &mut ClassifyCtx) -> Vec<TermMsg> {
|
||||
let rows = match output {
|
||||
Output::Event(v) => hive_sh4re::stream_enrich::classify_stream_value(v, ctx),
|
||||
Output::Stdout(line) => vec![TermMsg::new(Level::Debug, line.clone())],
|
||||
Output::Stderr(line) => vec![TermMsg::new(Level::Warn, format!("stderr: {line}"))],
|
||||
};
|
||||
let ts = iso8601_utc(ts);
|
||||
rows.into_iter().map(|m| m.at(ts.as_str())).collect()
|
||||
}
|
||||
|
||||
/// Offer one row. Every failure is terminal for that row and for nothing else.
|
||||
async fn publish(client: &async_nats::Client, subject: &str, msg: TermMsg) {
|
||||
// Unconnected, the client buffers rather than fails, and reports the
|
||||
// library's default payload limit rather than the server's.
|
||||
if let Err(e) = swarm_queue_client::ensure_connected(client) {
|
||||
tracing::warn!(error = %swarm_queue_client::chain(&e), "subagent terminal: publish skipped");
|
||||
return;
|
||||
}
|
||||
let limit = swarm_queue_client::max_payload(client).saturating_sub(HEADROOM);
|
||||
let summary = msg.summary.clone();
|
||||
let Some(msg) = fit(msg, limit) else {
|
||||
tracing::warn!(
|
||||
summary,
|
||||
limit,
|
||||
"subagent terminal: row does not fit even without its body, dropped"
|
||||
);
|
||||
return;
|
||||
};
|
||||
let payload = match serde_json::to_vec(&msg) {
|
||||
Ok(payload) => payload,
|
||||
Err(e) => {
|
||||
tracing::warn!(error = %e, "subagent terminal: serialising failed, row dropped");
|
||||
return;
|
||||
}
|
||||
};
|
||||
if let Err(e) = client.publish(subject.to_owned(), payload.into()).await {
|
||||
tracing::warn!(error = %e, "subagent terminal: publish failed, row dropped");
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
|
||||
fn env(pairs: &[(&str, &str)]) -> impl Fn(&str) -> Option<String> {
|
||||
let pairs: Vec<(String, String)> = pairs
|
||||
.iter()
|
||||
.map(|(k, v)| ((*k).to_owned(), (*v).to_owned()))
|
||||
.collect();
|
||||
move |k| pairs.iter().find(|(n, _)| n == k).map(|(_, v)| v.clone())
|
||||
}
|
||||
|
||||
const STORE: [(&str, &str); 3] = [
|
||||
("BAO_ADDR", "https://bao.t.local:8200"),
|
||||
("BAO_CLIENT_CERT", "/run/credentials/c"),
|
||||
("BAO_CLIENT_KEY", "/run/credentials/k"),
|
||||
];
|
||||
|
||||
fn full() -> Vec<(&'static str, &'static str)> {
|
||||
let mut pairs = STORE.to_vec();
|
||||
pairs.push((ENV_NATS_URL, "tls://nats.t.local:4222"));
|
||||
pairs.push((ENV_AGENT_NAME, "iris"));
|
||||
pairs
|
||||
}
|
||||
|
||||
/// Each subagent publishes on its own subject under this agent's family,
|
||||
/// the one the queue grants the agent.
|
||||
#[test]
|
||||
fn a_subagent_publishes_under_the_agents_own_family() {
|
||||
let target = Target::from_lookup(env(&full()))
|
||||
.expect("complete")
|
||||
.expect("configured");
|
||||
assert_eq!(target.subject("scout"), "$SWARM.term.iris.sub.scout");
|
||||
assert_eq!(target.subject("w-2"), "$SWARM.term.iris.sub.w-2");
|
||||
assert_eq!(target.store.path, "swarm/agents/iris/queue");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn without_a_queue_address_or_a_store_nothing_is_published() {
|
||||
let no_url: Vec<_> = full()
|
||||
.into_iter()
|
||||
.filter(|(k, _)| *k != ENV_NATS_URL)
|
||||
.collect();
|
||||
assert!(Target::from_lookup(env(&no_url)).expect("legal").is_none());
|
||||
let no_store: Vec<_> = full()
|
||||
.into_iter()
|
||||
.filter(|(k, _)| *k != "BAO_ADDR")
|
||||
.collect();
|
||||
assert!(
|
||||
Target::from_lookup(env(&no_store))
|
||||
.expect("legal")
|
||||
.is_none()
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn a_store_without_the_agent_name_is_an_error() {
|
||||
let no_name: Vec<_> = full()
|
||||
.into_iter()
|
||||
.filter(|(k, _)| *k != ENV_AGENT_NAME)
|
||||
.collect();
|
||||
assert!(Target::from_lookup(env(&no_name)).is_err());
|
||||
}
|
||||
|
||||
/// The sink must never wait on the publisher: a full queue drops the line
|
||||
/// and counts it.
|
||||
#[test]
|
||||
fn a_full_queue_drops_and_counts_rather_than_blocks() {
|
||||
let (publisher, _rx) = Publisher::channel(1);
|
||||
publisher.offer("scout", Output::Stdout("one".to_owned()));
|
||||
publisher.offer("scout", Output::Stdout("two".to_owned()));
|
||||
assert_eq!(publisher.dropped.load(Ordering::Relaxed), 1);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn every_kind_of_line_renders_as_stamped_rows() {
|
||||
let mut ctx = ClassifyCtx::default();
|
||||
let ts = 1_789_302_903;
|
||||
let event = serde_json::json!({
|
||||
"type": "assistant",
|
||||
"message": { "content": [{ "type": "text", "text": "looking" }] }
|
||||
});
|
||||
let rows = classify(&Output::Event(event), ts, &mut ctx);
|
||||
assert_eq!(rows.len(), 1);
|
||||
assert_eq!(rows[0].ts, "2026-09-13T12:35:03Z");
|
||||
|
||||
let rows = classify(&Output::Stderr("boom".to_owned()), ts, &mut ctx);
|
||||
assert_eq!(rows[0].level, Level::Warn);
|
||||
assert_eq!(rows[0].summary, "stderr: boom");
|
||||
|
||||
let rows = classify(&Output::Stdout("chatter".to_owned()), ts, &mut ctx);
|
||||
assert_eq!(rows[0].level, Level::Debug);
|
||||
assert_eq!(rows[0].ts, "2026-09-13T12:35:03Z");
|
||||
}
|
||||
}
|
||||
Loading…
Reference in a new issue