diff --git a/hive-c0re/src/job_queue/mod.rs b/hive-c0re/src/job_queue/mod.rs index 7ca0ce7f..d4afeea0 100644 --- a/hive-c0re/src/job_queue/mod.rs +++ b/hive-c0re/src/job_queue/mod.rs @@ -379,30 +379,34 @@ impl JobQueue { .collect() } - /// A node plus its live subtree, 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 + /// 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 the id up by identity alone — - /// no assumption that it names a DAG container or a root; "just show - /// whatever the backend sends" for whatever id the caller asks about. + /// `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. /// - /// Empty when `id` names no node in the graph. 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 emptiness. + /// 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_subtree(&self, id: u64) -> Vec { + pub fn node_subtrees(&self, ids: &[u64]) -> Vec { let inner = self.lock(); - let Some(node) = find_node(&inner, id) else { - return Vec::new(); - }; - inner.graph().wire_snapshot([node]) + let roots: Vec = ids.iter().filter_map(|id| find_node(&inner, *id)).collect(); + inner.graph().wire_snapshot(roots) } } diff --git a/hive-c0re/src/server.rs b/hive-c0re/src/server.rs index 4d22b800..5427a7f6 100644 --- a/hive-c0re/src/server.rs +++ b/hive-c0re/src/server.rs @@ -160,8 +160,8 @@ async fn dispatch(req: &HostRequest, coord: Arc) -> HostResponse { .collect(); HostResponse::dags(dags) } - HostRequest::QueueNodes { id } => { - HostResponse::nodes(coord.job_queue.node_subtree(*id)) + 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 diff --git a/hive-host-sock/src/lib.rs b/hive-host-sock/src/lib.rs index 363a3442..caf7e3d9 100644 --- a/hive-host-sock/src/lib.rs +++ b/hive-host-sock/src/lib.rs @@ -214,15 +214,19 @@ 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 job-queue node plus its live subtree, as generic - /// `hive-jobq-wire` nodes — `hivectl`'s wait/progress loop. Sibling of - /// [`Self::QueueDag`]: same graph, same `id`, through the generic + /// 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 `id` names a DAG container or root — whatever node - /// has that id, the backend hands back its subtree as-is. Result: + /// 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 { id: u64 }, + QueueNodes { ids: Vec }, /// List pending approval requests. Pending, /// Approve a pending request by id; the action runs immediately. @@ -548,13 +552,14 @@ 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 node plus its live subtree, as - /// generic `hive-jobq-wire` nodes. `None` for every other request kind. - /// An empty `Vec` means `id` names no live node in the graph (today: an - /// unknown id — see `JobQueue::node_subtree`'s doc comment for why a - /// *completed* DAG's nodes don't vanish the same way `QueueDag`'s do); - /// callers should read the root node's `state` for terminality, not - /// emptiness. + /// `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 diff --git a/hivectl/src/agents.rs b/hivectl/src/agents.rs index c66f7e5d..776ace6f 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_dags; +use crate::dag_progress::wait_for_nodes; 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_dags(socket, resp.queued_dags.unwrap_or_default(), no_wait).await + wait_for_nodes(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 917ba347..7d137b52 100644 --- a/hivectl/src/dag_progress.rs +++ b/hivectl/src/dag_progress.rs @@ -1,10 +1,10 @@ //! `hivectl` rebuild-queue progress rendering. //! //! Split out of `hivectl.rs` (which is already large): everything that -//! polls the daemon's DAG queue (`HostRequest::QueueNodes`) 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. +//! 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 comes from `payload.label`, and node-specific extras @@ -14,23 +14,31 @@ //! `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. +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 DAG ids (`HostRequest::QueueNodes`, ~1s interval) and -/// render progress until they all reach a terminal state. Exits non-zero -/// (via the returned `Err`) when any DAG ends `failed`; a `cancelled` DAG -/// terminates the wait but is an operator action, not an error. +/// 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. /// /// # 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, no_wait: bool) -> Result<()> { +/// 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<()> { use std::io::IsTerminal as _; if no_wait || ids.is_empty() { return Ok(()); @@ -38,51 +46,71 @@ pub(crate) async fn wait_for_dags(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_dags_animated(socket, ids).await + wait_for_nodes_animated(socket, ids).await } else { - wait_for_dags_plain(socket, ids).await + wait_for_nodes_plain(socket, ids).await } } -/// Find the DAG container among a `QueueNodes` response's nodes — the one -/// `GraphNode` with no structural parent. `wire_snapshot` hands back exactly -/// one per requested root, so this is a lookup, not a real search. -fn find_root(nodes: &[GraphNode]) -> Option<&GraphNode> { - nodes.iter().find(|n| n.parent.is_none()) +/// 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 } -/// Non-TTY progress: print a fresh line whenever a DAG's rendered state +/// 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 /// changes. No cursor tricks, so it's clean in pipes and CI logs. -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(); +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(); 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 resp = - crate::client::request(socket, hive_host_sock::HostRequest::QueueNodes { id }) - .await - .with_context(|| format!("connect to daemon socket {}", socket.display()))?; - let nodes = resp.nodes.unwrap_or_default(); - let Some(root) = find_root(&nodes) else { + 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_subtree`'s doc comment), so an id that - // resolves to no node never named a real DAG. A - // *completed* DAG's nodes keep riding here instead, with a + // `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. println!("job #{id}: gone from queue history"); pending.remove(&id); continue; }; - let line = render_dag_line(root, &nodes); + 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!("{} {}", dag_source(root), dag_agents(&nodes))); + failed.push(format!("{} {}", node_source(root), node_agents(nodes))); } pending.remove(&id); } @@ -95,11 +123,11 @@ async fn wait_for_dags_plain(socket: &Path, ids: Vec) -> Result<()> { } /// TTY progress: a live `indicatif` render — one braille-spinner line per -/// DAG node (grouped under a per-DAG header), each with its own elapsed +/// node (grouped under a per-root 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) -> Result<()> { +async fn wait_for_nodes_animated(socket: &Path, ids: Vec) -> Result<()> { use indicatif::{MultiProgress, ProgressBar, ProgressStyle}; let spinner = ProgressStyle::with_template(" {spinner} {msg}") @@ -110,32 +138,35 @@ async fn wait_for_dags_animated(socket: &Path, ids: Vec) -> Result<()> { 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 = - std::collections::HashMap::new(); - let mut node_bars: std::collections::HashMap<(u64, WireId), ProgressBar> = - std::collections::HashMap::new(); - let mut node_done: std::collections::HashSet<(u64, WireId)> = std::collections::HashSet::new(); + // 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(); - let mut pending: std::collections::BTreeSet = ids.into_iter().collect(); + let mut pending: 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 resp = - crate::client::request(socket, hive_host_sock::HostRequest::QueueNodes { id }) - .await - .with_context(|| format!("connect to daemon socket {}", socket.display()))?; - let nodes = resp.nodes.unwrap_or_default(); - let Some(root) = find_root(&nodes) else { + let nodes = groups.get(&id); + let root = nodes.and_then(|n| root_in_group(id, n)); + let (Some(nodes), Some(root)) = (nodes, root) else { mp.println(format!("job #{id}: gone from queue history")) .ok(); pending.remove(&id); continue; }; - let hdr = dag_bars.entry(id).or_insert_with(|| { + let hdr = root_bars.entry(id).or_insert_with(|| { let b = mp.add(ProgressBar::new_spinner()); b.set_style(plain.clone()); b @@ -143,8 +174,8 @@ async fn wait_for_dags_animated(socket: &Path, ids: Vec) -> Result<()> { hdr.set_message(format!( "{} {} {} · {}", state_glyph(root.state), - dag_source(root), - dag_agents(&nodes), + node_source(root), + node_agents(nodes), fmt_dur(node_elapsed(root, now)), )); for n in nodes.iter().filter(|n| n.id != root.id) { @@ -163,16 +194,16 @@ async fn wait_for_dags_animated(socket: &Path, ids: Vec) -> Result<()> { bar.finish_with_message(format!( " {} {}", state_glyph(n.state), - node_line(&nodes, n, now) + node_line(nodes, n, now) )); node_done.insert(key); } else { - bar.set_message(node_line(&nodes, n, now)); + bar.set_message(node_line(nodes, n, now)); } } if root.state.is_terminal() { if root.state == State::Failed { - failed.push(format!("{} {}", dag_source(root), dag_agents(&nodes))); + failed.push(format!("{} {}", node_source(root), node_agents(nodes))); } pending.remove(&id); } @@ -183,7 +214,7 @@ async fn wait_for_dags_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 dag_bars.values() { + for hdr in root_bars.values() { hdr.finish(); } mp.println(format!( @@ -205,12 +236,11 @@ fn finish_wait(mut failed: Vec) -> Result<()> { } } -/// The DAG's `source` tag (`"manual"`, `"meta_update"`, …), read from the -/// container node's `payload.data` — see `hive-c0re`'s `impl WireNode for -/// NodeKind`'s `NodeKind::Dag` arm, the one place it's put on the wire. -/// Falls back to the label when absent (defensive; every real DAG -/// container sets it). -fn dag_source(root: &GraphNode) -> &str { +/// A root's `source` tag (`"manual"`, `"meta_update"`, …), read from its +/// `payload.data` — see `hive-c0re`'s `impl WireNode for NodeKind`'s +/// `NodeKind::Dag` arm, the one place it's put on the wire. Falls back to +/// the label when absent (defensive; every real root sets it). +fn node_source(root: &GraphNode) -> &str { root.payload .data .get("source") @@ -218,11 +248,11 @@ fn dag_source(root: &GraphNode) -> &str { .unwrap_or(&root.payload.label) } -/// Distinct agents across a DAG's nodes, comma-joined for display. Each +/// 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 dag_agents(nodes: &[GraphNode]) -> String { +fn node_agents(nodes: &[GraphNode]) -> String { let mut seen: Vec<&str> = Vec::new(); for n in nodes { if let Some(agent) = n @@ -318,24 +348,23 @@ fn state_glyph(state: State) -> &'static str { } } -/// One progress line for a DAG: roll-up glyph, `source`, agents, then the +/// One progress line for a root: roll-up glyph, `source`, agents, then the /// node chain — the CLI twin of the dashboard's queue card. The glyph and -/// `source` come off `root` (the one entry in `nodes` with no parent); -/// everything else comes straight off the wire. Used by the plain -/// (non-TTY) path. +/// `source` 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_dag_line(root: &GraphNode, nodes: &[GraphNode]) -> String { +fn render_node_line(root: &GraphNode, nodes: &[GraphNode]) -> String { use std::fmt::Write as _; let mut out = format!( "{} {} {:<12}", state_glyph(root.state), - dag_source(root), - dag_agents(nodes) + node_source(root), + node_agents(nodes) ); let children: Vec<&GraphNode> = nodes.iter().filter(|n| n.id != root.id).collect(); for (i, n) in children.iter().enumerate() { @@ -355,9 +384,9 @@ mod tests { use hive_jobq_wire::{GraphNode, NodePayload}; use serde_json::json; - use super::render_dag_line; + use super::render_node_line; - fn dag_root(id: u64, source: &str, state: State) -> GraphNode { + fn root_node(id: u64, source: &str, state: State) -> GraphNode { GraphNode { id, parent: None, @@ -368,7 +397,7 @@ mod tests { finished_at: None, error: None, payload: NodePayload { - label: "dag".to_owned(), + label: "job".to_owned(), data: json!({ "source": source }), }, } @@ -392,12 +421,12 @@ mod tests { } #[test] - fn render_dag_line_shows_source_and_chain() { + fn render_node_line_shows_source_and_chain() { // The header shows what the backend sends: the glyph off `root`'s // own `state` (`Running` here) + `source` ("manual"); the operation // is read off the node chain, not a client-side label. `Done` // nodes stay in the chain — nothing here filters them off. - let root = dag_root(7, "manual", State::Running); + let root = root_node(7, "manual", State::Running); let nodes = vec![ root.clone(), work_node(0, 7, "alice", "prebuild", State::Done), @@ -405,7 +434,7 @@ mod tests { work_node(2, 7, "alice", "swap", State::Running), work_node(3, 7, "alice", "reconcile", State::Pending), ]; - let line = render_dag_line(&root, &nodes); + let line = render_node_line(&root, &nodes); assert!(line.starts_with("▶ manual alice"), "{line}"); assert!( line.contains("✔ prebuild → ✔ stop_for_update → ▶ swap → ⏸ reconcile"), @@ -414,13 +443,32 @@ mod tests { } #[test] - fn render_dag_line_surfaces_first_node_error() { - let root = dag_root(8, "manual", State::Failed); + 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); failed.error = Some("nix build exploded".to_owned()); let nodes = vec![root.clone(), failed]; - let line = render_dag_line(&root, &nodes); + let line = render_node_line(&root, &nodes); 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 8c68ac70..95c4751b 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 DAG progress rendering (`wait_for_dags` + the spinner / +/// Rebuild-queue node progress rendering (`wait_for_nodes` + 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 f082c529..4966e248 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_dags; +use crate::dag_progress::wait_for_nodes; 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_dags(socket, resp.queued_dags.unwrap_or_default(), no_wait).await?; + wait_for_nodes(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_dags(socket, resp.queued_dags.unwrap_or_default(), no_wait).await?; + wait_for_nodes(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_dags(socket, resp.queued_dags.unwrap_or_default(), false).await?; + wait_for_nodes(socket, resp.queued_dags.unwrap_or_default(), false).await?; rendered } diff --git a/hivectl/src/subvol.rs b/hivectl/src/subvol.rs index b1427732..f6740fb8 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_dags; +use crate::dag_progress::wait_for_nodes; 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 DAG now — the migration below snapshots + + // The stop is a queued job 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_dags(socket, stop_resp.queued_dags.unwrap_or_default(), false) + wait_for_nodes(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_dags(socket, start_resp.queued_dags.unwrap_or_default(), false) + wait_for_nodes(socket, start_resp.queued_dags.unwrap_or_default(), false) .await .with_context(|| { format!(