//! Hyperhive dashboard. Lists managed containers (with deep-links to each //! container's web UI), pending approvals (with unified diff vs the applied //! repo, plus approve/deny buttons), and the manager. use std::convert::Infallible; use std::net::SocketAddr; use std::path::Path; use std::sync::Arc; use anyhow::{Context, Result}; use axum::extract::Form; use axum::{ Router, extract::{Path as AxumPath, State}, http::{HeaderMap, StatusCode}, response::{ IntoResponse, Response, sse::{Event, KeepAlive, Sse}, }, routing::{get, post}, }; use hive_sh4re::Approval; use serde::{Deserialize, Serialize}; use tokio_stream::wrappers::BroadcastStream; use tokio_stream::{Stream, StreamExt}; use crate::container_view::{ContainerView, claude_has_session}; use crate::coordinator::Coordinator; use crate::lifecycle::{self, MANAGER_NAME}; use chrono::{DateTime, Utc}; mod approvals; mod build_logs; mod journal; mod lifecycle_ops; mod matrix_accounts; pub(crate) mod permissions; mod questions; mod reminders; mod schedules; mod state_files; mod topology; mod webhook; // Pre-computed at approval-submit time by the manager-socket handler // (`socket_server.rs`) and embedded in the `ApprovalAdded` event, so // re-exported at the module root to preserve the `crate::dashboard::approval_diff` // path across the submodule split. pub(crate) use approvals::approval_diff; // Run at broker-message ingest by the coordinator + the operator-msg path // (`main.rs`); re-exported to preserve the `crate::dashboard::scan_validated_paths` // path across the split. pub use state_files::scan_validated_paths; #[derive(Clone)] struct AppState { coord: Arc, } #[allow( clippy::too_many_lines, reason = "the body is dominated by the flat axum route table — one line \ per endpoint mapping a URL to its (now per-concern submodule) \ handler; splitting that exhaustive list across helpers would \ obscure the route map for no readability gain" )] pub async fn serve(port: u16, coord: Arc) -> Result<()> { // API-only: the gateway static-serves the dashboard dist and proxies // non-static requests here (see hive-gateway.nix). Unmatched paths 404. let app = Router::new() .route("/api/state", get(api_state)) .route("/api/journal/{name}", get(journal::get_journal)) .route("/api/journal-host", get(journal::get_journal_host)) .route("/api/approval-diff/{id}", get(approvals::get_approval_diff)) .route("/api/state-file", get(state_files::get_state_file)) .route( "/api/matrix-accounts", get(matrix_accounts::get_matrix_accounts), ) .route("/api/reminders", get(reminders::api_reminders)) .route("/api/operator-inbox", get(api_operator_inbox)) .route("/api/stats-hive", get(api_stats_hive)) .route("/api/container-resources", get(api_container_resources)) .route("/api/audit-log", get(api_audit_log)) .route("/api/build-logs", get(build_logs::get_build_logs_all)) .route( "/api/build-logs/{agent}", get(build_logs::get_build_logs_agent), ) .route( "/api/build-logs/id/{id}", get(build_logs::get_build_log_full), ) .route( "/api/build-logs/id/{id}/stream", get(build_logs::get_build_log_stream), ) .route( "/api/build-logs/id/{id}/raw", get(build_logs::get_build_log_raw), ) .route("/api/agent/{name}/mark-all-read", post(post_mark_all_read)) .route("/api/topology/set-parent", post(topology::post_set_parent)) .route( "/api/topology/set-parent-bulk", post(topology::post_set_parent_bulk), ) .route("/api/tool-groups", get(permissions::get_tool_groups)) .route( "/api/tool-groups/{agent}", post(permissions::post_tool_groups), ) .route("/api/capabilities", get(permissions::get_capabilities)) .route( "/api/capabilities/{agent}", post(permissions::post_capabilities), ) .route("/api/permissions", post(permissions::post_permissions)) .route( "/api/permissions/stale", get(permissions::get_stale_permissions), ) .route( "/api/permissions/{agent}", axum::routing::delete(permissions::delete_agent_permissions), ) .route( "/api/schedules", get(schedules::api_schedules).post(schedules::post_schedule_new), ) .route( "/api/schedules/{id}", axum::routing::patch(schedules::patch_schedule), ) .route( "/api/schedules/{id}/cancel", post(schedules::post_schedule_cancel), ) .route( "/api/schedules/{id}/pause", post(schedules::post_schedule_pause), ) .route( "/api/schedules/{id}/resume", post(schedules::post_schedule_resume), ) .route( "/api/schedules/{id}/fire-now", post(schedules::post_schedule_fire_now), ) .route( "/api/rebuild-queue/{id}/cancel", post(schedules::post_rebuild_queue_cancel), ) .route("/webhook/knowledge", post(webhook::post_webhook_knowledge)) // Backend routes — the frontend calls these `/api/` paths. The // transitional bare top-level aliases were removed once the // frontend migrated. `/webhook/knowledge` keeps its own prefix // (forge-driven, not the SPA). .route("/api/approve/{id}", post(approvals::post_approve)) .route("/api/deny/{id}", post(approvals::post_deny)) .route("/api/destroy/{name}", post(lifecycle_ops::post_destroy)) .route("/api/kill/{name}", post(lifecycle_ops::post_kill)) .route("/api/restart/{name}", post(lifecycle_ops::post_restart)) .route("/api/start/{name}", post(lifecycle_ops::post_start)) .route("/api/rebuild/{name}", post(lifecycle_ops::post_rebuild)) .route("/api/update-all", post(lifecycle_ops::post_update_all)) .route( "/api/answer-question/{id}", post(questions::post_answer_question), ) .route( "/api/cancel-question/{id}", post(questions::post_cancel_question), ) .route("/api/purge-tombstone/{name}", post(post_purge_tombstone)) .route( "/api/matrix-account-login", post(matrix_accounts::post_matrix_account_login), ) .route( "/api/cancel-reminder/{id}", post(reminders::post_cancel_reminder), ) .route( "/api/retry-reminder/{id}", post(reminders::post_retry_reminder), ) .route("/api/request-spawn", post(post_request_spawn)) .route("/api/op-send", post(post_op_send)) .route("/api/meta-update", post(post_meta_update)) .route("/api/dashboard/stream", get(dashboard_stream)) .route("/api/dashboard/history", get(dashboard_history)) // No static fallback — the gateway owns the dist; unmatched paths 404. .with_state(AppState { coord }); // Binds loopback-only; external access via gateway. // Rationale: docs/gateway.md::Firewall posture. let addr = SocketAddr::from(([127, 0, 0, 1], port)); let listener = bind_with_retry(addr).await?; tracing::info!(%addr, "dashboard listening"); axum::serve(listener, app).await?; Ok(()) } // SPA shape + SSE channels: docs/web-ui/shape.md. /// `SO_REUSEADDR` bind with retry. Retry mechanics, attempt-cap /// rationale, and log-level cadence: `docs/web-ui/shape.md::Listener bind`. async fn bind_with_retry(addr: SocketAddr) -> Result { let mut delay_ms = 250u64; let mut attempts = 0u32; loop { match try_bind(addr) { Ok(l) => { if attempts > 0 { tracing::info!( %addr, attempts, "dashboard: bind succeeded after retry" ); } return Ok(l); } Err(e) if e.kind() == std::io::ErrorKind::AddrInUse => { let attempt = attempts + 1; if attempt <= 12 { tracing::warn!( %addr, attempt, "dashboard: AddrInUse, retrying in {delay_ms}ms" ); } else { tracing::info!( %addr, attempt, "dashboard: AddrInUse still holding, retrying in {delay_ms}ms" ); } tokio::time::sleep(std::time::Duration::from_millis(delay_ms)).await; attempts += 1; delay_ms = (delay_ms * 2).min(2000); } Err(e) => { return Err(e).with_context(|| format!("bind dashboard on {addr}")); } } } } fn try_bind(addr: SocketAddr) -> std::io::Result { let sock = match addr { SocketAddr::V4(_) => tokio::net::TcpSocket::new_v4()?, SocketAddr::V6(_) => tokio::net::TcpSocket::new_v6()?, }; sock.set_reuseaddr(true)?; sock.bind(addr)?; sock.listen(1024) } #[allow(clippy::struct_excessive_bools)] #[derive(Serialize)] 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, manager_port: u16, any_stale: bool, containers: Vec, transients: Vec, approvals: Vec, /// Last 30 resolved approvals (approved / denied / failed), newest- /// first. Drives the "history" tab on the approvals section. approval_history: Vec, /// Pending operator-targeted questions (`target IS NULL`). Any /// agent can `ask` the operator and `ask` returns immediately with /// the id; on `/answer-question` we mark the row answered and /// fire `HelperEvent::QuestionAnswered` back into the asker's /// inbox. Peer-to-peer questions live in the same table but never /// surface here (see `OperatorQuestions::pending`). questions: Vec, /// Last 20 answered questions, newest-first. question_history: Vec, /// 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, /// Sub-agents whose FNV-1a hashed web UI port collides with at /// least one other agent. Operator resolves by renaming. The /// dashboard renders a banner at the top listing each cluster. port_conflicts: Vec, /// Inputs in `meta/flake.lock` the operator can selectively /// `nix flake update`. Hyperhive first, then `agent-` rows. meta_inputs: Vec, /// 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, /// Current state of the global rebuild queue — pending + running /// long-lived ops (rebuild / meta-update / spawn) plus the most /// recent few terminal entries the queue retains for history. /// Live transitions arrive via the `RebuildQueueChanged` event. /// See `rebuild_queue.rs`. rebuild_queue: Vec, /// 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.matrix.gui.enable` is on). The gateway /// (hive-gateway.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//` 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://:/` TCP links) is retained as a defensive /// fallback for the env being unset. See `docs/gateway.md::Vhost map`. gateway_enabled: bool, /// Public URL of the forge vhost served by hive-gateway (e.g. /// `"https://forge.pr1ma.darkest.space"`). Sourced from the /// `HIVE_FORGE_PUBLIC_URL` env var, which the c0re NixOS module /// sets when `forge.behindGateway = true`. `None` when absent — /// the frontend falls back to `http://:3000`. forge_public_url: Option, /// 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, /// Human name of the wider swarm this hive belongs to (e.g. /// `"constellat1on"`). Sourced from `HYPERHIVE_SWARM_NAME` env /// var, set from `services.hyperhive.swarmName`. `None` when /// unset — chrome omits the swarm segment of the breadcrumb. swarm_name: Option, /// Peer hives in the same swarm. Parsed from `HYPERHIVE_PEERS` /// (JSON array of `{domain,cert_fingerprint}` objects, emitted by /// the c0re NixOS module from `services.hyperhive.swarm.peers`). /// Empty on single-hive deploys. Feeds the P33RS dashboard tab. peer_hives: Vec, /// 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, } /// One peer hive for the P33RS dashboard tab. Derived from /// `HYPERHIVE_PEERS` env; `url` is the peer's HTTPS dashboard root. /// `cert_fingerprint` is `Some("sha256:")` when the peer uses a /// self-signed cert and the operator pinned its fingerprint in /// `services.hyperhive.swarm.peers`. #[derive(Serialize)] struct PeerHiveView { name: String, url: String, cert_fingerprint: Option, } /// `OpQuestion` + computed `question_refs` / `answer_refs`. Built /// from the snapshot read; the live channel attaches the same /// fields directly on `QuestionAdded` / `QuestionResolved`. #[derive(Serialize)] struct QuestionView { #[serde(flatten)] inner: crate::operator_questions::OpQuestion, #[serde(skip_serializing_if = "Vec::is_empty")] question_refs: Vec, #[serde(skip_serializing_if = "Vec::is_empty")] answer_refs: Vec, } impl QuestionView { fn from_question(q: crate::operator_questions::OpQuestion) -> Self { let question_refs = scan_validated_paths(&q.question); let answer_refs = q .answer .as_deref() .map(scan_validated_paths) .unwrap_or_default(); Self { inner: q, question_refs, answer_refs, } } } #[derive(Serialize)] struct PortConflict { port: u16, /// All agent names sharing this port (sorted, ≥2 entries). agents: Vec, } #[derive(Serialize, Clone, Debug)] pub struct TombstoneView { pub name: String, /// Bytes used by the state dir tree. Cheap-ish to compute; let the /// operator know how much they're holding onto. pub state_bytes: u64, /// Mtime (unix seconds) of the state dir; rough "last seen". pub last_seen: i64, pub has_creds: bool, } #[derive(Serialize)] struct TransientView { name: String, kind: &'static str, 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, /// `approved` / `denied` / `failed`. status: &'static str, /// RFC 3339 UTC. Renders as a relative time on the dashboard. resolved_at: DateTime, /// Operator-supplied deny reason (for `denied`) or build error /// (for `failed`). None on `approved`. #[serde(skip_serializing_if = "Option::is_none")] note: Option, } #[derive(Serialize)] struct ApprovalView { id: i64, agent: String, kind: &'static str, /// First 12 chars of the `commit_ref`, for `ApplyCommit` only. sha_short: Option, /// Raw unified diff text, for `ApplyCommit` only. The client splits /// on `\n` and per-line classifies (`+` / `-` / `@@` / `--- ` / `+++ ` /// → diff-add / diff-del / diff-hunk / diff-file). Shipping raw /// instead of pre-rendered HTML saves bytes on the wire (no /// per-line `` markup) and removes the only HTML-escape /// surface from the snapshot. diff: Option, /// Manager-supplied description shown on the approval card. #[serde(skip_serializing_if = "Option::is_none")] description: Option, /// Forge PR number, for `MergeConfigPr` only. Lets the frontend /// build a "review PR on forge" link /// (`{forgeBase}/agent-configs/{agent}/pulls/{pr_number}`) the same /// way it builds the `apply_commit` "commit on forge" link from the /// sha. `None` for every other kind. #[serde(skip_serializing_if = "Option::is_none")] pr_number: Option, /// Raw `commit_ref` payload for `UpdateMetaInputs` (JSON-encoded /// `Vec` 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, /// 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, } /// 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(what: &str, result: std::result::Result) -> 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); async fn api_state(headers: HeaderMap, State(state): State) -> axum::Json { 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).await; 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, &transient_snapshot); let port_conflicts = build_port_conflicts(&containers); // Both operator-targeted and peer threads surface on the dashboard // (the client filters by target). Each row is wrapped in QuestionView // so the snapshot carries the same file_refs the live event variants // attach. let questions: Vec = log_default("questions.pending_all", state.coord.questions.pending_all()) .into_iter() .map(QuestionView::from_question) .collect(); let question_history: Vec = log_default( "questions.recent_answered_all", state.coord.questions.recent_answered_all(20), ) .into_iter() .map(QuestionView::from_question) .collect(); // 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, manager_port: lifecycle::agent_web_port(MANAGER_NAME), any_stale, containers, transients, approvals, approval_history, meta_inputs: read_meta_inputs(), meta_update_running: state.coord.meta_update_in_progress(), questions, question_history, tombstones, port_conflicts, rebuild_queue: state.coord.rebuild_queue.snapshot(), 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()), peer_hives: parse_peer_hives(), server_warnings, }) } /// Parse `HYPERHIVE_PEERS` env var into dashboard-ready `PeerHiveView` /// entries. The env var is a JSON array of `{domain, cert_fingerprint}` /// objects emitted by the c0re NixOS module from /// `services.hyperhive.swarm.peers`. Each entry becomes /// `{ name: domain, url: "https://domain/" }` for the P33RS tab. /// Returns empty vec when unset (single-hive deploy). fn parse_peer_hives() -> Vec { #[derive(serde::Deserialize)] struct Raw { domain: String, cert_fingerprint: Option, } let Ok(json) = std::env::var("HYPERHIVE_PEERS") else { return Vec::new(); }; let Ok(raw): Result, _> = serde_json::from_str(&json) else { tracing::warn!("HYPERHIVE_PEERS is not valid JSON; ignoring"); return Vec::new(); }; raw.into_iter() .map(|r| { let cert_fingerprint = r.cert_fingerprint.and_then(|fp| { if validate_cert_fingerprint(&fp) { Some(fp) } else { tracing::warn!( domain = %r.domain, fingerprint = %fp, "HYPERHIVE_PEERS: invalid cert_fingerprint format \ (expected `sha256:<64 hex chars>`); ignoring fingerprint" ); None } }); PeerHiveView { name: r.domain.clone(), url: format!("https://{}/", r.domain), cert_fingerprint, } }) .collect() } /// Validate a TLS certificate fingerprint string from `HYPERHIVE_PEERS`. /// Accepts `sha256:<64 hex chars>` (upper or lower case). fn validate_cert_fingerprint(fp: &str) -> bool { let Some(hex) = fp.strip_prefix("sha256:") else { return false; }; hex.len() == 64 && hex.chars().all(|c| c.is_ascii_hexdigit()) } /// Group live containers by their assigned web UI port; clusters with /// more than one member are port-hash collisions the operator needs /// to resolve by renaming. Manager (fixed at 8000) and sub-agents /// (8100..8999) can't collide with each other — collisions are /// strictly between sub-agents. fn build_port_conflicts(containers: &[ContainerView]) -> Vec { let mut by_port: std::collections::BTreeMap> = 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() } #[derive(Serialize, Clone, Debug)] pub struct MetaInputView { /// Input key in meta's `flake.nix` — `hyperhive`, `agent-`, etc. pub name: String, /// Full locked sha. Not displayed verbatim; the dashboard /// truncates to the first 12 chars for the chip. pub rev: String, /// Unix seconds — `locked.lastModified`. Drives the relative /// "2h ago" timestamp on each input row. pub last_modified: i64, /// `original.url` if available, for the tooltip / row meta text. #[serde(skip_serializing_if = "Option::is_none")] pub url: Option, } /// Walk `flake.lock`'s `nodes` graph from `root` and emit one /// `MetaInputView` per fetched input, at **every** depth. That /// surfaces the direct meta inputs (`hyperhive`, `agent-`), the /// agent flakes' own inputs (`agent-dmatrix/mcp-matrix`, /// `hyperhive/nixpkgs`), and any deeper transitive inputs — so the /// operator can bump any of them individually. Names are /// slash-separated paths from root, the syntax `nix flake update` /// accepts for transitive inputs. /// /// Filtering: /// - Inputs that resolve via a `follows` chain (lock value is an /// array) are skipped — they alias another node, not their own /// fetched derivation, so updating them does nothing. /// - A node is emitted only when it carries a `locked.rev`. /// - Each fetched node is walked exactly once (a `visited` set): /// the lock graph shares nodes (many flakes reference one /// nixpkgs), so without this a shared subtree re-walks per parent /// and a cycle would recurse forever. The result is a spanning /// tree — every input shown once, at its shallowest path. fn read_meta_inputs() -> Vec { let mut out = Vec::new(); let Ok(raw) = std::fs::read_to_string("/var/lib/hyperhive/meta/flake.lock") else { return out; }; let Ok(json) = serde_json::from_str::(&raw) else { return out; }; let Some(nodes) = json.get("nodes").and_then(|v| v.as_object()) else { return out; }; let Some(root_name) = json.get("root").and_then(|v| v.as_str()) else { return out; }; let mut visited = std::collections::HashSet::new(); visited.insert(root_name.to_owned()); walk_meta_inputs(nodes, root_name, "", &mut visited, &mut out); // hyperhive first, then alphabetical. String-sorting the // slash-paths puts every node directly above its own children // (`agent-foo`, `agent-foo/bar`, `agent-foo/bar/baz`), so the // result is a pre-order traversal the tree renderer can consume. out.sort_by(|a, b| match (a.name.as_str(), b.name.as_str()) { ("hyperhive", _) => std::cmp::Ordering::Less, (_, "hyperhive") => std::cmp::Ordering::Greater, _ => a.name.cmp(&b.name), }); out } fn walk_meta_inputs( nodes: &serde_json::Map, node_name: &str, prefix: &str, visited: &mut std::collections::HashSet, out: &mut Vec, ) { let Some(node) = nodes.get(node_name) else { return; }; let Some(inputs_map) = node.get("inputs").and_then(|v| v.as_object()) else { return; }; // Two passes: claim (and emit) every direct input of this node // before descending into any of them. A shallow input that a // deeper flake also references then keeps its shallow path // rather than being captured first by the deep walk. let mut to_recurse: Vec<(String, String)> = Vec::new(); for (alias, target) in inputs_map { // Inputs map value is either a string (node name) or an // array (a `follows` chain). The latter just aliases another // node — we can't `nix flake update` it directly, so skip. let serde_json::Value::String(target_name) = target else { continue; }; // Walk each fetched node once — guards shared subtrees and // cycles, and keeps the panel free of duplicate rows. if !visited.insert(target_name.clone()) { continue; } let Some(target_node) = nodes.get(target_name) else { continue; }; let path = if prefix.is_empty() { alias.clone() } else { format!("{prefix}/{alias}") }; if let Some(rev) = target_node .get("locked") .and_then(|v| v.get("rev")) .and_then(|v| v.as_str()) { let last_modified = target_node .get("locked") .and_then(|v| v.get("lastModified")) .and_then(serde_json::Value::as_i64) .unwrap_or(0); let url = target_node .get("original") .and_then(|v| v.get("url")) .and_then(|v| v.as_str()) .map(str::to_owned); out.push(MetaInputView { name: path.clone(), rev: rev.to_owned(), last_modified, url, }); } to_recurse.push((target_name.clone(), path)); } // Recurse hyperhive's subtree before any agent's — without this, // when meta's top-level `nixpkgs` is a `follows` alias the // `String` check above skips it, and the alphabetical BTreeMap // iteration descends into `agent-*` first. The agent walk then // claims `nixpkgs` at `agent-X/nixpkgs` instead of // `hyperhive/nixpkgs`, which is where the operator expects it. // Sort by the same "hyperhive first, then alpha" // priority `read_meta_inputs` uses for the final output. to_recurse.sort_by(|(a, _), (b, _)| match (a.as_str(), b.as_str()) { ("hyperhive", _) => std::cmp::Ordering::Less, (_, "hyperhive") => std::cmp::Ordering::Greater, _ => a.cmp(b), }); for (target_name, path) in to_recurse { walk_meta_inputs(nodes, &target_name, &path, visited, out); } } /// 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, ) -> Vec { transient_snapshot .iter() .filter(|(name, _)| !containers.iter().any(|c| &c.name == *name)) .map(|(name, st)| TransientView { name: name.clone(), kind: transient_label(st.kind), secs: st.since.elapsed().as_secs(), }) .collect() } /// Render each pending approval into its dashboard view (short sha + /// unified diff for `ApplyCommit`, 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::ApprovalStatus::Approved => "approved", hive_sh4re::ApprovalStatus::Denied => "denied", hive_sh4re::ApprovalStatus::Failed => "failed", hive_sh4re::ApprovalStatus::Cancelled => "cancelled", // Pending shouldn't appear in recent_resolved, but be defensive. hive_sh4re::ApprovalStatus::Pending => "pending", }; let kind = match a.kind { hive_sh4re::ApprovalKind::ApplyCommit => "apply_commit", hive_sh4re::ApprovalKind::Spawn => "spawn", hive_sh4re::ApprovalKind::InitConfig => "init_config", hive_sh4re::ApprovalKind::UpdateMetaInputs => "update_meta_inputs", hive_sh4re::ApprovalKind::SchedulePrompt => "schedule_prompt", hive_sh4re::ApprovalKind::MergeConfigPr => "merge_config_pr", }; ApprovalHistoryView { id: a.id, agent: a.agent, kind, sha_short, status, resolved_at: a.resolved_at.unwrap_or_default(), note: a.note, } } async fn build_approval_views(approvals: Vec) -> Vec { let mut out = Vec::with_capacity(approvals.len()); for a in approvals { out.push(match a.kind { hive_sh4re::ApprovalKind::ApplyCommit => { // Prefer the canonical fetched sha from applied; // commit_ref is only the manager's claim and may be // amended out from under us. let displayed = a.fetched_sha.as_deref().unwrap_or(&a.commit_ref); let sha = displayed[..displayed.len().min(12)].to_owned(); let diff = approval_diff(&a.agent, a.id).await; ApprovalView { id: a.id, agent: a.agent.clone(), kind: "apply_commit", sha_short: Some(sha), diff: Some(diff), description: a.description, pr_number: None, commit_ref: None, requested_at: a.requested_at, } } hive_sh4re::ApprovalKind::Spawn => ApprovalView { id: a.id, agent: a.agent, kind: "spawn", sha_short: None, diff: None, description: a.description, pr_number: None, commit_ref: None, requested_at: a.requested_at, }, hive_sh4re::ApprovalKind::InitConfig => ApprovalView { id: a.id, agent: a.agent, kind: "init_config", sha_short: None, diff: None, description: a.description, pr_number: None, commit_ref: None, requested_at: a.requested_at, }, hive_sh4re::ApprovalKind::UpdateMetaInputs => ApprovalView { id: a.id, agent: a.agent, kind: "update_meta_inputs", sha_short: None, diff: None, description: a.description, pr_number: None, commit_ref: Some(a.commit_ref), requested_at: a.requested_at, }, hive_sh4re::ApprovalKind::SchedulePrompt => ApprovalView { id: a.id, agent: a.agent, kind: "schedule_prompt", sha_short: None, diff: None, description: a.description, pr_number: None, commit_ref: Some(a.commit_ref), requested_at: a.requested_at, }, hive_sh4re::ApprovalKind::MergeConfigPr => { // commit_ref = PR number; fetched_sha = the reviewed PR // head. Show the head sha; the forge PR diff surface is // a later phase of the PR-based config flow — None for now. 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::().ok(); ApprovalView { id: a.id, agent: a.agent, kind: "merge_config_pr", sha_short: sha, diff: None, description: a.description, pr_number, commit_ref: None, requested_at: a.requested_at, } } }); } out } /// State-dir names that don't appear in the live container list. Each /// one surfaces in the dashboard as a row with R3V1V3 + PURG3 actions. fn build_tombstone_views( coord: &Coordinator, containers: &[ContainerView], transient_snapshot: &std::collections::HashMap, ) -> Vec { let _ = coord; // kept_state_names is a free fn but takes &self by future plan let live: std::collections::HashSet<&str> = containers .iter() .map(|c| c.name.as_str()) .chain(transient_snapshot.keys().map(String::as_str)) .collect(); Coordinator::kept_state_names() .into_iter() .filter(|name| !live.contains(name.as_str())) .map(|name| { let root = Coordinator::agent_state_root(&name); let state_bytes = dir_size_bytes(&root); let last_seen = std::fs::metadata(&root) .and_then(|m| m.modified()) .ok() .and_then(|t| t.duration_since(std::time::UNIX_EPOCH).ok()) .and_then(|d| i64::try_from(d.as_secs()).ok()) .unwrap_or(0); let has_creds = claude_has_session(&Coordinator::agent_claude_dir(&name)); TombstoneView { name, state_bytes, last_seen, has_creds, } }) .collect() } /// Sum the byte size of every regular file under `root`. Cheap to compute /// for typical agent state (config repo + claude creds + notes file — /// usually a few MB); fine to do inline on each /api/state. Returns 0 on /// any error. fn dir_size_bytes(root: &Path) -> u64 { fn walk(p: &Path, acc: &mut u64) { let Ok(rd) = std::fs::read_dir(p) else { return }; for entry in rd.flatten() { let Ok(ft) = entry.file_type() else { continue }; if ft.is_dir() { walk(&entry.path(), acc); } else if ft.is_file() && let Ok(meta) = entry.metadata() { *acc += meta.len(); } } } let mut total = 0u64; walk(root, &mut total); total } async fn dashboard_history(State(state): State) -> 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 = messages .into_iter() .filter_map(|m| match m { crate::broker::MessageEvent::Sent { id, from, to, body, at, in_reply_to, } => { let file_refs = scan_validated_paths(&body); Some(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); Some(crate::dashboard_events::DashboardEvent::Delivered { seq: 0, id, from, to, body, at: hive_sh4re::wire_time::from_secs(at), in_reply_to, file_refs, }) } // Ping events are never persisted to sqlite — this arm is // unreachable in practice but required for exhaustiveness. crate::broker::MessageEvent::Ping { .. } => None, }) .collect(); axum::Json(serde_json::json!({ "seq": seq, "events": 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)] 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, } async fn dashboard_stream( State(state): State, axum::extract::Query(q): axum::extract::Query, ) -> Sse>> { 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> = q.kinds.and_then(|raw| { let set: std::collections::HashSet = 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()) } #[derive(Deserialize)] struct RequestSpawnForm { name: String, } #[cfg(test)] mod tests { use super::*; #[test] fn problem_details_carry_rfc9457_status_and_detail() { // Contract the frontend depends on: the problem_details crate // serialises the RFC 9457 members we rely on — `status` (numeric) // and `detail` (the caller message; the FE reads `.detail`). let pd = problem_details::ProblemDetails::from_status_code(StatusCode::BAD_REQUEST) .with_detail("bad input"); let v = serde_json::to_value(&pd).expect("problem details serialise"); assert_eq!(v["status"], 400); assert_eq!(v["detail"], "bad input"); // The 500 wrapper path carries the internal-error status. let five = problem_details::ProblemDetails::from_status_code(StatusCode::INTERNAL_SERVER_ERROR) .with_detail("boom"); let fv = serde_json::to_value(&five).expect("problem details serialise"); assert_eq!(fv["status"], 500); } #[test] fn walk_meta_inputs_keeps_nixpkgs_under_hyperhive_post_follows_refactor() { // Reproduce the shape where meta has // `nixpkgs.follows = "hyperhive/nixpkgs"` at the top level // (rendered as an array — `["hyperhive" "nixpkgs"]` — which // walk_meta_inputs skips because we can't `nix flake update` // a follows alias). The remaining top-level inputs are // `hyperhive` (string) and `agent-z` (string). Without the // hyperhive-first recursion sort, the BTreeMap alphabetical // order descends into `agent-z` first and claims // `nixpkgs` at `agent-z/nixpkgs`. let raw = r#"{ "root": "root", "version": 7, "nodes": { "root": { "inputs": { "hyperhive": "hyperhive", "nixpkgs": ["hyperhive", "nixpkgs"], "agent-z": "agent-z" } }, "hyperhive": { "inputs": { "nixpkgs": "nixpkgs" }, "locked": {"rev": "hhrev", "lastModified": 1}, "original": {"url": "git+file:///tmp/hyperhive"} }, "agent-z": { "inputs": { "nixpkgs": "nixpkgs" }, "locked": {"rev": "azrev", "lastModified": 2}, "original": {"url": "git+file:///tmp/agent-z"} }, "nixpkgs": { "locked": {"rev": "npkrev", "lastModified": 3}, "original": {"url": "github:NixOS/nixpkgs/nixos-26.05"} } } }"#; let json: serde_json::Value = serde_json::from_str(raw).unwrap(); let nodes = json.get("nodes").unwrap().as_object().unwrap(); let root_name = json.get("root").unwrap().as_str().unwrap(); let mut visited = std::collections::HashSet::new(); visited.insert(root_name.to_owned()); let mut out = Vec::new(); walk_meta_inputs(nodes, root_name, "", &mut visited, &mut out); let nixpkgs = out .iter() .find(|v| v.rev == "npkrev") .expect("nixpkgs node should be emitted exactly once"); assert_eq!( nixpkgs.name, "hyperhive/nixpkgs", "nixpkgs should be claimed under hyperhive, not under agent-z. \ got: {:?}", nixpkgs.name ); // And the agent-z path should NOT also carry a nixpkgs entry — // the spanning-tree visited set guarantees it's claimed once. assert!( !out.iter().any(|v| v.name == "agent-z/nixpkgs"), "agent-z/nixpkgs should not be emitted (already claimed under hyperhive)" ); } #[test] fn validate_agent_name_accepts_canonical_shapes() { assert!(validate_agent_name("damocles").is_none()); assert!(validate_agent_name("hm1nd").is_none()); assert!(validate_agent_name("agent-with-dashes").is_none()); assert!(validate_agent_name("snake_case").is_none()); assert!(validate_agent_name("mixed_2-3").is_none()); let max = "a".repeat(63); assert!( validate_agent_name(&max).is_none(), "63-char name should pass" ); } // The two-axis guard (`guard_agent_name`) wires `validate_agent_name` // + an async coordinator lookup. The lookup needs a populated // `Coordinator`, which needs sqlite + tokio runtime; rather than // build that scaffolding for an integration-flavoured test we cover // the format axis here (the existence axis is enforced by the // shared `containers_snapshot` API, tested in `coordinator.rs`'s // own suite). 9 cases below cover the boundary-length case and // other expected rejects to make the contract explicit. #[test] fn validate_agent_name_rejects_bad_input() { assert!(validate_agent_name("").is_some()); let too_long = "a".repeat(64); assert!(validate_agent_name(&too_long).is_some()); // Path-traversal attempts. assert!(validate_agent_name("../etc/passwd").is_some()); assert!(validate_agent_name("alice/bob").is_some()); // Uppercase rejected — canonical lowercase convention. assert!(validate_agent_name("Alice").is_some()); // No spaces, dots, special chars. assert!(validate_agent_name("alice bob").is_some()); assert!(validate_agent_name("alice.bob").is_some()); assert!(validate_agent_name("alice;DROP TABLE messages").is_some()); // Non-ASCII (incl. unicode homoglyphs of ASCII dash). assert!(validate_agent_name("damóclès").is_some()); assert!(validate_agent_name("alice\u{2013}bob").is_some()); // en-dash } } /// Snapshot the current tombstone list and emit a /// `TombstonesChanged` event. Call after any mutation that could /// add or remove a tombstone (`actions::destroy`, /// `post_purge_tombstone`, spawn finalisation). Cheap — the list /// is tiny. pub(crate) async fn emit_tombstones_snapshot(coord: &Arc) { let containers = coord.containers_snapshot().await; let transient_snapshot = coord.transient_snapshot(); let tombstones = build_tombstone_views(coord, &containers, &transient_snapshot); coord.emit_dashboard_event(crate::dashboard_events::DashboardEvent::TombstonesChanged { seq: coord.next_seq(), tombstones, }); } /// Snapshot meta/flake.lock's root inputs + emit /// `MetaInputsChanged`. Call after any mutation that bumps a lock /// (`run_meta_update`, `auto_update::rebuild_agent`). pub(crate) fn emit_meta_inputs_snapshot(coord: &Coordinator) { let inputs = read_meta_inputs(); coord.emit_dashboard_event(crate::dashboard_events::DashboardEvent::MetaInputsChanged { seq: coord.next_seq(), inputs, }); } /// Unread operator-directed messages for the dashboard's Y3R C4LL inbox. /// Returns messages addressed to `"operator"` that haven't been /// acked yet (the operator clears them via the existing /// `POST /api/agent/operator/mark-all-read`). Newest-first; path-shaped /// tokens are validated so the client renders file links like the /// terminal does. Shape: `{ "messages": [{ id, from, body, at, /// in_reply_to, file_refs }] }`. async fn api_operator_inbox(State(state): State) -> Response { const INBOX_LIMIT: u64 = 100; match state .coord .broker .unread_for_recipient("operator", INBOX_LIMIT) { Ok(messages) => { let items: Vec = messages .into_iter() .filter_map(|m| match m { crate::broker::MessageEvent::Sent { id, from, body, at, in_reply_to, .. } => { let file_refs = scan_validated_paths(&body); Some(serde_json::json!({ "id": id, "from": from, "body": body, "at": hive_sh4re::wire_time::from_secs(at), "in_reply_to": in_reply_to, "file_refs": file_refs, })) } crate::broker::MessageEvent::Delivered { .. } | crate::broker::MessageEvent::Ping { .. } => None, }) .collect(); axum::Json(serde_json::json!({ "messages": items })).into_response() } Err(e) => error_response(&format!("operator-inbox failed: {e:#}")), } } #[derive(Deserialize)] struct StatsHiveQuery { window: Option, } /// Hive-wide turn-stats rollup for the dashboard swarm-stats view. /// Aggregates every agent's `hyperhive-turn-stats.sqlite` read-only /// (skips missing/unreadable ones). Window defaults to `24h`. async fn api_stats_hive( State(state): State, axum::extract::Query(q): axum::extract::Query, ) -> Response { let window = crate::hive_stats::Window::parse(q.window.as_deref().unwrap_or("24h")); axum::Json(crate::hive_stats::hive_snapshot( window, &state.coord.model_prices, )) .into_response() } /// Live per-agent-container CPU + memory load from cgroup v2. Samples /// CPU over a short interval (~200 ms), so this call briefly awaits. async fn api_container_resources() -> Response { axum::Json(crate::container_stats::gather().await).into_response() } /// `GET /api/audit-log` — most-recent agent-initiated privileged-action /// audit entries, newest first (server-clamped to 500). Backs the /// operator dashboard's audit view. Returns /// `{ "entries": [AuditEntry…], "total": N }` so the UI can show /// "latest 500 of N" rather than silently capping. `ts_unix` is in /// **seconds**. async fn api_audit_log(State(state): State) -> Response { const LIMIT: usize = 500; let entries = match state.coord.audit_log.list_recent(LIMIT) { Ok(rows) => rows, Err(e) => return error_response(&format!("audit-log: {e:#}")), }; let total = match state.coord.audit_log.count_total() { Ok(n) => n, Err(e) => return error_response(&format!("audit-log count: {e:#}")), }; axum::Json(serde_json::json!({ "entries": entries, "total": total })).into_response() } /// Validate that a path-param agent name conforms to the hyperhive /// naming whitelist: 1-63 chars of `[a-z0-9_-]`. Rejects empty, /// uppercase, slashes, dots, and any non-ASCII (incl. unicode /// homoglyphs of dash/underscore). Returns `None` on accept, `Some(reason)` /// on reject — caller wraps the reason in a 400 response. Conservative /// whitelist matching `nixos-container` basename rules and the existing /// agent-name convention across the codebase. pub(crate) fn validate_agent_name(name: &str) -> Option<&'static str> { if name.is_empty() { return Some("agent name must not be empty"); } if name.len() > 63 { return Some("agent name must be 63 characters or fewer"); } if !name .bytes() .all(|b| b.is_ascii_lowercase() || b.is_ascii_digit() || b == b'-' || b == b'_') { return Some("agent name must contain only [a-z0-9_-]"); } None } /// Two-axis path-param guard for write routes. Combines: /// /// 1. **format validation** (`validate_agent_name`) — rejects path /// traversal / unicode homoglyphs / empty + too-long names with /// HTTP 400. /// 2. **existence check** — looks up `name` in the coordinator's /// container snapshot; unknown name → HTTP 404 with a clear /// "no such agent" message. catches the operator-typo case where /// a destructive POST would otherwise hit silently (mark-all-read /// returning 0) or hit downstream lifecycle code that fails with /// a confusing nspawn error. /// /// Returns `None` when both checks pass (caller proceeds), `Some(Response)` /// when the request should be rejected. Use at the top of every write /// handler taking a name path-param. Read-only GET handlers and /// handlers that legitimately operate on tombstoned agents (e.g. /// `mark-all-read` on broker rows for a destroyed agent) call /// `validate_agent_name` directly and skip the existence check. async fn guard_agent_name(state: &AppState, name: &str) -> Option { if let Some(reason) = validate_agent_name(name) { return Some( (StatusCode::BAD_REQUEST, format!("bad agent name: {reason}")).into_response(), ); } let snapshot = state.coord.containers_snapshot().await; if !snapshot.iter().any(|c| c.name == name) { return Some((StatusCode::NOT_FOUND, format!("no such agent: {name}")).into_response()); } None } /// Operator-driven "clear this agent's inbox" — backs the side-panel /// "mark all read" button. Marks every message addressed to the /// agent as acked (backfilling `delivered_at` for any still-pending /// rows so vacuum can collect them). Returns `{ "marked": N }` so the /// frontend can show "cleared N messages" feedback without an extra /// fetch. async fn post_mark_all_read( State(state): State, AxumPath(name): AxumPath, ) -> Response { if let Some(reason) = validate_agent_name(&name) { return (StatusCode::BAD_REQUEST, format!("bad agent name: {reason}")).into_response(); } match state.coord.broker.mark_all_read(&name) { Ok(n) => { tracing::info!(%name, marked = n, "operator marked all messages read"); axum::Json(serde_json::json!({ "marked": n })).into_response() } Err(e) => error_response(&format!("mark-all-read {name} failed: {e:#}")), } } async fn post_purge_tombstone( State(state): State, AxumPath(name): AxumPath, ) -> Response { // Format guard FIRST so a name like `..` can't traverse into the // parent of `/var/lib/hyperhive/agents/{name}` and have // `remove_dir_all` wipe `/var/lib/hyperhive/` itself. Existing // manager + live-container checks below don't catch `..` — only // the whitelist does. Existence check via // `containers_snapshot()` is deliberately NOT used here: // tombstoned agents are gone from the snapshot by design; that's // the whole point of this endpoint. if let Some(reason) = validate_agent_name(&name) { return (StatusCode::BAD_REQUEST, format!("bad agent name: {reason}")).into_response(); } // Sanity: refuse to purge if a live container still exists with this // name. The dashboard already filters tombstones to non-live names, // but the operator could send a stale POST. let live = lifecycle::list().await.unwrap_or_default(); if live .iter() .any(|c| c == &format!("{}{name}", lifecycle::AGENT_PREFIX) || c == &name) { return error_response(&format!( "refusing to purge {name}: container still exists — use DESTR0Y first" )); } let mut errors = Vec::new(); for dir in [ Coordinator::agent_state_root(&name), Coordinator::agent_applied_dir(&name), ] { if dir.exists() && let Err(e) = std::fs::remove_dir_all(&dir) { errors.push(format!("{}: {e}", dir.display())); } } let _ = state .coord .approvals .fail_pending_for_agent(&name, "agent state purged"); if errors.is_empty() { tracing::info!(%name, "tombstone purged"); // Fire the post-purge tombstones snapshot so dashboards // drop the row live; matching form carries // `data-no-refresh`. emit_tombstones_snapshot(&state.coord).await; (StatusCode::OK, "ok").into_response() } else { error_response(&format!("purge {name} partial: {}", errors.join(", "))) } } /// Operator-side compose form on the dashboard terminal. Drops a /// message into the broker as `{from: "operator", to, body}`. Same /// shape that per-agent web UIs use via `OperatorMsg`, but here the /// operator picks the recipient explicitly with `@name`. No /// validation that `to` resolves to a known agent — broker accepts /// arbitrary recipients (and the agent's inbox grows whether or not /// they exist, which is fine for spawn-then-greet flows). #[derive(Deserialize)] struct OpSendForm { to: String, body: String, } /// Form for `POST /meta-update`. Inputs ride in as a comma-separated /// list under the `inputs` field — the JS submitter joins the /// checked boxes since axum's `Form` extractor doesn't natively /// decode repeated keys without a helper. #[derive(Deserialize)] struct MetaUpdateForm { inputs: String, } /// Bulk-update selected meta flake inputs, then rebuild the affected /// agents in the background. Idempotent w.r.t. selection — choosing /// an input that's already at the latest sha is a no-op (no commit, /// no rebuild ripple). Returns immediately after queueing the work; /// dashboard polls for progress via container `pending` spinners + /// the meta-inputs row sha update. async fn post_meta_update( State(state): State, Form(form): Form, ) -> Response { let inputs: Vec = form .inputs .split(',') .map(|s| s.trim().to_owned()) .filter(|s| !s.is_empty()) .collect(); if inputs.is_empty() { return error_response("meta-update: no inputs selected"); } let inputs_label = inputs.join(", "); let parent_id = state.coord.rebuild_queue.enqueue_with_inputs( crate::rebuild_queue::QueueKind::MetaUpdate, "hyperhive".to_owned(), crate::rebuild_queue::QueueSource::Manual, format!("meta-update via dashboard ({inputs_label})"), None, inputs.clone(), ); // Pre-enqueue cascade rebuilds NOW so they're visible in the queue // alongside the parent. The worker's MetaUpdate arm // no longer enqueues children — it just runs the lock bump and // (on failure) cancels these pre-queued children. let cascade_agents = crate::rebuild_queue::meta_update_cascade_agents(&inputs).await; let cascade_reason = format!("meta-update cascade ({inputs_label})"); for name in cascade_agents { state.coord.rebuild_queue.enqueue( crate::rebuild_queue::QueueKind::Rebuild, name, crate::rebuild_queue::QueueSource::MetaUpdate, cascade_reason.clone(), Some(parent_id), ); } state.coord.emit_rebuild_queue_snapshot(); (StatusCode::OK, "ok").into_response() } async fn post_op_send(State(state): State, Form(form): Form) -> Response { let to = form.to.trim().to_owned(); let body = form.body.trim().to_owned(); if to.is_empty() { return error_response("op-send: `to` required"); } if body.is_empty() { return error_response("op-send: `body` required"); } if to == "*" { let errors = state .coord .broadcast_send(hive_sh4re::OPERATOR_RECIPIENT, &body); if !errors.is_empty() { return error_response(&format!( "op-send broadcast partial fail: {}", errors.join("; ") )); } } else if let Err(e) = state.coord.broker.send(&hive_sh4re::Message { from: hive_sh4re::OPERATOR_RECIPIENT.to_owned(), to: to.clone(), body, in_reply_to: None, }) { return error_response(&format!("op-send to {to} failed: {e:#}")); } // 200 instead of 303 → the client doesn't refetch /api/state. The // broker `send` already emitted a `MessageEvent` which the // dashboard channel forwarder mirrors as `DashboardEvent::Sent`, // and the page's terminal + inbox derive from that stream — so the // operator's send shows up the same way an agent's send does, with // no full-state refresh in between. (axum::http::StatusCode::OK, "ok").into_response() } async fn post_request_spawn( State(state): State, Form(form): Form, ) -> Response { let name = form.name.trim().to_owned(); if name.is_empty() { return error_response("spawn: `name` required"); } match state.coord.approvals.submit_kind( &name, hive_sh4re::ApprovalKind::Spawn, "", None, hive_sh4re::MANAGER_AGENT, ) { Ok(id) => { tracing::info!(%id, %name, "operator: spawn approval queued via dashboard"); // Phase 5b: notify the dashboard event channel so live // subscribers can append the row without a snapshot // refetch. Spawn approvals carry no diff/sha. state .coord .emit_approval_added(crate::coordinator::ApprovalAdded { id, agent: &name, approval_kind: "spawn", sha_short: None, diff: None, description: None, pr_number: None, }); (StatusCode::OK, "ok").into_response() } Err(e) => error_response(&format!("request-spawn {name} failed: {e:#}")), } } fn transient_label(k: crate::coordinator::TransientKind) -> &'static str { use crate::coordinator::TransientKind::{ Destroying, Rebuilding, Restarting, Spawning, Starting, Stopping, }; match k { Spawning => "spawning", Starting => "starting", Stopping => "stopping", Restarting => "restarting", Rebuilding => "rebuilding", Destroying => "destroying", } } /// Convert either a logical name or a container name back to the logical /// name. Sub-agents are `h-foo` → `foo`; manager stays `root`. fn strip_container_prefix(name: &str) -> String { name.strip_prefix(lifecycle::AGENT_PREFIX) .unwrap_or(name) .to_owned() } /// The common internal-error case as a `ProblemDetails`: a 500 RFC 9457 /// (`application/problem+json`) value via the `problem_details` crate. /// `from_status_code` sets `status` + `title` (the canonical reason phrase) /// and leaves `type` as the default `about:blank`; `with_detail` carries the /// caller message; the crate's axum `IntoResponse` emits the /// `application/problem+json` body the frontend parses (it reads `detail`). /// Handlers that surface client failures return `Result<_, ProblemDetails>` /// and hand this (or an inline `from_status_code(4xx)`) straight to `Err` — /// no manual `.into_response()`. fn error_problem(message: &str) -> problem_details::ProblemDetails { problem_details::ProblemDetails::from_status_code(StatusCode::INTERNAL_SERVER_ERROR) .with_detail(message) } /// `Response` wrapper around [`error_problem`] for the many handlers typed /// `-> Response` whose only failure mode is a 500 — they funnel errors /// through here rather than threading a `Result` return type. fn error_response(message: &str) -> Response { error_problem(message).into_response() }