diff --git a/hive-c0re/src/job_queue/mod.rs b/hive-c0re/src/job_queue/mod.rs index dfe916ec..2780ee6a 100644 --- a/hive-c0re/src/job_queue/mod.rs +++ b/hive-c0re/src/job_queue/mod.rs @@ -36,7 +36,7 @@ pub mod templates; #[cfg(test)] mod tests; -use std::sync::Mutex; +use std::sync::{Arc, Mutex}; use chrono::{DateTime, Utc}; use hive_host_sock::jobs::NodeView; @@ -109,27 +109,29 @@ struct DagMeta { created_at: DateTime, } -/// The mutable queue state behind the mutex: **just the crate scheduler**. +/// The crate scheduler, specialised to this host's node + resource types. +/// /// A **DAG is a single container node** ([`NodeKind::Dag`], `parent = None`) /// whose subtree is the DAG's work — so the container's `NodeId` is the DAG id, /// its rolled-up state is the DAG state, and there are no grouping side-tables: -/// membership + meta are graph queries ([`QueueInner::container`] / -/// [`QueueInner::dag_meta`] + the `hive_jobq::Graph` accessors). One shared -/// crate [`Graph`] holds every DAG. +/// membership + meta are graph queries ([`container`] / [`dag_meta`] + the +/// `hive_jobq::Graph` accessors). One shared crate [`Graph`] holds every DAG. /// -/// There is deliberately **no per-node side map** any more. The last one held -/// the `build_logs` row id; that link now lives on the log row itself -/// (`build_logs.node_id`), so it survives a restart and needs no lock held -/// alongside the scheduler's — which is what lets the scheduler's own lock be -/// the only one the run loop takes. -struct QueueInner { - sched: Scheduler, -} +/// There is deliberately **no wrapper struct and no per-node side map**. The +/// last map held the `build_logs` row id; that link now lives on the log row +/// itself (`build_logs.node_id`). With nothing else to guard, the mutex holds +/// the scheduler *directly* — which is what lets `hive_jobq` drive the run loop +/// (it takes `&Arc>>`, a type a host-side wrapper could not +/// satisfy). +type Sched = Scheduler; /// The queue. Lives on `Coordinator` (one per hive-c0re process); a single /// scheduler task ([`scheduler::run_worker`]) drives it. pub struct JobQueue { - inner: Mutex, + /// The scheduler, held directly rather than behind a host-side wrapper — + /// `hive_jobq`'s run-loop seam takes `&Arc>>`, so this + /// *is* the type the crate drives. + sched: Arc>, /// Wakes the scheduler when something new arrives or state changed. pub(crate) notify: Notify, } @@ -173,12 +175,11 @@ fn outcome_of(result: Result<(), String>) -> Outcome { /// # Errors /// Propagates a crate graph-insert error (malformed dep/parent / dep-scope). fn insert_group( - inner: &mut QueueInner, + inner: &mut Sched, declare: impl FnOnce(&Job), group_parent: Option, ) -> anyhow::Result<()> { inner - .sched .insert_job(group_parent, |b| { declare(b); // c0re names no handles: a DAG is addressed by its container node, @@ -199,15 +200,13 @@ impl JobQueue { u32::try_from(build_slots.max(1)).unwrap_or(u32::MAX), ); Self { - inner: Mutex::new(QueueInner { - sched: Scheduler::new(Graph::new(), table), - }), + sched: Arc::new(Mutex::new(Scheduler::new(Graph::new(), table))), notify: Notify::new(), } } - fn lock(&self) -> std::sync::MutexGuard<'_, QueueInner> { - self.inner.lock().expect("job_queue mutex poisoned") + fn lock(&self) -> std::sync::MutexGuard<'_, Sched> { + self.sched.lock().expect("job_queue mutex poisoned") } /// Submit a DAG: insert a [`NodeKind::Dag`] **container node** carrying the @@ -225,7 +224,6 @@ impl JobQueue { pub fn submit(&self, spec: DagSpec) -> anyhow::Result { let mut inner = self.lock(); let container = inner - .sched .append( NodeKind::Dag { source: spec.source, @@ -241,7 +239,7 @@ impl JobQueue { // `Finishing` and its children become runnable — it never needs claiming // or executing, and stays out of `claim_ready`. It rolls up terminal when // its whole subtree settles (that's the DAG-done signal). - inner.sched.complete(container, Outcome::Done); + inner.complete(container, Outcome::Done); drop(inner); self.notify.notify_one(); Ok(container.get()) @@ -255,15 +253,15 @@ impl JobQueue { pub fn claim_ready(&self) -> Vec { let mut inner = self.lock(); let inner = &mut *inner; - let started = inner.sched.settle(); + let started = inner.settle(); let mut claims = Vec::with_capacity(started.len()); for id in started { - let Some(node) = inner.sched.graph().node(id) else { + let Some(node) = inner.graph().node(id) else { continue; }; let kind = node.payload.clone(); let agent = node.payload.agent().to_owned(); - let Some(container) = inner.sched.graph().root_of(id) else { + let Some(container) = inner.graph().root_of(id) else { continue; }; claims.push(Claim { @@ -291,7 +289,7 @@ impl JobQueue { // complete) to express "grew nothing". The shared part is the outcome // mapping, and that's a free fn. let mut inner = self.lock(); - inner.sched.complete(node_id, outcome_of(result)); + inner.complete(node_id, outcome_of(result)); drop(inner); self.notify.notify_one(); } @@ -304,7 +302,7 @@ impl JobQueue { /// `Job::default()`. #[must_use] pub fn new_job(&self) -> Job { - self.lock().sched.new_job() + self.lock().new_job() } /// [`JobQueue::complete_node`] plus the work the node declared while it ran. @@ -318,10 +316,7 @@ impl JobQueue { // A rejected grown job is logged, not propagated: the node's own work // already ran, and refusing to complete it here would both misreport // that and wedge the DAG on a node stuck `Running`. - if let Err(e) = inner - .sched - .complete_growing(node_id, outcome_of(result), grown) - { + if let Err(e) = inner.complete_growing(node_id, outcome_of(result), grown) { tracing::error!( node = node_id.get(), error = %e, @@ -354,10 +349,10 @@ impl JobQueue { /// just that branch. Nothing here knows about DAGs. pub fn cancel(&self, id: u64) -> bool { let mut inner = self.lock(); - let Some(node) = inner.sched.graph().resolve_id(id) else { + let Some(node) = inner.graph().resolve_id(id) else { return false; }; - if !inner.sched.cancel_node(node) { + if !inner.cancel_node(node) { return false; } drop(inner); @@ -377,12 +372,8 @@ impl JobQueue { #[must_use] pub fn first_error(&self, dag_id: u64) -> Option { let inner = self.lock(); - let container = inner.container(dag_id)?; - inner - .sched - .graph() - .first_error(container) - .map(ToOwned::to_owned) + let container = container(&inner, dag_id)?; + inner.graph().first_error(container).map(ToOwned::to_owned) } /// `(agent, label, takes_container_down)` for the live transient-pill set, @@ -414,7 +405,6 @@ impl JobQueue { pub fn running_transients(&self) -> Vec { let inner = self.lock(); inner - .sched .graph() .nodes() .filter(|n| matches!(n.state, State::Running)) @@ -443,9 +433,11 @@ impl JobQueue { #[must_use] pub fn snapshot(&self) -> Vec { let inner = self.lock(); - let mut ids = inner.visible_dags(); + let mut ids = visible_dags(&inner); ids.sort_unstable_by_key(|c| c.get()); - ids.into_iter().filter_map(|c| inner.dag_view(c)).collect() + ids.into_iter() + .filter_map(|c| dag_view(&inner, c)) + .collect() } /// Number of live (non-terminal) DAGs — tests + diagnostics. @@ -453,191 +445,186 @@ impl JobQueue { #[must_use] pub fn live_count(&self) -> usize { let inner = self.lock(); - inner - .containers() + containers(&inner) .into_iter() - .filter(|&c| inner.sched.graph().is_settled(c) == Some(false)) + .filter(|&c| inner.graph().is_settled(c) == Some(false)) .count() } } -impl QueueInner { - /// The container node of `dag_id` — the `NodeKind::Dag` root whose id equals - /// `dag_id`. `NodeId` is un-fabricable from a raw `u64`, so this is a search. - fn container(&self, dag_id: u64) -> Option { - self.sched.graph().nodes().find_map(|n| { - (n.parent.is_none() - && n.id.get() == dag_id - && matches!(n.payload, NodeKind::Dag { .. })) +/// The container node of `dag_id` — the `NodeKind::Dag` root whose id equals +/// `dag_id`. `NodeId` is un-fabricable from a raw `u64`, so this is a search. +fn container(sched: &Sched, dag_id: u64) -> Option { + sched.graph().nodes().find_map(|n| { + (n.parent.is_none() && n.id.get() == dag_id && matches!(n.payload, NodeKind::Dag { .. })) .then_some(n.id) - }) - } + }) +} - /// The container's carried domain metadata as an owned read-view. The data - /// lives solely in the [`NodeKind::Dag`] payload — this is a derived read, - /// not a stored side-table. - fn dag_meta(&self, container: NodeId) -> Option { - let NodeKind::Dag { - source, - reason, - created_at, - } = &self.sched.graph().node(container)?.payload - else { - return None; +/// The container's carried domain metadata as an owned read-view. The data +/// lives solely in the [`NodeKind::Dag`] payload — this is a derived read, +/// not a stored side-table. +fn dag_meta(sched: &Sched, container: NodeId) -> Option { + let NodeKind::Dag { + source, + reason, + created_at, + } = &sched.graph().node(container)?.payload + else { + return None; + }; + Some(DagMeta { + source: *source, + reason: reason.clone(), + created_at: *created_at, + }) +} + +/// Project a DAG into its wire [`DagView`]: a near-raw view of the +/// container's work nodes, with `Done` nodes excluded. Lifecycle +/// (`state` / `started_at` / `finished_at` / `error`) is read straight +/// off each `hive_jobq::Node`; the client derives the DAG label, roll-up +/// state, and DAG timestamps from the node set. Non-derivable per-node +/// payload (`approval_id`, meta `inputs`) rides the owning node. Returns +/// `None` when every work node is `Done` or `Skipped` — a fully-settled +/// DAG drops out of the snapshot entirely (a `Failed` one lingers until +/// aged out). +fn dag_view(sched: &Sched, container: NodeId) -> Option { + let meta = dag_meta(sched, container)?; + let mut nodes = Vec::new(); + // Whether anything in this DAG still has an outcome worth showing. + // Kept separate from `nodes` being non-empty: skipped nodes ride the + // wire so the dashboard can mark the branches that weren't taken, but + // they must not by themselves hold a finished DAG in the snapshot. + let mut any_unsettled = false; + // DAG-level timestamps are taken over *all* subtree nodes (including the + // `Done` ones excluded from the wire) — the client can't derive them + // from a `Done`-filtered node set, so the host computes them here. + let mut started: Vec> = Vec::new(); + let mut finished: Vec> = Vec::new(); + for node in sched.graph().descendants(container) { + let id = node.id; + if let Some(s) = node.started_at { + started.push(s); + } + if let Some(f) = node.finished_at { + finished.push(f); + } + // `Done` nodes drop off the wire — a finished step isn't + // interesting. `Skipped` ones stay: which branch a run *didn't* + // take is the readable half of an outcome-branched DAG. + if matches!(node.state, State::Done) { + continue; + } + any_unsettled |= !matches!(node.state, State::Skipped); + let deps: Vec = node + .deps + .iter() + .filter_map(|d| match d { + Dep::Node { id, .. } => Some(id.get()), + Dep::Resource { .. } => None, + }) + .collect(); + // Non-derivable per-node payload rides the node that owns it. Every + // deploy phase carries the approval id, but only the subtree root + // projects it onto the wire — hanging the approval link off all of + // them would render the same card once per phase. + let approval_id = match &node.payload { + NodeKind::DeployWindow { approval_id, .. } => Some(*approval_id), + _ => None, }; - Some(DagMeta { - source: *source, - reason: reason.clone(), - created_at: *created_at, - }) + let inputs = match &node.payload { + NodeKind::MetaLock { inputs, .. } => inputs.clone(), + _ => Vec::new(), + }; + // Looked up from the log row itself (`build_logs.node_id`), not a + // host-side map. One indexed query per node in the snapshot; the + // node set is bounded by `MAX_HISTORY_DAGS` and the store is a + // local sqlite file, so this is cheaper than the lock contention + // a second shared map would reintroduce. + let build_log_id = crate::build_logs::global().and_then(|h| h.id_for_node(id.get())); + // `node.parent` is the structural jobq parent. Top-level nodes + // have `parent == Some(container)` (direct children of the Dag + // container); those become `parent: None` on the wire since the + // container itself is not part of the work-node payload. Sub-nodes + // carry the id of their containing parent work-node. + let parent = node + .parent + .filter(|&p| p != container) + .map(hive_jobq::NodeId::get); + nodes.push(NodeView { + id: id.get(), + agent: node.payload.agent().to_owned(), + kind: node.payload.as_str().to_owned(), + deps, + state: node.state, + started_at: node.started_at, + finished_at: node.finished_at, + error: node.error.clone(), + approval_id, + inputs, + build_log_id, + parent, + }); } + if !any_unsettled { + return None; + } + let is_terminal = sched.graph().is_settled(container) == Some(true); + Some(DagView { + id: container.get(), + source: meta.source, + reason: meta.reason.clone(), + created_at: meta.created_at, + started_at: started.into_iter().min(), + finished_at: is_terminal.then(|| finished.into_iter().max()).flatten(), + nodes, + }) +} - /// Project a DAG into its wire [`DagView`]: a near-raw view of the - /// container's work nodes, with `Done` nodes excluded. Lifecycle - /// (`state` / `started_at` / `finished_at` / `error`) is read straight - /// off each `hive_jobq::Node`; the client derives the DAG label, roll-up - /// state, and DAG timestamps from the node set. Non-derivable per-node - /// payload (`approval_id`, meta `inputs`) rides the owning node. Returns - /// `None` when every work node is `Done` or `Skipped` — a fully-settled - /// DAG drops out of the snapshot entirely (a `Failed` one lingers until - /// aged out). - fn dag_view(&self, container: NodeId) -> Option { - let meta = self.dag_meta(container)?; - let mut nodes = Vec::new(); - // Whether anything in this DAG still has an outcome worth showing. - // Kept separate from `nodes` being non-empty: skipped nodes ride the - // wire so the dashboard can mark the branches that weren't taken, but - // they must not by themselves hold a finished DAG in the snapshot. - let mut any_unsettled = false; - // DAG-level timestamps are taken over *all* subtree nodes (including the - // `Done` ones excluded from the wire) — the client can't derive them - // from a `Done`-filtered node set, so the host computes them here. - let mut started: Vec> = Vec::new(); - let mut finished: Vec> = Vec::new(); - for node in self.sched.graph().descendants(container) { - let id = node.id; - if let Some(s) = node.started_at { - started.push(s); - } - if let Some(f) = node.finished_at { - finished.push(f); - } - // `Done` nodes drop off the wire — a finished step isn't - // interesting. `Skipped` ones stay: which branch a run *didn't* - // take is the readable half of an outcome-branched DAG. - if matches!(node.state, State::Done) { - continue; - } - any_unsettled |= !matches!(node.state, State::Skipped); - let deps: Vec = node - .deps - .iter() - .filter_map(|d| match d { - Dep::Node { id, .. } => Some(id.get()), - Dep::Resource { .. } => None, - }) - .collect(); - // Non-derivable per-node payload rides the node that owns it. Every - // deploy phase carries the approval id, but only the subtree root - // projects it onto the wire — hanging the approval link off all of - // them would render the same card once per phase. - let approval_id = match &node.payload { - NodeKind::DeployWindow { approval_id, .. } => Some(*approval_id), - _ => None, - }; - let inputs = match &node.payload { - NodeKind::MetaLock { inputs, .. } => inputs.clone(), - _ => Vec::new(), - }; - // Looked up from the log row itself (`build_logs.node_id`), not a - // host-side map. One indexed query per node in the snapshot; the - // node set is bounded by `MAX_HISTORY_DAGS` and the store is a - // local sqlite file, so this is cheaper than the lock contention - // a second shared map would reintroduce. - let build_log_id = crate::build_logs::global().and_then(|h| h.id_for_node(id.get())); - // `node.parent` is the structural jobq parent. Top-level nodes - // have `parent == Some(container)` (direct children of the Dag - // container); those become `parent: None` on the wire since the - // container itself is not part of the work-node payload. Sub-nodes - // carry the id of their containing parent work-node. - let parent = node - .parent - .filter(|&p| p != container) - .map(hive_jobq::NodeId::get); - nodes.push(NodeView { - id: id.get(), - agent: node.payload.agent().to_owned(), - kind: node.payload.as_str().to_owned(), - deps, - state: node.state, - started_at: node.started_at, - finished_at: node.finished_at, - error: node.error.clone(), - approval_id, - inputs, - build_log_id, - parent, - }); +/// When a DAG's work node finishes on `finished_at` — the max over its +/// subtree (read off the graph `Node`, as unix seconds), for the history +/// cap ordering. +fn dag_finished_at(sched: &Sched, container: NodeId) -> i64 { + sched + .graph() + .descendants(container) + .filter_map(|n| n.finished_at) + .map(|t| t.timestamp()) + .max() + .unwrap_or(0) +} + +/// Every DAG container node id in the graph. +fn containers(sched: &Sched) -> Vec { + sched + .graph() + .nodes() + .filter(|n| n.parent.is_none() && matches!(n.payload, NodeKind::Dag { .. })) + .map(|n| n.id) + .collect() +} + +/// The **visible** DAG set for the snapshot: every live (non-terminal) DAG, +/// plus the newest [`MAX_HISTORY_DAGS`] terminal ones. Crate nodes for +/// evicted DAGs linger in the graph (bounded-prune is a Stage-C follow-up); +/// this filter is what bounds what the dashboard sees. +fn visible_dags(sched: &Sched) -> Vec { + let mut live: Vec = Vec::new(); + let mut terminal: Vec<(NodeId, i64)> = Vec::new(); + for c in containers(sched) { + if sched.graph().is_settled(c) == Some(true) { + terminal.push((c, dag_finished_at(sched, c))); + } else { + live.push(c); } - if !any_unsettled { - return None; - } - let is_terminal = self.sched.graph().is_settled(container) == Some(true); - Some(DagView { - id: container.get(), - source: meta.source, - reason: meta.reason.clone(), - created_at: meta.created_at, - started_at: started.into_iter().min(), - finished_at: is_terminal.then(|| finished.into_iter().max()).flatten(), - nodes, - }) - } - - /// When a DAG's work node finishes on `finished_at` — the max over its - /// subtree (read off the graph `Node`, as unix seconds), for the history - /// cap ordering. - fn dag_finished_at(&self, container: NodeId) -> i64 { - self.sched - .graph() - .descendants(container) - .filter_map(|n| n.finished_at) - .map(|t| t.timestamp()) - .max() - .unwrap_or(0) - } - - /// Every DAG container node id in the graph. - fn containers(&self) -> Vec { - self.sched - .graph() - .nodes() - .filter(|n| n.parent.is_none() && matches!(n.payload, NodeKind::Dag { .. })) - .map(|n| n.id) - .collect() - } - - /// The **visible** DAG set for the snapshot: every live (non-terminal) DAG, - /// plus the newest [`MAX_HISTORY_DAGS`] terminal ones. Crate nodes for - /// evicted DAGs linger in the graph (bounded-prune is a Stage-C follow-up); - /// this filter is what bounds what the dashboard sees. - fn visible_dags(&self) -> Vec { - let mut live: Vec = Vec::new(); - let mut terminal: Vec<(NodeId, i64)> = Vec::new(); - for c in self.containers() { - if self.sched.graph().is_settled(c) == Some(true) { - terminal.push((c, self.dag_finished_at(c))); - } else { - live.push(c); - } - } - // Newest first, so truncating to the cap keeps the most recent. - terminal.sort_by(|a, b| b.1.cmp(&a.1).then(b.0.get().cmp(&a.0.get()))); - terminal.truncate(MAX_HISTORY_DAGS); - let mut kept = live; - kept.extend(terminal.into_iter().map(|(c, _)| c)); - kept } + // Newest first, so truncating to the cap keeps the most recent. + terminal.sort_by(|a, b| b.1.cmp(&a.1).then(b.0.get().cmp(&a.0.get()))); + terminal.truncate(MAX_HISTORY_DAGS); + let mut kept = live; + kept.extend(terminal.into_iter().map(|(c, _)| c)); + kept } /// Truncate a node error to [`MAX_ERROR_LEN`] on a char boundary, appending `…`. diff --git a/hive-c0re/src/job_queue/tests.rs b/hive-c0re/src/job_queue/tests.rs index 82065941..802372a2 100644 --- a/hive-c0re/src/job_queue/tests.rs +++ b/hive-c0re/src/job_queue/tests.rs @@ -117,7 +117,6 @@ fn settle_rebuild_tail(q: &JobQueue, agent: &str, expect_ok: bool) { fn declared_resources(q: &JobQueue, node_id: hive_jobq::NodeId) -> Vec { let inner = q.lock(); inner - .sched .graph() .node(node_id) .expect("node exists")