refactor(#2352): extract standalone hivectl crate, hive-c0re daemon-only

This commit is contained in:
damocles 2026-07-15 22:36:13 +02:00
commit cc67a05974
13 changed files with 236 additions and 147 deletions

File diff suppressed because it is too large Load diff

View file

@ -1,404 +0,0 @@
//! `hivectl` rebuild-queue progress rendering.
//!
//! Split out of `hivectl.rs` (which is already large): everything that
//! polls the daemon's DAG queue (`HostRequest::QueueDag`) and renders the
//! per-DAG / per-node progress lives here. [`wait_for_dags`] is the entry
//! point the command handlers call; it dispatches to a live `indicatif`
//! animation on a TTY and a plain line-on-change stream otherwise.
use std::path::Path;
use anyhow::{Context as _, Result, bail};
/// Poll the submitted DAG ids (`HostRequest::QueueDag`, ~1s interval) and
/// render progress until they all reach a terminal state. Exits non-zero
/// (via the returned `Err`) when any DAG (or fan-out child) ends `failed`;
/// a `cancelled` DAG terminates the wait but is an operator action, not an
/// error.
///
/// # Errors
///
/// Returns an error if the daemon socket can't be reached, or if any polled
/// DAG finished in the `failed` state.
pub(crate) async fn wait_for_dags(socket: &Path, ids: Vec<u64>, no_wait: bool) -> Result<()> {
use std::io::IsTerminal as _;
if no_wait || ids.is_empty() {
return Ok(());
}
// Animate only on a real terminal. Piped / CI output falls back to the
// plain line-on-change stream so logs stay free of spinner redraw noise.
if std::io::stderr().is_terminal() {
wait_for_dags_animated(socket, ids).await
} else {
wait_for_dags_plain(socket, ids).await
}
}
/// Non-TTY progress: print a fresh line whenever a DAG's rendered state
/// changes. No cursor tricks, so it's clean in pipes and CI logs.
async fn wait_for_dags_plain(socket: &Path, ids: Vec<u64>) -> Result<()> {
let mut pending: std::collections::BTreeSet<u64> = ids.into_iter().collect();
let mut last: std::collections::HashMap<u64, String> = std::collections::HashMap::new();
let mut failed: Vec<String> = Vec::new();
while !pending.is_empty() {
for id in pending.clone() {
let resp =
hive_c0re::client::request(socket, hive_host_sock::HostRequest::QueueDag { id })
.await
.with_context(|| format!("connect to daemon socket {}", socket.display()))?;
let dags = resp.dags.unwrap_or_default();
if dags.is_empty() {
// Evicted from the queue's history tail — it finished a
// while ago; nothing left to report on.
println!("job #{id}: gone from queue history");
pending.remove(&id);
continue;
}
let mut all_terminal = true;
for d in &dags {
let line = render_dag_line(d);
if last.get(&d.id) != Some(&line) {
println!("{line}");
last.insert(d.id, line);
}
// Node-level terminality, not the roll-up: a DAG rolls
// up `failed` the moment one node fails while its
// after-any recovery node (rebuild's tail Reconcile)
// may still be running — keep watching so the operator
// sees whether the agent came back.
if d.nodes.iter().all(|n| n.state.is_terminal()) {
if d.state == hive_sh4re::jobs::State::Failed {
failed.push(format!("{} {}", d.kind.as_str(), dag_agents(d)));
}
} else {
all_terminal = false;
}
}
if all_terminal {
pending.remove(&id);
}
}
if !pending.is_empty() {
tokio::time::sleep(std::time::Duration::from_secs(1)).await;
}
}
finish_wait(failed)
}
/// TTY progress: a live `indicatif` render — one braille-spinner line per
/// DAG node (grouped under a per-DAG header), each with its own elapsed
/// timer, plus an overall-elapsed footer. A node with more than one
/// dependency (fan-in) gets its own row with an `(after …)` marker rather
/// than being crammed onto a chain line.
async fn wait_for_dags_animated(socket: &Path, ids: Vec<u64>) -> Result<()> {
use indicatif::{MultiProgress, ProgressBar, ProgressStyle};
let spinner = ProgressStyle::with_template(" {spinner} {msg}")
.unwrap_or_else(|_| ProgressStyle::default_spinner())
.tick_strings(&["", "", "", "", "", "", "", "", "", "", "·"]);
let plain =
ProgressStyle::with_template("{msg}").unwrap_or_else(|_| ProgressStyle::default_bar());
let mp = MultiProgress::new();
let started = std::time::Instant::now();
// Per-DAG header bars + per-node bars, keyed so we update in place. The
// insertion order groups a DAG's nodes right under its header.
let mut dag_bars: std::collections::HashMap<u64, ProgressBar> =
std::collections::HashMap::new();
let mut node_bars: std::collections::HashMap<(u64, hive_sh4re::jobs::NodeId), ProgressBar> =
std::collections::HashMap::new();
let mut node_done: std::collections::HashSet<(u64, hive_sh4re::jobs::NodeId)> =
std::collections::HashSet::new();
let mut pending: std::collections::BTreeSet<u64> = ids.into_iter().collect();
let mut failed: Vec<String> = Vec::new();
while !pending.is_empty() {
let now = now_unix();
for id in pending.clone() {
let resp =
hive_c0re::client::request(socket, hive_host_sock::HostRequest::QueueDag { id })
.await
.with_context(|| format!("connect to daemon socket {}", socket.display()))?;
let dags = resp.dags.unwrap_or_default();
if dags.is_empty() {
mp.println(format!("job #{id}: gone from queue history"))
.ok();
pending.remove(&id);
continue;
}
let mut all_terminal = true;
for d in &dags {
let hdr = dag_bars.entry(d.id).or_insert_with(|| {
let b = mp.add(ProgressBar::new_spinner());
b.set_style(plain.clone());
b
});
hdr.set_message(format!(
"{} {} {} · {}",
state_glyph(d.state),
d.kind.as_str(),
dag_agents(d),
fmt_dur(dag_elapsed(d, now)),
));
for n in &d.nodes {
let key = (d.id, n.id);
if node_done.contains(&key) {
continue;
}
let bar = node_bars.entry(key).or_insert_with(|| {
let b = mp.add(ProgressBar::new_spinner());
b.set_style(spinner.clone());
b.enable_steady_tick(std::time::Duration::from_millis(120));
b
});
if n.state.is_terminal() {
bar.set_style(plain.clone());
bar.finish_with_message(format!(
" {} {}",
state_glyph(n.state),
node_line(d, n, now)
));
node_done.insert(key);
} else {
bar.set_message(node_line(d, n, now));
}
}
if d.nodes.iter().all(|n| n.state.is_terminal()) {
if d.state == hive_sh4re::jobs::State::Failed {
failed.push(format!("{} {}", d.kind.as_str(), dag_agents(d)));
}
} else {
all_terminal = false;
}
}
if all_terminal {
pending.remove(&id);
}
}
if !pending.is_empty() {
tokio::time::sleep(std::time::Duration::from_secs(1)).await;
}
}
// Leave the final node lines on screen; drop the still-spinning headers
// into their terminal state and print an overall-elapsed footer.
for hdr in dag_bars.values() {
hdr.finish();
}
mp.println(format!(
"done · {} elapsed",
fmt_dur(started.elapsed().as_secs().try_into().unwrap_or(0))
))
.ok();
finish_wait(failed)
}
/// Shared tail: succeed when nothing failed, else bail listing the failures.
fn finish_wait(mut failed: Vec<String>) -> Result<()> {
if failed.is_empty() {
Ok(())
} else {
failed.sort();
failed.dedup();
bail!("queued job(s) failed: {}", failed.join(", "))
}
}
/// Distinct agents across a DAG's nodes, comma-joined for display — the
/// per-node replacement for the old DAG-level `agent` field. Single-agent
/// DAGs render one name; a hive-wide DAG lists each.
fn dag_agents(d: &hive_sh4re::jobs::DagView) -> String {
let mut seen: Vec<&str> = Vec::new();
for n in &d.nodes {
if !seen.contains(&n.agent.as_str()) {
seen.push(&n.agent);
}
}
seen.join(",")
}
/// Current unix time in seconds (0 on the impossible pre-epoch error).
fn now_unix() -> i64 {
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.ok()
.and_then(|d| i64::try_from(d.as_secs()).ok())
.unwrap_or(0)
}
/// Elapsed seconds for a DAG: `started_at` (falling back to `enqueued_at`)
/// through `finished_at` or `now`.
fn dag_elapsed(d: &hive_sh4re::jobs::DagView, now: i64) -> i64 {
let start = d.started_at.unwrap_or(d.enqueued_at);
(d.finished_at.unwrap_or(now) - start).max(0)
}
/// Elapsed seconds for a node: `started_at` → `finished_at`/`now`, or 0
/// when it hasn't started.
fn node_elapsed(n: &hive_sh4re::jobs::NodeView, now: i64) -> i64 {
match n.started_at {
Some(start) => (n.finished_at.unwrap_or(now) - start).max(0),
None => 0,
}
}
/// Compact duration: `45s` under a minute, else `1m03s`.
fn fmt_dur(secs: i64) -> String {
let s = secs.max(0);
if s < 60 {
format!("{s}s")
} else {
format!("{}m{:02}s", s / 60, s % 60)
}
}
/// One animated node line: kind, live step, an `(after …)` marker for a
/// fan-in node (>1 dep), its elapsed timer, and a truncated error tail.
fn node_line(d: &hive_sh4re::jobs::DagView, n: &hive_sh4re::jobs::NodeView, now: i64) -> String {
use std::fmt::Write as _;
let mut s = n.kind.clone();
if n.state == hive_sh4re::jobs::State::Running
&& let Some(step) = &n.step
{
let _ = write!(s, " ({step})");
}
if n.deps.len() > 1 {
let after: Vec<&str> = n
.deps
.iter()
.filter_map(|dep| d.nodes.iter().find(|m| m.id == *dep))
.map(|m| m.kind.as_str())
.collect();
if !after.is_empty() {
let _ = write!(s, " (after {})", after.join(", "));
}
}
let el = node_elapsed(n, now);
if el > 0 {
let _ = write!(s, " · {}", fmt_dur(el));
}
if let Some(err) = &n.error {
let short: String = err.chars().take(100).collect();
let _ = write!(s, " — {short}");
}
s
}
fn state_glyph(state: hive_sh4re::jobs::State) -> &'static str {
match state {
hive_sh4re::jobs::State::Queued => "",
hive_sh4re::jobs::State::Running => "",
hive_sh4re::jobs::State::Done => "",
hive_sh4re::jobs::State::Failed => "",
hive_sh4re::jobs::State::Cancelled => "",
}
}
/// One progress line for a DAG: roll-up glyph, template, agent, then
/// the node chain with the running node's live step label — the CLI
/// twin of the dashboard's queue card. Used by the plain (non-TTY) path.
fn render_dag_line(d: &hive_sh4re::jobs::DagView) -> String {
use std::fmt::Write as _;
let mut out = format!(
"{} {} {:<12}",
state_glyph(d.state),
d.kind.as_str(),
dag_agents(d)
);
for (i, n) in d.nodes.iter().enumerate() {
let sep = if i == 0 { " " } else { "" };
let _ = write!(out, "{sep}{} {}", state_glyph(n.state), n.kind);
if n.state == hive_sh4re::jobs::State::Running
&& let Some(step) = &n.step
{
let _ = write!(out, " ({step})");
}
}
if let Some(err) = d.nodes.iter().find_map(|n| n.error.as_deref()) {
let short: String = err.chars().take(120).collect();
let _ = write!(out, " — {short}");
}
out
}
#[cfg(test)]
mod tests {
use hive_sh4re::jobs::{DagView, NodeView, Source, State, Template};
use super::render_dag_line;
fn node(id: u32, agent: &str, kind: &str, state: State, step: Option<&str>) -> NodeView {
NodeView {
id,
agent: agent.to_owned(),
kind: kind.to_owned(),
deps: if id == 0 { vec![] } else { vec![id - 1] },
state,
step: step.map(str::to_owned),
build_log_id: None,
started_at: None,
finished_at: None,
error: None,
}
}
#[test]
fn render_dag_line_shows_chain_and_running_step() {
let dag = DagView {
id: 7,
kind: Template::Rebuild,
state: State::Running,
source: Source::Manual,
reason: "manual".to_owned(),
enqueued_at: 0,
started_at: Some(1),
finished_at: None,
inputs: vec![],
approval_id: None,
perm_payload: None,
nodes: vec![
node(0, "alice", "prebuild", State::Done, None),
node(1, "alice", "stop_for_update", State::Done, None),
node(
2,
"alice",
"swap",
State::Running,
Some("nixos-container update"),
),
node(3, "alice", "reconcile", State::Queued, None),
],
};
let line = render_dag_line(&dag);
assert!(line.starts_with("▶ rebuild alice"), "{line}");
assert!(
line.contains(
"✔ prebuild → ✔ stop_for_update → ▶ swap (nixos-container update) → ⏸ reconcile"
),
"{line}"
);
}
#[test]
fn render_dag_line_surfaces_first_node_error() {
let mut failed = node(0, "bob", "prebuild", State::Failed, None);
failed.error = Some("nix build exploded".to_owned());
let dag = DagView {
id: 8,
kind: Template::Rebuild,
state: State::Failed,
source: Source::Manual,
reason: "manual".to_owned(),
enqueued_at: 0,
started_at: Some(1),
finished_at: Some(2),
inputs: vec![],
approval_id: None,
perm_payload: None,
nodes: vec![failed],
};
let line = render_dag_line(&dag);
assert!(line.contains("✖ rebuild"), "{line}");
assert!(line.contains("— nix build exploded"), "{line}");
}
}

View file

@ -1,27 +0,0 @@
use std::path::Path;
use anyhow::{Context, Result, bail};
use hive_host_sock::{HostRequest, HostResponse};
use tokio::io::{AsyncBufReadExt, AsyncWriteExt, BufReader};
use tokio::net::UnixStream;
pub async fn request(socket: &Path, req: HostRequest) -> Result<HostResponse> {
let stream = UnixStream::connect(socket)
.await
.with_context(|| format!("connect to {}", socket.display()))?;
let (read, mut write) = stream.into_split();
let mut payload = serde_json::to_string(&req)?;
payload.push('\n');
write.write_all(payload.as_bytes()).await?;
write.flush().await?;
let mut reader = BufReader::new(read);
let mut line = String::new();
reader.read_line(&mut line).await?;
if line.is_empty() {
bail!("server closed connection without responding");
}
let resp: HostResponse = serde_json::from_str(line.trim())?;
Ok(resp)
}

View file

@ -1,15 +1,15 @@
//! `hive-c0re` library — module surface shared by the `hive-c0re`
//! daemon binary and the `hivectl` operator CLI.
//! `hive-c0re` library — the coordinator daemon's module surface
//! (coordinator, broker, axum dashboard, admin/manager/agent unix
//! sockets, background sweepers). Consumed by the `hive-c0re` daemon
//! binary (`src/main.rs`).
//!
//! `hive-c0re` (daemon) keeps the systemd service shape it always had:
//! coordinator, broker, axum dashboard, admin/manager/agent unix
//! sockets, background sweepers. `hivectl` (sibling bin under
//! `src/bin/hivectl.rs`) reuses a thin subset (`forge`, `matrix`,
//! `lifecycle`) to expose host-side administration verbs — manually
//! provisioning forge / matrix users for an agent, etc.
//! The operator CLI lives in the **standalone `hivectl` crate**, which
//! talks to the daemon over the host admin socket (`hive-host-sock` wire
//! types) rather than linking this crate — so `hive-c0re` is daemon-only,
//! not a library shared with a CLI.
//!
//! Every module is re-exported `pub` so anything in the crate is
//! addressable from either binary; the lib doesn't have a curated
//! addressable from the daemon binary; the lib doesn't have a curated
//! surface beyond "this is where the modules live".
//!
//! Cohesive clusters live in directory submodules (`stores`, `stats`,
@ -19,7 +19,6 @@
pub mod actions;
pub mod agent_config;
pub mod client;
pub mod container_view;
pub mod coordinator;
pub mod dashboard;

View file

@ -1,9 +1,8 @@
use std::path::PathBuf;
use std::sync::Arc;
use anyhow::{Context as _, Result, bail};
use anyhow::{Context as _, Result};
use clap::{Parser, Subcommand};
use hive_host_sock::{HostRequest, HostResponse};
// Every module hangs off the `hive_c0re` library (see `src/lib.rs`).
// The daemon and the `hivectl` sibling binary share the same module
@ -12,7 +11,7 @@ use hive_host_sock::{HostRequest, HostResponse};
// explicit (any new daemon entry point reads off the next add).
use hive_c0re::coordinator::{Coordinator, HiveEnv, ServeConfig};
use hive_c0re::{
agent_sockets, auto_update, broker, client, crash_watch, dashboard, dashboard_events, forge,
agent_sockets, auto_update, broker, crash_watch, dashboard, dashboard_events, forge,
host_stats, job_queue, knowledge, matrix, mcp_sockets, migrate, reminder_scheduler,
scheduled_prompts_worker, server, socket_server, sweep_health, warnings,
};
@ -91,52 +90,6 @@ enum Cmd {
#[arg(long)]
build_slots: Option<usize>,
},
/// Spawn a new agent container directly (`hive-agent-<name>`). Bypasses
/// the approval queue — use only as an operator on the host. For
/// approval-gated spawns, use `request-spawn` instead.
Spawn { name: String },
/// Queue a spawn request as an approval. The container is created on
/// `approve <id>` (CLI) or the dashboard's APPR0VE button.
RequestSpawn { name: String },
/// Stop a managed container (graceful).
Kill { name: String },
/// Tear down a sub-agent container. Container is removed; persistent
/// state (config repos + Claude credentials) is kept by default. Pass
/// `--purge` to also wipe the agent's state dirs (config + creds +
/// notes). No undo.
Destroy {
name: String,
#[arg(long)]
purge: bool,
},
/// Apply pending config to a managed container.
Rebuild { name: String },
/// List managed containers.
List,
/// List pending approval requests submitted by the manager.
Pending,
/// Approve a pending request by id; the action runs immediately.
Approve { id: i64 },
/// Deny a pending request by id.
Deny { id: i64 },
/// Move an agent in the topology tree. Set `--parent` to a new
/// parent agent name; pass `--root` to promote the agent to root
/// (no parent). Refuses cycles and unknown agents. The manager
/// is reparentable like any other agent — its privileges come
/// from the privileged MCP socket, not its tree position.
SetParent {
child: String,
/// New parent agent name. Mutually exclusive with `--root`.
/// Exactly one of `--parent` / `--root` is required — clap
/// rejects both-absent calls so a fat-fingered
/// `hive-c0re set-parent alice` doesn't silently promote
/// alice to root.
#[arg(long, conflicts_with = "root", required_unless_present = "root")]
parent: Option<String>,
/// Promote `child` to root (no parent).
#[arg(long)]
root: bool,
},
}
#[tokio::main]
@ -206,37 +159,6 @@ async fn main() -> Result<()> {
}
cmd_serve(sc.env, sc.model_prices, sc.build_slots, db, &cli.socket).await
}
Cmd::Spawn { name } => {
render(client::request(&cli.socket, HostRequest::Spawn { name }).await?)
}
Cmd::RequestSpawn { name } => {
render(client::request(&cli.socket, HostRequest::RequestSpawn { name }).await?)
}
Cmd::Kill { name } => {
render(client::request(&cli.socket, HostRequest::Kill { name }).await?)
}
Cmd::Destroy { name, purge } => {
render(client::request(&cli.socket, HostRequest::Destroy { name, purge }).await?)
}
Cmd::Rebuild { name } => {
render(client::request(&cli.socket, HostRequest::Rebuild { name }).await?)
}
Cmd::List => render(client::request(&cli.socket, HostRequest::List).await?),
Cmd::Pending => render(client::request(&cli.socket, HostRequest::Pending).await?),
Cmd::Approve { id } => {
render(client::request(&cli.socket, HostRequest::Approve { id }).await?)
}
Cmd::Deny { id } => render(client::request(&cli.socket, HostRequest::Deny { id }).await?),
Cmd::SetParent {
child,
parent,
root,
} => {
let new_parent = if root { None } else { parent };
render(
client::request(&cli.socket, HostRequest::SetParent { child, new_parent }).await?,
)
}
}
}
@ -648,11 +570,3 @@ fn spawn_broker_to_dashboard_forwarder(coord: Arc<Coordinator>) {
}
});
}
fn render(resp: HostResponse) -> Result<()> {
println!("{}", serde_json::to_string_pretty(&resp)?);
if !resp.ok {
bail!(resp.error.unwrap_or_else(|| "request failed".to_owned()));
}
Ok(())
}