//! Publishing this agent's turn-state header onto the swarm queue. //! //! The per-agent web UI's `/api/state` already carries what a header bar //! wants — what the turn loop is doing, since when, on which model, how much //! context and cost the last turn spent. None of it reaches the swarm: the //! `agent-status` KV bucket republishes once a minute, which is fine for //! "is this agent alive" and useless for "is it thinking right now". This //! offers the same values upward on their own subject so a swarm-level //! header can render an agent without reaching into its hive. //! //! **Published on transition, not on a timer.** The whole reason this //! exists is that the once-a-minute bucket is too stale, so a second //! periodic publisher would reproduce the problem it is here to fix. This //! subscribes to the event bus and republishes whenever the snapshot it //! builds differs from the one it last sent — see [`run`] for why the //! comparison is on the serialized bytes rather than on an event allow-list. //! //! **A core subject, not `JetStream`**, for the reason //! [`crate::swarm_term`] gives at more length: nothing here is worth the //! durability of the notices stream, and it keeps the agent's grant to a //! plain publish. It does cost this subject the one thing a header would //! like and a terminal would not — a late subscriber sees nothing until the //! next transition, rather than the current value. Accepted, because the //! alternative is a KV bucket written on every turn-state flip, and the //! swarm already has one of those at the resolution it can afford. //! //! **Best-effort in every direction.** The turn loop and the web UI must not //! notice whether the queue exists, so the bus receiver is the only thing //! this task blocks on; a failed publish is a log line and the next //! transition is still attempted. use serde::Serialize; use tokio::sync::broadcast; use swarm_queue_client::wanted::AgentState; use crate::events::{Bus, BusEvent, LiveEvent, TurnState}; use crate::swarm_queue::{Connection, Presented}; use crate::term_msg::iso8601_utc; /// Subject family carrying agent turn-state headers, the swarm-wide /// agreement this publisher holds up its end of. One leaf subject per /// agent, matching `$SWARM.term`'s shape, so a subscriber can follow one /// agent without filtering the swarm's whole header traffic. const SUBJECT_PREFIX: &str = "$SWARM.agent-state"; /// What the queue's client id looks like either side of the hive's own name. /// /// These mirror the auth-callout responder's `--hive-client-prefix` and /// `--agent-client-suffix`, which is where the authoritative pair lives: the /// responder parses the hive back out of the presented client id and grants /// exactly the configured subjects for *that* hive. The agent is not told /// those flags, so it restates their defaults — a deployment that retunes /// either one has to change them here too, and the symptom of not doing so /// is every publish refused rather than a wrong subject being accepted. const CLIENT_ID_PREFIX: &str = "hive-"; const CLIENT_ID_SUFFIX: &str = "-agent"; /// How much room to leave under the announced limit for everything the /// publish adds around the payload — subject, headers, protocol framing. /// Same margin and same reasoning as `swarm_term::HEADROOM`. const HEADROOM: usize = 1024; /// One turn-state header, as a swarm-level reader receives it. /// /// The field names here are the published contract; two of them /// deliberately do **not** match the per-agent web UI's `StateSnapshot`, /// which is the other place these same values are served from: /// /// - `turn_state_since` is an ISO 8601 / RFC 3339 UTC string, where /// `StateSnapshot` sends unix seconds. The sibling subject's /// [`crate::term_msg::TermMsg::ts`] is already spelled that way, and two /// messages on adjacent subjects disagreeing about how to write a time is /// a cost paid by every reader of both. /// - `agent_state` replaces `StateSnapshot`'s `paused: bool` with the /// swarm's own [`AgentState`] — the same enum the *wanted* state is /// declared in, so a reader can compare actual against wanted directly /// instead of translating one into the other's vocabulary first. /// /// `turn_state` and `agent_state` are two axes, not one: [`AgentState`] has /// no notion of `thinking` and [`TurnState`] has no notion of `offline`, so /// a single field would have to drop one of them. #[derive(Debug, Clone, PartialEq, Eq, Serialize)] struct AgentStateMsg { /// What the turn loop is doing right now. turn_state: TurnState, /// When it entered that state, ISO 8601 / RFC 3339 UTC — see the struct /// doc. A reader ticks the elapsed time off this rather than being sent /// an age that is wrong the moment it is serialized. turn_state_since: String, /// What this agent can honestly say it *is* — see [`observed_state`], /// which documents why only two of the four variants can ever appear /// here. agent_state: AgentState, /// The alias `claude --model` was invoked with (e.g. `"sonnet"`). model: String, /// The concrete model id the most recent completed turn actually ran /// on, resolved from that turn's assistant events. `None` before any /// turn has completed. resolved_model: Option, /// Effective context-window token budget for the current model. context_window_tokens: u64, /// Last-inference token usage from the most recent completed turn — the /// current context-window occupancy. `None` until the first turn ends. ctx_usage: Option, /// Cumulative token usage across the most recent turn's inferences, the /// cost signal. `None` until the first turn ends. cost_usage: Option, } /// The hive a queue client id names. /// /// The hive in the subject has to be the string the **responder** parses out /// of this same client id, because that is what its grant is built from. The /// harness knows a hive display name too, from a different source and with /// no rule tying the two together — deriving the subject from that one /// instead produces a publish the broker refuses, which surfaces as a header /// that is simply never populated. /// /// `None` for an id that does not have the expected shape: a publisher that /// guessed at a subject would be asking for a grant it cannot have, so the /// caller disables itself instead. fn hive_from_client_id(client_id: &str) -> Option<&str> { let hive = client_id .strip_prefix(CLIENT_ID_PREFIX)? .strip_suffix(CLIENT_ID_SUFFIX)?; (!hive.is_empty()).then_some(hive) } /// What this process can honestly report about the agent-state axis. /// /// Only two of [`AgentState`]'s four variants are knowable from inside the /// container, and the split is not arbitrary: /// /// - [`AgentState::Up`] and [`AgentState::Paused`] are the agent's own /// facts. Running is self-evident — this code is executing — and the /// pause marker is a file in the harness dir that this process reads on /// every turn-loop iteration anyway (`crate::paths::paused_marker`). /// - [`AgentState::Offline`] and [`AgentState::Destroyed`] are **not /// reportable from here, ever**. A stopped agent publishes nothing and a /// destroyed one does not exist, so a message claiming either would be a /// message from a process contradicting itself. Those two remain /// hive-c0re's to observe and the `agent-status` bucket's to carry; a /// reader wanting the full four-state picture needs both sources, and /// this one's silence is the only evidence it can offer for the other /// two. /// /// Read fresh per publish rather than cached: the marker is written by /// hive-c0re (through hive-priv) and by this process's own auto-pause, so /// nothing in here observes every write. fn observed_state() -> AgentState { if crate::paths::paused_marker().exists() { AgentState::Paused } else { AgentState::Up } } /// Build the current header from the bus plus the pause marker. fn snapshot(bus: &Bus) -> AgentStateMsg { let (turn_state, since_unix) = bus.state_snapshot(); let model = bus.model(); AgentStateMsg { turn_state, turn_state_since: iso8601_utc(since_unix), agent_state: observed_state(), context_window_tokens: bus.effective_context_window(&model), model, resolved_model: bus.last_resolved_model(), ctx_usage: bus.last_ctx_usage(), cost_usage: bus.last_cost_usage(), } } /// The subject this agent's header goes to under the credential it connected /// with, or `None` when a hive client id names no hive. fn subject(presented: &Presented, agent: &str) -> Option { match presented { Presented::Agent => Some(format!("{SUBJECT_PREFIX}.{agent}")), Presented::Hive { client_id } => { let Some(hive) = hive_from_client_id(client_id) else { tracing::warn!( %client_id, expected = format!("{CLIENT_ID_PREFIX}{CLIENT_ID_SUFFIX}"), "queue client id does not name a hive; not publishing turn state upward" ); return None; }; Some(format!("{SUBJECT_PREFIX}.{hive}.{agent}")) } } } /// Start the publish task, if this agent has a queue credential and a label /// to name its subject with. /// /// Returns without spawning in every other case — no queue or no label — each /// of which is a legal state for an agent rather than an error, and each /// logged once here rather than per transition. pub fn spawn(bus: &Bus) { if !crate::swarm_queue::configured() { // `swarm_queue::init` already said why at boot; repeating it here // would be the same fact logged twice. return; } let agent = crate::identity::label(); if agent.is_empty() { tracing::warn!("this agent has no label; not publishing turn state upward"); return; } tokio::spawn(run(bus.subscribe(), bus.clone(), agent)); } /// Watch the bus and publish whenever the header actually changed. /// /// **Why the change test is on the serialized bytes rather than on which /// event arrived.** The header is assembled from six independent pieces of /// `Bus` state, only some of which announce themselves with an event of /// their own — `set_resolved_model`, for one, emits nothing and lands /// alongside the `TokenUsageChanged` of the same turn. An allow-list of /// interesting variants would therefore have to encode which event happens /// to be adjacent to each field's write, and would go quietly stale the /// first time a field moved. Rebuilding on (almost) any event and comparing /// the result is the same publish traffic with none of that coupling: a /// rebuild is a handful of mutex reads and one `stat`. /// /// `Stream` is the exception, skipped before the rebuild: it is the only /// high-rate variant — one per `stream-json` line, so hundreds per turn — /// and it carries nothing this header reads. Every other variant, including /// any added later, funnels into the comparison and costs nothing when it /// changes nothing. async fn run(mut rx: broadcast::Receiver, bus: Bus, agent: String) { let Some(Connection { client, presented }) = crate::swarm_queue::client().await else { return; }; let Some(subject) = subject(&presented, &agent) else { return; }; tracing::info!(subject, "publishing agent turn state to the swarm queue"); // The last payload actually sent, so a rebuild that changed nothing is // dropped here instead of on the wire. `None` until the first publish, // which is what makes the boot-time header go out at all: an agent that // sat idle from boot would otherwise never announce itself. let mut last: Option> = None; loop { match rx.recv().await { // Skipped before the rebuild — see this function's doc. Ok(BusEvent { event: LiveEvent::Stream(_), .. }) => continue, Ok(_) => {} // The bus drops events for a subscriber that falls behind. A // header is a current value rather than a log, so a missed // event costs nothing here: the rebuild below reads the state // that the missed events led to, not the events themselves. Err(broadcast::error::RecvError::Lagged(missed)) => { tracing::debug!(missed, "swarm agent state: lagged, rebuilding anyway"); } // Every sender is gone, so the harness is shutting down. Err(broadcast::error::RecvError::Closed) => return, } publish(&client, &subject, &snapshot(&bus), &mut last).await; } } /// Offer one header, if it differs from the last one sent. Every failure is /// terminal for that publish and for nothing else. /// /// `last` is only updated on a publish the client accepted, so a transition /// lost to a transport failure is re-sent by the next one rather than /// deduped away against a value the swarm never saw. async fn publish( client: &async_nats::Client, subject: &str, msg: &AgentStateMsg, last: &mut Option>, ) { let payload = match serde_json::to_vec(msg) { Ok(payload) => payload, Err(e) => { tracing::warn!(error = %e, "swarm agent state: serialising failed, header dropped"); return; } }; if last.as_ref() == Some(&payload) { return; } // An unconnected client does not fail a publish, it buffers it — and the // limit read below is the library's pre-connect default until the server // has announced its own, which is smaller than any deployment sets. if let Err(e) = swarm_queue_client::ensure_connected(client) { tracing::warn!(error = %swarm_queue_client::chain(&e), "swarm agent state: publish skipped"); return; } // No degrade path, unlike `swarm_term`'s: every field here is a number, // an enum, or a model name, so there is no arbitrary-length body worth // spending to get under the limit. The check stays because being one // byte over does not truncate the message, it gets it refused and the // connection closed — which would cost every header racing behind it // through the reconnect, not just this one. let limit = swarm_queue_client::max_payload(client).saturating_sub(HEADROOM); if payload.len() > limit { tracing::warn!( len = payload.len(), limit, "swarm agent state: header over the payload limit, dropped" ); return; } if let Err(e) = client .publish(subject.to_owned(), payload.clone().into()) .await { tracing::warn!(error = %e, "swarm agent state: publish failed, header dropped"); return; } *last = Some(payload); } #[cfg(test)] mod tests { use super::{AgentStateMsg, Presented, hive_from_client_id, subject}; use crate::events::TurnState; use swarm_queue_client::wanted::AgentState; /// `TokenUsage` is `#[non_exhaustive]`, so a downstream crate cannot /// write one as a struct expression at all — see `turn_stats`'s note on /// the same restriction. Deserializing one is the shortest way to a /// populated fixture here, and it keeps the numbers next to the field /// names they belong to. /// /// ⚠️ All four fields, every time: none of them carries a serde default, /// so a block written with only the field a given test reads fails to /// deserialize at all rather than zero-filling the rest. fn usage(json: &serde_json::Value) -> hive_claude::TokenUsage { serde_json::from_value(json.clone()).expect("a usage block") } /// A header with every optional field populated, so a test asserting on /// the serialized shape sees the widest form a reader can receive. fn msg() -> AgentStateMsg { AgentStateMsg { turn_state: TurnState::Thinking, // As `iso8601_utc` renders a transition time. turn_state_since: "2026-09-13T12:35:03Z".to_owned(), agent_state: AgentState::Paused, model: "sonnet".to_owned(), resolved_model: Some("claude-sonnet-4-5-20260805".to_owned()), context_window_tokens: 200_000, ctx_usage: Some(usage(&serde_json::json!({ "input_tokens": 11, "output_tokens": 22, "cache_read_input_tokens": 33, "cache_creation_input_tokens": 44, }))), // Distinct from `ctx_usage`'s numbers in every field, so a // renderer served one block where it asked for the other shows // up as a wrong value rather than as a coincidence. cost_usage: Some(usage(&serde_json::json!({ "input_tokens": 55, "output_tokens": 66, "cache_read_input_tokens": 77, "cache_creation_input_tokens": 88, }))), } } #[test] fn a_client_id_names_the_hive_between_the_prefix_and_the_suffix() { assert_eq!(hive_from_client_id("hive-alpha-agent"), Some("alpha")); // A hive whose own name contains the suffix still resolves: only the // trailing one is stripped, so the responder and this agree. assert_eq!( hive_from_client_id("hive-alpha-agent-agent"), Some("alpha-agent") ); } /// Each way an id can fail to name a hive. A guessed subject would be /// refused by the broker, and a refusal reaches an operator as a header /// that never fills in rather than as an error, so none of these may /// fall back. #[test] fn an_id_of_another_shape_names_no_hive() { for id in [ // The hive's own id, not an agent's. "hive-alpha", // Missing the prefix the responder keys on. "alpha-agent", // Prefix and suffix but nothing between them: the subject would // be `$SWARM.agent-state..`, whose empty token matches no // grant. "hive--agent", // Prefix and suffix overlapping with no hive at all. "hive-agent", "", ] { assert_eq!(hive_from_client_id(id), None, "{id} must name no hive"); } } /// The subject the swarm side subscribes to, built end to end from the /// two things that identify this publisher. `swarm-controller`'s relay /// spells the same prefix in its own constant and the two never meet, so /// this pins the half that lives here. #[test] fn the_subject_is_the_prefix_then_the_hive_then_the_agent() { let presented = Presented::Hive { client_id: "hive-alpha-agent".to_owned(), }; assert_eq!( subject(&presented, "mara").as_deref(), Some("$SWARM.agent-state.alpha.mara") ); } /// An agent that connected with its own credential publishes on the /// subject that credential is granted, which names no hive. #[test] fn an_agent_on_its_own_credential_publishes_on_its_hive_free_subject() { assert_eq!( subject(&Presented::Agent, "mara").as_deref(), Some("$SWARM.agent-state.mara") ); } /// The published contract, asserted on the JSON a subscriber parses /// rather than on the Rust struct: a renderer is written against these /// key names, and renaming a field in Rust without noticing is exactly /// the failure this catches. #[test] fn the_payload_carries_the_contracts_field_names() { let value: serde_json::Value = serde_json::from_slice(&serde_json::to_vec(&msg()).expect("serialises")) .expect("valid json"); let obj = value.as_object().expect("an object"); let mut keys: Vec<&str> = obj.keys().map(String::as_str).collect(); keys.sort_unstable(); assert_eq!( keys, [ "agent_state", "context_window_tokens", "cost_usage", "ctx_usage", "model", "resolved_model", "turn_state", "turn_state_since", ] ); } /// Both enums cross the wire as the snake-case strings their `serde` /// attributes promise, not as Rust variant names — and the two stay /// separate fields, since neither vocabulary contains the other's values. #[test] fn the_two_state_axes_serialise_as_their_own_snake_case_strings() { let value = serde_json::to_value(msg()).expect("serialises"); assert_eq!(value["turn_state"], serde_json::json!("thinking")); assert_eq!(value["agent_state"], serde_json::json!("paused")); } /// The departure from `StateSnapshot` that a reader is most likely to /// get wrong: this field is a string, and a number here would parse as a /// valid — and completely wrong — date on the other side. #[test] fn turn_state_since_is_an_iso_8601_string_not_unix_seconds() { let value = serde_json::to_value(msg()).expect("serialises"); assert_eq!( value["turn_state_since"], serde_json::json!("2026-09-13T12:35:03Z") ); } /// The token blocks ship whole rather than pre-summed: a header renderer /// wanting the context percentage needs the same three fields /// `TokenUsage::context_tokens` adds up, and one of them alone reads as /// near-zero once prompt caching is on. #[test] fn the_usage_blocks_ship_their_own_fields() { let value = serde_json::to_value(msg()).expect("serialises"); assert_eq!(value["ctx_usage"]["cache_read_input_tokens"], 33); assert_eq!(value["cost_usage"]["input_tokens"], 55); assert_eq!(value["context_window_tokens"], 200_000); } /// A pre-first-turn header, which is what a freshly booted agent /// actually publishes. The nullable fields have to be present and null /// rather than absent — a renderer that reads `ctx_usage` off a header /// missing the key gets `undefined`, which is not the same thing as /// "this agent has not finished a turn yet". #[test] fn a_header_from_before_the_first_turn_still_carries_every_key() { let fresh = AgentStateMsg { turn_state: TurnState::Idle, agent_state: AgentState::Up, resolved_model: None, ctx_usage: None, cost_usage: None, ..msg() }; let value = serde_json::to_value(fresh).expect("serialises"); assert_eq!(value["resolved_model"], serde_json::Value::Null); assert_eq!(value["ctx_usage"], serde_json::Value::Null); assert_eq!(value["cost_usage"], serde_json::Value::Null); assert_eq!(value["turn_state"], serde_json::json!("idle")); assert_eq!(value["agent_state"], serde_json::json!("up")); } }