diff --git a/Cargo.lock b/Cargo.lock index 4b0ded9a..4e71f1ae 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -2036,6 +2036,7 @@ name = "hive-subagent-mcp" version = "0.1.0" dependencies = [ "anyhow", + "async-nats", "axum", "clap", "hive-agent-sock", @@ -2049,6 +2050,8 @@ dependencies = [ "schemars", "serde", "serde_json", + "swarm-queue-client", + "swarm-secret-client", "tokio", "tracing", "tracing-subscriber", diff --git a/docs/swarm/README.md b/docs/swarm/README.md index e62f7d58..88ec8360 100644 --- a/docs/swarm/README.md +++ b/docs/swarm/README.md @@ -169,6 +169,16 @@ agent talks about itself here and reads nothing. Rows aren't retained — a subscriber that wasn't listening missed them, the same as on the agent's own live stream. +An agent's **subagents** publish their terminals too, from the agent's +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: the daemon creates the stream `term-sub-` on +first use, holding rows for 24 hours, and the swarm lists an agent's subagents +from that stream's subjects. The grant covers `CREATE` and `INFO` on that one stream name and no +other `JetStream` subject. Subagents publish output only and read nothing. + The queue would refuse a row too large for its `max_payload` outright and take the connection down with it, so the harness drops such a row's body before sending and leaves a marker in its place; the summary, level and icon still diff --git a/docs/swarm/ui.md b/docs/swarm/ui.md index 8267603f..43dc0744 100644 --- a/docs/swarm/ui.md +++ b/docs/swarm/ui.md @@ -7,18 +7,45 @@ domain, covers host-level detail for one hive. ## What it shows -| route | what | -| ------------------------- | --------------------------------------------------------------------------------------------------- | -| `/` | the hive directory, each hive with its last reported status | -| `/agents` | every agent: status, config PR, wanted state; create agents, link forge, matrix and GitHub accounts | -| `/agents//terminal` | one agent's live terminal | -| `/jobs` | the controller's job graph — where agent creation and credential mints show progress | -| `/issues` | a cross-repo issue report | +| route | what | +| ---------------------------------------------- | --------------------------------------------------------------------------------------------------- | +| `/` | the hive directory, each hive with its last reported status | +| `/agents` | every agent: status, config PR, wanted state; create agents, link forge, matrix and GitHub accounts | +| `/agents//terminal` | one agent's live terminal | +| `/agents//subagents//terminal` | one of that agent's subagents' live terminal | +| `/jobs` | the controller's job graph — where agent creation and credential mints show progress | +| `/issues` | a cross-repo issue report | Everything it shows comes from [`swarm-controller`](../../swarm-controller/README.md). An agent created here or with `swarmctl agent create` starts `paused`; set it `up` from its card. +### Agent and subagent terminals + +An agent's detail panel on `/agents` shows a small live preview of its +terminal, and below it the agent's **subagents**: every subagent the swarm +queue holds terminal rows for from the last 24 hours. Picking one shows its +terminal in the same preview, and **expand** opens it in its own tab, as it +does for the agent. The list refreshes every 10 seconds, so a subagent the +agent spawns shows up once it writes output. + +Both terminals are live views with no input box: a preview shows rows +published while it's open. Subagents take no input from the swarm. A subagent +terminal has no turn-state badges. + +
Where the subagent list and rows come from + +The agent's subagent daemon publishes each subagent's terminal rows on +`$SWARM.term..sub.` with the agent's own queue credential, +into a stream named `term-sub-` that it creates on first use. The +stream keeps rows for 24 hours. swarm-controller serves the list as +`GET /api/agents//subagents`, the subagents named by the subjects in that +stream, and relays one subagent's rows as SSE on +`GET /api/agents//subagents//term/stream`. An agent with no +store identity, or no queue address, publishes no subagent terminals. + +
+ ### Linking external accounts Each agent on `/agents` opens three dialogs that write a credential for it diff --git a/frontend/packages/swarm-ui/src/App.tsx b/frontend/packages/swarm-ui/src/App.tsx index 8c1a482f..9df8630d 100644 --- a/frontend/packages/swarm-ui/src/App.tsx +++ b/frontend/packages/swarm-ui/src/App.tsx @@ -25,6 +25,10 @@ export function App() { + diff --git a/frontend/packages/swarm-ui/src/pages/agents/AgentTermPreview.tsx b/frontend/packages/swarm-ui/src/pages/agents/AgentTermPreview.tsx index fdd0ff48..bcf0e43f 100644 --- a/frontend/packages/swarm-ui/src/pages/agents/AgentTermPreview.tsx +++ b/frontend/packages/swarm-ui/src/pages/agents/AgentTermPreview.tsx @@ -2,9 +2,10 @@ // page, a small read-only preview embedded below `AgentsPage`'s detail // panel fields. No input — sending input back to the agent is explicitly // follow-up scope, not this issue's. Consumes swarm-controller's -// `GET /api/agents/{name}/term/stream` via `useSwarmTermStream`, -// rendering through the same `@hive/shared` `Row` component the per-hive -// local terminal uses — "share components between hive and swarm level +// `GET /api/agents/{name}/term/stream`, or for a subagent +// `GET /api/agents/{name}/subagents/{subagent}/term/stream`, via +// `useSwarmTermStream`, rendering through the same `@hive/shared` `Row` +// component the per-hive local terminal uses — "share components between hive and swarm level // term as much as possible" (mara). // // **Floating header badges** — turn state / model / ctx / cost parity @@ -100,21 +101,50 @@ function fmtTokens(n: number): string { return String(n); } +// The paths one terminal is reached at. A subagent's terminal sits under its +// agent's, and has no turn-state header: the swarm relays none for it. +interface TermPaths { + label: string; + page: string; + stream: string; + state: string | null; +} + +export function termPaths(agentName: string, subagent?: string): TermPaths { + const agent = `/agents/${encodeURIComponent(agentName)}`; + if (subagent === undefined) { + return { + label: agentName, + page: `${agent}/terminal`, + stream: `/api${agent}/term/stream`, + state: `/api${agent}/state/stream`, + }; + } + const sub = `${agent}/subagents/${encodeURIComponent(subagent)}`; + return { + label: `${agentName}/${subagent}`, + page: `${sub}/terminal`, + stream: `/api${sub}/term/stream`, + state: null, + }; +} + export function AgentTermPreview({ agentName, + subagent, showHeaderBadges = true, fullHeight = false, }: { agentName: string; + // Set: this is one of `agentName`'s subagents, output only like the + // agent's own. + subagent?: string; showHeaderBadges?: boolean; fullHeight?: boolean; }) { - const { rows, connection } = useSwarmTermStream( - `/api/agents/${encodeURIComponent(agentName)}/term/stream`, - ); - const { header } = useSwarmAgentStateStream( - `/api/agents/${encodeURIComponent(agentName)}/state/stream`, - ); + const paths = termPaths(agentName, subagent); + const { rows, connection } = useSwarmTermStream(paths.stream); + const { header } = useSwarmAgentStateStream(paths.state); const { open: openTab } = useDynamicTabs(); const logRef = useRef(null); const [stickToBottom, setStickToBottom] = useState(true); @@ -156,13 +186,8 @@ export function AgentTermPreview({ variant="quiet" icon={} value="expand" - title={`open ${agentName}'s terminal in its own tab`} - onClick={() => - openTab({ - href: `/agents/${encodeURIComponent(agentName)}/terminal`, - label: agentName, - }) - } + title={`open ${paths.label}'s terminal in its own tab`} + onClick={() => openTab({ href: paths.page, label: paths.label })} /> )} diff --git a/frontend/packages/swarm-ui/src/pages/agents/AgentTerminalPage.tsx b/frontend/packages/swarm-ui/src/pages/agents/AgentTerminalPage.tsx index 4b1f8b9b..149b6a4f 100644 --- a/frontend/packages/swarm-ui/src/pages/agents/AgentTerminalPage.tsx +++ b/frontend/packages/swarm-ui/src/pages/agents/AgentTerminalPage.tsx @@ -13,30 +13,42 @@ // idempotent (dedupes by href), so this is safe to call every mount. import { useEffect } from "preact/hooks"; import type { RouteComponentProps } from "wouter-preact"; -import { AgentTermPreview } from "./AgentTermPreview.js"; +import { AgentTermPreview, termPaths } from "./AgentTermPreview.js"; import { useDynamicTabs } from "../../shell/useDynamicTabs.js"; import "./AgentTerminalPage.css"; +// Also mounted at `/agents/:name/subagents/:subagent/terminal`, for one of +// the agent's subagents. export function AgentTerminalPage({ params, -}: RouteComponentProps<{ name: string }>) { +}: RouteComponentProps<{ name: string; subagent?: string }>) { const { open } = useDynamicTabs(); const name = decodeURIComponent(params.name); + const subagent = + params.subagent === undefined + ? undefined + : decodeURIComponent(params.subagent); + const { page, label } = termPaths(name, subagent); - // Runs once per mounted agent name — re-registering on every render + // Runs once per mounted terminal — re-registering on every render // would just be redundant work against an already-deduped list, but // there's no reason to pay it more than once per identity. useEffect(() => { - open({ href: `/agents/${encodeURIComponent(name)}/terminal`, label: name }); + open({ href: page, label }); // eslint-disable-next-line react-hooks/exhaustive-deps -- `open` is a // fresh function identity every render (`useDynamicTabs` doesn't // memoize it); including it would re-fire this effect every render - // instead of once per agent `name`. - }, [name]); + // instead of once per terminal `page`. + }, [page]); return (
- +
); } diff --git a/frontend/packages/swarm-ui/src/pages/agents/AgentsPage.tsx b/frontend/packages/swarm-ui/src/pages/agents/AgentsPage.tsx index 85e624ad..911cada4 100644 --- a/frontend/packages/swarm-ui/src/pages/agents/AgentsPage.tsx +++ b/frontend/packages/swarm-ui/src/pages/agents/AgentsPage.tsx @@ -32,6 +32,7 @@ import { Badge } from "@hive/shared/badge.js"; import { LinkIcon } from "@hive/shared/icons.js"; import { AgentCard } from "./AgentCard.js"; import { AgentTermPreview } from "./AgentTermPreview.js"; +import { SubagentTerms } from "./SubagentTerms.js"; import { FRESHNESS, type AgentRow } from "./AgentTypes.js"; import { LinkedAccounts } from "./LinkedAccounts.js"; import { Button } from "../../ui/button/Button.js"; @@ -654,6 +655,10 @@ export function AgentsPage() { preview, no header/no input — the full terminal (+ sending input back) is explicit follow-up scope. */} + ) : (

select an agent to see its details

diff --git a/frontend/packages/swarm-ui/src/pages/agents/SubagentTerms.css b/frontend/packages/swarm-ui/src/pages/agents/SubagentTerms.css new file mode 100644 index 00000000..3d2247a8 --- /dev/null +++ b/frontend/packages/swarm-ui/src/pages/agents/SubagentTerms.css @@ -0,0 +1,17 @@ +/* — one row of subagent names under the agent's own + terminal preview; the picked subagent's preview sits below it. */ +.ui-subagent-terms { + margin-block-start: 1rem; +} + +.ui-subagent-terms-list { + display: flex; + flex-wrap: wrap; + align-items: center; + gap: 0.5em; +} + +.ui-subagent-terms-heading, +.ui-subagent-terms-empty { + color: var(--muted); +} diff --git a/frontend/packages/swarm-ui/src/pages/agents/SubagentTerms.tsx b/frontend/packages/swarm-ui/src/pages/agents/SubagentTerms.tsx new file mode 100644 index 00000000..7c1bfd67 --- /dev/null +++ b/frontend/packages/swarm-ui/src/pages/agents/SubagentTerms.tsx @@ -0,0 +1,76 @@ +// — an agent's subagents, under its own terminal preview in +// `AgentsPage`'s detail panel. The list is `GET /api/agents/{name}/subagents`: +// the subagents the swarm queue holds terminal rows for. Subagents are +// spawned on demand, so it is refetched every `REFRESH_MS`. Picking one +// shows its terminal in `AgentTermPreview`, expand-to-tab included. Read-only +// like the agent's own: subagents take no input. +// +// Mount it keyed by agent name: the refresh timer is armed once per mount. +import { useState } from "preact/hooks"; +import { ApiErrorPanel } from "@hive/shared/api-error-panel.js"; +import { readApiError, type ProblemDetails } from "@hive/shared/api-error.js"; +import { Badge } from "@hive/shared/badge.js"; +import { AgentTermPreview } from "./AgentTermPreview.js"; +import { useRefreshInterval } from "../../ui/refresh-interval/RefreshInterval.js"; +import "./SubagentTerms.css"; + +const REFRESH_MS = 10_000; + +export function SubagentTerms({ agentName }: { agentName: string }) { + const [names, setNames] = useState(null); + const [error, setError] = useState(null); + const [selected, setSelected] = useState(null); + + async function refresh() { + const res = await fetch( + `/api/agents/${encodeURIComponent(agentName)}/subagents`, + ); + if (!res.ok) { + setError(await readApiError(res)); + return; + } + setNames((await res.json()) as string[]); + setError(null); + } + + useRefreshInterval(REFRESH_MS, () => { + refresh().catch((e: unknown) => setError({ detail: String(e) })); + }); + + return ( +
+
+ subagents + {names === null && !error && } + {names?.length === 0 && ( + + none with output on record + + )} + {names?.map((name) => ( + setSelected(name === selected ? null : name)} + /> + ))} +
+ {error && ( + + )} + {selected && ( + + )} +
+ ); +} diff --git a/frontend/packages/swarm-ui/src/pages/agents/useSwarmAgentStateStream.ts b/frontend/packages/swarm-ui/src/pages/agents/useSwarmAgentStateStream.ts index 3050e7cb..8fbd31f0 100644 --- a/frontend/packages/swarm-ui/src/pages/agents/useSwarmAgentStateStream.ts +++ b/frontend/packages/swarm-ui/src/pages/agents/useSwarmAgentStateStream.ts @@ -47,7 +47,7 @@ export interface UseSwarmAgentStateStreamResult { } export function useSwarmAgentStateStream( - streamUrl: string, + streamUrl: string | null, ): UseSwarmAgentStateStreamResult { const [header, setHeader] = useState(null); const [connection, setConnection] = useState("connecting"); @@ -55,6 +55,8 @@ export function useSwarmAgentStateStream( useEffect(() => { setHeader(null); setConnection("connecting"); + // `null`: a terminal with no header to follow (a subagent's). + if (streamUrl === null) return; const es = new EventSource(streamUrl); es.onopen = () => setConnection("open"); es.onmessage = (e) => { diff --git a/frontend/packages/swarm-ui/src/shell/Shell.tsx b/frontend/packages/swarm-ui/src/shell/Shell.tsx index 9dcbbc9b..6459c0b8 100644 --- a/frontend/packages/swarm-ui/src/shell/Shell.tsx +++ b/frontend/packages/swarm-ui/src/shell/Shell.tsx @@ -62,13 +62,15 @@ const NAV_ITEMS: { href: string; label: string; accent: string }[] = [ // its own table view (plus the list+detail split) hits the same cap. const WIDE_BODY_ROUTES = new Set(["/issues", "/agents"]); -// `/agents/:name/terminal` needs the same widening but can't join +// `/agents/:name/terminal` (and a subagent's, under +// `/agents/:name/subagents/:subagent/`) needs the same widening but can't join // the `Set` above — it's a parameterized path, one literal string can't // match every agent name. A 60em-capped terminal pane reads as cramped // (this is a full tab, not a sidebar preview), same reasoning as the // table routes above just for a different reason (a wide pane, not a // wide table). -const WIDE_BODY_ROUTE_PATTERN = /^\/agents\/[^/]+\/terminal$/; +const WIDE_BODY_ROUTE_PATTERN = + /^\/agents\/[^/]+(\/subagents\/[^/]+)?\/terminal$/; function isWideBodyRoute(location: string): boolean { return ( diff --git a/hive-agent/src/main.rs b/hive-agent/src/main.rs index a60120e0..06d5d0bf 100644 --- a/hive-agent/src/main.rs +++ b/hive-agent/src/main.rs @@ -28,7 +28,6 @@ mod reminders; mod serve_common; mod state_entry_watch; mod stats; -mod stream_enrich; mod swarm_agent_icon; mod swarm_agent_state; mod swarm_queue; diff --git a/hive-agent/src/swarm_term.rs b/hive-agent/src/swarm_term.rs index 620b81d4..4a8c686f 100644 --- a/hive-agent/src/swarm_term.rs +++ b/hive-agent/src/swarm_term.rs @@ -29,6 +29,7 @@ use tokio::sync::broadcast; use crate::events::BusEvent; use crate::swarm_queue::{Connection, Presented}; use crate::term_msg::{ClassifyCtx, TermMsg, classify}; +use hive_sh4re::term_msg::fit; /// Subject family carrying agent terminal rows, the swarm-wide agreement this /// publisher holds up its end of. One leaf subject per agent, so a subscriber @@ -47,9 +48,6 @@ const SUBJECT_PREFIX: &str = "$SWARM.term"; const CLIENT_ID_PREFIX: &str = "hive-"; const CLIENT_ID_SUFFIX: &str = "-agent"; -/// What a row that lost its body to the payload limit carries instead. -const DROPPED_BODY: &str = "[body dropped: over the queue's payload limit]"; - /// How much room to leave under the announced limit for everything the /// publish adds around the payload — subject, headers, protocol framing. /// @@ -78,41 +76,6 @@ fn hive_from_client_id(client_id: &str) -> Option<&str> { (!hive.is_empty()).then_some(hive) } -/// Bring `msg` under `limit` serialized bytes, or report that it cannot be. -/// -/// The body is the only field that carries arbitrary length — a diff, a whole -/// tool result — so it is the only one worth spending: replacing it keeps the -/// row's identity, level, summary and time, which is what makes the row -/// readable at all, and a reader sees that something was there rather than -/// seeing nothing. -/// -/// `None` means even the degraded row does not fit, so the caller reports it -/// instead of publishing. That matters more than it looks: an oversize publish -/// is not truncated by the server, it is refused and the connection is closed, -/// which costs the row *and* every row racing behind it through the reconnect. -fn fit(msg: TermMsg, limit: usize) -> Option { - if serialized_len(&msg)? <= limit { - return Some(msg); - } - // Rebuilt rather than mutated in place: `body_format` describes the body, - // and a marker left tagged `Diff` renders as a broken diff downstream. - let degraded = TermMsg { - body: Some(DROPPED_BODY.to_owned()), - body_format: None, - ..msg - }; - (serialized_len(°raded)? <= limit).then_some(degraded) -} - -/// Serialized size of a row, or `None` if it does not serialize at all. -/// -/// Measured by serializing rather than estimated from field lengths: the -/// payload is what the server measures, and JSON escaping makes the two differ -/// by an unbounded factor on exactly the rows that are already near the limit. -fn serialized_len(msg: &TermMsg) -> Option { - serde_json::to_vec(msg).ok().map(|v| v.len()) -} - /// The subject this agent's rows go to under the credential it connected /// with, or `None` when a hive client id names no hive. fn subject(presented: &Presented, agent: &str) -> Option { @@ -217,12 +180,12 @@ async fn publish(client: &async_nats::Client, subject: &str, msg: TermMsg) { #[cfg(test)] mod tests { - use super::{DROPPED_BODY, Presented, fit, hive_from_client_id, subject}; + use super::{Presented, hive_from_client_id, subject}; use crate::events::LiveEvent; - use crate::term_msg::{BodyFormat, ClassifyCtx, Level, TermMsg, classify}; + use crate::term_msg::{ClassifyCtx, classify}; + use hive_sh4re::term_msg::fit; - /// A row's time as `classify` renders it, for the tests below that need - /// one without classifying an event to get it. + /// The time `classify` renders for the event the test below classifies. const ROW_TS: &str = "2026-09-13T12:35:03Z"; #[test] @@ -278,82 +241,6 @@ mod tests { assert_eq!(subject(&unparseable, "mara"), None); } - #[test] - fn a_row_that_already_fits_is_published_unchanged() { - let msg = TermMsg::new(Level::Info, "turn ok") - .icon("✅") - .body("a short body", Some(BodyFormat::Markdown)); - let fitted = fit(msg, 4096).expect("a small row fits"); - assert_eq!(fitted.body.as_deref(), Some("a short body")); - assert_eq!(fitted.body_format, Some(BodyFormat::Markdown)); - } - - /// The case the degrade exists for: a body larger than the limit costs the - /// body and nothing else, and what comes back is actually under the limit - /// rather than merely smaller. - #[test] - fn an_oversize_body_is_replaced_and_the_result_fits() { - let limit = 512; - let msg = TermMsg::new(Level::Info, "Edit(src/main.rs)") - .icon("🔧") - .body("x".repeat(limit * 4), Some(BodyFormat::Diff)) - .coalesce("tool-1") - .at(ROW_TS); - let fitted = fit(msg, limit).expect("dropping the body brings this under the limit"); - assert_eq!(fitted.body.as_deref(), Some(DROPPED_BODY)); - // The tag describes a body that is no longer there; left set, a reader - // renders the marker as a diff. - assert_eq!(fitted.body_format, None); - // The fields that make the row readable survive. - assert_eq!(fitted.summary, "Edit(src/main.rs)"); - assert_eq!(fitted.icon.as_deref(), Some("🔧")); - assert_eq!(fitted.coalesce_key.as_deref(), Some("tool-1")); - // Including the time: a degraded row that lost it would be a row a - // subscriber cannot place, and nothing downstream could tell. - assert_eq!(fitted.ts, ROW_TS); - assert!( - serde_json::to_vec(&fitted).expect("serialises").len() <= limit, - "the degraded row must be under the limit, not merely smaller" - ); - } - - /// A row whose summary alone exceeds the limit cannot be degraded into - /// one, and publishing it anyway would cost the connection rather than the - /// row. Reported by the caller, never sent. - #[test] - fn a_row_too_large_even_without_its_body_is_refused() { - let limit = 256; - let msg = TermMsg::new(Level::Warn, "s".repeat(limit * 4)).body("x".repeat(limit), None); - assert!(fit(msg, limit).is_none()); - } - - /// Control for the test above: the same oversize summary with a body that - /// would fit still refuses, so the refusal is about total size and not - /// about the body having been present. - #[test] - fn the_refusal_is_about_size_rather_than_the_body_being_present() { - let limit = 256; - let msg = TermMsg::new(Level::Warn, "s".repeat(limit * 4)); - assert!(fit(msg, limit).is_none()); - } - - /// JSON escaping is why the size is measured by serializing: a body of - /// quotes serializes to twice its own length, so a row that fits by - /// character count can still be refused on the wire. - #[test] - fn the_limit_is_measured_on_the_serialized_bytes() { - let body = "\"".repeat(200); - let msg = TermMsg::new(Level::Info, "quotes").body(body.clone(), None); - let limit = body.len() + 64; - // Well under the limit as characters, over it once escaped. - assert!( - serde_json::to_vec(&msg).expect("serialises").len() > limit, - "this fixture must be oversize only after escaping" - ); - let fitted = fit(msg, limit).expect("dropping the body fits"); - assert_eq!(fitted.body.as_deref(), Some(DROPPED_BODY)); - } - /// What a subscriber actually receives: the payload this module hands /// `publish` is the bare row, so the time has to be *in* it. A `ts` that /// existed in Rust but never serialized would leave the queue exactly as diff --git a/hive-agent/src/term_msg.rs b/hive-agent/src/term_msg.rs index 1965a0fe..caeddcc8 100644 --- a/hive-agent/src/term_msg.rs +++ b/hive-agent/src/term_msg.rs @@ -1,5 +1,5 @@ -//! Terminal-message wire shape: what the per-agent web UI's live/history -//! endpoints actually serve for the "terminal" event stream, as opposed to +//! The harness's events as terminal rows: what the per-agent web UI's +//! live/history endpoints serve for the "terminal" event stream, as opposed to //! agent-state changes (`StatusChanged`/`ModelChanged`/`EffortChanged`/ //! `TokenUsageChanged`/`TurnStateChanged`), which have never rendered as //! terminal rows (the header/badges poll `/api/state`, not this stream) and @@ -9,167 +9,17 @@ //! map to exactly one, but `LiveEvent::Stream` (one raw claude //! `stream-json` line) can expand to several: an `assistant` message with //! both a text block and a `tool_use` block produces two rows. -//! -//! Seven fields, and classification belongs **here**, not in the client: the -//! web UI renders what it is handed and owns no per-tool dispatch table, so -//! a new tool needs no frontend change. `level` carries styling, so a row -//! never names a CSS class. `kind`, `unread`, `from` and `expanded_default` -//! are deliberately absent — adding one back is a design change, not an -//! oversight. -//! -//! The seventh field is `ts`, and it is the row's own: [`classify`] stamps -//! every row it returns with the time of the event it classified, not the -//! time it ran. The swarm queue publishes the bare row with no envelope -//! around it, so a subscriber that reads a row has nowhere else to learn -//! when the thing happened; carrying the source event's time means a row -//! replayed out of sqlite months later still says when it happened rather -//! than when it was read. -use chrono::{DateTime, SecondsFormat}; -use serde::Serialize; -use std::collections::HashMap; +pub use hive_sh4re::term_msg::{ClassifyCtx, Level, TermMsg, iso8601_utc}; use crate::events::LiveEvent; -#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)] -#[serde(rename_all = "snake_case")] -pub enum Level { - /// Low-signal / ambient chatter — thinking, progress ticks, harness - /// housekeeping. The client's default rendering can dim/de-emphasize - /// these without hiding them outright. - Debug, - /// Routine substantive content — turn boundaries, assistant text, tool - /// calls/results, message bodies. - Info, - /// Heads-up, not necessarily broken — stderr lines, an unclassified - /// event shape landing (the old `.sys` catch-all), API retries. - Warn, - /// Something actually failed — a turn ending non-ok, a tool result with - /// `is_error: true`. - Error, -} - -#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)] -#[serde(rename_all = "snake_case")] -pub enum BodyFormat { - Markdown, - Diff, -} - -/// One terminal row. `body_format: None` with `body: Some(_)` means plain -/// text (the common case — no explicit tag on the wire for it, same logic -/// as `body` itself being absent meaning "nothing to expand"). -#[derive(Debug, Clone, Serialize)] -pub struct TermMsg { - /// When the classified event happened, ISO 8601 / RFC 3339 UTC. - /// - /// Empty on a freshly built row and filled by [`classify`], which is - /// the only path a row reaches either wire by — a builder that had to - /// be handed the time would repeat the same value across the several - /// rows one `stream-json` line expands into. - pub ts: String, - #[serde(skip_serializing_if = "Option::is_none")] - pub icon: Option, - pub level: Level, - pub summary: String, - #[serde(skip_serializing_if = "Option::is_none")] - pub body: Option, - #[serde(skip_serializing_if = "Option::is_none")] - pub body_format: Option, - #[serde(skip_serializing_if = "Option::is_none")] - pub coalesce_key: Option, -} - -impl TermMsg { - pub fn new(level: Level, summary: impl Into) -> Self { - Self { - ts: String::new(), - icon: None, - level, - summary: summary.into(), - body: None, - body_format: None, - coalesce_key: None, - } - } - - #[must_use] - pub fn icon(mut self, icon: impl Into) -> Self { - self.icon = Some(icon.into()); - self - } - - #[must_use] - pub fn body(mut self, body: impl Into, format: Option) -> Self { - self.body = Some(body.into()); - self.body_format = format; - self - } - - #[must_use] - pub fn coalesce(mut self, key: impl Into) -> Self { - self.coalesce_key = Some(key.into()); - self - } - - /// Stamp the row with the time of the event it came from. - #[must_use] - pub fn at(mut self, ts: impl Into) -> Self { - self.ts = ts.into(); - self - } -} - -/// Render an event's unix-seconds stamp as ISO 8601 / RFC 3339 UTC. -/// -/// Seconds rather than milliseconds because that is the resolution the -/// event bus and the sqlite `events.ts` column actually carry -/// (`crate::events`) — a `.000` on every row would be precision the source -/// does not have. UTC rather than local: the reader of a swarm-published -/// row is not on the machine that wrote it, and an offset-carrying stamp -/// would make two agents' rows sort by string differently than by time. -/// -/// A stamp outside the representable range renders as the epoch: a row -/// whose only defect is an absurd clock is still worth reading. -#[must_use] -pub fn iso8601_utc(unix_seconds: i64) -> String { - DateTime::from_timestamp(unix_seconds, 0) - .unwrap_or_default() - .to_rfc3339_opts(SecondsFormat::Secs, true) -} - -/// Per-connection/per-request classification state. A live SSE stream keeps -/// one of these alive for the connection's lifetime — `tool_use` id → name -/// correlation, so a `tool_result` can tell it's answering a `recv` call and -/// render its body as markdown instead of plain text. The history endpoint -/// uses a fresh one per page: correlation only works within the page -/// actually returned, not across the live/history boundary. Accepted -/// degradation — the only user-visible effect is a `recv` result whose -/// `tool_use` fell on the other side of a page/reconnect boundary rendering -/// its body as plain text instead of markdown. -#[derive(Default)] -pub struct ClassifyCtx { - tool_name_by_id: HashMap, -} - -impl ClassifyCtx { - pub fn record_tool_use(&mut self, id: &str, name: &str) { - self.tool_name_by_id.insert(id.to_owned(), name.to_owned()); - } - - #[must_use] - pub fn tool_name(&self, id: &str) -> Option<&str> { - self.tool_name_by_id.get(id).map(String::as_str) - } -} - /// Classify one [`LiveEvent`] into zero or more terminal rows, each /// stamped with `ts` — the event's own unix-seconds time, from the bus on /// the live path and from the `events` table on replay. /// /// The stamp is applied here, over the rows the classifiers return, so that -/// every row reaching a wire carries one: a row built anywhere else has an -/// empty `ts` and cannot get out without passing through this function. +/// every row the harness puts on a wire carries one. pub fn classify(ev: &LiveEvent, ts: i64, ctx: &mut ClassifyCtx) -> Vec { let ts = iso8601_utc(ts); classify_rows(ev, ctx) @@ -202,7 +52,7 @@ fn classify_rows(ev: &LiveEvent, ctx: &mut ClassifyCtx) -> Vec { vec![msg] } LiveEvent::Note { text } => vec![classify_note(text)], - LiveEvent::Stream(v) => crate::stream_enrich::classify_stream_value(v, ctx), + LiveEvent::Stream(v) => hive_sh4re::stream_enrich::classify_stream_value(v, ctx), // Agent-state transitions never render as terminal rows — the // header/badges read `/api/state`, not this stream (see module doc). LiveEvent::StatusChanged { .. } diff --git a/hive-sh4re/Cargo.toml b/hive-sh4re/Cargo.toml index 2efaebd9..04e98e91 100644 --- a/hive-sh4re/Cargo.toml +++ b/hive-sh4re/Cargo.toml @@ -13,11 +13,10 @@ hive-priv-sock.workspace = true hive-types.workspace = true schemars.workspace = true serde.workspace = true +# `stream_enrich` walks raw claude `stream-json` values into `term_msg` rows. +serde_json.workspace = true strum.workspace = true # Facade only, for the one warn in `permissions::ToolGroup::parse_list`: an # unknown tool-group name is skipped rather than fatal, so the log line is the # only trace it leaves. tracing.workspace = true - -[dev-dependencies] -serde_json.workspace = true diff --git a/hive-sh4re/src/lib.rs b/hive-sh4re/src/lib.rs index 37129bbb..85270971 100644 --- a/hive-sh4re/src/lib.rs +++ b/hive-sh4re/src/lib.rs @@ -9,4 +9,6 @@ pub mod manager; pub mod paths; pub mod permissions; pub mod schedule; +pub mod stream_enrich; +pub mod term_msg; pub mod wire_time; diff --git a/hive-agent/src/stream_enrich.rs b/hive-sh4re/src/stream_enrich.rs similarity index 99% rename from hive-agent/src/stream_enrich.rs rename to hive-sh4re/src/stream_enrich.rs index b73994c8..0b5fbe70 100644 --- a/hive-agent/src/stream_enrich.rs +++ b/hive-sh4re/src/stream_enrich.rs @@ -1,13 +1,13 @@ //! Classify raw claude stream-json values into [`crate::term_msg::TermMsg`] -//! rows before SSE delivery. +//! rows. //! //! [`classify_stream_value`] is the entry point — it walks one raw claude -//! `stream-json` line (the payload of a [`crate::events::LiveEvent::Stream`]) -//! and returns zero or more terminal rows. Applied at SSE-emit time in -//! `crate::web_ui::stream` so both the live tail (`events/stream`) and the -//! history replay (`events/history`) endpoints deliver the same classified -//! shape — the sqlite event log stores the raw, unclassified event, so the -//! DB never needs migration when classification logic changes. +//! `stream-json` line and returns zero or more terminal rows. Two callers: +//! `hive-agent` applies it at SSE-emit time, so both the live tail +//! (`events/stream`) and the history replay (`events/history`) deliver the +//! same classified shape while the sqlite event log keeps the raw event; and +//! `hive-subagent-mcp` applies it to each subagent's output before publishing +//! it to the swarm queue. //! //! The per-tool icon/summary formatting below (`tool_icon`, `fmt_tool_use` //! and its per-family helpers, `rich_tool_body`) is reused as-is from diff --git a/hive-sh4re/src/term_msg.rs b/hive-sh4re/src/term_msg.rs new file mode 100644 index 00000000..d896934c --- /dev/null +++ b/hive-sh4re/src/term_msg.rs @@ -0,0 +1,277 @@ +//! Terminal-message wire shape: one classified terminal row, as the +//! per-agent web UI's live/history endpoints serve it and as the swarm queue +//! carries it from an agent and from its subagents. +//! +//! Seven fields, and classification belongs **here**, not in the client: the +//! web UI renders what it is handed and owns no per-tool dispatch table, so +//! a new tool needs no frontend change. `level` carries styling, so a row +//! never names a CSS class. `kind`, `unread`, `from` and `expanded_default` +//! are deliberately absent — adding one back is a design change, not an +//! oversight. +//! +//! The seventh field is `ts`, and it is the row's own: whoever classifies an +//! event stamps every row it produced with the time of that event, not the +//! time it ran. The swarm queue publishes the bare row with no envelope +//! around it, so a subscriber that reads a row has nowhere else to learn +//! when the thing happened; carrying the source event's time means a row +//! replayed out of sqlite months later still says when it happened rather +//! than when it was read. + +use chrono::{DateTime, SecondsFormat}; +use serde::Serialize; +use std::collections::HashMap; + +#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)] +#[serde(rename_all = "snake_case")] +pub enum Level { + /// Low-signal / ambient chatter — thinking, progress ticks, harness + /// housekeeping. The client's default rendering can dim/de-emphasize + /// these without hiding them outright. + Debug, + /// Routine substantive content — turn boundaries, assistant text, tool + /// calls/results, message bodies. + Info, + /// Heads-up, not necessarily broken — stderr lines, an unclassified + /// event shape landing (the old `.sys` catch-all), API retries. + Warn, + /// Something actually failed — a turn ending non-ok, a tool result with + /// `is_error: true`. + Error, +} + +#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)] +#[serde(rename_all = "snake_case")] +pub enum BodyFormat { + Markdown, + Diff, +} + +/// One terminal row. `body_format: None` with `body: Some(_)` means plain +/// text (the common case — no explicit tag on the wire for it, same logic +/// as `body` itself being absent meaning "nothing to expand"). +#[derive(Debug, Clone, Serialize)] +pub struct TermMsg { + /// When the classified event happened, ISO 8601 / RFC 3339 UTC. + /// + /// Empty on a freshly built row and filled with [`TermMsg::at`] by the + /// caller that classified the event — a builder that had to be handed + /// the time would repeat the same value across the several rows one + /// `stream-json` line expands into. + pub ts: String, + #[serde(skip_serializing_if = "Option::is_none")] + pub icon: Option, + pub level: Level, + pub summary: String, + #[serde(skip_serializing_if = "Option::is_none")] + pub body: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub body_format: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub coalesce_key: Option, +} + +impl TermMsg { + pub fn new(level: Level, summary: impl Into) -> Self { + Self { + ts: String::new(), + icon: None, + level, + summary: summary.into(), + body: None, + body_format: None, + coalesce_key: None, + } + } + + #[must_use] + pub fn icon(mut self, icon: impl Into) -> Self { + self.icon = Some(icon.into()); + self + } + + #[must_use] + pub fn body(mut self, body: impl Into, format: Option) -> Self { + self.body = Some(body.into()); + self.body_format = format; + self + } + + #[must_use] + pub fn coalesce(mut self, key: impl Into) -> Self { + self.coalesce_key = Some(key.into()); + self + } + + /// Stamp the row with the time of the event it came from. + #[must_use] + pub fn at(mut self, ts: impl Into) -> Self { + self.ts = ts.into(); + self + } +} + +/// Render an event's unix-seconds stamp as ISO 8601 / RFC 3339 UTC. +/// +/// Seconds rather than milliseconds because that is the resolution the +/// event bus and the sqlite `events.ts` column actually carry +/// (`hive-agent`'s `events`) — a `.000` on every row would be precision the source +/// does not have. UTC rather than local: the reader of a swarm-published +/// row is not on the machine that wrote it, and an offset-carrying stamp +/// would make two agents' rows sort by string differently than by time. +/// +/// A stamp outside the representable range renders as the epoch: a row +/// whose only defect is an absurd clock is still worth reading. +#[must_use] +pub fn iso8601_utc(unix_seconds: i64) -> String { + DateTime::from_timestamp(unix_seconds, 0) + .unwrap_or_default() + .to_rfc3339_opts(SecondsFormat::Secs, true) +} + +/// Per-connection/per-request classification state. A live SSE stream keeps +/// one of these alive for the connection's lifetime — `tool_use` id → name +/// correlation, so a `tool_result` can tell it's answering a `recv` call and +/// render its body as markdown instead of plain text. The history endpoint +/// uses a fresh one per page: correlation only works within the page +/// actually returned, not across the live/history boundary. Accepted +/// degradation — the only user-visible effect is a `recv` result whose +/// `tool_use` fell on the other side of a page/reconnect boundary rendering +/// its body as plain text instead of markdown. +#[derive(Default)] +pub struct ClassifyCtx { + tool_name_by_id: HashMap, +} + +impl ClassifyCtx { + pub fn record_tool_use(&mut self, id: &str, name: &str) { + self.tool_name_by_id.insert(id.to_owned(), name.to_owned()); + } + + #[must_use] + pub fn tool_name(&self, id: &str) -> Option<&str> { + self.tool_name_by_id.get(id).map(String::as_str) + } +} + +/// What a row that lost its body to the payload limit carries instead. +pub const DROPPED_BODY: &str = "[body dropped: over the queue's payload limit]"; + +/// Bring `msg` under `limit` serialized bytes, or report that it cannot be. +/// +/// The body is the only field that carries arbitrary length — a diff, a whole +/// tool result — so it is the only one worth spending: replacing it keeps the +/// row's identity, level, summary and time, which is what makes the row +/// readable at all, and a reader sees that something was there rather than +/// seeing nothing. +/// +/// `None` means even the degraded row does not fit, so the caller reports it +/// instead of publishing. That matters more than it looks: an oversize publish +/// is not truncated by the server, it is refused and the connection is closed, +/// which costs the row *and* every row racing behind it through the reconnect. +#[must_use] +pub fn fit(msg: TermMsg, limit: usize) -> Option { + if serialized_len(&msg)? <= limit { + return Some(msg); + } + // Rebuilt rather than mutated in place: `body_format` describes the body, + // and a marker left tagged `Diff` renders as a broken diff downstream. + let degraded = TermMsg { + body: Some(DROPPED_BODY.to_owned()), + body_format: None, + ..msg + }; + (serialized_len(°raded)? <= limit).then_some(degraded) +} + +/// Serialized size of a row, or `None` if it does not serialize at all. +/// +/// Measured by serializing rather than estimated from field lengths: the +/// payload is what the server measures, and JSON escaping makes the two differ +/// by an unbounded factor on exactly the rows that are already near the limit. +fn serialized_len(msg: &TermMsg) -> Option { + serde_json::to_vec(msg).ok().map(|v| v.len()) +} + +#[cfg(test)] +mod tests { + use super::{BodyFormat, DROPPED_BODY, Level, TermMsg, fit}; + + /// A row's time, for the tests below that need one. + const ROW_TS: &str = "2026-09-13T12:35:03Z"; + + #[test] + fn a_row_that_already_fits_is_published_unchanged() { + let msg = TermMsg::new(Level::Info, "turn ok") + .icon("✅") + .body("a short body", Some(BodyFormat::Markdown)); + let fitted = fit(msg, 4096).expect("a small row fits"); + assert_eq!(fitted.body.as_deref(), Some("a short body")); + assert_eq!(fitted.body_format, Some(BodyFormat::Markdown)); + } + + /// The case the degrade exists for: a body larger than the limit costs the + /// body and nothing else, and what comes back is actually under the limit + /// rather than merely smaller. + #[test] + fn an_oversize_body_is_replaced_and_the_result_fits() { + let limit = 512; + let msg = TermMsg::new(Level::Info, "Edit(src/main.rs)") + .icon("🔧") + .body("x".repeat(limit * 4), Some(BodyFormat::Diff)) + .coalesce("tool-1") + .at(ROW_TS); + let fitted = fit(msg, limit).expect("dropping the body brings this under the limit"); + assert_eq!(fitted.body.as_deref(), Some(DROPPED_BODY)); + // The tag describes a body that is no longer there; left set, a reader + // renders the marker as a diff. + assert_eq!(fitted.body_format, None); + // The fields that make the row readable survive. + assert_eq!(fitted.summary, "Edit(src/main.rs)"); + assert_eq!(fitted.icon.as_deref(), Some("🔧")); + assert_eq!(fitted.coalesce_key.as_deref(), Some("tool-1")); + // Including the time: a degraded row that lost it would be a row a + // subscriber cannot place, and nothing downstream could tell. + assert_eq!(fitted.ts, ROW_TS); + assert!( + serde_json::to_vec(&fitted).expect("serialises").len() <= limit, + "the degraded row must be under the limit, not merely smaller" + ); + } + + /// A row whose summary alone exceeds the limit cannot be degraded into + /// one, and publishing it anyway would cost the connection rather than the + /// row. Reported by the caller, never sent. + #[test] + fn a_row_too_large_even_without_its_body_is_refused() { + let limit = 256; + let msg = TermMsg::new(Level::Warn, "s".repeat(limit * 4)).body("x".repeat(limit), None); + assert!(fit(msg, limit).is_none()); + } + + /// Control for the test above: the same oversize summary with a body that + /// would fit still refuses, so the refusal is about total size and not + /// about the body having been present. + #[test] + fn the_refusal_is_about_size_rather_than_the_body_being_present() { + let limit = 256; + let msg = TermMsg::new(Level::Warn, "s".repeat(limit * 4)); + assert!(fit(msg, limit).is_none()); + } + + /// JSON escaping is why the size is measured by serializing: a body of + /// quotes serializes to twice its own length, so a row that fits by + /// character count can still be refused on the wire. + #[test] + fn the_limit_is_measured_on_the_serialized_bytes() { + let body = "\"".repeat(200); + let msg = TermMsg::new(Level::Info, "quotes").body(body.clone(), None); + let limit = body.len() + 64; + // Well under the limit as characters, over it once escaped. + assert!( + serde_json::to_vec(&msg).expect("serialises").len() > limit, + "this fixture must be oversize only after escaping" + ); + let fitted = fit(msg, limit).expect("dropping the body fits"); + assert_eq!(fitted.body.as_deref(), Some(DROPPED_BODY)); + } +} diff --git a/hive-subagent-mcp/Cargo.toml b/hive-subagent-mcp/Cargo.toml index 6ea7c82e..fcc059d2 100644 --- a/hive-subagent-mcp/Cargo.toml +++ b/hive-subagent-mcp/Cargo.toml @@ -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 diff --git a/hive-subagent-mcp/src/lib.rs b/hive-subagent-mcp/src/lib.rs index 1bed0f8f..bc84e8ef 100644 --- a/hive-subagent-mcp/src/lib.rs +++ b/hive-subagent-mcp/src/lib.rs @@ -21,3 +21,4 @@ pub mod mcp_config; pub mod paths; pub mod role; pub mod session; +pub mod swarm_term; diff --git a/hive-subagent-mcp/src/main.rs b/hive-subagent-mcp/src/main.rs index c11ef75e..54356392 100644 --- a/hive-subagent-mcp/src/main.rs +++ b/hive-subagent-mcp/src/main.rs @@ -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 diff --git a/hive-subagent-mcp/src/session.rs b/hive-subagent-mcp/src/session.rs index c601deae..fcd9ca74 100644 --- a/hive-subagent-mcp/src/session.rs +++ b/hive-subagent-mcp/src/session.rs @@ -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, } /// 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) -> 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, 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()); + } } } diff --git a/hive-subagent-mcp/src/swarm_term.rs b/hive-subagent-mcp/src/swarm_term.rs new file mode 100644 index 00000000..9934e0e5 --- /dev/null +++ b/hive-subagent-mcp/src/swarm_term.rs @@ -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..sub.` (`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-`, 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, + dropped: Arc, +} + +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 { + 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) { + 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) -> anyhow::Result> { + 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 { + 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 = 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) -> anyhow::Result { + 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, + connect_after: Option, + stream_open: bool, + stream_after: Option, +} + +impl Queue { + /// A client to publish on, connecting and opening the stream first when + /// they are due. + async fn ready(&mut self, target: &Arc) -> Option { + 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, target: Target, dropped: Arc) { + let target = Arc::new(target); + let mut ctxs: HashMap = 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 { + 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 { + 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"); + } +} diff --git a/hivectl/src/watch.rs b/hivectl/src/watch.rs index c259ae31..d7dcccf1 100644 --- a/hivectl/src/watch.rs +++ b/hivectl/src/watch.rs @@ -122,7 +122,7 @@ fn render_frame(frame: &str) { "note" => format!("· {}", str_field("text").unwrap_or("")), "status_changed" => format!("· status: {}", str_field("status").unwrap_or("")), // Claude stream-json lines, enriched server-side - // (`hive-agent/src/stream_enrich.rs`) with `_icon`/`_summary` so + // (`hive-sh4re/src/stream_enrich.rs`) with `_icon`/`_summary` so // clients don't need to reimplement the web UI's dispatch logic. "stream" => match (str_field("_icon"), str_field("_summary")) { (Some(icon), Some(summary)) => format!("{icon} {summary}"), diff --git a/nix/agent-modules/mcp.nix b/nix/agent-modules/mcp.nix index e41b69df..1a02336c 100644 --- a/nix/agent-modules/mcp.nix +++ b/nix/agent-modules/mcp.nix @@ -31,11 +31,14 @@ let # (several agents that rarely compile at the same time) working. subagentMemoryHigh = containerMemoryMaxBytes * 2 / 3; # Set on the harness only for an ACP agent on the opencode preset - # (./agent-service.nix). The subagent daemon then reads that key from the - # store as this agent, so it is handed the same store identity - # ./queue-identity.nix hands the harness. + # (./agent-service.nix). acpApiKeyEnv = config.systemd.services.hive-agent.environment.HIVE_ACP_API_KEY_ENV or null; - subagentReadsStore = acpApiKeyEnv != null && config.services.hyperhive.agent.bao.addr != null; + # The subagent daemon reads this agent's own queue credential from the + # store, to publish its subagents' terminals as the agent, and on the + # opencode preset the provider key too. Either way it is handed the same + # store identity ./queue-identity.nix hands the harness, under the same + # switch. + subagentReadsStore = config.services.hyperhive.agent.bao.addr != null; in { options.services.hyperhive.agent.allowedRecipients = lib.mkOption { diff --git a/nix/host-modules/swarm-nats.nix b/nix/host-modules/swarm-nats.nix index 2111355e..f5cb2fb9 100644 --- a/nix/host-modules/swarm-nats.nix +++ b/nix/host-modules/swarm-nats.nix @@ -844,6 +844,16 @@ in # on a hive presents that one, so it cannot be scoped to one # agent's key. "--agent-token-publish-subject ${lib.escapeShellArg "\$\$KV.agent-icons.{agent}"}" + # Its subagents' terminal rows, and the one stream that keeps + # them (`swarm_queue_client::subagent_term`). The agent's + # subagent daemon creates the stream on first use, so it gets + # `CREATE` and `INFO` on that stream name alone: no `UPDATE`, + # no `DELETE`, no consumer, no other stream. A `CREATE`'s + # config travels in its payload, which no subject grant can + # narrow. + "--agent-token-publish-subject ${lib.escapeShellArg "\$\$SWARM.term.{agent}.sub.>"}" + "--agent-token-publish-subject ${lib.escapeShellArg "\$\$JS.API.STREAM.CREATE.term-sub-{agent}"}" + "--agent-token-publish-subject ${lib.escapeShellArg "\$\$JS.API.STREAM.INFO.term-sub-{agent}"}" "--store-cert-role ${lib.escapeShellArg authCertRole}" ]; # Every credential arrives by `LoadCredential` and is named diff --git a/nix/module-eval/agent-runtime.nix b/nix/module-eval/agent-runtime.nix index b69b9374..66629f66 100644 --- a/nix/module-eval/agent-runtime.nix +++ b/nix/module-eval/agent-runtime.nix @@ -86,13 +86,32 @@ let && (subagent opencode).environment.HIVE_ACP_API_KEY_ENV == "T_PROVIDER_KEY"; } { - # Only the opencode preset has a provider key variable; any other ACP - # command reads nothing from the store. + # Only the opencode preset has a provider key variable. name = "an ACP agent off the opencode preset is told no provider key variable"; ok = !((harness acpBao).environment ? HIVE_ACP_API_KEY_ENV) - && (subagent acpBao).environment.HIVE_ACP_API_KEY_ENV or null == null - && !((subagent acpBao).serviceConfig ? LoadCredential); + && (subagent acpBao).environment.HIVE_ACP_API_KEY_ENV or null == null; + } + { + # Every agent with a store, whatever its runtime: the daemon reads the + # agent's queue credential with it to publish subagent terminals. + name = "a subagent daemon gets the agent's store identity whenever the agent has a store"; + ok = + lib.all + ( + u: + u.serviceConfig.LoadCredential == [ + "hive-agent-bao-cert" + "hive-agent-bao-key" + "hive-agent-bao-server-ca" + ] + && u.environment.HIVE_AGENT_NAME == "a1" + && u.environment.BAO_ADDR == "https://bao.t.local:8200" + ) + [ + (subagent acpBao) + (subagent opencodeBao) + ]; } { name = "an opencode agent's subagent daemon gets the agent's store identity"; diff --git a/nix/module-eval/nats-authelia.nix b/nix/module-eval/nats-authelia.nix index 12ae6dce..142807ea 100644 --- a/nix/module-eval/nats-authelia.nix +++ b/nix/module-eval/nats-authelia.nix @@ -161,6 +161,32 @@ let in lib.hasInfix "--agent-token-publish-subject '$$KV.agent-icons.{agent}'" exec; } + { + # The whole per-agent grant, compared as a list rather than by infix: a + # wider subject added beside these (`$$JS.API.>`, another stream's name) + # would pass every presence check above. + name = "a verified agent's grant is exactly its own subjects and its subagent stream"; + ok = + let + args = lib.splitString " " (responderOf natsOldPath).serviceConfig.ExecStart; + granted = lib.concatLists ( + lib.imap0 ( + i: a: + lib.optional (a == "--agent-token-publish-subject" && i + 1 < lib.length args) ( + lib.elemAt args (i + 1) + ) + ) args + ); + in + granted == [ + "'$$SWARM.term.{agent}'" + "'$$SWARM.agent-state.{agent}'" + "'$$KV.agent-icons.{agent}'" + "'$$SWARM.term.{agent}.sub.>'" + "'$$JS.API.STREAM.CREATE.term-sub-{agent}'" + "'$$JS.API.STREAM.INFO.term-sub-{agent}'" + ]; + } { # Each credential by `LoadCredential`, and the environment naming where # the unit sees it: a DynamicUser cannot read the copies directly. diff --git a/swarm-controller/src/main.rs b/swarm-controller/src/main.rs index 066c3c86..2a89d872 100644 --- a/swarm-controller/src/main.rs +++ b/swarm-controller/src/main.rs @@ -58,6 +58,7 @@ mod queue_identity; mod read_policy; mod status; mod store; +mod subagent_term; mod term_stream; mod vcs_metrics; mod wanted; @@ -2992,6 +2993,8 @@ fn build_app(state: AppState) -> axum::Router { .routes(routes!(linked_accounts::get_linked_accounts)) .routes(routes!(get_hive_wanted)) .routes(routes!(term_stream::stream_agent_term)) + .routes(routes!(subagent_term::list_subagents)) + .routes(routes!(subagent_term::stream_subagent_term)) .routes(routes!(agent_state_stream::stream_agent_state)) .routes(routes!(issue_report::get_repos)) .routes(routes!(issue_report::get_issue_report_all)) @@ -4981,6 +4984,55 @@ mod tests { "{problem:?}" ); } + + #[tokio::test] + async fn listing_subagents_needs_a_queue_and_an_identifier() { + let (state, _sched) = state_with_roster(); + let no_queue = super::subagent_term::list_subagents( + axum::extract::State(state.clone()), + axum::extract::Path("iris".to_owned()), + ) + .await + .expect_err("no queue is configured"); + assert_eq!( + no_queue.status, + Some(axum::http::StatusCode::SERVICE_UNAVAILABLE), + "{no_queue:?}" + ); + let bad_name = super::subagent_term::list_subagents( + axum::extract::State(state), + axum::extract::Path("iris.sub".to_owned()), + ) + .await + .expect_err("a dot would widen the subject"); + assert_eq!( + bad_name.status, + Some(axum::http::StatusCode::BAD_REQUEST), + "{bad_name:?}" + ); + } + + /// A subagent name is a subject token, so a wildcard or a dot in it is + /// refused before anything subscribes. + #[tokio::test] + async fn a_subagent_stream_refuses_a_name_that_is_not_one_token() { + let (state, _sched) = state_with_roster(); + for subagent in ["*", ">", "a.b"] { + let Err(problem) = super::subagent_term::stream_subagent_term( + axum::extract::State(state.clone()), + axum::extract::Path(("iris".to_owned(), subagent.to_owned())), + ) + .await + else { + panic!("{subagent:?} must be refused"); + }; + assert_eq!( + problem.status, + Some(axum::http::StatusCode::BAD_REQUEST), + "{subagent:?}: {problem:?}" + ); + } + } } #[cfg(test)] diff --git a/swarm-controller/src/subagent_term.rs b/swarm-controller/src/subagent_term.rs new file mode 100644 index 00000000..40ac5e9b --- /dev/null +++ b/swarm-controller/src/subagent_term.rs @@ -0,0 +1,172 @@ +//! An agent's subagents: which ones the swarm has rows for, and a live SSE +//! relay of one subagent's terminal. +//! +//! The agent's subagent daemon publishes each subagent's classified rows on +//! `$SWARM.term.{agent}.sub.{subagent}` under the agent's own queue credential, +//! into the stream `term-sub-{agent}` it creates itself +//! (`swarm_queue_client::subagent_term`). Subagents are spawned on demand +//! under caller-chosen names, so the list is the set of subjects that stream +//! holds rows for, read with `STREAM.INFO`. This daemon never creates the +//! stream: a missing one is an agent with no subagent output yet, and lists +//! nothing. +//! +//! **No hive lookup.** Unlike [`crate::term_stream`], there is no +//! hive-scoped subject to follow: only the agent's own credential is granted +//! these subjects. +//! +//! **Live tail only, output only.** The relay is a core subscription, the +//! same as the agent's own terminal route, so a reader sees rows published +//! while it is attached. Nothing is sent toward a subagent: subagents take no +//! input from the swarm. + +use std::convert::Infallible; + +use axum::Json; +use axum::extract::{Path, State}; +use axum::http::StatusCode; +use axum::response::sse::{Event, KeepAlive, Sse}; +use futures_util::{Stream, StreamExt as _, TryStreamExt as _}; +use swarm_queue_client::subagent_term; + +use crate::AppState; + +#[utoipa::path( + get, + path = "/api/agents/{name}/subagents", + params(("name" = String, Path, description = "agent whose subagents to list")), + responses( + (status = 200, description = "names of the agent's subagents with terminal rows \ + in the swarm queue, sorted; empty when the agent has published none", + body = Vec), + (status = 400, description = "the agent name is not shaped like an identifier \ + (problem+json)", body = String), + (status = 503, description = "no swarm queue is configured on this host, or the \ + queue could not be reached (problem+json)", body = String), + (status = 500, description = "reading the agent's subagent stream failed \ + (problem+json)", body = String), + ), + tag = "agents" +)] +pub(crate) async fn list_subagents( + State(state): State, + Path(agent): Path, +) -> Result>, problem_details::ProblemDetails> { + let agent = ident(&agent).map_err(bad_request)?; + let client = connected_client(&state).map_err(|e| unavailable(&e))?; + let js = async_nats::jetstream::new(client); + let stream_name = subagent_term::stream_name(&agent); + let stream = match js.get_stream(&stream_name).await { + Ok(stream) => stream, + Err(e) if is_stream_not_found(&e) => return Ok(Json(Vec::new())), + Err(e) => return Err(read_failed(&agent, &e)), + }; + let subjects: Vec<(String, usize)> = stream + .info_with_subjects(subagent_term::stream_subjects(&agent)) + .await + .map_err(|e| read_failed(&agent, &e))? + .try_collect() + .await + .map_err(|e| read_failed(&agent, &e))?; + Ok(Json(subagent_term::subagent_names( + &agent, + subjects.iter().map(|(s, _)| s.as_str()), + ))) +} + +#[utoipa::path( + get, + path = "/api/agents/{name}/subagents/{subagent}/term/stream", + params( + ("name" = String, Path, description = "agent the subagent belongs to"), + ("subagent" = String, Path, description = "subagent whose terminal to stream"), + ), + responses( + (status = 200, description = "server-sent event stream; each event's `data` is \ + one already-classified TermMsg row, JSON as the agent's subagent daemon \ + published it (opaque to this daemon) — live only, no replay", + body = String, content_type = "text/event-stream"), + (status = 400, description = "the agent or subagent name is not shaped like an \ + identifier (problem+json)", body = String), + (status = 503, description = "no swarm queue is configured on this host, or the \ + queue could not be reached (problem+json)", body = String), + (status = 500, description = "the subscribe itself failed (problem+json)", + body = String), + ), + tag = "agents" +)] +pub(crate) async fn stream_subagent_term( + State(state): State, + Path((agent, subagent)): Path<(String, String)>, +) -> Result>>, problem_details::ProblemDetails> { + let agent = ident(&agent).map_err(bad_request)?; + let subagent = ident(&subagent).map_err(bad_request)?; + let client = connected_client(&state).map_err(|e| unavailable(&e))?; + let subscriber = crate::term_stream::subscribe_all( + &client, + &[subagent_term::subject(&agent, &subagent)], + "subagent 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 queue client, or why there is none: the 503 every queue-backed route +/// answers. +fn connected_client(state: &AppState) -> Result { + let status = state + .status + .as_ref() + .ok_or_else(|| "no swarm queue is configured on this host".to_owned())?; + let client = status.queue_client(); + swarm_queue_client::ensure_connected(&client).map_err(|e| swarm_queue_client::chain(&e))?; + Ok(client) +} + +/// A path segment as a single subject token, or why it is not one. +fn ident(name: &str) -> Result { + hive_types::Ident::parse(name).map(hive_types::Ident::into_string) +} + +fn bad_request(reason: &str) -> problem_details::ProblemDetails { + crate::error_problem(StatusCode::BAD_REQUEST, reason) +} + +fn unavailable(detail: &str) -> problem_details::ProblemDetails { + crate::error_problem(StatusCode::SERVICE_UNAVAILABLE, detail) +} + +fn is_stream_not_found(e: &async_nats::jetstream::context::GetStreamError) -> bool { + matches!( + e.kind(), + async_nats::jetstream::context::GetStreamErrorKind::JetStream(js) + if js.error_code() == async_nats::jetstream::ErrorCode::STREAM_NOT_FOUND + ) +} + +fn read_failed(agent: &str, e: &dyn std::error::Error) -> problem_details::ProblemDetails { + let detail = swarm_queue_client::chain(e); + tracing::warn!(%agent, error = %detail, "listing subagents: reading the stream failed"); + crate::error_problem(StatusCode::INTERNAL_SERVER_ERROR, &detail) +} + +#[cfg(test)] +mod tests { + use swarm_queue_client::subagent_term::subagent_names; + + /// The list is what `STREAM.INFO` reports for the stream's subjects, by + /// name: a subagent that published once and one that published many rows + /// list the same way, and the order is the name's, not the stream's. + #[test] + fn the_list_is_the_streams_subjects_by_name() { + let reported = [ + ("$SWARM.term.atlas.sub.scout".to_owned(), 12), + ("$SWARM.term.atlas.sub.builder".to_owned(), 1), + ]; + assert_eq!( + subagent_names("atlas", reported.iter().map(|(s, _)| s.as_str())), + ["builder", "scout"] + ); + } +} diff --git a/swarm-nats-auth/src/main.rs b/swarm-nats-auth/src/main.rs index 65673f73..a1607d6c 100644 --- a/swarm-nats-auth/src/main.rs +++ b/swarm-nats-auth/src/main.rs @@ -361,12 +361,18 @@ mod tests { "$SWARM.agent-state.{agent}", "--agent-token-publish-subject", "$KV.agent-icons.{agent}", + "--agent-token-publish-subject", + "$SWARM.term.{agent}.sub.>", + "--agent-token-publish-subject", + "$JS.API.STREAM.CREATE.term-sub-{agent}", + "--agent-token-publish-subject", + "$JS.API.STREAM.INFO.term-sub-{agent}", "--store-cert-role", "swarm-nats-auth", ]) .expect("the unit's own argument vector must parse"); assert_eq!(args.agent_client_suffix, "-agent"); - assert_eq!(args.agent_token_publish_subjects.len(), 3); + assert_eq!(args.agent_token_publish_subjects.len(), 6); } /// The control for the case above: an ordinary value parses through the diff --git a/swarm-nats-auth/src/policy.rs b/swarm-nats-auth/src/policy.rs index 8c45997a..6bbe9716 100644 --- a/swarm-nats-auth/src/policy.rs +++ b/swarm-nats-auth/src/policy.rs @@ -1219,6 +1219,37 @@ mod tests { ); } + /// The subagent templates as `swarm-nats.nix` spells them expand to the + /// subjects and stream name the subagent daemon uses, and to `CREATE` and + /// `INFO` on that one stream only. + #[test] + fn an_agents_subagent_grant_is_its_own_stream_and_nothing_wider() { + use swarm_queue_client::subagent_term::{stream_name, stream_subjects, subject}; + + let p = policy_with_agent_subject() + .with_agent_token_subjects(vec![ + "$SWARM.term.{agent}.sub.>".to_owned(), + "$JS.API.STREAM.CREATE.term-sub-{agent}".to_owned(), + "$JS.API.STREAM.INFO.term-sub-{agent}".to_owned(), + ]) + .expect("per-agent templates are valid"); + let g = p.agent_token_permissions("atlas").expect("configured"); + assert_eq!( + g.publish, + vec![ + stream_subjects("atlas"), + format!("$JS.API.STREAM.CREATE.{}", stream_name("atlas")), + format!("$JS.API.STREAM.INFO.{}", stream_name("atlas")), + ] + ); + assert!(subject("atlas", "scout").starts_with(g.publish[0].trim_end_matches('>'))); + let argus = p.agent_token_permissions("argus").expect("configured"); + assert!( + g.publish.iter().all(|s| !argus.publish.contains(s)), + "atlas and argus share a subject: {g:?} {argus:?}" + ); + } + #[test] fn an_agent_token_subject_without_the_placeholder_is_refused() { let err = policy() diff --git a/swarm-queue-client/Cargo.toml b/swarm-queue-client/Cargo.toml index b4dd258a..aa62f848 100644 --- a/swarm-queue-client/Cargo.toml +++ b/swarm-queue-client/Cargo.toml @@ -30,6 +30,9 @@ kv = ["async-nats/kv"] # `cargo check -p swarm-nats-auth` — no `kv` anywhere in that build — # surfaced it as `cannot find jetstream in async_nats`). notices = ["async-nats/jetstream"] +# `subagent_term::open_or_create`, for the one publisher that creates its own +# stream. Explicit for the same reason as `notices`. +subagent-term = ["async-nats/jetstream"] [dependencies] # Bare (no `kv`/`jetstream`) unless a consumer opts into the `kv` feature diff --git a/swarm-queue-client/src/lib.rs b/swarm-queue-client/src/lib.rs index 4ad9d858..2316a6a9 100644 --- a/swarm-queue-client/src/lib.rs +++ b/swarm-queue-client/src/lib.rs @@ -125,6 +125,15 @@ pub enum Error { #[source] source: async_nats::jetstream::context::CreateStreamError, }, + + #[cfg(feature = "subagent-term")] + #[error("creating the {stream} stream")] + CreateAgentStream { + // Owned: the name carries the agent's (`subagent_term::stream_name`). + stream: String, + #[source] + source: async_nats::jetstream::context::CreateStreamError, + }, } /// Render an error and its source chain on one line. @@ -248,6 +257,11 @@ pub struct DeployRequest { #[cfg(feature = "notices")] pub mod notices; +/// The subject an agent's subagents' terminal rows go to, and the per-agent +/// stream that lists them. Only opening the stream needs the `subagent-term` +/// feature; the names are unconditional, like [`agent_icon`]'s. +pub mod subagent_term; + /// Only the fields this needs; authelia returns several. #[derive(serde::Deserialize)] struct TokenResponse { diff --git a/swarm-queue-client/src/subagent_term.rs b/swarm-queue-client/src/subagent_term.rs new file mode 100644 index 00000000..69d25915 --- /dev/null +++ b/swarm-queue-client/src/subagent_term.rs @@ -0,0 +1,161 @@ +//! Subagent terminals: the subject each subagent's rows go to, and the +//! stream that keeps them. +//! +//! An agent's subagent daemon publishes each subagent's classified terminal +//! rows on `$SWARM.term..sub.`, under the agent's own queue +//! credential. The queue grants that credential this subject family and +//! `CREATE`/`INFO` on the one stream [`crate::subagent_term::stream_name`] +//! names, nothing else of `JetStream`, so the agent creates its own stream on +//! first use. +//! +//! The stream exists so the swarm can list an agent's subagents: subagents +//! are spawned on demand under caller-chosen names, and the set of subjects +//! the stream holds is the list. A reader takes it from the stream's subject +//! counts ([`crate::subagent_term::subagent_names`]) and never creates the +//! stream. +//! +//! The name and the subject live here because three crates must agree on +//! them: the publisher in `hive-subagent-mcp`, the reader in +//! `swarm-controller`, and the queue grant in `swarm-nats.nix`, whose +//! `{agent}` templates spell the same strings. + +/// The subject family every agent terminal shares. +const TERM_PREFIX: &str = "$SWARM.term"; + +/// The token between the agent and the subagent name. A subagent subject has +/// two more tokens than the agent's own `$SWARM.term.`, so neither +/// family's subscriber receives the other's rows. +const SUBAGENT_TOKEN: &str = "sub"; + +/// How long the stream keeps a row. The subagent list is the subjects with a +/// row inside this window, so a subagent drops off the list this long after +/// its last output. +pub const MAX_AGE: std::time::Duration = std::time::Duration::from_hours(24); + +/// The stream holding `agent`'s subagent rows: `term-sub-`. +#[must_use] +pub fn stream_name(agent: &str) -> String { + format!("term-sub-{agent}") +} + +/// The subject `subagent`'s rows are published on. +#[must_use] +pub fn subject(agent: &str, subagent: &str) -> String { + format!("{TERM_PREFIX}.{agent}.{SUBAGENT_TOKEN}.{subagent}") +} + +/// Every subject [`stream_name`]'s stream captures, which is also the publish +/// grant's shape. +#[must_use] +pub fn stream_subjects(agent: &str) -> String { + format!("{TERM_PREFIX}.{agent}.{SUBAGENT_TOKEN}.>") +} + +/// The subagent a stream subject names, or `None` for a subject that is not +/// exactly one subagent token under `agent`'s family. +#[must_use] +pub fn subagent_from_subject<'a>(agent: &str, subject: &'a str) -> Option<&'a str> { + let name = subject + .strip_prefix(TERM_PREFIX)? + .strip_prefix('.')? + .strip_prefix(agent)? + .strip_prefix('.')? + .strip_prefix(SUBAGENT_TOKEN)? + .strip_prefix('.')?; + (!name.is_empty() && !name.contains('.')).then_some(name) +} + +/// The subagents named by a stream's subjects, sorted and deduplicated. +/// Subjects outside `agent`'s family are skipped. +#[must_use] +pub fn subagent_names<'a>(agent: &str, subjects: impl IntoIterator) -> Vec { + let mut names: Vec = subjects + .into_iter() + .filter_map(|s| subagent_from_subject(agent, s)) + .map(str::to_owned) + .collect(); + names.sort_unstable(); + names.dedup(); + names +} + +/// Open `agent`'s subagent stream, creating it if it does not exist yet. +/// +/// Called by the publisher only. The reader opens the stream with `get_stream` +/// and reads a missing one as "no subagents". +#[cfg(feature = "subagent-term")] +pub async fn open_or_create( + client: &async_nats::Client, + agent: &str, +) -> Result { + let js = async_nats::jetstream::new(client.clone()); + let name = stream_name(agent); + match js.get_stream(&name).await { + 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, + }) + } + } +} + +#[cfg(test)] +mod tests { + use super::*; + + /// The names the grant in `swarm-nats.nix` spells with `{agent}`. + #[test] + fn names_match_the_queue_grant() { + assert_eq!(stream_name("iris"), "term-sub-iris"); + assert_eq!(subject("iris", "scout"), "$SWARM.term.iris.sub.scout"); + assert_eq!(stream_subjects("iris"), "$SWARM.term.iris.sub.>"); + } + + #[test] + fn a_subject_reads_back_as_its_subagent() { + assert_eq!( + subagent_from_subject("iris", &subject("iris", "scout")), + Some("scout") + ); + } + + /// The agent's own terminal, another agent's subagent, and anything + /// deeper than one token are not subagents of `iris`. + #[test] + fn other_subjects_name_no_subagent() { + for s in [ + "$SWARM.term.iris", + "$SWARM.term.h1.iris", + "$SWARM.term.iris.sub", + "$SWARM.term.iris.sub.", + "$SWARM.term.iris.sub.a.b", + "$SWARM.term.iris-2.sub.scout", + "$SWARM.term.argus.sub.scout", + "$SWARM.agent-state.iris.sub.scout", + ] { + assert_eq!(subagent_from_subject("iris", s), None, "{s}"); + } + } + + #[test] + fn the_list_is_sorted_unique_and_scoped_to_the_agent() { + let subjects = [ + "$SWARM.term.iris.sub.zed", + "$SWARM.term.iris.sub.alpha", + "$SWARM.term.argus.sub.other", + "$SWARM.term.iris.sub.alpha", + ]; + assert_eq!(subagent_names("iris", subjects), ["alpha", "zed"]); + } +}