Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
b04e7d985d | ||
|
|
0c4d56a585 | ||
|
|
309c1c91c6 | ||
|
|
10b0f640af |
11 changed files with 361 additions and 221 deletions
2
Cargo.lock
generated
2
Cargo.lock
generated
|
|
@ -1725,6 +1725,7 @@ version = "0.1.0"
|
||||||
dependencies = [
|
dependencies = [
|
||||||
"chrono",
|
"chrono",
|
||||||
"hive-jobq",
|
"hive-jobq",
|
||||||
|
"hive-jobq-wire",
|
||||||
"hive-sh4re",
|
"hive-sh4re",
|
||||||
"hive-types",
|
"hive-types",
|
||||||
"serde",
|
"serde",
|
||||||
|
|
@ -1861,6 +1862,7 @@ dependencies = [
|
||||||
"clap-markdown",
|
"clap-markdown",
|
||||||
"clap_complete",
|
"clap_complete",
|
||||||
"hive-host-sock",
|
"hive-host-sock",
|
||||||
|
"hive-jobq-wire",
|
||||||
"hive-sh4re",
|
"hive-sh4re",
|
||||||
"hive-types",
|
"hive-types",
|
||||||
"http-body-util",
|
"http-body-util",
|
||||||
|
|
|
||||||
|
|
@ -289,8 +289,8 @@ impl JobQueue {
|
||||||
#[must_use]
|
#[must_use]
|
||||||
pub fn first_error(&self, dag_id: u64) -> Option<String> {
|
pub fn first_error(&self, dag_id: u64) -> Option<String> {
|
||||||
let inner = self.lock();
|
let inner = self.lock();
|
||||||
let container = container(&inner, dag_id)?;
|
let node = find_node(&inner, dag_id)?;
|
||||||
inner.graph().first_error(container).map(ToOwned::to_owned)
|
inner.graph().first_error(node).map(ToOwned::to_owned)
|
||||||
}
|
}
|
||||||
|
|
||||||
/// `(agent, label, takes_container_down)` for the live transient-pill set,
|
/// `(agent, label, takes_container_down)` for the live transient-pill set,
|
||||||
|
|
@ -378,15 +378,45 @@ impl JobQueue {
|
||||||
.filter_map(|c| dag_view(&inner, c))
|
.filter_map(|c| dag_view(&inner, c))
|
||||||
.collect()
|
.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<GraphNode> {
|
||||||
|
let inner = self.lock();
|
||||||
|
let roots: Vec<NodeId> = ids.iter().filter_map(|id| find_node(&inner, *id)).collect();
|
||||||
|
inner.graph().wire_snapshot(roots)
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
/// The container node of `dag_id` — the `NodeKind::Dag` root whose id equals
|
/// The graph node whose id equals `id`, whatever its kind or depth.
|
||||||
/// `dag_id`. `NodeId` is un-fabricable from a raw `u64`, so this is a search.
|
/// `NodeId` is un-fabricable from a raw `u64`, so this is a search.
|
||||||
fn container(sched: &Sched, dag_id: u64) -> Option<NodeId> {
|
fn find_node(sched: &Sched, id: u64) -> Option<NodeId> {
|
||||||
sched.graph().nodes().find_map(|n| {
|
sched
|
||||||
(n.parent.is_none() && n.id.get() == dag_id && matches!(n.payload, NodeKind::Dag { .. }))
|
.graph()
|
||||||
.then_some(n.id)
|
.nodes()
|
||||||
})
|
.find_map(|n| (n.id.get() == id).then_some(n.id))
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Project a DAG into its wire [`DagView`]: a near-raw view of the
|
/// Project a DAG into its wire [`DagView`]: a near-raw view of the
|
||||||
|
|
|
||||||
|
|
@ -160,6 +160,9 @@ async fn dispatch(req: &HostRequest, coord: Arc<Coordinator>) -> HostResponse {
|
||||||
.collect();
|
.collect();
|
||||||
HostResponse::dags(dags)
|
HostResponse::dags(dags)
|
||||||
}
|
}
|
||||||
|
HostRequest::QueueNodes { ids } => {
|
||||||
|
HostResponse::nodes(coord.job_queue.node_subtrees(ids))
|
||||||
|
}
|
||||||
HostRequest::List => HostResponse::list(lifecycle::list().await?),
|
HostRequest::List => HostResponse::list(lifecycle::list().await?),
|
||||||
// The agents root is ours and not world-traversable, so this
|
// The agents root is ours and not world-traversable, so this
|
||||||
// question is only answerable on this side of the socket —
|
// question is only answerable on this side of the socket —
|
||||||
|
|
|
||||||
|
|
@ -10,6 +10,7 @@ workspace = true
|
||||||
[dependencies]
|
[dependencies]
|
||||||
chrono.workspace = true
|
chrono.workspace = true
|
||||||
hive-jobq.workspace = true
|
hive-jobq.workspace = true
|
||||||
|
hive-jobq-wire.workspace = true
|
||||||
hive-sh4re.workspace = true
|
hive-sh4re.workspace = true
|
||||||
hive-types.workspace = true
|
hive-types.workspace = true
|
||||||
serde.workspace = true
|
serde.workspace = true
|
||||||
|
|
|
||||||
|
|
@ -214,6 +214,19 @@ pub enum HostRequest {
|
||||||
/// `hivectl`'s wait/progress loop. A multi-step op is a single DAG
|
/// `hivectl`'s wait/progress loop. A multi-step op is a single DAG
|
||||||
/// (its whole graph in `nodes`). Result: [`HostResponse::dags`].
|
/// (its whole graph in `nodes`). Result: [`HostResponse::dags`].
|
||||||
QueueDag { id: u64 },
|
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<u64> },
|
||||||
/// List pending approval requests.
|
/// List pending approval requests.
|
||||||
Pending,
|
Pending,
|
||||||
/// Approve a pending request by id; the action runs immediately.
|
/// Approve a pending request by id; the action runs immediately.
|
||||||
|
|
@ -539,6 +552,16 @@ pub struct HostResponse {
|
||||||
/// been evicted from the queue's history tail.
|
/// been evicted from the queue's history tail.
|
||||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||||||
pub dags: Option<Vec<jobs::DagView>>,
|
pub dags: Option<Vec<jobs::DagView>>,
|
||||||
|
/// `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<Vec<hive_jobq_wire::GraphNode>>,
|
||||||
/// Free-form operator-facing output lines the client prints verbatim
|
/// Free-form operator-facing output lines the client prints verbatim
|
||||||
/// (one per line). Carries results a request produced daemon-side that
|
/// (one per line). Carries results a request produced daemon-side that
|
||||||
/// have no structured home — e.g. a freshly-minted matrix token, a
|
/// have no structured home — e.g. a freshly-minted matrix token, a
|
||||||
|
|
@ -639,6 +662,17 @@ impl HostResponse {
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// `QueueNodes` result — the polled node + its subtree, as generic
|
||||||
|
/// wire nodes.
|
||||||
|
#[must_use]
|
||||||
|
pub fn nodes(nodes: Vec<hive_jobq_wire::GraphNode>) -> Self {
|
||||||
|
Self {
|
||||||
|
ok: true,
|
||||||
|
nodes: Some(nodes),
|
||||||
|
..Self::default()
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
/// A success carrying operator-facing output lines the client prints
|
/// A success carrying operator-facing output lines the client prints
|
||||||
/// verbatim — the result shape for the `Matrix*` provisioning requests.
|
/// verbatim — the result shape for the `Matrix*` provisioning requests.
|
||||||
#[must_use]
|
#[must_use]
|
||||||
|
|
|
||||||
|
|
@ -21,6 +21,7 @@ clap.workspace = true
|
||||||
clap_complete.workspace = true
|
clap_complete.workspace = true
|
||||||
clap-markdown = "0.1"
|
clap-markdown = "0.1"
|
||||||
hive-host-sock.workspace = true
|
hive-host-sock.workspace = true
|
||||||
|
hive-jobq-wire.workspace = true
|
||||||
hive-sh4re.workspace = true
|
hive-sh4re.workspace = true
|
||||||
hive-types.workspace = true
|
hive-types.workspace = true
|
||||||
http-body-util.workspace = true
|
http-body-util.workspace = true
|
||||||
|
|
|
||||||
|
|
@ -13,7 +13,7 @@ use anyhow::{Context as _, Result, bail};
|
||||||
use hive_host_sock::HostRequest;
|
use hive_host_sock::HostRequest;
|
||||||
|
|
||||||
use crate::cli::{AgentCmd, AgentQuotaCmd};
|
use crate::cli::{AgentCmd, AgentQuotaCmd};
|
||||||
use crate::dag_progress::wait_for_dags;
|
use crate::dag_progress::wait_for_nodes;
|
||||||
use crate::util::render;
|
use crate::util::render;
|
||||||
|
|
||||||
async fn agents_restart(socket: &Path, name: &str, no_wait: bool) -> Result<()> {
|
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()))?;
|
.with_context(|| format!("connect to daemon socket {}", socket.display()))?;
|
||||||
if resp.ok {
|
if resp.ok {
|
||||||
println!("restart queued: {name}");
|
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 {
|
} else {
|
||||||
bail!(
|
bail!(
|
||||||
"restart {name}: {}",
|
"restart {name}: {}",
|
||||||
|
|
|
||||||
|
|
@ -1,26 +1,45 @@
|
||||||
//! `hivectl` rebuild-queue progress rendering.
|
//! `hivectl` rebuild-queue progress rendering.
|
||||||
//!
|
//!
|
||||||
//! Split out of `hivectl.rs` (which is already large): everything that
|
//! Split out of `hivectl.rs` (which is already large): everything that
|
||||||
//! polls the daemon's DAG queue (`HostRequest::QueueDag`) and renders the
|
//! polls the daemon's node queue (`HostRequest::QueueNodes`) and renders
|
||||||
//! per-DAG / per-node progress lives here. [`wait_for_dags`] is the entry
|
//! progress lives here. [`wait_for_nodes`] is the entry point the command
|
||||||
//! point the command handlers call; it dispatches to a live `indicatif`
|
//! handlers call; it dispatches to a live `indicatif` animation on a TTY
|
||||||
//! animation on a TTY and a plain line-on-change stream otherwise.
|
//! 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.
|
||||||
|
|
||||||
|
use std::collections::{BTreeSet, HashMap, HashSet};
|
||||||
use std::path::Path;
|
use std::path::Path;
|
||||||
|
|
||||||
use anyhow::{Context as _, Result, bail};
|
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::QueueDag`, ~1s interval) and
|
/// Poll the submitted ids (`HostRequest::QueueNodes`, ~1s interval, all
|
||||||
/// render progress until they all reach a terminal state. Exits non-zero
|
/// ids in one request per tick) and render progress until they all reach
|
||||||
/// (via the returned `Err`) when any DAG (or fan-out child) ends `failed`;
|
/// a terminal state. Exits non-zero (via the returned `Err`) when any
|
||||||
/// a `cancelled` DAG terminates the wait but is an operator action, not an
|
/// node ends `failed`; a `cancelled` node terminates the wait but is an
|
||||||
/// error.
|
/// operator action, not an error.
|
||||||
///
|
///
|
||||||
/// # Errors
|
/// # Errors
|
||||||
///
|
///
|
||||||
/// Returns an error if the daemon socket can't be reached, or if any polled
|
/// Returns an error if the daemon socket can't be reached, or if any
|
||||||
/// DAG finished in the `failed` state.
|
/// polled node finished in the `failed` state.
|
||||||
pub(crate) async fn wait_for_dags(socket: &Path, ids: Vec<u64>, no_wait: bool) -> Result<()> {
|
pub(crate) async fn wait_for_nodes(socket: &Path, ids: Vec<u64>, no_wait: bool) -> Result<()> {
|
||||||
use std::io::IsTerminal as _;
|
use std::io::IsTerminal as _;
|
||||||
if no_wait || ids.is_empty() {
|
if no_wait || ids.is_empty() {
|
||||||
return Ok(());
|
return Ok(());
|
||||||
|
|
@ -28,52 +47,72 @@ pub(crate) async fn wait_for_dags(socket: &Path, ids: Vec<u64>, no_wait: bool) -
|
||||||
// Animate only on a real terminal. Piped / CI output falls back to the
|
// 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.
|
// plain line-on-change stream so logs stay free of spinner redraw noise.
|
||||||
if std::io::stderr().is_terminal() {
|
if std::io::stderr().is_terminal() {
|
||||||
wait_for_dags_animated(socket, ids).await
|
wait_for_nodes_animated(socket, ids).await
|
||||||
} else {
|
} else {
|
||||||
wait_for_dags_plain(socket, ids).await
|
wait_for_nodes_plain(socket, ids).await
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Non-TTY progress: print a fresh line whenever a DAG's rendered state
|
/// 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<GraphNode>) -> HashMap<u64, Vec<GraphNode>> {
|
||||||
|
let mut groups: HashMap<u64, Vec<GraphNode>> = 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
|
||||||
/// changes. No cursor tricks, so it's clean in pipes and CI logs.
|
/// 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<()> {
|
async fn wait_for_nodes_plain(socket: &Path, ids: Vec<u64>) -> Result<()> {
|
||||||
let mut pending: std::collections::BTreeSet<u64> = ids.into_iter().collect();
|
let mut pending: BTreeSet<u64> = ids.into_iter().collect();
|
||||||
let mut last: std::collections::HashMap<u64, String> = std::collections::HashMap::new();
|
let mut last: HashMap<u64, String> = HashMap::new();
|
||||||
let mut failed: Vec<String> = Vec::new();
|
let mut failed: Vec<String> = Vec::new();
|
||||||
while !pending.is_empty() {
|
while !pending.is_empty() {
|
||||||
|
let batch: Vec<u64> = 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() {
|
for id in pending.clone() {
|
||||||
let resp = crate::client::request(socket, hive_host_sock::HostRequest::QueueDag { id })
|
let nodes = groups.get(&id);
|
||||||
.await
|
let root = nodes.and_then(|n| root_in_group(id, n));
|
||||||
.with_context(|| format!("connect to daemon socket {}", socket.display()))?;
|
let (Some(nodes), Some(root)) = (nodes, root) else {
|
||||||
let dags = resp.dags.unwrap_or_default();
|
// Unknown id — nothing prunes the graph yet (see
|
||||||
if dags.is_empty() {
|
// `JobQueue::node_subtrees`'s doc comment), so an id that
|
||||||
// Evicted from the queue's history tail — it finished a
|
// resolves to no node never named a real job. A
|
||||||
// while ago; nothing left to report on.
|
// *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");
|
println!("job #{id}: gone from queue history");
|
||||||
pending.remove(&id);
|
pending.remove(&id);
|
||||||
continue;
|
continue;
|
||||||
|
};
|
||||||
|
let line = render_node_line(root, nodes);
|
||||||
|
if last.get(&id) != Some(&line) {
|
||||||
|
println!("{line}");
|
||||||
|
last.insert(id, line);
|
||||||
}
|
}
|
||||||
let mut all_terminal = true;
|
if root.state.is_terminal() {
|
||||||
for d in &dags {
|
if root.state == State::Failed {
|
||||||
let line = render_dag_line(d);
|
failed.push(format!("{} {}", root.payload.label, node_agents(nodes)));
|
||||||
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);
|
pending.remove(&id);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
@ -85,11 +124,11 @@ async fn wait_for_dags_plain(socket: &Path, ids: Vec<u64>) -> Result<()> {
|
||||||
}
|
}
|
||||||
|
|
||||||
/// TTY progress: a live `indicatif` render — one braille-spinner line per
|
/// 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
|
/// timer, plus an overall-elapsed footer. A node with more than one
|
||||||
/// dependency (fan-in) gets its own row with an `(after …)` marker rather
|
/// dependency (fan-in) gets its own row with an `(after …)` marker rather
|
||||||
/// than being crammed onto a chain line.
|
/// than being crammed onto a chain line.
|
||||||
async fn wait_for_dags_animated(socket: &Path, ids: Vec<u64>) -> Result<()> {
|
async fn wait_for_nodes_animated(socket: &Path, ids: Vec<u64>) -> Result<()> {
|
||||||
use indicatif::{MultiProgress, ProgressBar, ProgressStyle};
|
use indicatif::{MultiProgress, ProgressBar, ProgressStyle};
|
||||||
|
|
||||||
let spinner = ProgressStyle::with_template(" {spinner} {msg}")
|
let spinner = ProgressStyle::with_template(" {spinner} {msg}")
|
||||||
|
|
@ -100,77 +139,73 @@ async fn wait_for_dags_animated(socket: &Path, ids: Vec<u64>) -> Result<()> {
|
||||||
|
|
||||||
let mp = MultiProgress::new();
|
let mp = MultiProgress::new();
|
||||||
let started = std::time::Instant::now();
|
let started = std::time::Instant::now();
|
||||||
// Per-DAG header bars + per-node bars, keyed so we update in place. The
|
// Per-root header bars + per-node bars, keyed so we update in place.
|
||||||
// insertion order groups a DAG's nodes right under its header.
|
// The insertion order groups a root's nodes right under its header.
|
||||||
let mut dag_bars: std::collections::HashMap<u64, ProgressBar> =
|
let mut root_bars: HashMap<u64, ProgressBar> = HashMap::new();
|
||||||
std::collections::HashMap::new();
|
let mut node_bars: HashMap<(u64, WireId), ProgressBar> = HashMap::new();
|
||||||
let mut node_bars: std::collections::HashMap<(u64, hive_host_sock::jobs::NodeId), ProgressBar> =
|
let mut node_done: HashSet<(u64, WireId)> = HashSet::new();
|
||||||
std::collections::HashMap::new();
|
|
||||||
let mut node_done: std::collections::HashSet<(u64, hive_host_sock::jobs::NodeId)> =
|
|
||||||
std::collections::HashSet::new();
|
|
||||||
|
|
||||||
let mut pending: std::collections::BTreeSet<u64> = ids.into_iter().collect();
|
let mut pending: BTreeSet<u64> = ids.into_iter().collect();
|
||||||
let mut failed: Vec<String> = Vec::new();
|
let mut failed: Vec<String> = Vec::new();
|
||||||
|
|
||||||
while !pending.is_empty() {
|
while !pending.is_empty() {
|
||||||
let now = now_unix();
|
let now = now_unix();
|
||||||
|
let batch: Vec<u64> = 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() {
|
for id in pending.clone() {
|
||||||
let resp = crate::client::request(socket, hive_host_sock::HostRequest::QueueDag { id })
|
let nodes = groups.get(&id);
|
||||||
.await
|
let root = nodes.and_then(|n| root_in_group(id, n));
|
||||||
.with_context(|| format!("connect to daemon socket {}", socket.display()))?;
|
let (Some(nodes), Some(root)) = (nodes, root) else {
|
||||||
let dags = resp.dags.unwrap_or_default();
|
|
||||||
if dags.is_empty() {
|
|
||||||
mp.println(format!("job #{id}: gone from queue history"))
|
mp.println(format!("job #{id}: gone from queue history"))
|
||||||
.ok();
|
.ok();
|
||||||
pending.remove(&id);
|
pending.remove(&id);
|
||||||
continue;
|
continue;
|
||||||
}
|
};
|
||||||
let mut all_terminal = true;
|
let hdr = root_bars.entry(id).or_insert_with(|| {
|
||||||
for d in &dags {
|
let b = mp.add(ProgressBar::new_spinner());
|
||||||
let hdr = dag_bars.entry(d.id).or_insert_with(|| {
|
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 b = mp.add(ProgressBar::new_spinner());
|
let b = mp.add(ProgressBar::new_spinner());
|
||||||
b.set_style(plain.clone());
|
b.set_style(spinner.clone());
|
||||||
|
b.enable_steady_tick(std::time::Duration::from_millis(120));
|
||||||
b
|
b
|
||||||
});
|
});
|
||||||
hdr.set_message(format!(
|
if n.state.is_terminal() {
|
||||||
"{} {} {} · {}",
|
bar.set_style(plain.clone());
|
||||||
state_glyph(d.rollup_state()),
|
bar.finish_with_message(format!(
|
||||||
d.source.as_str(),
|
" {} {}",
|
||||||
dag_agents(d),
|
state_glyph(n.state),
|
||||||
fmt_dur(dag_elapsed(d, now)),
|
node_line(nodes, n, now)
|
||||||
));
|
));
|
||||||
for n in &d.nodes {
|
node_done.insert(key);
|
||||||
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 {
|
} else {
|
||||||
all_terminal = false;
|
bar.set_message(node_line(nodes, n, now));
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
if all_terminal {
|
if root.state.is_terminal() {
|
||||||
|
if root.state == State::Failed {
|
||||||
|
failed.push(format!("{} {}", root.payload.label, node_agents(nodes)));
|
||||||
|
}
|
||||||
pending.remove(&id);
|
pending.remove(&id);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
@ -180,7 +215,7 @@ async fn wait_for_dags_animated(socket: &Path, ids: Vec<u64>) -> Result<()> {
|
||||||
}
|
}
|
||||||
// Leave the final node lines on screen; drop the still-spinning headers
|
// Leave the final node lines on screen; drop the still-spinning headers
|
||||||
// into their terminal state and print an overall-elapsed footer.
|
// into their terminal state and print an overall-elapsed footer.
|
||||||
for hdr in dag_bars.values() {
|
for hdr in root_bars.values() {
|
||||||
hdr.finish();
|
hdr.finish();
|
||||||
}
|
}
|
||||||
mp.println(format!(
|
mp.println(format!(
|
||||||
|
|
@ -202,14 +237,21 @@ fn finish_wait(mut failed: Vec<String>) -> Result<()> {
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Distinct agents across a DAG's nodes, comma-joined for display — the
|
/// Distinct agents across a root's nodes, comma-joined for display. Each
|
||||||
/// per-node replacement for the old DAG-level `agent` field. Single-agent
|
/// node's `agent` (when it targets one) rides in `payload.data["agent"]` —
|
||||||
/// DAGs render one name; a hive-wide DAG lists each.
|
/// an opaque kvp, not a typed field, since `GraphNode` carries nothing
|
||||||
fn dag_agents(d: &hive_host_sock::jobs::DagView) -> String {
|
/// domain-specific (see `hive_jobq_wire::WireNode::data`'s doc comment).
|
||||||
|
fn node_agents(nodes: &[GraphNode]) -> String {
|
||||||
let mut seen: Vec<&str> = Vec::new();
|
let mut seen: Vec<&str> = Vec::new();
|
||||||
for n in &d.nodes {
|
for n in nodes {
|
||||||
if !seen.contains(&n.agent.as_str()) {
|
if let Some(agent) = n
|
||||||
seen.push(&n.agent);
|
.payload
|
||||||
|
.data
|
||||||
|
.get("agent")
|
||||||
|
.and_then(serde_json::Value::as_str)
|
||||||
|
&& !seen.contains(&agent)
|
||||||
|
{
|
||||||
|
seen.push(agent);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
seen.join(",")
|
seen.join(",")
|
||||||
|
|
@ -224,23 +266,13 @@ fn now_unix() -> i64 {
|
||||||
.unwrap_or(0)
|
.unwrap_or(0)
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Elapsed seconds for a DAG: `started_at` (falling back to `created_at`)
|
/// Elapsed seconds for a node: `started_at` → `finished_at`/`now` once
|
||||||
/// through `finished_at` or `now`. The wire carries these as RFC3339
|
/// it's run; `created_at` → `now` while it's still queued. `GraphNode`
|
||||||
/// `DateTime<Utc>`; compare in unix seconds against `now`.
|
/// carries no group-level timestamp separate from its own, so a header
|
||||||
fn dag_elapsed(d: &hive_host_sock::jobs::DagView, now: i64) -> i64 {
|
/// line just reads this off whichever node it's summarizing.
|
||||||
let start = d
|
fn node_elapsed(n: &GraphNode, now: i64) -> i64 {
|
||||||
.started_at
|
let start = n.started_at.unwrap_or(n.created_at);
|
||||||
.map_or_else(|| d.created_at.timestamp(), |t| t.timestamp());
|
(n.finished_at.map_or(now, |t| t.timestamp()) - start.timestamp()).max(0)
|
||||||
(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`.
|
/// Compact duration: `45s` under a minute, else `1m03s`.
|
||||||
|
|
@ -254,20 +286,26 @@ fn fmt_dur(secs: i64) -> String {
|
||||||
}
|
}
|
||||||
|
|
||||||
/// One animated node line: kind, an `(after …)` marker for a fan-in node
|
/// One animated node line: kind, an `(after …)` marker for a fan-in node
|
||||||
/// (>1 dep), its elapsed timer, and a truncated error tail.
|
/// (>1 node-dependency), its elapsed timer, and a truncated error tail.
|
||||||
fn node_line(
|
/// Resource deps (`GraphDep::Resource`) don't name another node, so they're
|
||||||
d: &hive_host_sock::jobs::DagView,
|
/// filtered out of the fan-in count the same way `hive-c0re`'s own
|
||||||
n: &hive_host_sock::jobs::NodeView,
|
/// `dag_view` projection already does.
|
||||||
now: i64,
|
fn node_line(nodes: &[GraphNode], n: &GraphNode, now: i64) -> String {
|
||||||
) -> String {
|
|
||||||
use std::fmt::Write as _;
|
use std::fmt::Write as _;
|
||||||
let mut s = n.kind.clone();
|
let mut s = n.payload.label.clone();
|
||||||
if n.deps.len() > 1 {
|
let dep_ids: Vec<WireId> = n
|
||||||
let after: Vec<&str> = n
|
.deps
|
||||||
.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
|
||||||
.iter()
|
.iter()
|
||||||
.filter_map(|dep| d.nodes.iter().find(|m| m.id == *dep))
|
.filter_map(|id| nodes.iter().find(|m| m.id == *id))
|
||||||
.map(|m| m.kind.as_str())
|
.map(|m| m.payload.label.as_str())
|
||||||
.collect();
|
.collect();
|
||||||
if !after.is_empty() {
|
if !after.is_empty() {
|
||||||
let _ = write!(s, " (after {})", after.join(", "));
|
let _ = write!(s, " (after {})", after.join(", "));
|
||||||
|
|
@ -284,38 +322,45 @@ fn node_line(
|
||||||
s
|
s
|
||||||
}
|
}
|
||||||
|
|
||||||
fn state_glyph(state: hive_host_sock::jobs::State) -> &'static str {
|
fn state_glyph(state: State) -> &'static str {
|
||||||
match state {
|
match state {
|
||||||
hive_host_sock::jobs::State::Pending => "⏸",
|
State::Pending => "⏸",
|
||||||
// `Finishing` is own-work-done with sub-nodes still going — in flight,
|
// `Finishing` is own-work-done with sub-nodes still going — in flight,
|
||||||
// so it reads the same as running.
|
// so it reads the same as running.
|
||||||
hive_host_sock::jobs::State::Running | hive_host_sock::jobs::State::Finishing => "▶",
|
State::Running | State::Finishing => "▶",
|
||||||
hive_host_sock::jobs::State::Done => "✔",
|
State::Done => "✔",
|
||||||
hive_host_sock::jobs::State::Failed => "✖",
|
State::Failed => "✖",
|
||||||
hive_host_sock::jobs::State::Cancelled => "⊘",
|
State::Cancelled => "⊘",
|
||||||
// Distinct from cancelled: nothing went wrong, this branch just
|
// Distinct from cancelled: nothing went wrong, this branch just
|
||||||
// wasn't the one the run took.
|
// wasn't the one the run took.
|
||||||
hive_host_sock::jobs::State::Skipped => "·",
|
State::Skipped => "·",
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
/// One progress line for a DAG: roll-up glyph, `source`, agents, then the
|
/// 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 header shows
|
/// node chain — the CLI twin of the dashboard's queue card. The glyph and
|
||||||
/// what the backend sends (`source` + the raw node kinds); only the roll-up
|
/// label come off `root`; everything else comes straight off the wire.
|
||||||
/// state glyph is derived from the node set. Used by the plain (non-TTY) path.
|
/// Used by the plain (non-TTY) path.
|
||||||
fn render_dag_line(d: &hive_host_sock::jobs::DagView) -> String {
|
///
|
||||||
|
/// 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 {
|
||||||
use std::fmt::Write as _;
|
use std::fmt::Write as _;
|
||||||
let mut out = format!(
|
let mut out = format!(
|
||||||
"{} {} {:<12}",
|
"{} {} {:<12}",
|
||||||
state_glyph(d.rollup_state()),
|
state_glyph(root.state),
|
||||||
d.source.as_str(),
|
root.payload.label,
|
||||||
dag_agents(d)
|
node_agents(nodes)
|
||||||
);
|
);
|
||||||
for (i, n) in d.nodes.iter().enumerate() {
|
let children: Vec<&GraphNode> = nodes.iter().filter(|n| n.id != root.id).collect();
|
||||||
|
for (i, n) in children.iter().enumerate() {
|
||||||
let sep = if i == 0 { " " } else { " → " };
|
let sep = if i == 0 { " " } else { " → " };
|
||||||
let _ = write!(out, "{sep}{} {}", state_glyph(n.state), n.kind);
|
let _ = write!(out, "{sep}{} {}", state_glyph(n.state), n.payload.label);
|
||||||
}
|
}
|
||||||
if let Some(err) = d.nodes.iter().find_map(|n| n.error.as_deref()) {
|
if let Some(err) = children.iter().find_map(|n| n.error.as_deref()) {
|
||||||
let short: String = err.chars().take(120).collect();
|
let short: String = err.chars().take(120).collect();
|
||||||
let _ = write!(out, " — {short}");
|
let _ = write!(out, " — {short}");
|
||||||
}
|
}
|
||||||
|
|
@ -324,50 +369,62 @@ fn render_dag_line(d: &hive_host_sock::jobs::DagView) -> String {
|
||||||
|
|
||||||
#[cfg(test)]
|
#[cfg(test)]
|
||||||
mod tests {
|
mod tests {
|
||||||
use hive_host_sock::jobs::{DagView, NodeView, Source, State};
|
use hive_host_sock::jobs::State;
|
||||||
use hive_sh4re::wire_time::from_secs;
|
use hive_jobq_wire::{GraphNode, NodePayload};
|
||||||
|
use serde_json::json;
|
||||||
|
|
||||||
use super::render_dag_line;
|
use super::render_node_line;
|
||||||
|
|
||||||
fn node(id: u64, agent: &str, kind: &str, state: State) -> NodeView {
|
fn root_node(id: u64, label: &str, state: State) -> GraphNode {
|
||||||
NodeView {
|
GraphNode {
|
||||||
id,
|
id,
|
||||||
parent: None,
|
parent: None,
|
||||||
agent: agent.to_owned(),
|
|
||||||
kind: kind.to_owned(),
|
|
||||||
deps: if id == 0 { vec![] } else { vec![id - 1] },
|
|
||||||
state,
|
state,
|
||||||
|
deps: Vec::new(),
|
||||||
|
created_at: hive_sh4re::wire_time::from_secs(0),
|
||||||
started_at: None,
|
started_at: None,
|
||||||
finished_at: None,
|
finished_at: None,
|
||||||
error: None,
|
error: None,
|
||||||
approval_id: None,
|
payload: NodePayload {
|
||||||
inputs: vec![],
|
label: label.to_owned(),
|
||||||
build_log_id: None,
|
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 }),
|
||||||
|
},
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
fn render_dag_line_shows_source_and_chain() {
|
fn render_node_line_shows_label_and_chain() {
|
||||||
// The header shows what the backend sends: the roll-up state glyph
|
// The header shows what the backend sends: the glyph off `root`'s
|
||||||
// (derived — Running here) + the DAG `source` ("manual"); the operation
|
// own `state` (`Running` here) + `root`'s own `payload.label`
|
||||||
// is read off the node chain, not a client-side label. (`Done` nodes
|
// ("manual"); the operation is read off the node chain, not any
|
||||||
// are included here to exercise glyph rendering; production filters
|
// extra metadata. `Done` nodes stay in the chain — nothing here
|
||||||
// them off.)
|
// filters them off.
|
||||||
let dag = DagView {
|
let root = root_node(7, "manual", State::Running);
|
||||||
id: 7,
|
let nodes = vec![
|
||||||
source: Source::Manual,
|
root.clone(),
|
||||||
reason: "manual".to_owned(),
|
work_node(0, 7, "alice", "prebuild", State::Done),
|
||||||
created_at: from_secs(0),
|
work_node(1, 7, "alice", "stop_for_update", State::Done),
|
||||||
started_at: Some(from_secs(1)),
|
work_node(2, 7, "alice", "swap", State::Running),
|
||||||
finished_at: None,
|
work_node(3, 7, "alice", "reconcile", State::Pending),
|
||||||
nodes: vec![
|
];
|
||||||
node(0, "alice", "prebuild", State::Done),
|
let line = render_node_line(&root, &nodes);
|
||||||
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.starts_with("▶ manual alice"), "{line}");
|
||||||
assert!(
|
assert!(
|
||||||
line.contains("✔ prebuild → ✔ stop_for_update → ▶ swap → ⏸ reconcile"),
|
line.contains("✔ prebuild → ✔ stop_for_update → ▶ swap → ⏸ reconcile"),
|
||||||
|
|
@ -376,20 +433,32 @@ mod tests {
|
||||||
}
|
}
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
fn render_dag_line_surfaces_first_node_error() {
|
fn render_node_line_surfaces_first_node_error() {
|
||||||
let mut failed = node(0, "bob", "prebuild", State::Failed);
|
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());
|
failed.error = Some("nix build exploded".to_owned());
|
||||||
let dag = DagView {
|
let nodes = vec![root.clone(), failed];
|
||||||
id: 8,
|
let line = render_node_line(&root, &nodes);
|
||||||
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("✖ manual"), "{line}");
|
||||||
assert!(line.contains("— nix build exploded"), "{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));
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -24,7 +24,7 @@ mod cli;
|
||||||
/// The host admin socket client (`request`), split out so it lives with
|
/// The host admin socket client (`request`), split out so it lives with
|
||||||
/// hivectl rather than in the daemon crate.
|
/// hivectl rather than in the daemon crate.
|
||||||
mod client;
|
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.
|
/// plain renderers), split out to keep this file manageable.
|
||||||
mod dag_progress;
|
mod dag_progress;
|
||||||
use cli::{Cli, Cmd, ForgeCmd, GatewayCmd, GithubCmd, WgCmd};
|
use cli::{Cli, Cmd, ForgeCmd, GatewayCmd, GithubCmd, WgCmd};
|
||||||
|
|
|
||||||
|
|
@ -5,7 +5,7 @@ use std::path::Path;
|
||||||
|
|
||||||
use anyhow::{Context as _, Result};
|
use anyhow::{Context as _, Result};
|
||||||
|
|
||||||
use crate::dag_progress::wait_for_dags;
|
use crate::dag_progress::wait_for_nodes;
|
||||||
use crate::util::render_lifecycle;
|
use crate::util::render_lifecycle;
|
||||||
|
|
||||||
pub(crate) async fn stop(
|
pub(crate) async fn stop(
|
||||||
|
|
@ -24,7 +24,7 @@ pub(crate) async fn stop(
|
||||||
// watch the already-queued agent DAGs before surfacing the error —
|
// watch the already-queued agent DAGs before surfacing the error —
|
||||||
// they run regardless.
|
// they run regardless.
|
||||||
let rendered = render_lifecycle(&resp, "stop queued");
|
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
|
rendered
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -37,7 +37,7 @@ pub(crate) async fn start(
|
||||||
.await
|
.await
|
||||||
.with_context(|| format!("connect to daemon socket {}", socket.display()))?;
|
.with_context(|| format!("connect to daemon socket {}", socket.display()))?;
|
||||||
let rendered = render_lifecycle(&resp, "start queued");
|
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
|
rendered
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -66,6 +66,6 @@ pub(crate) async fn restart(
|
||||||
.await
|
.await
|
||||||
.with_context(|| format!("connect to daemon socket {}", socket.display()))?;
|
.with_context(|| format!("connect to daemon socket {}", socket.display()))?;
|
||||||
let rendered = render_lifecycle(&resp, "restart queued");
|
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
|
rendered
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -7,7 +7,7 @@ use std::path::Path;
|
||||||
use anyhow::{Context as _, Result, bail};
|
use anyhow::{Context as _, Result, bail};
|
||||||
|
|
||||||
use crate::cli::{SnapshotCmd, SubvolCmd};
|
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};
|
use crate::util::{agent_exists, daemon_request};
|
||||||
|
|
||||||
/// A [`LifecycleScope`](hive_host_sock::LifecycleScope) targeting exactly one
|
/// 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")
|
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
|
// 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 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
|
.await
|
||||||
.with_context(|| format!("waiting for {name} to stop before the migration"))?;
|
.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")
|
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
|
.await
|
||||||
.with_context(|| {
|
.with_context(|| {
|
||||||
format!(
|
format!(
|
||||||
|
|
|
||||||
Loading…
Reference in a new issue