diff --git a/Cargo.lock b/Cargo.lock index 22171566..93b98034 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1725,7 +1725,6 @@ version = "0.1.0" dependencies = [ "chrono", "hive-jobq", - "hive-jobq-wire", "hive-sh4re", "hive-types", "serde", @@ -1862,7 +1861,6 @@ dependencies = [ "clap-markdown", "clap_complete", "hive-host-sock", - "hive-jobq-wire", "hive-sh4re", "hive-types", "http-body-util", diff --git a/hive-c0re/src/job_queue/mod.rs b/hive-c0re/src/job_queue/mod.rs index d4afeea0..6b2411df 100644 --- a/hive-c0re/src/job_queue/mod.rs +++ b/hive-c0re/src/job_queue/mod.rs @@ -289,8 +289,8 @@ impl JobQueue { #[must_use] pub fn first_error(&self, dag_id: u64) -> Option { let inner = self.lock(); - let node = find_node(&inner, dag_id)?; - inner.graph().first_error(node).map(ToOwned::to_owned) + let container = container(&inner, dag_id)?; + inner.graph().first_error(container).map(ToOwned::to_owned) } /// `(agent, label, takes_container_down)` for the live transient-pill set, @@ -378,45 +378,15 @@ impl JobQueue { .filter_map(|c| dag_view(&inner, c)) .collect() } - - /// One or more nodes plus their live subtrees, as generic wire nodes — - /// the `QueueNodes` polling surface behind `hivectl`'s wait/progress - /// loop. Sibling of [`Self::snapshot`] (which serves the same graph - /// through the typed `DagView`/`NodeView` projection for the - /// dashboard's `/api/state.rebuild_queue`), this one goes through - /// [`GraphWire::wire_snapshot`] instead — no `Done`-node filtering, no - /// roll-up field (a node's own `state` answers that, see - /// `hive_jobq_wire`'s doc comment). Looks each id up by identity - /// alone — no assumption that it names a DAG container or a root; - /// "just show whatever the backend sends" for whatever ids the caller - /// asks about. Multiple ids in one call is the normal shape for a - /// batch op (e.g. restarting every agent submits one root per agent) — - /// callers should request the whole batch together rather than poll - /// one id per round-trip. - /// - /// An id with no matching node in the graph is silently dropped from - /// the result rather than erroring the whole batch — some ids in a - /// batch may already be evicted while others are still live. Today - /// that only happens for a genuinely unknown id: nothing prunes the - /// graph yet (bounded-prune is a Stage-C follow-up, see - /// [`visible_dags`]), so a *completed* DAG's nodes keep riding here - /// with a terminal `state` rather than disappearing — callers - /// watching for "done" should read the root's `state`, not absence. - #[must_use] - pub fn node_subtrees(&self, ids: &[u64]) -> Vec { - let inner = self.lock(); - let roots: Vec = ids.iter().filter_map(|id| find_node(&inner, *id)).collect(); - inner.graph().wire_snapshot(roots) - } } -/// The graph node whose id equals `id`, whatever its kind or depth. -/// `NodeId` is un-fabricable from a raw `u64`, so this is a search. -fn find_node(sched: &Sched, id: u64) -> Option { - sched - .graph() - .nodes() - .find_map(|n| (n.id.get() == id).then_some(n.id)) +/// The container node of `dag_id` — the `NodeKind::Dag` root whose id equals +/// `dag_id`. `NodeId` is un-fabricable from a raw `u64`, so this is a search. +fn container(sched: &Sched, dag_id: u64) -> Option { + sched.graph().nodes().find_map(|n| { + (n.parent.is_none() && n.id.get() == dag_id && matches!(n.payload, NodeKind::Dag { .. })) + .then_some(n.id) + }) } /// Project a DAG into its wire [`DagView`]: a near-raw view of the diff --git a/hive-c0re/src/server.rs b/hive-c0re/src/server.rs index 5427a7f6..26803aa2 100644 --- a/hive-c0re/src/server.rs +++ b/hive-c0re/src/server.rs @@ -160,9 +160,6 @@ async fn dispatch(req: &HostRequest, coord: Arc) -> HostResponse { .collect(); HostResponse::dags(dags) } - HostRequest::QueueNodes { ids } => { - HostResponse::nodes(coord.job_queue.node_subtrees(ids)) - } HostRequest::List => HostResponse::list(lifecycle::list().await?), // The agents root is ours and not world-traversable, so this // question is only answerable on this side of the socket — diff --git a/hive-host-sock/Cargo.toml b/hive-host-sock/Cargo.toml index dd7fc555..ffd069fd 100644 --- a/hive-host-sock/Cargo.toml +++ b/hive-host-sock/Cargo.toml @@ -10,7 +10,6 @@ workspace = true [dependencies] chrono.workspace = true hive-jobq.workspace = true -hive-jobq-wire.workspace = true hive-sh4re.workspace = true hive-types.workspace = true serde.workspace = true diff --git a/hive-host-sock/src/lib.rs b/hive-host-sock/src/lib.rs index caf7e3d9..2363e3db 100644 --- a/hive-host-sock/src/lib.rs +++ b/hive-host-sock/src/lib.rs @@ -214,19 +214,6 @@ pub enum HostRequest { /// `hivectl`'s wait/progress loop. A multi-step op is a single DAG /// (its whole graph in `nodes`). Result: [`HostResponse::dags`]. QueueDag { id: u64 }, - /// Fetch one or more job-queue nodes plus their live subtrees, as - /// generic `hive-jobq-wire` nodes — `hivectl`'s wait/progress loop. - /// Sibling of [`Self::QueueDag`]: same graph, through the generic - /// projection instead of the typed `DagView`/`NodeView` (kept for - /// `QueueDag`'s other consumer, `/api/state.rebuild_queue`). No - /// assumption that an id names a DAG container or root — whatever - /// node has that id, the backend hands back its subtree as-is. A - /// batch op that submits several independent roots (e.g. one per - /// agent on a hive-wide restart) is a single request naming all of - /// them, not one request per id — the caller owns bundling `ids`, - /// this request just answers whatever it's asked. Result: - /// [`HostResponse::nodes`]. - QueueNodes { ids: Vec }, /// List pending approval requests. Pending, /// Approve a pending request by id; the action runs immediately. @@ -552,16 +539,6 @@ pub struct HostResponse { /// been evicted from the queue's history tail. #[serde(default, skip_serializing_if = "Option::is_none")] pub dags: Option>, - /// `QueueNodes` result — the requested nodes plus their live subtrees, - /// as generic `hive-jobq-wire` nodes, all roots' subtrees combined in - /// one flat list. `None` for every other request kind. An id with no - /// live node in the graph is silently dropped rather than erroring - /// the whole batch — see `JobQueue::node_subtrees`'s doc comment for - /// why a *completed* DAG's nodes don't vanish the same way - /// `QueueDag`'s do; callers should read each root node's `state` for - /// terminality, not absence. - #[serde(default, skip_serializing_if = "Option::is_none")] - pub nodes: Option>, /// Free-form operator-facing output lines the client prints verbatim /// (one per line). Carries results a request produced daemon-side that /// have no structured home — e.g. a freshly-minted matrix token, a @@ -662,17 +639,6 @@ impl HostResponse { } } - /// `QueueNodes` result — the polled node + its subtree, as generic - /// wire nodes. - #[must_use] - pub fn nodes(nodes: Vec) -> Self { - Self { - ok: true, - nodes: Some(nodes), - ..Self::default() - } - } - /// A success carrying operator-facing output lines the client prints /// verbatim — the result shape for the `Matrix*` provisioning requests. #[must_use] diff --git a/hivectl/Cargo.toml b/hivectl/Cargo.toml index baf32cf2..e282912d 100644 --- a/hivectl/Cargo.toml +++ b/hivectl/Cargo.toml @@ -21,7 +21,6 @@ clap.workspace = true clap_complete.workspace = true clap-markdown = "0.1" hive-host-sock.workspace = true -hive-jobq-wire.workspace = true hive-sh4re.workspace = true hive-types.workspace = true http-body-util.workspace = true diff --git a/hivectl/src/agents.rs b/hivectl/src/agents.rs index 776ace6f..c66f7e5d 100644 --- a/hivectl/src/agents.rs +++ b/hivectl/src/agents.rs @@ -13,7 +13,7 @@ use anyhow::{Context as _, Result, bail}; use hive_host_sock::HostRequest; use crate::cli::{AgentCmd, AgentQuotaCmd}; -use crate::dag_progress::wait_for_nodes; +use crate::dag_progress::wait_for_dags; use crate::util::render; async fn agents_restart(socket: &Path, name: &str, no_wait: bool) -> Result<()> { @@ -27,7 +27,7 @@ async fn agents_restart(socket: &Path, name: &str, no_wait: bool) -> Result<()> .with_context(|| format!("connect to daemon socket {}", socket.display()))?; if resp.ok { println!("restart queued: {name}"); - wait_for_nodes(socket, resp.queued_dags.unwrap_or_default(), no_wait).await + wait_for_dags(socket, resp.queued_dags.unwrap_or_default(), no_wait).await } else { bail!( "restart {name}: {}", diff --git a/hivectl/src/dag_progress.rs b/hivectl/src/dag_progress.rs index fa25516b..aaca5710 100644 --- a/hivectl/src/dag_progress.rs +++ b/hivectl/src/dag_progress.rs @@ -1,45 +1,26 @@ //! `hivectl` rebuild-queue progress rendering. //! //! Split out of `hivectl.rs` (which is already large): everything that -//! polls the daemon's node queue (`HostRequest::QueueNodes`) and renders -//! progress lives here. [`wait_for_nodes`] 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. -//! -//! Consumes `hive-jobq-wire`'s generic [`GraphNode`]/`NodePayload`: a -//! node's kind (and, for a root, what submitted it) comes straight off -//! `payload.label`; node-specific extras (`agent`) ride in -//! `payload.data`'s opaque kvps — the shape `hive-c0re`'s -//! `impl WireNode for NodeKind` produces (see its doc comment). There's -//! no separate roll-up field on the wire: a node's own -//! `state` already reflects everything below it (see `hive_jobq_wire`'s -//! doc comment), so nothing here re-derives a roll-up or waits for every -//! node to go terminal before reading a result. -//! -//! `ids` is always a **batch**, not a single id: a hive-wide op (e.g. -//! restarting every agent) submits one root per agent, and the whole -//! batch is polled together in a single `QueueNodes` request per tick — -//! [`group_by_root`] splits the combined response back into per-root -//! groups rather than issuing one round-trip per id. +//! 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::collections::{BTreeSet, HashMap, HashSet}; use std::path::Path; use anyhow::{Context as _, Result, bail}; -use hive_host_sock::jobs::State; -use hive_jobq_wire::{GraphDep, GraphNode, WireId}; -/// Poll the submitted ids (`HostRequest::QueueNodes`, ~1s interval, all -/// ids in one request per tick) and render progress until they all reach -/// a terminal state. Exits non-zero (via the returned `Err`) when any -/// node ends `failed`; a `cancelled` node terminates the wait but is an -/// operator action, not an error. +/// 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 node finished in the `failed` state. -pub(crate) async fn wait_for_nodes(socket: &Path, ids: Vec, no_wait: bool) -> Result<()> { +/// 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, no_wait: bool) -> Result<()> { use std::io::IsTerminal as _; if no_wait || ids.is_empty() { return Ok(()); @@ -47,72 +28,52 @@ pub(crate) async fn wait_for_nodes(socket: &Path, ids: Vec, no_wait: bool) // 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_nodes_animated(socket, ids).await + wait_for_dags_animated(socket, ids).await } else { - wait_for_nodes_plain(socket, ids).await + wait_for_dags_plain(socket, ids).await } } -/// Split a combined `QueueNodes` response back into per-root groups, keyed -/// by each root's own id. A node with `parent: None` is a root and starts -/// its own group (keyed by its own id); everything else joins its direct -/// parent's group. Today's job shapes are root + flat children (no deeper -/// nesting), so a direct-parent lookup is enough — matches the flat -/// child-iteration every render helper below already assumed. -fn group_by_root(nodes: Vec) -> HashMap> { - let mut groups: HashMap> = HashMap::new(); - for n in nodes { - let key = n.parent.unwrap_or(n.id); - groups.entry(key).or_default().push(n); - } - groups -} - -/// This id's own root node within its group, found by id rather than -/// "no parent" — `group_by_root` already keyed each group by its root's -/// id, so the root is just the member whose id matches the key. -fn root_in_group(id: u64, nodes: &[GraphNode]) -> Option<&GraphNode> { - nodes.iter().find(|n| n.id == id) -} - -/// Non-TTY progress: print a fresh line whenever a node's rendered state +/// 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_nodes_plain(socket: &Path, ids: Vec) -> Result<()> { - let mut pending: BTreeSet = ids.into_iter().collect(); - let mut last: HashMap = HashMap::new(); +async fn wait_for_dags_plain(socket: &Path, ids: Vec) -> Result<()> { + let mut pending: std::collections::BTreeSet = ids.into_iter().collect(); + let mut last: std::collections::HashMap = std::collections::HashMap::new(); let mut failed: Vec = Vec::new(); while !pending.is_empty() { - let batch: Vec = pending.iter().copied().collect(); - let resp = crate::client::request( - socket, - hive_host_sock::HostRequest::QueueNodes { ids: batch }, - ) - .await - .with_context(|| format!("connect to daemon socket {}", socket.display()))?; - let groups = group_by_root(resp.nodes.unwrap_or_default()); for id in pending.clone() { - let nodes = groups.get(&id); - let root = nodes.and_then(|n| root_in_group(id, n)); - let (Some(nodes), Some(root)) = (nodes, root) else { - // Unknown id — nothing prunes the graph yet (see - // `JobQueue::node_subtrees`'s doc comment), so an id that - // resolves to no node never named a real job. A - // *completed* job's nodes keep riding here instead, with a - // terminal root `state`, which is what the check below - // watches for. + let resp = crate::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 line = render_node_line(root, nodes); - if last.get(&id) != Some(&line) { - println!("{line}"); - last.insert(id, line); } - if root.state.is_terminal() { - if root.state == State::Failed { - failed.push(format!("{} {}", root.payload.label, node_agents(nodes))); + 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.rollup_state() == hive_host_sock::jobs::State::Failed { + failed.push(format!("{} {}", d.source.as_str(), dag_agents(d))); + } + } else { + all_terminal = false; + } + } + if all_terminal { pending.remove(&id); } } @@ -124,11 +85,11 @@ async fn wait_for_nodes_plain(socket: &Path, ids: Vec) -> Result<()> { } /// TTY progress: a live `indicatif` render — one braille-spinner line per -/// node (grouped under a per-root header), each with its own elapsed +/// 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_nodes_animated(socket: &Path, ids: Vec) -> Result<()> { +async fn wait_for_dags_animated(socket: &Path, ids: Vec) -> Result<()> { use indicatif::{MultiProgress, ProgressBar, ProgressStyle}; let spinner = ProgressStyle::with_template(" {spinner} {msg}") @@ -139,73 +100,77 @@ async fn wait_for_nodes_animated(socket: &Path, ids: Vec) -> Result<()> { let mp = MultiProgress::new(); let started = std::time::Instant::now(); - // Per-root header bars + per-node bars, keyed so we update in place. - // The insertion order groups a root's nodes right under its header. - let mut root_bars: HashMap = HashMap::new(); - let mut node_bars: HashMap<(u64, WireId), ProgressBar> = HashMap::new(); - let mut node_done: HashSet<(u64, WireId)> = HashSet::new(); + // 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 = + std::collections::HashMap::new(); + let mut node_bars: std::collections::HashMap<(u64, hive_host_sock::jobs::NodeId), ProgressBar> = + std::collections::HashMap::new(); + let mut node_done: std::collections::HashSet<(u64, hive_host_sock::jobs::NodeId)> = + std::collections::HashSet::new(); - let mut pending: BTreeSet = ids.into_iter().collect(); + let mut pending: std::collections::BTreeSet = ids.into_iter().collect(); let mut failed: Vec = Vec::new(); while !pending.is_empty() { let now = now_unix(); - let batch: Vec = pending.iter().copied().collect(); - let resp = crate::client::request( - socket, - hive_host_sock::HostRequest::QueueNodes { ids: batch }, - ) - .await - .with_context(|| format!("connect to daemon socket {}", socket.display()))?; - let groups = group_by_root(resp.nodes.unwrap_or_default()); for id in pending.clone() { - let nodes = groups.get(&id); - let root = nodes.and_then(|n| root_in_group(id, n)); - let (Some(nodes), Some(root)) = (nodes, root) else { + let resp = crate::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 hdr = root_bars.entry(id).or_insert_with(|| { - let b = mp.add(ProgressBar::new_spinner()); - b.set_style(plain.clone()); - b - }); - hdr.set_message(format!( - "{} {} {} · {}", - state_glyph(root.state), - root.payload.label, - node_agents(nodes), - fmt_dur(node_elapsed(root, now)), - )); - for n in nodes.iter().filter(|n| n.id != root.id) { - let key = (id, n.id); - if node_done.contains(&key) { - continue; - } - let bar = node_bars.entry(key).or_insert_with(|| { + } + 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(spinner.clone()); - b.enable_steady_tick(std::time::Duration::from_millis(120)); + b.set_style(plain.clone()); b }); - if n.state.is_terminal() { - bar.set_style(plain.clone()); - bar.finish_with_message(format!( - " {} {}", - state_glyph(n.state), - node_line(nodes, n, now) - )); - node_done.insert(key); + hdr.set_message(format!( + "{} {} {} · {}", + state_glyph(d.rollup_state()), + d.source.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.rollup_state() == hive_host_sock::jobs::State::Failed { + failed.push(format!("{} {}", d.source.as_str(), dag_agents(d))); + } } else { - bar.set_message(node_line(nodes, n, now)); + all_terminal = false; } } - if root.state.is_terminal() { - if root.state == State::Failed { - failed.push(format!("{} {}", root.payload.label, node_agents(nodes))); - } + if all_terminal { pending.remove(&id); } } @@ -215,7 +180,7 @@ async fn wait_for_nodes_animated(socket: &Path, ids: Vec) -> Result<()> { } // 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 root_bars.values() { + for hdr in dag_bars.values() { hdr.finish(); } mp.println(format!( @@ -237,21 +202,14 @@ fn finish_wait(mut failed: Vec) -> Result<()> { } } -/// Distinct agents across a root's nodes, comma-joined for display. Each -/// node's `agent` (when it targets one) rides in `payload.data["agent"]` — -/// an opaque kvp, not a typed field, since `GraphNode` carries nothing -/// domain-specific (see `hive_jobq_wire::WireNode::data`'s doc comment). -fn node_agents(nodes: &[GraphNode]) -> String { +/// 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_host_sock::jobs::DagView) -> String { let mut seen: Vec<&str> = Vec::new(); - for n in nodes { - if let Some(agent) = n - .payload - .data - .get("agent") - .and_then(serde_json::Value::as_str) - && !seen.contains(&agent) - { - seen.push(agent); + for n in &d.nodes { + if !seen.contains(&n.agent.as_str()) { + seen.push(&n.agent); } } seen.join(",") @@ -266,13 +224,23 @@ fn now_unix() -> i64 { .unwrap_or(0) } -/// Elapsed seconds for a node: `started_at` → `finished_at`/`now` once -/// it's run; `created_at` → `now` while it's still queued. `GraphNode` -/// carries no group-level timestamp separate from its own, so a header -/// line just reads this off whichever node it's summarizing. -fn node_elapsed(n: &GraphNode, now: i64) -> i64 { - let start = n.started_at.unwrap_or(n.created_at); - (n.finished_at.map_or(now, |t| t.timestamp()) - start.timestamp()).max(0) +/// Elapsed seconds for a DAG: `started_at` (falling back to `created_at`) +/// through `finished_at` or `now`. The wire carries these as RFC3339 +/// `DateTime`; compare in unix seconds against `now`. +fn dag_elapsed(d: &hive_host_sock::jobs::DagView, now: i64) -> i64 { + let start = d + .started_at + .map_or_else(|| d.created_at.timestamp(), |t| t.timestamp()); + (d.finished_at.map_or(now, |t| t.timestamp()) - 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_host_sock::jobs::NodeView, now: i64) -> i64 { + match n.started_at { + Some(start) => (n.finished_at.map_or(now, |t| t.timestamp()) - start.timestamp()).max(0), + None => 0, + } } /// Compact duration: `45s` under a minute, else `1m03s`. @@ -286,26 +254,20 @@ fn fmt_dur(secs: i64) -> String { } /// One animated node line: kind, an `(after …)` marker for a fan-in node -/// (>1 node-dependency), its elapsed timer, and a truncated error tail. -/// Resource deps (`GraphDep::Resource`) don't name another node, so they're -/// filtered out of the fan-in count the same way `hive-c0re`'s own -/// `dag_view` projection already does. -fn node_line(nodes: &[GraphNode], n: &GraphNode, now: i64) -> String { +/// (>1 dep), its elapsed timer, and a truncated error tail. +fn node_line( + d: &hive_host_sock::jobs::DagView, + n: &hive_host_sock::jobs::NodeView, + now: i64, +) -> String { use std::fmt::Write as _; - let mut s = n.payload.label.clone(); - let dep_ids: Vec = n - .deps - .iter() - .filter_map(|d| match d { - GraphDep::Node { id, .. } => Some(*id), - GraphDep::Resource { .. } => None, - }) - .collect(); - if dep_ids.len() > 1 { - let after: Vec<&str> = dep_ids + let mut s = n.kind.clone(); + if n.deps.len() > 1 { + let after: Vec<&str> = n + .deps .iter() - .filter_map(|id| nodes.iter().find(|m| m.id == *id)) - .map(|m| m.payload.label.as_str()) + .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(", ")); @@ -322,45 +284,38 @@ fn node_line(nodes: &[GraphNode], n: &GraphNode, now: i64) -> String { s } -fn state_glyph(state: State) -> &'static str { +fn state_glyph(state: hive_host_sock::jobs::State) -> &'static str { match state { - State::Pending => "⏸", + hive_host_sock::jobs::State::Pending => "⏸", // `Finishing` is own-work-done with sub-nodes still going — in flight, // so it reads the same as running. - State::Running | State::Finishing => "▶", - State::Done => "✔", - State::Failed => "✖", - State::Cancelled => "⊘", + hive_host_sock::jobs::State::Running | hive_host_sock::jobs::State::Finishing => "▶", + hive_host_sock::jobs::State::Done => "✔", + hive_host_sock::jobs::State::Failed => "✖", + hive_host_sock::jobs::State::Cancelled => "⊘", // Distinct from cancelled: nothing went wrong, this branch just // wasn't the one the run took. - State::Skipped => "·", + hive_host_sock::jobs::State::Skipped => "·", } } -/// One progress line for a root: roll-up glyph, label, agents, then the -/// node chain — the CLI twin of the dashboard's queue card. The glyph and -/// label come off `root`; everything else comes straight off the wire. -/// Used by the plain (non-TTY) path. -/// -/// The chain shows every entry `wire_snapshot` sends, `Done` ones -/// included — that filtering was a dashboard-history-bounding concern, -/// not relevant to a single actively-watched job, so a completed step -/// keeps its checkmark instead of vanishing from the line, matching how -/// the animated path already behaves. -fn render_node_line(root: &GraphNode, nodes: &[GraphNode]) -> String { +/// One progress line for a DAG: roll-up glyph, `source`, agents, then the +/// node chain — the CLI twin of the dashboard's queue card. The header shows +/// what the backend sends (`source` + the raw node kinds); only the roll-up +/// state glyph is derived from the node set. Used by the plain (non-TTY) path. +fn render_dag_line(d: &hive_host_sock::jobs::DagView) -> String { use std::fmt::Write as _; let mut out = format!( "{} {} {:<12}", - state_glyph(root.state), - root.payload.label, - node_agents(nodes) + state_glyph(d.rollup_state()), + d.source.as_str(), + dag_agents(d) ); - let children: Vec<&GraphNode> = nodes.iter().filter(|n| n.id != root.id).collect(); - for (i, n) in children.iter().enumerate() { + for (i, n) in d.nodes.iter().enumerate() { let sep = if i == 0 { " " } else { " → " }; - let _ = write!(out, "{sep}{} {}", state_glyph(n.state), n.payload.label); + let _ = write!(out, "{sep}{} {}", state_glyph(n.state), n.kind); } - if let Some(err) = children.iter().find_map(|n| n.error.as_deref()) { + 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}"); } @@ -369,62 +324,50 @@ fn render_node_line(root: &GraphNode, nodes: &[GraphNode]) -> String { #[cfg(test)] mod tests { - use hive_host_sock::jobs::State; - use hive_jobq_wire::{GraphNode, NodePayload}; - use serde_json::json; + use hive_host_sock::jobs::{DagView, NodeView, Source, State}; + use hive_sh4re::wire_time::from_secs; - use super::render_node_line; + use super::render_dag_line; - fn root_node(id: u64, label: &str, state: State) -> GraphNode { - GraphNode { + fn node(id: u64, agent: &str, kind: &str, state: State) -> NodeView { + NodeView { id, parent: None, + agent: agent.to_owned(), + kind: kind.to_owned(), + deps: if id == 0 { vec![] } else { vec![id - 1] }, state, - deps: Vec::new(), - created_at: hive_sh4re::wire_time::from_secs(0), started_at: None, finished_at: None, error: None, - payload: NodePayload { - label: label.to_owned(), - data: json!({}), - }, - } - } - - fn work_node(id: u64, root: u64, agent: &str, label: &str, state: State) -> GraphNode { - GraphNode { - id, - parent: Some(root), - state, - deps: Vec::new(), - created_at: hive_sh4re::wire_time::from_secs(0), - started_at: None, - finished_at: None, - error: None, - payload: NodePayload { - label: label.to_owned(), - data: json!({ "agent": agent }), - }, + approval_id: None, + inputs: vec![], + build_log_id: None, } } #[test] - fn render_node_line_shows_label_and_chain() { - // The header shows what the backend sends: the glyph off `root`'s - // own `state` (`Running` here) + `root`'s own `payload.label` - // ("manual"); the operation is read off the node chain, not any - // extra metadata. `Done` nodes stay in the chain — nothing here - // filters them off. - let root = root_node(7, "manual", State::Running); - let nodes = vec![ - root.clone(), - work_node(0, 7, "alice", "prebuild", State::Done), - work_node(1, 7, "alice", "stop_for_update", State::Done), - work_node(2, 7, "alice", "swap", State::Running), - work_node(3, 7, "alice", "reconcile", State::Pending), - ]; - let line = render_node_line(&root, &nodes); + fn render_dag_line_shows_source_and_chain() { + // The header shows what the backend sends: the roll-up state glyph + // (derived — Running here) + the DAG `source` ("manual"); the operation + // is read off the node chain, not a client-side label. (`Done` nodes + // are included here to exercise glyph rendering; production filters + // them off.) + let dag = DagView { + id: 7, + source: Source::Manual, + reason: "manual".to_owned(), + created_at: from_secs(0), + started_at: Some(from_secs(1)), + finished_at: None, + nodes: vec![ + node(0, "alice", "prebuild", State::Done), + node(1, "alice", "stop_for_update", State::Done), + node(2, "alice", "swap", State::Running), + node(3, "alice", "reconcile", State::Pending), + ], + }; + let line = render_dag_line(&dag); assert!(line.starts_with("▶ manual alice"), "{line}"); assert!( line.contains("✔ prebuild → ✔ stop_for_update → ▶ swap → ⏸ reconcile"), @@ -433,32 +376,20 @@ mod tests { } #[test] - fn render_node_line_surfaces_first_node_error() { - let root = root_node(8, "manual", State::Failed); - let mut failed = work_node(0, 8, "bob", "prebuild", State::Failed); + fn render_dag_line_surfaces_first_node_error() { + let mut failed = node(0, "bob", "prebuild", State::Failed); failed.error = Some("nix build exploded".to_owned()); - let nodes = vec![root.clone(), failed]; - let line = render_node_line(&root, &nodes); + let dag = DagView { + id: 8, + source: Source::Manual, + reason: "manual".to_owned(), + created_at: from_secs(0), + started_at: Some(from_secs(1)), + finished_at: Some(from_secs(2)), + nodes: vec![failed], + }; + let line = render_dag_line(&dag); assert!(line.contains("✖ manual"), "{line}"); assert!(line.contains("— nix build exploded"), "{line}"); } - - #[test] - fn group_by_root_splits_a_combined_batch_response() { - // Two independent roots' subtrees riding the same QueueNodes - // response (the whole point of batching): each node must land in - // its own root's group, not get mixed into the other's. - let nodes = vec![ - root_node(1, "manual", State::Running), - work_node(10, 1, "alice", "prebuild", State::Running), - root_node(2, "manual", State::Done), - work_node(20, 2, "bob", "prebuild", State::Done), - ]; - let groups = super::group_by_root(nodes); - assert_eq!(groups.len(), 2); - assert_eq!(groups[&1].len(), 2, "root 1's group must have its child"); - assert_eq!(groups[&2].len(), 2, "root 2's group must have its child"); - assert!(groups[&1].iter().any(|n| n.id == 10)); - assert!(groups[&2].iter().any(|n| n.id == 20)); - } } diff --git a/hivectl/src/main.rs b/hivectl/src/main.rs index 95c4751b..8c68ac70 100644 --- a/hivectl/src/main.rs +++ b/hivectl/src/main.rs @@ -24,7 +24,7 @@ mod cli; /// The host admin socket client (`request`), split out so it lives with /// hivectl rather than in the daemon crate. mod client; -/// Rebuild-queue node progress rendering (`wait_for_nodes` + the spinner / +/// Rebuild-queue DAG progress rendering (`wait_for_dags` + the spinner / /// plain renderers), split out to keep this file manageable. mod dag_progress; use cli::{Cli, Cmd, ForgeCmd, GatewayCmd, GithubCmd, WgCmd}; diff --git a/hivectl/src/power.rs b/hivectl/src/power.rs index 4966e248..f082c529 100644 --- a/hivectl/src/power.rs +++ b/hivectl/src/power.rs @@ -5,7 +5,7 @@ use std::path::Path; use anyhow::{Context as _, Result}; -use crate::dag_progress::wait_for_nodes; +use crate::dag_progress::wait_for_dags; use crate::util::render_lifecycle; pub(crate) async fn stop( @@ -24,7 +24,7 @@ pub(crate) async fn stop( // watch the already-queued agent DAGs before surfacing the error — // they run regardless. let rendered = render_lifecycle(&resp, "stop queued"); - wait_for_nodes(socket, resp.queued_dags.unwrap_or_default(), no_wait).await?; + wait_for_dags(socket, resp.queued_dags.unwrap_or_default(), no_wait).await?; rendered } @@ -37,7 +37,7 @@ pub(crate) async fn start( .await .with_context(|| format!("connect to daemon socket {}", socket.display()))?; let rendered = render_lifecycle(&resp, "start queued"); - wait_for_nodes(socket, resp.queued_dags.unwrap_or_default(), no_wait).await?; + wait_for_dags(socket, resp.queued_dags.unwrap_or_default(), no_wait).await?; rendered } @@ -66,6 +66,6 @@ pub(crate) async fn restart( .await .with_context(|| format!("connect to daemon socket {}", socket.display()))?; let rendered = render_lifecycle(&resp, "restart queued"); - wait_for_nodes(socket, resp.queued_dags.unwrap_or_default(), false).await?; + wait_for_dags(socket, resp.queued_dags.unwrap_or_default(), false).await?; rendered } diff --git a/hivectl/src/subvol.rs b/hivectl/src/subvol.rs index f6740fb8..b1427732 100644 --- a/hivectl/src/subvol.rs +++ b/hivectl/src/subvol.rs @@ -7,7 +7,7 @@ use std::path::Path; use anyhow::{Context as _, Result, bail}; use crate::cli::{SnapshotCmd, SubvolCmd}; -use crate::dag_progress::wait_for_nodes; +use crate::dag_progress::wait_for_dags; use crate::util::{agent_exists, daemon_request}; /// A [`LifecycleScope`](hive_host_sock::LifecycleScope) targeting exactly one @@ -78,10 +78,10 @@ async fn subvol_upgrade(socket: &Path, name: &str, yes: bool) -> Result<()> { stop_resp.error.as_deref().unwrap_or("unknown error") ); } - // The stop is a queued job now — the migration below snapshots + + // The stop is a queued DAG now — the migration below snapshots + // swaps the state dir and MUST NOT run under a live bind mount, so // wait for the stop to actually execute before touching anything. - wait_for_nodes(socket, stop_resp.queued_dags.unwrap_or_default(), false) + wait_for_dags(socket, stop_resp.queued_dags.unwrap_or_default(), false) .await .with_context(|| format!("waiting for {name} to stop before the migration"))?; @@ -130,7 +130,7 @@ async fn subvol_upgrade(socket: &Path, name: &str, yes: bool) -> Result<()> { start_resp.error.as_deref().unwrap_or("unknown error") ); } - wait_for_nodes(socket, start_resp.queued_dags.unwrap_or_default(), false) + wait_for_dags(socket, start_resp.queued_dags.unwrap_or_default(), false) .await .with_context(|| { format!(