hivectl: migrate dag_progress to hive-jobq-wire's generic GraphNode

This commit is contained in:
damocles 2026-08-03 18:47:22 +02:00 committed by mara
commit 10b0f640af
8 changed files with 272 additions and 175 deletions

2
Cargo.lock generated
View file

@ -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",

View file

@ -378,6 +378,30 @@ impl JobQueue {
.filter_map(|c| dag_view(&inner, c)) .filter_map(|c| dag_view(&inner, c))
.collect() .collect()
} }
/// One DAG's container 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 [`GraphWire::wire_snapshot`] instead — no `Done`-node
/// filtering, no roll-up field (the root's own `state` answers that,
/// see `hive_jobq_wire`'s doc comment).
///
/// Empty when `dag_id` names no DAG container 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.
#[must_use]
pub fn dag_nodes(&self, dag_id: u64) -> Vec<GraphNode> {
let inner = self.lock();
let Some(root) = container(&inner, dag_id) else {
return Vec::new();
};
inner.graph().wire_snapshot([root])
}
} }
/// The container node of `dag_id` — the `NodeKind::Dag` root whose id equals /// The container node of `dag_id` — the `NodeKind::Dag` root whose id equals

View file

@ -308,6 +308,15 @@ impl hive_jobq_wire::WireNode for NodeKind {
{ {
data.insert("inputs".to_owned(), inputs.clone().into()); data.insert("inputs".to_owned(), inputs.clone().into());
} }
// The DAG container's own metadata — nowhere else on the wire, since
// `GraphNode` carries no DAG-level fields (a group root is an
// ordinary node). `hivectl` needs `source` for its progress line;
// `reason` rides along for free rather than adding a second variant
// later for the one field the first pass missed.
if let NodeKind::Dag { source, reason, .. } = self {
data.insert("source".to_owned(), source.as_str().into());
data.insert("reason".to_owned(), reason.clone().into());
}
// Not in the payload at all — the build log is keyed on node identity // Not in the payload at all — the build log is keyed on node identity
// in a side table, which is why `data` is handed the id. // in a side table, which is why `data` is handed the id.
if let Some(log) = crate::build_logs::global().and_then(|h| h.id_for_node(id)) { if let Some(log) = crate::build_logs::global().and_then(|h| h.id_for_node(id)) {

View file

@ -160,6 +160,7 @@ async fn dispatch(req: &HostRequest, coord: Arc<Coordinator>) -> HostResponse {
.collect(); .collect();
HostResponse::dags(dags) HostResponse::dags(dags)
} }
HostRequest::QueueNodes { id } => HostResponse::nodes(coord.job_queue.dag_nodes(*id)),
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 —

View file

@ -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

View file

@ -214,6 +214,14 @@ 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 job-queue DAG's container node plus its live subtree, as
/// generic `hive-jobq-wire` nodes — `hivectl`'s wait/progress loop.
/// Sibling of [`Self::QueueDag`]: same DAG, same
/// `id` (the container's own node id, what `queued_dags` already
/// carries), through the generic projection instead of the typed
/// `DagView`/`NodeView` (kept for `QueueDag`'s other consumer,
/// `/api/state.rebuild_queue`). Result: [`HostResponse::nodes`].
QueueNodes { id: 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 +547,15 @@ 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 DAG's container 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 DAG in
/// the graph (today: an unknown id — see `JobQueue::dag_nodes`'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.
#[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 +656,17 @@ impl HostResponse {
} }
} }
/// `QueueNodes` result — the polled DAG's container + 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]

View file

@ -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

View file

@ -1,20 +1,30 @@
//! `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 DAG queue (`HostRequest::QueueNodes`) and renders the
//! per-DAG / per-node progress lives here. [`wait_for_dags`] is the entry //! per-DAG / per-node progress lives here. [`wait_for_dags`] is the entry
//! point the command handlers call; it dispatches to a live `indicatif` //! point the command handlers call; it dispatches to a live `indicatif`
//! animation on a TTY and a plain line-on-change stream otherwise. //! 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
//! (`source`, `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.
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 DAG ids (`HostRequest::QueueNodes`, ~1s interval) and
/// render progress until they all reach a terminal state. Exits non-zero /// 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`; /// (via the returned `Err`) when any DAG ends `failed`; a `cancelled` DAG
/// a `cancelled` DAG terminates the wait but is an operator action, not an /// terminates the wait but is an operator action, not an error.
/// error.
/// ///
/// # Errors /// # Errors
/// ///
@ -34,6 +44,13 @@ pub(crate) async fn wait_for_dags(socket: &Path, ids: Vec<u64>, no_wait: bool) -
} }
} }
/// 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())
}
/// Non-TTY progress: print a fresh line whenever a DAG'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. /// 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_dags_plain(socket: &Path, ids: Vec<u64>) -> Result<()> {
@ -42,38 +59,31 @@ async fn wait_for_dags_plain(socket: &Path, ids: Vec<u64>) -> Result<()> {
let mut failed: Vec<String> = Vec::new(); let mut failed: Vec<String> = Vec::new();
while !pending.is_empty() { while !pending.is_empty() {
for id in pending.clone() { for id in pending.clone() {
let resp = crate::client::request(socket, hive_host_sock::HostRequest::QueueDag { id }) let resp =
.await crate::client::request(socket, hive_host_sock::HostRequest::QueueNodes { id })
.with_context(|| format!("connect to daemon socket {}", socket.display()))?; .await
let dags = resp.dags.unwrap_or_default(); .with_context(|| format!("connect to daemon socket {}", socket.display()))?;
if dags.is_empty() { let nodes = resp.nodes.unwrap_or_default();
// Evicted from the queue's history tail — it finished a let Some(root) = find_root(&nodes) else {
// while ago; nothing left to report on. // Unknown id — nothing prunes the graph yet (see
// `JobQueue::dag_nodes`'s doc comment), so an id that
// resolves to no container never named a real DAG. A
// *completed* DAG'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_dag_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!("{} {}", dag_source(root), dag_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);
} }
} }
@ -104,10 +114,9 @@ async fn wait_for_dags_animated(socket: &Path, ids: Vec<u64>) -> Result<()> {
// insertion order groups a DAG's nodes right under its header. // insertion order groups a DAG's nodes right under its header.
let mut dag_bars: std::collections::HashMap<u64, ProgressBar> = let mut dag_bars: std::collections::HashMap<u64, ProgressBar> =
std::collections::HashMap::new(); std::collections::HashMap::new();
let mut node_bars: std::collections::HashMap<(u64, hive_host_sock::jobs::NodeId), ProgressBar> = let mut node_bars: std::collections::HashMap<(u64, WireId), ProgressBar> =
std::collections::HashMap::new(); std::collections::HashMap::new();
let mut node_done: std::collections::HashSet<(u64, hive_host_sock::jobs::NodeId)> = let mut node_done: std::collections::HashSet<(u64, WireId)> = std::collections::HashSet::new();
std::collections::HashSet::new();
let mut pending: std::collections::BTreeSet<u64> = ids.into_iter().collect(); let mut pending: std::collections::BTreeSet<u64> = ids.into_iter().collect();
let mut failed: Vec<String> = Vec::new(); let mut failed: Vec<String> = Vec::new();
@ -115,62 +124,56 @@ async fn wait_for_dags_animated(socket: &Path, ids: Vec<u64>) -> Result<()> {
while !pending.is_empty() { while !pending.is_empty() {
let now = now_unix(); let now = now_unix();
for id in pending.clone() { for id in pending.clone() {
let resp = crate::client::request(socket, hive_host_sock::HostRequest::QueueDag { id }) let resp =
.await crate::client::request(socket, hive_host_sock::HostRequest::QueueNodes { id })
.with_context(|| format!("connect to daemon socket {}", socket.display()))?; .await
let dags = resp.dags.unwrap_or_default(); .with_context(|| format!("connect to daemon socket {}", socket.display()))?;
if dags.is_empty() { let nodes = resp.nodes.unwrap_or_default();
let Some(root) = find_root(&nodes) else {
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 = dag_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),
dag_source(root),
dag_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!("{} {}", dag_source(root), dag_agents(&nodes)));
}
pending.remove(&id); pending.remove(&id);
} }
} }
@ -202,14 +205,34 @@ fn finish_wait(mut failed: Vec<String>) -> Result<()> {
} }
} }
/// Distinct agents across a DAG's nodes, comma-joined for display — the /// The DAG's `source` tag (`"manual"`, `"meta_update"`, …), read from the
/// per-node replacement for the old DAG-level `agent` field. Single-agent /// container node's `payload.data` — see `hive-c0re`'s `impl WireNode for
/// DAGs render one name; a hive-wide DAG lists each. /// NodeKind`'s `NodeKind::Dag` arm, the one place it's put on the wire.
fn dag_agents(d: &hive_host_sock::jobs::DagView) -> String { /// Falls back to the label when absent (defensive; every real DAG
/// container sets it).
fn dag_source(root: &GraphNode) -> &str {
root.payload
.data
.get("source")
.and_then(serde_json::Value::as_str)
.unwrap_or(&root.payload.label)
}
/// Distinct agents across a DAG'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 {
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 +247,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 +267,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 +303,46 @@ 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 DAG: roll-up glyph, `source`, 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 /// `source` come off `root` (the one entry in `nodes` with no parent);
/// state glyph is derived from the node set. Used by the plain (non-TTY) path. /// everything else comes straight off the wire. Used by the plain
fn render_dag_line(d: &hive_host_sock::jobs::DagView) -> String { /// (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 {
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(), dag_source(root),
dag_agents(d) dag_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 +351,61 @@ 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_dag_line;
fn node(id: u64, agent: &str, kind: &str, state: State) -> NodeView { fn dag_root(id: u64, source: &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: "dag".to_owned(),
build_log_id: None, data: json!({ "source": source }),
},
}
}
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_dag_line_shows_source_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) + `source` ("manual"); the operation
// is read off the node chain, not a client-side label. (`Done` nodes // is read off the node chain, not a client-side label. `Done`
// are included here to exercise glyph rendering; production filters // nodes stay in the chain — nothing here filters them off.
// them off.) let root = dag_root(7, "manual", State::Running);
let dag = DagView { let nodes = vec![
id: 7, root.clone(),
source: Source::Manual, work_node(0, 7, "alice", "prebuild", State::Done),
reason: "manual".to_owned(), work_node(1, 7, "alice", "stop_for_update", State::Done),
created_at: from_secs(0), work_node(2, 7, "alice", "swap", State::Running),
started_at: Some(from_secs(1)), work_node(3, 7, "alice", "reconcile", State::Pending),
finished_at: None, ];
nodes: vec![ let line = render_dag_line(&root, &nodes);
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.starts_with("▶ manual alice"), "{line}");
assert!( assert!(
line.contains("✔ prebuild → ✔ stop_for_update → ▶ swap → ⏸ reconcile"), line.contains("✔ prebuild → ✔ stop_for_update → ▶ swap → ⏸ reconcile"),
@ -377,18 +415,11 @@ mod tests {
#[test] #[test]
fn render_dag_line_surfaces_first_node_error() { fn render_dag_line_surfaces_first_node_error() {
let mut failed = node(0, "bob", "prebuild", State::Failed); let root = dag_root(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_dag_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}");
} }