hyperhive/hive-c0re/src/dashboard/state_snapshot.rs
atlas 7003d14d2c deploy: split the homeserver's host decisions out of swarm.matrix
`swarm.*` is what a hive needs to be a *client* of the swarm. For the
homeserver that is what it IS from anywhere: its package, the name it
answers to, the ports and URLs it is reached on, and the client id it is
registered under. Whether it is exposed, which peers it trusts, how large
a request it accepts and where its host-local secrets sit are decisions
of the machine running it, so openFirewall, trustedServers,
maxRequestSize, registrationTokenFile, gui.enable and
sso.clientSecretFile move to `deploy.matrix.*`.

Two sub-blocks split rather than moving whole, on their own evidence.
`gui.enable` is whether THIS host serves the web client; `gui.package` is
which client, an artifact identity, and stays. `sso.clientSecretFile` is a
path on one host; `clientId` must match the id in authelia's register, so
it is swarm-wide. Each half now points at the other, because the rendered
docs put them on separate pages.

hive-gateway passed the whole `swarm.matrix` attrset into vhosts.nix, so
that file read a moving option through an argument with no option path
anywhere in it. It now takes `matrixDeployCfg` beside `matrixCfg` — the
only shape that carries a split namespace across that boundary.

While there: vhosts.nix read `matrixCfg.enable`, which has been a rename
alias for `deploy.matrix.enable` since the enable moved. Reading it made
the module system print `Obsolete option services.hyperhive.swarm.matrix.
enable is used` on EVERY evaluation of every host — a deprecation warning
no operator could silence, because the config tripping it was ours. That
shim lives in hive-matrix.nix rather than in this file's table, which is
why deploy.nix's header claim to be their single home is now qualified
in the new block's comment.

glue-matrix-bao-token.nix read the registration token through its own
`matrixCfg` alias; with that read repointed, the binding had no reader
left, so it goes, and the comment naming it is reworded.

module-eval gains a case configuring a hive through all six OLD paths and
asserting two rendered effects — the host firewall's port list and the
container's bind-mount table — because the new paths evaluate fine
without the shims. `gui.enable` is set to the opposite of its default so
the definition has to land rather than agreeing with it by accident.
2026-09-07 14:24:52 +02:00

901 lines
37 KiB
Rust

//! `/api/state` cold-load snapshot plus the dashboard's live read side:
//! the `StateSnapshot` shape and its view builders, the
//! `/api/dashboard/stream` SSE channel, and the `/api/dashboard/history`
//! backfill. SPA shape + SSE channels: docs/web-ui/shape.md.
use std::convert::Infallible;
use axum::{
extract::State,
http::HeaderMap,
response::{
IntoResponse, Response,
sse::{Event, KeepAlive, Sse},
},
};
use chrono::{DateTime, Utc};
use hive_sh4re::approvals::Approval;
use serde::{Deserialize, Serialize};
use tokio_stream::wrappers::BroadcastStream;
use tokio_stream::{Stream, StreamExt};
use utoipa::IntoParams;
use crate::container_view::ContainerView;
use super::meta_inputs::{MetaInputView, read_meta_inputs};
use super::tombstones::{TombstoneView, build_tombstone_views};
use super::{AppState, approvals, error_response, scan_validated_paths};
#[allow(clippy::struct_excessive_bools)]
#[derive(Serialize)]
pub(super) struct StateSnapshot {
/// Broker seq at the moment this snapshot was assembled. Clients
/// dedupe their buffered SSE traffic against this value: any
/// `MessageEvent` with `seq <= snapshot.seq` is already reflected in
/// the snapshot (or pre-dates it); anything with `seq > snapshot.seq`
/// is post-snapshot and should be applied. Set to 0 in the
/// pre-emit case (no events ever fired) — clients treat that as
/// "apply everything you've buffered".
seq: u64,
hostname: String,
/// Current rev of the pinned `hyperhive` flake input, resolved the
/// same way `get_agent_meta`'s per-agent `hyperhive_rev` is (see
/// `auto_update::current_flake_rev`). `None` when the flake ref
/// isn't a local path pin (e.g. a bare `github:` url) or the rev
/// can't be resolved. Feeds the dashboard start page so an operator
/// can tell what build a hive is running without shelling in.
hyperhive_rev: Option<String>,
any_stale: bool,
containers: Vec<ContainerView>,
transients: Vec<TransientView>,
approvals: Vec<ApprovalView>,
/// Last 30 resolved approvals (approved / denied / failed), newest-
/// first. Drives the "history" tab on the approvals section.
approval_history: Vec<ApprovalHistoryView>,
/// State dirs (config history + claude creds + /state/ notes) that
/// survive after a destroy-without-purge. The operator can re-spawn
/// with the same name to resume, or PURG3 to wipe them.
tombstones: Vec<TombstoneView>,
/// Agents whose FNV-1a hashed web UI port collides with at least
/// one other. Operator resolves by renaming. The dashboard renders
/// a banner at the top listing each cluster.
port_conflicts: Vec<PortConflict>,
/// Inputs in `meta/flake.lock` the operator can selectively
/// `nix flake update`. Hyperhive first, then `agent-<n>` rows.
meta_inputs: Vec<MetaInputView>,
/// True while a dashboard-triggered `meta-update` (flake lock bump +
/// agent rebuild ripple) is running in the background. Lets a
/// client that cold-loads mid-update render the META INPUTS panel's
/// disabled "updating…" state; live transitions arrive via the
/// `MetaUpdateRunning` event.
meta_update_running: bool,
/// Whether the hive-forge container is up. When true the dashboard
/// links each container's config + each approval's commit into the
/// forge's `agent-configs` repos.
forge_present: bool,
/// Whether the matrix GUI is reachable at `/matrix/`. Sourced from
/// `HIVE_MATRIX_GUI_ENABLED` env var (set by the c0re NixOS module
/// when `services.hyperhive.deploy.matrix.gui.enable` is on). The gateway
/// (`nix/host-modules/hive-gateway/vhosts.nix`) does the actual `/matrix/` static serving;
/// this flag is just an availability signal for iris's dashboard
/// chrome so the `M4TR1X →` tab doesn't flash when the GUI is off.
matrix_gui_enabled: bool,
/// Whether `hive-gateway` is in front of this dashboard. Sourced
/// from the `HIVE_GATEWAY_ENABLED` env var, which the c0re NixOS
/// module now always sets (the gateway runs unconditionally
/// alongside hyperhive), so this is effectively always true: the
/// dashboard frontend builds same-origin `/agent/<name>/` links to
/// the per-agent web UI (the gateway routes them via the
/// runtime-generated `agents.conf` include file — see
/// `gateway_nginx.rs`). The `false` branch (direct
/// `http://<hostname>:<port>/` TCP links) is retained as a defensive
/// fallback for the env being unset. See `docs/networking/gateway.md::Vhost map`.
gateway_enabled: bool,
/// Public URL of the forge (e.g.
/// `"https://forge.pr1ma.darkest.space"`). Sourced from the
/// `HIVE_FORGE_PUBLIC_URL` env var, which the c0re NixOS module
/// sets from `services.hyperhive.swarm.forge.publicUrl` (defaults to
/// the gateway vhost URL when `deploy.forgejo.behindGateway = true`,
/// `null` otherwise). `None` when absent — the frontend **hides**
/// forge links rather than guessing `http://<hostname>:3000`,
/// which is only right by accident on deployments that aren't
/// plain localhost.
forge_public_url: Option<String>,
/// Human name of this single-host hive instance (e.g. `"pr1ma"`).
/// Sourced from `HYPERHIVE_HIVE_NAME` env var, set by the c0re
/// NixOS module from `services.hyperhive.hiveName`. `None` when
/// the option is unset — chrome falls back to `hostname`.
hive_name: Option<String>,
/// Human name of the wider swarm this hive belongs to (e.g.
/// `"constellat1on"`). Sourced from `HYPERHIVE_SWARM_NAME` env
/// var, set from `services.hyperhive.swarm.name`. `None` when
/// unset — chrome omits the swarm segment of the breadcrumb.
swarm_name: Option<String>,
/// Server-level warnings for the dashboard's top-of-page banner
/// (currently host disk-pressure; more producers can be added
/// backend-side). Empty when all clear. Built by
/// `host_stats::server_warnings`; the frontend renders this list
/// generically, so new warning kinds need no frontend change.
server_warnings: Vec<crate::host_stats::ServerWarning>,
}
#[derive(Serialize)]
struct PortConflict {
port: u16,
/// All agent names sharing this port (sorted, ≥2 entries).
agents: Vec<String>,
}
#[derive(Serialize)]
struct TransientView {
name: String,
/// Owned: the label is the running node's wire tag, not one of a fixed set.
kind: String,
secs: u64,
}
#[derive(Serialize)]
struct ApprovalHistoryView {
id: i64,
agent: String,
kind: &'static str,
/// First 12 chars of the canonical sha (preferred) or
/// manager-supplied ref. None for resolved spawn approvals.
sha_short: Option<String>,
/// `approved` / `denied` / `failed`.
status: &'static str,
/// RFC 3339 UTC. Renders as a relative time on the dashboard.
resolved_at: DateTime<Utc>,
/// Operator-supplied deny reason (for `denied`) or build error
/// (for `failed`). None on `approved`.
#[serde(skip_serializing_if = "Option::is_none")]
note: Option<String>,
}
#[derive(Serialize)]
struct ApprovalView {
id: i64,
agent: String,
kind: &'static str,
/// First 12 chars of the reviewed PR head sha, for `MergeConfigPr`
/// only. Display-only (the short chip on the card).
sha_short: Option<String>,
/// Manager-supplied description shown on the approval card.
#[serde(skip_serializing_if = "Option::is_none")]
description: Option<String>,
/// Forge PR number, for `MergeConfigPr` only. Lets the frontend
/// build a "review PR on forge" link
/// (`{forgeBase}/agent-configs/{agent}/pulls/{pr_number}`). `None`
/// for every other kind.
#[serde(skip_serializing_if = "Option::is_none")]
pr_number: Option<u64>,
/// Raw `commit_ref` payload for `UpdateMetaInputs` (JSON-encoded
/// `Vec<String>` of input names; `"[]"` = all inputs) and
/// `SchedulePrompt` (JSON-encoded `SchedulePromptPayload`). The
/// frontend parses this to render a human-readable card body.
/// `None` for every other kind.
#[serde(skip_serializing_if = "Option::is_none")]
commit_ref: Option<String>,
/// RFC 3339 UTC time the approval was queued. Rendered as a
/// relative time on the card so the operator can spot a stale
/// request.
requested_at: DateTime<Utc>,
}
/// Replace silent `.unwrap_or_default()` on the data sources behind
/// `/api/state` so that whichever query degrades surfaces in journald
/// instead of leaving the operator staring at an empty list. The
/// dashboard still degrades to a sensible default value; the warn
/// is just the diagnostic breadcrumb the old code swallowed.
fn log_default<T, E>(what: &str, result: std::result::Result<T, E>) -> T
where
T: Default,
E: std::fmt::Debug,
{
match result {
Ok(v) => v,
Err(e) => {
tracing::warn!(target: "api_state", source = %what, error = ?e, "snapshot source failed; using default");
T::default()
}
}
}
/// Window over which container crashes count toward the `agents_crashing`
/// banner warning. Wide enough that a crash-looping container (restarted
/// by `Restart=on-failure` every few seconds) keeps the warning lit
/// between flaps, short enough that a single recovered crash clears within
/// minutes.
const CRASH_WARNING_WINDOW: std::time::Duration = std::time::Duration::from_mins(10);
/// Cold-load snapshot of the whole dashboard.
///
/// Includes the roster, approvals (+ history),
/// tombstones, job queue, meta inputs, and more. Live clients then
/// follow `/api/dashboard/stream` (SSE) for incremental updates keyed
/// off `seq`.
// `StateSnapshot` is a large tree of nested view types (`ContainerView`,
// `ApprovalView`, ...) with no `ToSchema` anywhere in that
// graph; wiring it up is a schema-modelling project of its own, well past
// "annotate what's reachable". `serde_json::Value` placeholder for now —
// see the batch report.
#[utoipa::path(
get,
path = "/api/state",
responses((status = 200, description = "full dashboard snapshot", body = serde_json::Value)),
tag = "state_snapshot"
)]
pub(super) async fn api_state(
headers: HeaderMap,
State(state): State<AppState>,
) -> axum::Json<StateSnapshot> {
let host = headers
.get("host")
.and_then(|h| h.to_str().ok())
.unwrap_or("localhost");
let hostname = host.split(':').next().unwrap_or(host).to_owned();
// Capture the unified dashboard-channel seq *before* any read so the
// dedupe contract is "events with seq > snapshot.seq are
// post-snapshot, never missed." An event landing during snapshot
// construction may be doubly applied (snapshot caught the write +
// client also applies the SSE frame) — that's a renderer's problem
// to make idempotent, not ours to avoid here.
let seq = state.coord.current_seq();
// Refresh the coordinator's cached container snapshot before
// reading. Cold-load clients then see whatever the latest rescan
// produced; live clients converge via the matching
// `ContainerStateChanged` / `ContainerRemoved` events the rescan
// emits.
//
// Bound the rescan: it shells out (`nixos-container list` etc.), so a
// saturated/wedged build backend — e.g. hive-c0re mid-startup-sweep
// hammering slow `nixos-container update` subprocesses — can stall it
// long enough that `/api/state` hangs for the whole request (the
// ~minute-long /state reported in the field). On timeout we skip the
// fresh rescan and serve the last cached snapshot instead; live
// clients still converge via the SSE events a later successful rescan
// emits, and the next /state call retries the refresh. Introspection
// stays responsive regardless of the build backend's health.
if tokio::time::timeout(
std::time::Duration::from_secs(3),
state.coord.rescan_containers_and_emit(),
)
.await
.is_err()
{
tracing::warn!(
"api_state: container rescan exceeded 3s (build backend likely saturated); \
serving last cached snapshot"
);
}
let containers = state.coord.containers_snapshot().await;
let any_stale = containers.iter().any(|c| c.needs_update);
let transient_snapshot = state.coord.transient_snapshot();
let pending_approvals = approvals::gc_orphans(
&state.coord,
log_default("approvals.pending", state.coord.approvals.pending()),
);
let transients = build_transient_views(&containers, &transient_snapshot);
let approvals = build_approval_views(pending_approvals);
let approval_history = log_default(
"approvals.recent_resolved",
state.coord.approvals.recent_resolved(30),
)
.into_iter()
.map(history_view)
.collect();
let tombstones = build_tombstone_views(&state.coord, &containers);
let port_conflicts = build_port_conflicts(&containers);
// Banner warnings: host probes (disk) + agent-state (pending logins,
// crashing agents). Built before the response struct because the
// agent-state producer borrows `containers`, which moves in below.
let server_warnings = {
let mut w = crate::host_stats::server_warnings();
w.extend(crate::host_stats::agent_state_warnings(
&containers,
&state.coord.recent_crash_counts(CRASH_WARNING_WINDOW),
));
w
};
axum::Json(StateSnapshot {
seq,
hostname,
hyperhive_rev: crate::auto_update::current_flake_rev(&state.coord.hyperhive_flake),
any_stale,
containers,
transients,
approvals,
approval_history,
meta_inputs: read_meta_inputs(),
meta_update_running: state.coord.meta_update_in_progress(),
tombstones,
port_conflicts,
forge_present: crate::forge::is_present().await,
matrix_gui_enabled: std::env::var_os("HIVE_MATRIX_GUI_ENABLED").is_some_and(|v| {
// Accept any truthy string ("1", "true", "yes") since the
// env var is set by NixOS module wiring with the literal
// "1"; defensive parse so manual overrides also work.
let s = v.to_string_lossy().to_ascii_lowercase();
matches!(s.as_str(), "1" | "true" | "yes")
}),
gateway_enabled: std::env::var_os("HIVE_GATEWAY_ENABLED").is_some_and(|v| {
// Same truthy-string parse as `matrix_gui_enabled`; the
// env var is set by the c0re NixOS module to the literal
// "1" — the gateway always runs alongside hyperhive.
let s = v.to_string_lossy().to_ascii_lowercase();
matches!(s.as_str(), "1" | "true" | "yes")
}),
forge_public_url: std::env::var("HIVE_FORGE_PUBLIC_URL")
.ok()
.filter(|s| !s.is_empty()),
hive_name: std::env::var("HYPERHIVE_HIVE_NAME")
.ok()
.filter(|s| !s.is_empty()),
swarm_name: std::env::var("HYPERHIVE_SWARM_NAME")
.ok()
.filter(|s| !s.is_empty()),
server_warnings,
})
}
/// Group containers by their assigned web UI port; clusters with more
/// than one member are port-hash collisions the operator resolves by
/// renaming an agent.
///
/// Every container's port comes from [`crate::lifecycle::agent_web_port`],
/// which hashes the name into one range for all of them — so the manager
/// is a candidate for a collision exactly like any sub-agent.
fn build_port_conflicts(containers: &[ContainerView]) -> Vec<PortConflict> {
let mut by_port: std::collections::BTreeMap<u16, Vec<String>> =
std::collections::BTreeMap::new();
for c in containers {
by_port.entry(c.port).or_default().push(c.name.clone());
}
by_port
.into_iter()
.filter(|(_, agents)| agents.len() > 1)
.map(|(port, mut agents)| {
agents.sort();
PortConflict { port, agents }
})
.collect()
}
/// Transient state for agents whose container does NOT yet exist
/// (`Spawning`). Lifecycle ops on existing containers surface as
/// `ContainerView.pending` inline; this list only catches pre-creation.
fn build_transient_views(
containers: &[ContainerView],
transient_snapshot: &std::collections::HashMap<String, Vec<crate::job_queue::RunningTransient>>,
) -> Vec<TransientView> {
let mut out: Vec<TransientView> = transient_snapshot
.iter()
.filter(|(name, _)| !containers.iter().any(|c| &c.name == *name))
// One row per running node, so an agent with several shows several
// rather than one of them arbitrarily.
.flat_map(|(name, sts)| sts.iter().map(move |st| (name, st)))
.map(|(name, st)| TransientView {
name: name.clone(),
kind: st.label.clone(),
// Clamped at 0: `since` is wall-clock now (the node's own
// `started_at`), so a backwards clock adjustment could otherwise
// render a negative age.
secs: (chrono::Utc::now() - st.since)
.num_seconds()
.max(0)
.cast_unsigned(),
})
.collect();
// The caller rebuilds this map on every poll, so its iteration order
// differs between renders and the rows would shuffle under an operator
// watching an operation run. Sorted for the same reason
// `build_port_conflicts` sorts its clusters.
out.sort_by(|a, b| {
a.name
.cmp(&b.name)
.then_with(|| a.kind.cmp(&b.kind))
.then_with(|| a.secs.cmp(&b.secs))
});
out
}
/// Render each pending approval into its dashboard view (short sha for
/// `MergeConfigPr`, just the name for `Spawn`).
/// Project a resolved sqlite row into the lean shape the dashboard
/// history tab consumes — no `diff_html` (rendering 30 of them
/// per /api/state poll would mean 30 git diffs per refresh).
fn history_view(a: Approval) -> ApprovalHistoryView {
let displayed = a.fetched_sha.as_deref().unwrap_or(&a.commit_ref);
let sha_short = if displayed.is_empty() {
None
} else {
Some(displayed[..displayed.len().min(12)].to_owned())
};
let status = match a.status {
hive_sh4re::approvals::ApprovalStatus::Approved => "approved",
hive_sh4re::approvals::ApprovalStatus::Denied => "denied",
hive_sh4re::approvals::ApprovalStatus::Failed => "failed",
hive_sh4re::approvals::ApprovalStatus::Cancelled => "cancelled",
// Pending shouldn't appear in recent_resolved, but be defensive.
hive_sh4re::approvals::ApprovalStatus::Pending => "pending",
};
let kind = a.kind.as_str();
ApprovalHistoryView {
id: a.id,
agent: a.agent.to_string(),
kind,
sha_short,
status,
resolved_at: a.resolved_at.unwrap_or_default(),
note: a.note,
}
}
fn build_approval_views(approvals: Vec<Approval>) -> Vec<ApprovalView> {
let mut out = Vec::with_capacity(approvals.len());
for a in approvals {
out.push(match a.kind {
hive_sh4re::approvals::ApprovalKind::Spawn => ApprovalView {
id: a.id,
agent: a.agent.to_string(),
kind: "spawn",
sha_short: None,
description: a.description,
pr_number: None,
commit_ref: None,
requested_at: a.requested_at,
},
hive_sh4re::approvals::ApprovalKind::InitConfig => ApprovalView {
id: a.id,
agent: a.agent.to_string(),
kind: "init_config",
sha_short: None,
description: a.description,
pr_number: None,
commit_ref: None,
requested_at: a.requested_at,
},
hive_sh4re::approvals::ApprovalKind::UpdateMetaInputs => ApprovalView {
id: a.id,
agent: a.agent.to_string(),
kind: "update_meta_inputs",
sha_short: None,
description: a.description,
pr_number: None,
commit_ref: Some(a.commit_ref),
requested_at: a.requested_at,
},
hive_sh4re::approvals::ApprovalKind::SchedulePrompt => ApprovalView {
id: a.id,
agent: a.agent.to_string(),
kind: "schedule_prompt",
sha_short: None,
description: a.description,
pr_number: None,
commit_ref: Some(a.commit_ref),
requested_at: a.requested_at,
},
hive_sh4re::approvals::ApprovalKind::MergeConfigPr => {
// commit_ref = PR number; fetched_sha = the reviewed PR
// head. Show the head sha; the config diff surface lives
// on the forge PR itself.
let sha = a
.fetched_sha
.as_deref()
.map(|s| s[..s.len().min(12)].to_owned());
// Surface the PR number so the frontend can link to the
// PR on the forge. commit_ref holds the number as text.
let pr_number = a.commit_ref.parse::<u64>().ok();
ApprovalView {
id: a.id,
agent: a.agent.to_string(),
kind: "merge_config_pr",
sha_short: sha,
description: a.description,
pr_number,
commit_ref: None,
requested_at: a.requested_at,
}
}
});
}
out
}
/// `/api/jobq/graph` query string — the generic jobq/wire query shape (any
/// host serving a `GraphWire` projection over HTTP takes the same one; see
/// `hive_jobq_wire::parse_states`, which does the actual parsing this side
/// of the query string). Today's only field is `states`: a comma-separated
/// allow-list of `hive_jobq::State` names (`"Pending"`, `"Running"`, ...).
/// Empty / absent ⇒ no filter (current behaviour, every visible root). Set
/// ⇒ only **root** groups whose own state is named are served — a root's
/// state is already its subtree's rolled-up answer (see `hive_jobq_wire`'s
/// doc), so filtering the root filters the whole group. Unknown tokens are
/// silently ignored (an unrecognised name matches nothing rather than
/// erroring the whole request), mirroring `DashboardStreamQuery::kinds`
/// above.
#[derive(Deserialize, Default, IntoParams)]
pub(super) struct JobqGraphQuery {
states: Option<String>,
}
#[utoipa::path(
get,
path = "/api/jobq/graph",
params(JobqGraphQuery),
responses(
(status = 200, description = "every node of every retained job group, \
as generic `hive_jobq` graph nodes: identity, the parent tree, \
dependency edges with their accepted-outcome sets, lifecycle, and \
one opaque per-node payload. This is the generic jobq/wire shape \
(`hive_jobq_wire::GraphWire::wire_snapshot`), not a hive-c0re-only \
projection — a `swarm-controller` serving its own graph responds \
with the same shape at the same path. Group roots ride as \
ordinary nodes (`parent: null`) and `Done` nodes are not \
filtered by default — a consumer renders the graph without \
knowing what any node means. `?states=` narrows to root groups \
in the named states.",
body = Vec<hive_jobq_wire::GraphNode>),
),
tag = "state_snapshot"
)]
pub(super) async fn jobq_graph(
State(state): State<AppState>,
axum::extract::Query(q): axum::extract::Query<JobqGraphQuery>,
) -> axum::Json<Vec<hive_jobq_wire::GraphNode>> {
let states = hive_jobq_wire::parse_states(q.states.as_deref());
axum::Json(state.coord.job_queue.graph_snapshot(states.as_deref()))
}
#[utoipa::path(
get,
path = "/api/jobq/rollup",
responses(
(status = 200, description = "counts by lifecycle state over the same \
groups `/api/jobq/graph` serves, as `(state, nodes, roots)` \
triples — the generic jobq/wire roll-up shape \
(`hive_jobq_wire::state_rollup`), same as `/api/jobq/graph` \
above. Every state is present, zero counts included, in a fixed \
order — a consumer renders a summary (\"3 running · 2 queued\") \
without fetching the graph and without re-deriving the tally. \
`roots` counts groups, `nodes` counts every step at any depth: one \
rebuild is 1 root and ~7 nodes, so a summary meaning *operations* \
reads `roots` and one meaning *steps* reads `nodes`.",
body = Vec<hive_jobq_wire::StateCount>),
),
tag = "state_snapshot"
)]
pub(super) async fn jobq_rollup(
State(state): State<AppState>,
) -> axum::Json<Vec<hive_jobq_wire::StateCount>> {
axum::Json(state.coord.job_queue.state_rollup())
}
/// Response body for `/api/dashboard/history`. No `ToSchema` — its
/// `events` field wraps [`crate::dashboard_events::DashboardEvent`],
/// which doesn't derive `ToSchema` either (a large enum with many
/// variants; see that type's doc comment for why annotating it is
/// out of scope here). The `responses(...)` doc below spells out the
/// shape in prose instead of a `body = ...` reference.
#[derive(Serialize)]
struct DashboardHistoryBody {
seq: u64,
events: Vec<crate::dashboard_events::DashboardEvent>,
}
#[utoipa::path(
get,
path = "/api/dashboard/history",
responses(
(status = 200, description = "`{ seq, events }` — up to the last 200 \
broker messages as `DashboardEvent::Sent`/`Delivered` JSON, plus \
`seq`: the dashboard channel's high-water mark at fetch time \
(used by clients to dedupe against buffered live SSE frames)"),
),
tag = "state_snapshot"
)]
pub(super) async fn dashboard_history(State(state): State<AppState>) -> Response {
// Backfill source for the dashboard terminal. Returns up to ~200
// historical broker messages (no other event kinds are persisted)
// converted to `DashboardEvent::Sent` JSON so the client can replay
// through the same dispatch path as live frames. Wrapped in
// `{ seq, events }`: the seq is the dashboard channel's high-water
// mark at fetch time. Clients use it to dedupe their buffered live
// SSE traffic (drop anything with `seq <= history_seq`) so a frame
// that lands between SSE-subscribe and history-fetch isn't shown
// twice and isn't lost. Historical rows carry `seq = 0`; the
// boundary seq is what closes the dedupe window.
const HISTORY_LIMIT: u64 = 200;
let seq = state.coord.current_seq();
match state.coord.broker.recent_all(HISTORY_LIMIT) {
Ok(mut messages) => {
messages.reverse();
let events: Vec<crate::dashboard_events::DashboardEvent> = messages
.into_iter()
.map(|m| match m {
crate::broker::MessageEvent::Sent {
id,
from,
to,
body,
at,
in_reply_to,
} => {
let file_refs = scan_validated_paths(&body);
crate::dashboard_events::DashboardEvent::Sent {
seq: 0,
id,
from,
to,
body,
at: hive_sh4re::wire_time::from_secs(at),
in_reply_to,
file_refs,
}
}
crate::broker::MessageEvent::Delivered {
id,
from,
to,
body,
at,
in_reply_to,
} => {
let file_refs = scan_validated_paths(&body);
crate::dashboard_events::DashboardEvent::Delivered {
seq: 0,
id,
from,
to,
body,
at: hive_sh4re::wire_time::from_secs(at),
in_reply_to,
file_refs,
}
}
})
.collect();
axum::Json(DashboardHistoryBody { seq, events }).into_response()
}
Err(e) => error_response(&format!("dashboard/history failed: {e:#}")),
}
}
/// `/dashboard/stream` query string. Today's only field is `kinds`:
/// a comma-separated allow-list of event-`kind` strings.
/// Empty / absent ⇒ no filter (current behaviour, all variants
/// forwarded). Set ⇒ only the named kinds reach the subscriber,
/// non-matches are skipped before the JSON serialise cost.
///
/// Useful for narrow pages (e.g. `flow.js` only cares about `sent`
/// / `delivered` / `container_state_changed` / `container_removed`)
/// that want to drop the dispatch overhead on every unrelated mutation.
#[derive(Deserialize, Default, IntoParams)]
pub(super) struct DashboardStreamQuery {
/// Comma-separated event kinds to forward. Each token is
/// trimmed; unknown kinds are silently ignored on lookup
/// (subscriber sees nothing instead of an error).
kinds: Option<String>,
}
#[utoipa::path(
get,
path = "/api/dashboard/stream",
params(DashboardStreamQuery),
responses(
(status = 200, description = "server-sent event stream; each event's \
`data` is a JSON-serialised `DashboardEvent` (seq-tagged; pair \
with `/api/dashboard/history` to backfill + dedupe on connect)",
body = String, content_type = "text/event-stream"),
),
tag = "state_snapshot"
)]
pub(super) async fn dashboard_stream(
State(state): State<AppState>,
axum::extract::Query(q): axum::extract::Query<DashboardStreamQuery>,
) -> Sse<impl Stream<Item = Result<Event, Infallible>>> {
let rx = state.coord.dashboard_subscribe();
// Pre-parse the allow-list once at subscription time, so the
// per-event hot path is just a `HashSet::contains` on a
// `&'static str` — no string churn per frame.
let kind_filter: Option<std::collections::HashSet<String>> = q.kinds.and_then(|raw| {
let set: std::collections::HashSet<String> = raw
.split(',')
.map(str::trim)
.filter(|s| !s.is_empty())
.map(str::to_owned)
.collect();
if set.is_empty() { None } else { Some(set) }
});
let stream = BroadcastStream::new(rx).filter_map(move |res| {
// Drop lagged frames. Browsers reconnect; the seq dedupe on
// reconnect skips any frame already reflected in the snapshot.
let event = res.ok()?;
if let Some(filter) = kind_filter.as_ref()
&& !filter.contains(event.kind_tag())
{
return None;
}
let json = serde_json::to_string(&event).ok()?;
Some(Ok(Event::default().data(json)))
});
Sse::new(stream).keep_alive(KeepAlive::default())
}
#[cfg(test)]
mod tests {
use std::collections::HashMap;
use chrono::{Duration, Utc};
use super::{ContainerView, build_port_conflicts, build_transient_views};
use crate::job_queue::RunningTransient;
fn rt(agent: &str, label: &str, age_secs: i64) -> RunningTransient {
RunningTransient {
agent: agent.to_owned(),
label: label.to_owned(),
takes_container_down: false,
since: Utc::now() - Duration::seconds(age_secs),
}
}
fn snapshot(rows: Vec<RunningTransient>) -> HashMap<String, Vec<RunningTransient>> {
let mut m: HashMap<String, Vec<RunningTransient>> = HashMap::new();
for r in rows {
m.entry(r.agent.clone()).or_default().push(r);
}
m
}
#[test]
fn an_agent_whose_container_exists_is_not_transient() {
let snap = snapshot(vec![rt("alice", "spawn", 5)]);
let live = [cv("alice", 8101, true)];
assert!(build_transient_views(&live, &snap).is_empty());
// Control: with no container by that name it *is* reported, so the
// emptiness above is the filter and not an inert builder.
assert_eq!(build_transient_views(&[], &snap).len(), 1);
}
#[test]
fn an_agent_with_several_running_nodes_gets_a_row_each() {
let snap = snapshot(vec![rt("alice", "spawn", 5), rt("alice", "update", 3)]);
assert_eq!(build_transient_views(&[], &snap).len(), 2);
}
/// The snapshot is a fresh `HashMap` per poll, so its iteration order is
/// not stable across renders. Sorting is what keeps a row still.
#[test]
fn rows_are_ordered_by_name_then_kind() {
let snap = snapshot(vec![
rt("zoe", "spawn", 1),
rt("adam", "update", 1),
rt("adam", "spawn", 1),
]);
let got: Vec<(String, String)> = build_transient_views(&[], &snap)
.into_iter()
.map(|v| (v.name, v.kind))
.collect();
assert_eq!(
got,
[
("adam".to_owned(), "spawn".to_owned()),
("adam".to_owned(), "update".to_owned()),
("zoe".to_owned(), "spawn".to_owned()),
]
);
}
/// A `since` in the future yields a negative duration, and casting that
/// unsigned would render an age near `u64::MAX` rather than a small one.
#[test]
fn a_future_since_clamps_to_zero_instead_of_wrapping() {
let snap = snapshot(vec![rt("alice", "spawn", -600)]);
let got = build_transient_views(&[], &snap);
assert_eq!(got[0].secs, 0, "clamped, not wrapped");
}
#[test]
fn an_empty_snapshot_produces_no_rows() {
assert!(build_transient_views(&[], &snapshot(vec![])).is_empty());
}
fn cv(name: &str, port: u16, running: bool) -> ContainerView {
ContainerView {
name: name.to_owned(),
container: format!("h-{name}"),
port,
running,
failed: false,
needs_update: false,
needs_login: false,
deployed_sha: None,
parent: None,
active_model: None,
status_text: None,
status_set_at: None,
paused: false,
cpu_quota: "200%".to_owned(),
memory_max: "4G".to_owned(),
}
}
#[test]
fn distinct_ports_produce_no_conflicts() {
let cs = [cv("alice", 8101, true), cv("bob", 8102, true)];
assert!(build_port_conflicts(&cs).is_empty());
}
#[test]
fn a_lone_container_does_not_conflict_with_itself() {
assert!(build_port_conflicts(&[cv("alice", 8101, true)]).is_empty());
}
#[test]
fn two_agents_on_one_port_are_reported_together() {
let cs = [cv("alice", 8101, true), cv("bob", 8101, true)];
let got = build_port_conflicts(&cs);
assert_eq!(got.len(), 1);
assert_eq!(got[0].port, 8101);
assert_eq!(got[0].agents, ["alice", "bob"]);
}
/// The manager takes its port from the same hash as everyone else, so
/// it collides like any other name. A doc comment here used to claim it
/// sat at a fixed 8000 and could not — no such special case exists.
#[test]
fn the_manager_can_collide_with_a_sub_agent() {
let cs = [cv("manager", 8101, true), cv("alice", 8101, true)];
let got = build_port_conflicts(&cs);
assert_eq!(got.len(), 1);
assert_eq!(got[0].agents, ["alice", "manager"]);
}
/// A stopped agent keeps its hashed port, so the clash is real the
/// moment it starts. Reporting it while it is down is the point.
#[test]
fn a_stopped_container_still_claims_its_port() {
let cs = [cv("alice", 8101, true), cv("bob", 8101, false)];
assert_eq!(build_port_conflicts(&cs).len(), 1);
}
#[test]
fn names_are_sorted_within_a_conflict() {
let cs = [cv("zoe", 8101, true), cv("adam", 8101, true)];
assert_eq!(build_port_conflicts(&cs)[0].agents, ["adam", "zoe"]);
}
#[test]
fn more_than_two_agents_land_in_one_cluster() {
let cs = [
cv("c", 8101, true),
cv("a", 8101, true),
cv("b", 8101, true),
];
let got = build_port_conflicts(&cs);
assert_eq!(got.len(), 1, "one cluster, not one row per pair");
assert_eq!(got[0].agents, ["a", "b", "c"]);
}
#[test]
fn separate_clusters_are_ordered_by_port() {
let cs = [
cv("d", 8300, true),
cv("c", 8300, true),
cv("b", 8200, true),
cv("a", 8200, true),
];
let ports: Vec<u16> = build_port_conflicts(&cs).iter().map(|p| p.port).collect();
assert_eq!(ports, [8200, 8300]);
}
#[test]
fn no_containers_is_not_a_conflict() {
assert!(build_port_conflicts(&[]).is_empty());
}
}