address review: switch hive-c0re to hive-jobq-wire's shared parse_states/filter_nodes_by_state, note the generic shape in endpoint docs
This commit is contained in:
parent
08efd7875e
commit
eae04ac2c5
3 changed files with 29 additions and 108 deletions
|
|
@ -580,35 +580,23 @@ fn build_approval_views(approvals: Vec<Approval>) -> Vec<ApprovalView> {
|
||||||
out
|
out
|
||||||
}
|
}
|
||||||
|
|
||||||
/// `/api/jobq/graph` query string. Today's only field is `states`: a
|
/// `/api/jobq/graph` query string — the generic jobq/wire query shape (any
|
||||||
/// comma-separated allow-list of `hive_jobq::State` names (`"Pending"`,
|
/// host serving a `GraphWire` projection over HTTP takes the same one; see
|
||||||
/// `"Running"`, ...). Empty / absent ⇒ no filter (current behaviour, every
|
/// `hive_jobq_wire::parse_states`, which does the actual parsing this side
|
||||||
/// visible root). Set ⇒ only **root** groups whose own state is named are
|
/// of the query string). Today's only field is `states`: a comma-separated
|
||||||
/// served — a root's state is already its subtree's rolled-up answer (see
|
/// allow-list of `hive_jobq::State` names (`"Pending"`, `"Running"`, ...).
|
||||||
/// `hive_jobq_wire`'s doc), so filtering the root filters the whole group.
|
/// Empty / absent ⇒ no filter (current behaviour, every visible root). Set
|
||||||
/// Unknown tokens are silently ignored (an unrecognised name matches
|
/// ⇒ only **root** groups whose own state is named are served — a root's
|
||||||
/// nothing rather than erroring the whole request), mirroring
|
/// state is already its subtree's rolled-up answer (see `hive_jobq_wire`'s
|
||||||
/// `DashboardStreamQuery::kinds` above.
|
/// doc), so filtering the root filters the whole group. Unknown tokens are
|
||||||
|
/// silently ignored (an unrecognised name matches nothing rather than
|
||||||
|
/// erroring the whole request), mirroring `DashboardStreamQuery::kinds`
|
||||||
|
/// above.
|
||||||
#[derive(Deserialize, Default, IntoParams)]
|
#[derive(Deserialize, Default, IntoParams)]
|
||||||
pub(super) struct JobqGraphQuery {
|
pub(super) struct JobqGraphQuery {
|
||||||
states: Option<String>,
|
states: Option<String>,
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Parses [`JobqGraphQuery::states`] into the list
|
|
||||||
/// [`crate::job_queue::JobQueue::graph_snapshot`] wants. `None` when absent
|
|
||||||
/// or when every token failed to parse — both mean "no filter" rather than
|
|
||||||
/// "match nothing", so an empty/garbled query reads as the unfiltered call
|
|
||||||
/// it replaces rather than an empty result set.
|
|
||||||
fn parse_states(raw: Option<&str>) -> Option<Vec<hive_jobq::State>> {
|
|
||||||
let states: Vec<hive_jobq::State> = raw?
|
|
||||||
.split(',')
|
|
||||||
.map(str::trim)
|
|
||||||
.filter(|s| !s.is_empty())
|
|
||||||
.filter_map(|s| serde_json::from_value(serde_json::Value::String(s.to_owned())).ok())
|
|
||||||
.collect();
|
|
||||||
(!states.is_empty()).then_some(states)
|
|
||||||
}
|
|
||||||
|
|
||||||
#[utoipa::path(
|
#[utoipa::path(
|
||||||
get,
|
get,
|
||||||
path = "/api/jobq/graph",
|
path = "/api/jobq/graph",
|
||||||
|
|
@ -617,10 +605,14 @@ fn parse_states(raw: Option<&str>) -> Option<Vec<hive_jobq::State>> {
|
||||||
(status = 200, description = "every node of every retained job group, \
|
(status = 200, description = "every node of every retained job group, \
|
||||||
as generic `hive_jobq` graph nodes: identity, the parent tree, \
|
as generic `hive_jobq` graph nodes: identity, the parent tree, \
|
||||||
dependency edges with their accepted-outcome sets, lifecycle, and \
|
dependency edges with their accepted-outcome sets, lifecycle, and \
|
||||||
one opaque per-node payload. Group roots ride as ordinary nodes \
|
one opaque per-node payload. This is the generic jobq/wire shape \
|
||||||
(`parent: null`) and `Done` nodes are not filtered by default — a \
|
(`hive_jobq_wire::GraphWire::wire_snapshot`), not a hive-c0re-only \
|
||||||
consumer renders the graph without knowing what any node means. \
|
projection — a `swarm-controller` serving its own graph responds \
|
||||||
`?states=` narrows to root groups in the named states.",
|
with the same shape at the same path. Group roots ride as \
|
||||||
|
ordinary nodes (`parent: null`) and `Done` nodes are not \
|
||||||
|
filtered by default — a consumer renders the graph without \
|
||||||
|
knowing what any node means. `?states=` narrows to root groups \
|
||||||
|
in the named states.",
|
||||||
body = Vec<hive_jobq_wire::GraphNode>),
|
body = Vec<hive_jobq_wire::GraphNode>),
|
||||||
),
|
),
|
||||||
tag = "state_snapshot"
|
tag = "state_snapshot"
|
||||||
|
|
@ -629,7 +621,7 @@ pub(super) async fn jobq_graph(
|
||||||
State(state): State<AppState>,
|
State(state): State<AppState>,
|
||||||
axum::extract::Query(q): axum::extract::Query<JobqGraphQuery>,
|
axum::extract::Query(q): axum::extract::Query<JobqGraphQuery>,
|
||||||
) -> axum::Json<Vec<hive_jobq_wire::GraphNode>> {
|
) -> axum::Json<Vec<hive_jobq_wire::GraphNode>> {
|
||||||
let states = parse_states(q.states.as_deref());
|
let states = hive_jobq_wire::parse_states(q.states.as_deref());
|
||||||
axum::Json(state.coord.job_queue.graph_snapshot(states.as_deref()))
|
axum::Json(state.coord.job_queue.graph_snapshot(states.as_deref()))
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -639,7 +631,9 @@ pub(super) async fn jobq_graph(
|
||||||
responses(
|
responses(
|
||||||
(status = 200, description = "counts by lifecycle state over the same \
|
(status = 200, description = "counts by lifecycle state over the same \
|
||||||
groups `/api/jobq/graph` serves, as `(state, nodes, roots)` \
|
groups `/api/jobq/graph` serves, as `(state, nodes, roots)` \
|
||||||
triples. Every state is present, zero counts included, in a fixed \
|
triples — the generic jobq/wire roll-up shape \
|
||||||
|
(`hive_jobq_wire::state_rollup`), same as `/api/jobq/graph` \
|
||||||
|
above. Every state is present, zero counts included, in a fixed \
|
||||||
order — a consumer renders a summary (\"3 running · 2 queued\") \
|
order — a consumer renders a summary (\"3 running · 2 queued\") \
|
||||||
without fetching the graph and without re-deriving the tally. \
|
without fetching the graph and without re-deriving the tally. \
|
||||||
`roots` counts groups, `nodes` counts every step at any depth: one \
|
`roots` counts groups, `nodes` counts every step at any depth: one \
|
||||||
|
|
|
||||||
|
|
@ -351,15 +351,16 @@ impl JobQueue {
|
||||||
/// still-matching descendant one level higher rather than hiding or
|
/// still-matching descendant one level higher rather than hiding or
|
||||||
/// orphaning it.
|
/// orphaning it.
|
||||||
///
|
///
|
||||||
/// The projection itself is [`hive_jobq_wire`]'s; all this layer supplies
|
/// The projection and state filter are both [`hive_jobq_wire`]'s; this
|
||||||
/// is *which* nodes to show — see [`visible_roots`] for why the graph
|
/// layer only supplies *which roots* are in view (see [`visible_roots`]
|
||||||
/// can't decide the root-visibility half of that for itself.
|
/// for why the graph can't decide that itself) — an older, different
|
||||||
|
/// question than state filtering, and neither replaces the other.
|
||||||
#[must_use]
|
#[must_use]
|
||||||
pub fn graph_snapshot(&self, states: Option<&[State]>) -> Vec<GraphNode> {
|
pub fn graph_snapshot(&self, states: Option<&[State]>) -> Vec<GraphNode> {
|
||||||
let inner = self.lock();
|
let inner = self.lock();
|
||||||
let roots = visible_roots(&inner);
|
let roots = visible_roots(&inner);
|
||||||
let nodes = inner.graph().wire_snapshot(roots);
|
let nodes = inner.graph().wire_snapshot(roots);
|
||||||
filter_nodes_by_state(nodes, states)
|
hive_jobq_wire::filter_nodes_by_state(nodes, states)
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Per-state counts over the **same** groups [`JobQueue::graph_snapshot`]
|
/// Per-state counts over the **same** groups [`JobQueue::graph_snapshot`]
|
||||||
|
|
@ -414,26 +415,6 @@ fn find_node(sched: &Sched, id: u64) -> Option<NodeId> {
|
||||||
.find_map(|n| (n.id.get() == id).then_some(n.id))
|
.find_map(|n| (n.id.get() == id).then_some(n.id))
|
||||||
}
|
}
|
||||||
|
|
||||||
/// [`JobQueue::graph_snapshot`]'s `states` ask, applied to the already-
|
|
||||||
/// projected node list: keeps every node — root or descendant — whose own
|
|
||||||
/// `state` is named.
|
|
||||||
///
|
|
||||||
/// Applied *after* [`GraphWire::wire_snapshot`] rather than as a root
|
|
||||||
/// pre-filter — narrowing which roots are visible at all is
|
|
||||||
/// [`visible_roots`]'s job (a different question: how much settled work is
|
|
||||||
/// retained, full stop); this is "of what's retained and live, which
|
|
||||||
/// individual nodes does the caller want shown right now." `None` (or an
|
|
||||||
/// unrecognised/absent query) is the identity filter.
|
|
||||||
fn filter_nodes_by_state(nodes: Vec<GraphNode>, states: Option<&[State]>) -> Vec<GraphNode> {
|
|
||||||
let Some(states) = states else {
|
|
||||||
return nodes;
|
|
||||||
};
|
|
||||||
nodes
|
|
||||||
.into_iter()
|
|
||||||
.filter(|n| states.contains(&n.state))
|
|
||||||
.collect()
|
|
||||||
}
|
|
||||||
|
|
||||||
/// The visible **group** set for [`JobQueue::graph_snapshot`]: every live group
|
/// The visible **group** set for [`JobQueue::graph_snapshot`]: every live group
|
||||||
/// root, plus the newest [`MAX_HISTORY_DAGS`] settled ones.
|
/// root, plus the newest [`MAX_HISTORY_DAGS`] settled ones.
|
||||||
///
|
///
|
||||||
|
|
|
||||||
|
|
@ -1527,60 +1527,6 @@ fn graph_snapshot_states_filters_by_node_state() {
|
||||||
);
|
);
|
||||||
}
|
}
|
||||||
|
|
||||||
/// One hand-built node, bypassing the scheduler entirely — `filter_nodes_by_state`
|
|
||||||
/// is a pure `Vec<GraphNode> -> Vec<GraphNode>` transform, so this exercises it
|
|
||||||
/// directly rather than trying (and failing, per this module's own rule) to
|
|
||||||
/// drive a live group into a genuinely mixed state through the real queue.
|
|
||||||
fn node(id: u64, parent: Option<u64>, state: State) -> GraphNode {
|
|
||||||
GraphNode {
|
|
||||||
id,
|
|
||||||
parent,
|
|
||||||
state,
|
|
||||||
deps: Vec::new(),
|
|
||||||
created_at: Utc::now(),
|
|
||||||
started_at: None,
|
|
||||||
finished_at: None,
|
|
||||||
error: None,
|
|
||||||
payload: hive_jobq_wire::NodePayload {
|
|
||||||
label: id.to_string(),
|
|
||||||
data: serde_json::Value::Null,
|
|
||||||
},
|
|
||||||
}
|
|
||||||
}
|
|
||||||
|
|
||||||
/// The case `graph_snapshot_states_filters_by_node_state` above can't reach:
|
|
||||||
/// a still-live group (root `Running`) holding a mix of already-`Done` and
|
|
||||||
/// still-`Pending` steps. Filtering out `Done` must drop exactly the `Done`
|
|
||||||
/// node and nothing else — the root and the `Pending` sibling both stay,
|
|
||||||
/// even though the root itself isn't in the requested state set.
|
|
||||||
#[test]
|
|
||||||
fn filter_nodes_by_state_keeps_matching_nodes_from_a_mixed_state_tree() {
|
|
||||||
let nodes = vec![
|
|
||||||
node(1, None, State::Running), // root: whole group still live
|
|
||||||
node(2, Some(1), State::Done), // finished step, should be hidden
|
|
||||||
node(3, Some(1), State::Pending), // not-yet-run step, should stay
|
|
||||||
];
|
|
||||||
|
|
||||||
let mut kept: Vec<u64> =
|
|
||||||
filter_nodes_by_state(nodes.clone(), Some(&[State::Running, State::Pending]))
|
|
||||||
.into_iter()
|
|
||||||
.map(|n| n.id)
|
|
||||||
.collect();
|
|
||||||
kept.sort_unstable();
|
|
||||||
assert_eq!(
|
|
||||||
kept,
|
|
||||||
vec![1, 3],
|
|
||||||
"Done is filtered out even though its still-live parent (root) isn't"
|
|
||||||
);
|
|
||||||
|
|
||||||
let mut unfiltered: Vec<u64> = filter_nodes_by_state(nodes, None)
|
|
||||||
.into_iter()
|
|
||||||
.map(|n| n.id)
|
|
||||||
.collect();
|
|
||||||
unfiltered.sort_unstable();
|
|
||||||
assert_eq!(unfiltered, vec![1, 2, 3], "None is the identity filter");
|
|
||||||
}
|
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
fn error_truncation_cuts_on_a_char_boundary() {
|
fn error_truncation_cuts_on_a_char_boundary() {
|
||||||
// `truncate_error` is a pure `&str -> String`. This used to submit a DAG,
|
// `truncate_error` is a pure `&str -> String`. This used to submit a DAG,
|
||||||
|
|
|
||||||
Loading…
Reference in a new issue