diff --git a/hive-c0re/src/job_queue/exec.rs b/hive-c0re/src/job_queue/exec.rs index 7dcdf4dd..78451968 100644 --- a/hive-c0re/src/job_queue/exec.rs +++ b/hive-c0re/src/job_queue/exec.rs @@ -10,10 +10,10 @@ use std::sync::Arc; use anyhow::{Context as _, Result}; -use super::Claim; +use super::{Claim, Job}; use hive_jobq::TerminalState; -use super::model::{NodeKind, NodeSpec}; +use super::model::NodeKind; use crate::coordinator::Coordinator; use crate::power::{ReconcileAction, reconcile_action}; @@ -30,19 +30,19 @@ pub const GRACEFUL_STOP_TIMEOUT: std::time::Duration = std::time::Duration::from #[derive(Debug, Default)] pub struct NodeOutput { /// Whole per-agent *subgraphs* to append into *this same* DAG at - /// runtime — the single in-DAG-growth channel. Each inner - /// `Vec` is one independent subgraph whose `deps` are local - /// (0-based within that subgraph); the scheduler appends each via - /// [`super::JobQueue::append_subgraph`], which rebases the deps onto the DAG's - /// node-id space and roots the subgraph on the emitting node. Used both - /// for the multi-node case (`MetaLock` growing one rebuild subgraph per - /// agent — the startup sweep's stale agents, the meta-update cascade's - /// affected agents) and the single-node case (a `Reconcile` planner - /// emitting its mechanical `Start` / `Stop` as a one-node subgraph). The - /// scheduler applies these *before* the emitting node's completion so the - /// DAG never rolls terminal with the appended work still pending — keeping - /// the lease-window transient held across the sub-step. - pub append_subgraph: Vec>, + /// runtime — the single in-DAG-growth channel. Each [`Job`] is one + /// independent subgraph, declared but not yet inserted: an executor cannot + /// reach the queue, so it hands the declaration back and the scheduler + /// inserts it via [`super::JobQueue::append_subgraph`] under its own lock, + /// rooted on the emitting node. Used both for the multi-node case + /// (`MetaLock` growing one rebuild subgraph per agent — the startup + /// sweep's stale agents, the meta-update cascade's affected agents) and + /// the single-node case (a `Reconcile` planner emitting its mechanical + /// `Start` / `Stop` as a one-node subgraph). The scheduler applies these + /// *before* the emitting node's completion so the DAG never rolls terminal + /// with the appended work still pending — keeping the lease-window + /// transient held across the sub-step. + pub append_subgraph: Vec, } /// Build-log sink for one claimed node. @@ -342,14 +342,17 @@ async fn run_meta_lock( .unwrap_or_default() .iter() .map(|agent| { + let job = Job::new(); super::templates::rebuild_nodes( + &job, agent, super::templates::RebuildOpts { relock: true, graceful: true, }, - 0, - ) + None, + ); + job }) .collect(); return Ok(NodeOutput { append_subgraph }); @@ -371,14 +374,17 @@ async fn run_meta_lock( let append_subgraph = cascade .iter() .map(|agent| { + let job = Job::new(); super::templates::rebuild_nodes( + &job, agent, super::templates::RebuildOpts { relock: false, graceful: false, }, - 0, - ) + None, + ); + job }) .collect(); Ok(NodeOutput { append_subgraph }) @@ -398,7 +404,11 @@ async fn run_reconcile(coord: &Arc, claim: &Claim) -> Result sub(NodeKind::Start { agent: name.clone(), diff --git a/hive-c0re/src/job_queue/mod.rs b/hive-c0re/src/job_queue/mod.rs index 0f0bc009..ba06fa86 100644 --- a/hive-c0re/src/job_queue/mod.rs +++ b/hive-c0re/src/job_queue/mod.rs @@ -47,9 +47,19 @@ use hive_jobq::{Dep, Graph, NodeId}; use tokio::sync::Notify; pub use hive_jobq::TerminalState; -pub use model::{DagSpec, DagView, NodeKind, NodeSpec, PermPayload, Source, State}; +pub use model::{DagSpec, DagView, NodeKind, PermPayload, Source, State}; use resource::Resource; +/// A job under construction: `hive_jobq`'s builder over this queue's payload +/// ([`NodeKind`]) and resource ([`Resource`]) types. Templates declare into +/// one of these; [`JobQueue::submit`] inserts it. +pub type Job = hive_jobq::JobBuilder; + +/// A handle to one node a template declared — where its edges, grouping and +/// resources are declared. `Copy`; naming a node as a dependency does not +/// consume the ability to name it again. +pub type Handle<'a> = hive_jobq::NodeRef<'a, NodeKind, Resource>; + /// How many terminal DAGs (`Done` / `Failed` / `Cancelled`) the snapshot /// retains, newest first. A flat cap over the whole sorted list: the /// dashboard renders one recent-builds list, so one number bounds it. @@ -143,56 +153,33 @@ impl Default for JobQueue { } } -/// Insert `nodes` into the shared graph, honouring the spec's explicit **parent -/// axis**: a node with `parent = None` is a top-level group root (re-parented to -/// `group_parent`, which is `None` for `submit` and the emitting node for -/// `append_subgraph`); a node with `parent = Some(idx)` becomes a child of the -/// already-inserted node at spec index `idx`. `deps` are translated to crate -/// `Dep::Node` edges verbatim — templates declare the parent axis + sibling -/// ordering directly, so there is no dep-on-root to drop and no lease to hoist: -/// each node declares its own `Dep::Resource`, and the crate's borrow model -/// keeps a resource continuous across a subtree (a root owns it, descendants -/// borrow it). Independent group roots (multiple `parent = None` nodes) carry no -/// cross-links, so a multi-agent DAG's per-agent subgraphs run concurrently, each -/// on its own lease. Records per-node `node_rt`. Returns the inserted ids -/// (index-aligned with `nodes`). A node with `parent = None` is re-parented to -/// `group_parent` (the DAG container for a template, or the emitting node for a -/// runtime-appended subgraph); a node's `parent` / dep targets must precede it -/// in `nodes` (submit-time `validate` enforces density + acyclicity). +/// Insert a declared `job` into the shared graph and record its per-node +/// `node_rt`, returning the inserted ids. +/// +/// A node that declared no parent hangs under `group_parent` — the DAG +/// container for a template, the emitting node for a runtime-appended +/// subgraph. Templates declare the parent axis + sibling ordering directly, so +/// there is no dep-on-root to drop and no lease to hoist: each node declares +/// its own resources, and the crate's borrow model keeps a resource continuous +/// across a subtree (a root owns it, descendants borrow it). Independent group +/// roots carry no cross-links, so a multi-agent DAG's per-agent subgraphs run +/// concurrently, each on its own lease. /// /// # Errors /// Propagates a crate graph-insert error (malformed dep/parent / dep-scope). fn insert_group( inner: &mut QueueInner, - nodes: &[NodeSpec], + job: Job, group_parent: Option, ) -> anyhow::Result> { - let mut ids: Vec = Vec::with_capacity(nodes.len()); - for ns in nodes { - let payload = ns.kind.clone(); - let mut deps: Vec> = payload - .resource_deps() - .into_iter() - .map(|(name, count)| Dep::Resource { name, count }) - .collect(); - for d in &ns.deps { - deps.push(Dep::Node { - id: ids[dep_index(d.on)], - when: d.when, - }); - } - let parent = match ns.parent { - Some(idx) => Some(ids[dep_index(idx)]), - None => group_parent, - }; - let id = inner - .sched - .append(payload, deps, parent) - .map_err(|e| anyhow::anyhow!("job_queue: graph insert failed: {e}"))?; - ids.push(id); + let ids = inner + .sched + .insert_job(job, group_parent) + .map_err(|e| anyhow::anyhow!("job_queue: graph insert failed: {e}"))?; + for &id in ids.values() { inner.node_rt.insert(id, NodeRuntime::default()); } - Ok(ids) + Ok(ids.into_values().collect()) } impl JobQueue { @@ -225,7 +212,6 @@ impl JobQueue { /// Propagates the spec-validation error (empty / cyclic / bad parent) or a /// graph-insert error (dependencies that aren't dependency-topological). pub fn submit(&self, spec: DagSpec) -> anyhow::Result { - templates::validate(&spec)?; let mut inner = self.lock(); let container = inner .sched @@ -240,7 +226,7 @@ impl JobQueue { ) .map_err(|e| anyhow::anyhow!("job_queue: container insert failed: {e}"))?; inner.node_rt.insert(container, NodeRuntime::default()); - insert_group(&mut inner, &spec.nodes, Some(container))?; + insert_group(&mut inner, spec.job, Some(container))?; // Settle the container's own (no-op) logic immediately so it parks in // `Finishing` and its children become runnable — it never needs claiming // or executing, and stays out of `claim_ready`. It rolls up terminal when @@ -261,8 +247,8 @@ impl JobQueue { /// the DAG's terminal node deps on the top root, roll-up keeps the DAG from /// settling early with no explicit wiring. Returns the new node ids; empty if /// the DAG is gone or `nodes` is empty. - pub fn append_subgraph(&self, dag_id: u64, nodes: &[NodeSpec], dep_on: NodeId) -> Vec { - if nodes.is_empty() { + pub fn append_subgraph(&self, dag_id: u64, job: Job, dep_on: NodeId) -> Vec { + if job.is_empty() { return Vec::new(); } let mut inner = self.lock(); @@ -275,7 +261,7 @@ impl JobQueue { // emitter stays `Finishing` until this appended subtree settles, and the // container node rolls up terminal only once its whole subtree (incl. this // appended work) has settled, so the DAG hook waits for free. - let ids = match insert_group(&mut inner, nodes, Some(dep_on)) { + let ids = match insert_group(&mut inner, job, Some(dep_on)) { Ok(ids) => ids, Err(e) => { tracing::error!( @@ -679,12 +665,6 @@ impl QueueInner { } } -/// A spec dependency index (`Dep.on`, a wire `u64`) as a `usize` for indexing -/// into the node/id vectors. `templates::validate` guarantees it's in range. -fn dep_index(on: u64) -> usize { - usize::try_from(on).unwrap_or(usize::MAX) -} - /// Truncate a node error to [`MAX_ERROR_LEN`] on a char boundary, appending `…`. fn truncate_error(e: &str) -> String { if e.len() <= MAX_ERROR_LEN { diff --git a/hive-c0re/src/job_queue/model.rs b/hive-c0re/src/job_queue/model.rs index 77ba0757..79a40d14 100644 --- a/hive-c0re/src/job_queue/model.rs +++ b/hive-c0re/src/job_queue/model.rs @@ -13,18 +13,10 @@ //! full design. use chrono::{DateTime, Utc}; -pub use hive_host_sock::jobs::{DagView, NodeId, PermPayload, Source, State}; +pub use hive_host_sock::jobs::{DagView, PermPayload, Source, State}; use serde::Serialize; -use hive_jobq::{DepWhen, TerminalState}; - -/// A dependency edge (intra-DAG only — cross-DAG ordering comes from -/// the per-agent lease + dedup, never from edges between DAGs). -#[derive(Debug, Clone, Copy, Serialize)] -pub struct Dep { - pub on: NodeId, - pub when: DepWhen, -} +use hive_jobq::TerminalState; /// The primitive operations — each kind maps to one executor fn in /// `exec.rs`, a thin wrapper over existing `lifecycle.rs` / `meta.rs` @@ -469,35 +461,24 @@ impl NodeKind { } } -/// Submit-time spec for one node. -#[derive(Debug, Clone)] -pub struct NodeSpec { - /// The node's payload — [`NodeKind`] is the queue's payload type directly, - /// and each variant carries the agent it targets (a DAG can span agents; - /// the queue derives per-agent leasing from [`NodeKind::agent`]). - pub kind: NodeKind, - pub deps: Vec, - /// The **structural parent** axis — the spec-local index of this node's - /// group parent, or `None` for a top-level (group-root) node. Independent - /// of `deps`: `deps` order execution, `parent` groups nodes into a subtree - /// whose resource the whole subtree borrows (the agent lease is owned by a - /// group root and re-entered by its descendants for continuity). A child - /// runs once its parent reaches `Finishing` (the parent gate), so a child - /// never `deps` on its own parent (that would deadlock — dep-scope - /// validation rejects it). - pub parent: Option, -} - -/// Submit-time spec for a whole DAG. Built by `templates.rs`; validated -/// (cycle rejection) by `JobQueue::submit`. No DAG-level `agent` — every -/// node carries its own (a DAG can span agents), and the queue derives -/// per-agent leasing from [`NodeKind::agent`]. Type-specific payloads -/// (`PermChange`'s file payload) ride the node that consumes them -/// ([`NodeKind::WritePermFile`]), not this generic spec. -#[derive(Debug, Clone)] +/// Submit-time spec for a whole DAG: the group's metadata plus the declared — +/// not yet inserted — nodes. Built by `templates.rs`, inserted by +/// `JobQueue::submit`. +/// +/// No DAG-level `agent` — every node carries its own (a DAG can span agents), +/// and the queue derives per-agent leasing from [`NodeKind::agent`]. +/// Type-specific payloads (`PermChange`'s file payload) ride the node that +/// consumes them ([`NodeKind::WritePermFile`]), not this generic spec. +/// +/// There is no separate per-node spec type: the nodes live in the builder, +/// which inserts them itself. A shape that has been declared is therefore +/// always insertable — a dangling edge or a cycle cannot be expressed, so +/// there is nothing left for a submit-time validation pass to reject. +#[derive(Debug)] pub struct DagSpec { pub source: Source, /// Free-form "why". pub reason: String, - pub nodes: Vec, + /// The DAG's declared nodes, with their edges, grouping and resources. + pub job: super::Job, } diff --git a/hive-c0re/src/job_queue/scheduler.rs b/hive-c0re/src/job_queue/scheduler.rs index 105d5326..3efdf913 100644 --- a/hive-c0re/src/job_queue/scheduler.rs +++ b/hive-c0re/src/job_queue/scheduler.rs @@ -125,7 +125,7 @@ fn handle_completion(coord: &Arc, done: NodeDone) { // `Done` just below — covers both the multi-node case (a `MetaLock` // growing per-agent rebuild subgraphs) and the single-node case (a // `Reconcile` planner's `Start` / `Stop`). - for subgraph in &output.append_subgraph { + for subgraph in output.append_subgraph { coord .job_queue .append_subgraph(claim.dag_id, subgraph, claim.node_id); diff --git a/hive-c0re/src/job_queue/submit.rs b/hive-c0re/src/job_queue/submit.rs index 43955ec9..cddf2db9 100644 --- a/hive-c0re/src/job_queue/submit.rs +++ b/hive-c0re/src/job_queue/submit.rs @@ -6,8 +6,8 @@ //! state, which needs an async `lifecycle::is_running` read that a pure/sync //! template can't do. So these fns are async — they read each agent's state, //! assemble a per-agent subgraph out of the shared pure primitives -//! (`templates::{node, after_ok, rebuild_nodes}`), and concatenate them into -//! ONE DAG (independent per-agent roots, concurrent on their own leases). +//! (`templates::{node, rebuild_nodes}`), all declaring into ONE job +//! (independent per-agent roots, concurrent on their own leases). //! //! Dynamic shape rule: `stop`/`start` carry a head `SetWanted(w)` (durable //! intent write) — `restart` does NOT (it bounces the container but leaves @@ -24,9 +24,9 @@ use std::sync::Arc; -use super::model::{DagSpec, Dep, NodeKind, NodeSpec}; -use super::templates::{RebuildOpts, after_ok, child, node, rebuild_nodes}; -use super::{Source, templates}; +use super::model::{DagSpec, NodeKind}; +use super::templates::{RebuildOpts, node, rebuild_nodes}; +use super::{Job, Source, templates}; use crate::coordinator::Coordinator; use crate::lifecycle; @@ -51,71 +51,75 @@ pub fn rebuild(coord: &Arc, agent: &str, source: Source, reason: St // The pure per-agent chain builders below take `running` (and `stale`) // explicitly so they stay pure + unit-testable without a live container; // the async `*_many` fns read the real state via `lifecycle::is_running` -// then hand it in. Each chain uses LOCAL (0-based) deps; `concat_subgraphs` -// rebases them into one DAG. +// then hand it in. Each chain declares into the shared job it is handed, and +// names the nodes it depends on — so there is nothing to rebase. /// One agent's **stop** subgraph. `SetWanted(Off)` head + `Reconcile` tail /// always; the graceful `Signal → Drain` quiesce only when the agent is /// actually running (nothing to drain on a down container). The `Reconcile` /// stays even for a down agent so a race-up between the state read and exec /// is still stopped in-DAG. -fn stop_chain(agent: &str, graceful: bool, running: bool) -> Vec { +fn stop_chain(b: &Job, agent: &str, graceful: bool, running: bool) { // `SetWanted` is the group root and owns the agent lease; the mechanical // steps are its children (borrow the lease, run once it reaches `Finishing`, // dep-ordered among themselves). let a = || agent.to_owned(); - let mut n = vec![node( + let wanted = node( + b, NodeKind::SetWanted { agent: a(), up: false, }, - Vec::new(), - )]; + ); + // Declaration order is dependency order: the quiesce steps come first so + // the `Reconcile` that waits on them can name them. if graceful && running { - n.push(child(0, NodeKind::Signal { agent: a() }, Vec::new())); - n.push(child(0, NodeKind::Drain { agent: a() }, after_ok(1))); - n.push(child(0, NodeKind::Reconcile { agent: a() }, after_ok(2))); + let signal = node(b, NodeKind::Signal { agent: a() }).part_of(wanted); + let drain = node(b, NodeKind::Drain { agent: a() }) + .part_of(wanted) + .after_ok(signal); + let _ = node(b, NodeKind::Reconcile { agent: a() }) + .part_of(wanted) + .after_ok(drain); } else { - n.push(child(0, NodeKind::Reconcile { agent: a() }, Vec::new())); + let _ = node(b, NodeKind::Reconcile { agent: a() }).part_of(wanted); } - n } /// One agent's **start** subgraph. `SetWanted(Up)` head; a down + stale-rev /// agent gets the rebuild subgraph (its tail `Reconcile` starts it on /// current derivations), otherwise a plain `Reconcile` (which starts a down /// agent and noops an already-running one). -fn start_chain(agent: &str, running: bool, stale: bool) -> Vec { - let mut n = vec![node( +fn start_chain(b: &Job, agent: &str, running: bool, stale: bool) { + let wanted = node( + b, NodeKind::SetWanted { agent: agent.to_owned(), up: true, }, - Vec::new(), - )]; + ); if !running && stale { - // Rebuild subtree after the SetWanted head (base = 1, so the rebuild's - // `MetaSync` root deps `after_ok(0)` = the head). `MetaSync`, + // Rebuild subtree chained behind the `SetWanted` head. `MetaSync`, // `Prebuild` + `Reconcile` are their own group roots (top-level, per // `rebuild_nodes`). - n.extend(rebuild_nodes( + rebuild_nodes( + b, agent, RebuildOpts { relock: true, graceful: false, }, - 1, - )); + Some(wanted), + ); } else { - n.push(child( - 0, + let _ = node( + b, NodeKind::Reconcile { agent: agent.to_owned(), }, - Vec::new(), - )); + ) + .part_of(wanted); } - n } /// One agent's **restart** subgraph. Restart NEVER rewrites `wanted` @@ -127,82 +131,49 @@ fn start_chain(agent: &str, running: bool, stale: bool) -> Vec { /// before `Reconcile`; a down agent gets just `Reconcile`, which /// converges to intent — a stopped (`wanted = Off`) agent stays stopped, /// a crashed (`wanted = Up`) agent comes back up. -fn restart_chain(agent: &str, graceful: bool, running: bool) -> Vec { +fn restart_chain(b: &Job, agent: &str, graceful: bool, running: bool) { let a = || agent.to_owned(); if !running { // Nothing to bounce — a lone Reconcile converges to intent. - return vec![node(NodeKind::Reconcile { agent: a() }, Vec::new())]; + let _ = node(b, NodeKind::Reconcile { agent: a() }); + return; } // Running: mechanical stop then Reconcile. The first stop node is the group // ROOT (no SetWanted head) and owns the agent lease; the rest are its // children (borrow the lease, dep-ordered), so the bounce holds one // continuous lease and `Reconcile` cancel-cascades if a stop step fails. - let mut n = vec![if graceful { - node(NodeKind::Signal { agent: a() }, Vec::new()) - } else { - node(NodeKind::StopForUpdate { agent: a() }, Vec::new()) - }]; + // + // `Reconcile` gates on the last mechanical step. For a non-graceful bounce + // that step *is* the root, and the parent gate already orders it — a child + // must NOT dep on its own parent (dep-scope), so it takes no sibling edge. if graceful { - n.push(child(0, NodeKind::Drain { agent: a() }, Vec::new())); - n.push(child( - 0, - NodeKind::StopForUpdate { agent: a() }, - after_ok(1), - )); - } - // `Reconcile` gates on the last mechanical step. When the only step is the - // root itself (non-graceful, `StopForUpdate` == index 0), the parent gate - // already orders `Reconcile` after it — a child must NOT dep on its own - // parent (dep-scope). So the sibling dep is added only for a graceful - // bounce, where the last step is a sibling child. - let deps = if n.len() > 1 { - after_ok(u64::try_from(n.len() - 1).unwrap_or(0)) + let signal = node(b, NodeKind::Signal { agent: a() }); + let drain = node(b, NodeKind::Drain { agent: a() }).part_of(signal); + let stop = node(b, NodeKind::StopForUpdate { agent: a() }) + .part_of(signal) + .after_ok(drain); + let _ = node(b, NodeKind::Reconcile { agent: a() }) + .part_of(signal) + .after_ok(stop); } else { - Vec::new() - }; - n.push(child(0, NodeKind::Reconcile { agent: a() }, deps)); - n -} - -/// Concatenate per-agent subgraphs (each with LOCAL 0-based deps) into one -/// node list, rebasing each subgraph's internal deps by its offset. A -/// subgraph root (empty deps — the `SetWanted` head) stays a root, so the -/// per-agent subgraphs are independent and run concurrently, each on its -/// own lease. -fn concat_subgraphs(chains: Vec>) -> Vec { - let mut out: Vec = Vec::new(); - for chain in chains { - let base = u64::try_from(out.len()).unwrap_or(u64::MAX); - for spec in chain { - let deps = spec - .deps - .into_iter() - .map(|d| Dep { - on: base + d.on, - when: d.when, - }) - .collect(); - out.push(NodeSpec { - kind: spec.kind, - deps, - // Rebase the structural parent by the same offset (a subgraph - // root keeps `parent = None`, so the per-agent groups stay - // independent + concurrent). - parent: spec.parent.map(|p| base + p), - }); - } + let stop = node(b, NodeKind::StopForUpdate { agent: a() }); + let _ = node(b, NodeKind::Reconcile { agent: a() }).part_of(stop); } - out } -/// Wrap assembled power-op `nodes` in a `DagSpec`. No tail node: a power op's +/// Wrap the per-agent subgraphs in a `DagSpec`. No tail node: a power op's /// effect is its nodes (`SetWanted` + `Reconcile`), with nothing left to do once /// they settle. -fn power_dag(source: Source, reason: String, nodes: Vec) -> DagSpec { +/// +/// There is no concatenation step: every chain declares into the same builder +/// and each keeps its own root, so the per-agent subgraphs are independent and +/// run concurrently, each on its own lease. Rebasing one subgraph's indices +/// onto another's used to be a function. +fn power_dag(source: Source, reason: String, job: Job) -> DagSpec { DagSpec { source, reason, - nodes, + job, } } @@ -218,11 +189,11 @@ pub(crate) fn stop_spec( source: Source, reason: String, ) -> DagSpec { - let chains = targets - .iter() - .map(|(agent, running)| stop_chain(agent, graceful, *running)) - .collect(); - power_dag(source, reason, concat_subgraphs(chains)) + let job = Job::new(); + for (agent, running) in targets { + stop_chain(&job, agent, graceful, *running); + } + power_dag(source, reason, job) } /// Assemble the start DAG from explicit `(agent, running, stale)` targets. @@ -236,11 +207,11 @@ pub(crate) fn start_spec( source: Source, reason: String, ) -> DagSpec { - let chains = targets - .iter() - .map(|(agent, running, stale)| start_chain(agent, *running, *stale)) - .collect(); - power_dag(source, reason, concat_subgraphs(chains)) + let job = Job::new(); + for (agent, running, stale) in targets { + start_chain(&job, agent, *running, *stale); + } + power_dag(source, reason, job) } /// Assemble the restart DAG from explicit `(agent, running)` targets. @@ -250,11 +221,11 @@ pub(crate) fn restart_spec( source: Source, reason: String, ) -> DagSpec { - let chains = targets - .iter() - .map(|(agent, running)| restart_chain(agent, graceful, *running)) - .collect(); - power_dag(source, reason, concat_subgraphs(chains)) + let job = Job::new(); + for (agent, running) in targets { + restart_chain(&job, agent, graceful, *running); + } + power_dag(source, reason, job) } /// Restart a single agent. Thin wrapper over [`restart_many`]. diff --git a/hive-c0re/src/job_queue/templates.rs b/hive-c0re/src/job_queue/templates.rs index 70c6d6a9..78f6550b 100644 --- a/hive-c0re/src/job_queue/templates.rs +++ b/hive-c0re/src/job_queue/templates.rs @@ -1,21 +1,19 @@ //! DAG shape builders — every operation as a template over the shared -//! node primitives — plus submit-time cycle validation (petgraph is -//! confined to this validation; the runtime store stays the plain -//! `Vec` + `deps`). +//! node primitives. //! //! Every node carries its own `agent` (there is no DAG-level agent) — the -//! `node` helper stamps each node's agent. This module holds the *pure* -//! shape builders (no I/O). The hive-wide **power ops** (`stop` / `start` / -//! `restart`) are NOT here: their per-agent shape depends on each agent's -//! live running state (an async `lifecycle::is_running` read), so they are -//! assembled dynamically in `submit.rs` out of the shared pure primitives -//! this module exports (`node`, `after_ok`, `rebuild_nodes`) — one -//! independent per-agent subgraph each, concurrent on its own lease, ONE -//! DAG for the whole hive-wide op. `stop`/`start` write the durable `wanted` -//! intent via a head `SetWanted(w)` node (holding the agent lease, so -//! intent+reconcile is atomic per-agent); `restart` writes no intent — it -//! bounces the container and lets the tail `Reconcile` converge to the -//! agent's existing `wanted`. +//! [`node`] helper stamps each node's agent and declares the resources that +//! node's kind needs. This module holds the *pure* shape builders (no I/O). +//! The hive-wide **power ops** (`stop` / `start` / `restart`) are NOT here: +//! their per-agent shape depends on each agent's live running state (an async +//! `lifecycle::is_running` read), so they are assembled dynamically in +//! `submit.rs` out of the shared pure primitives this module exports +//! ([`node`], [`rebuild_nodes`]) — one independent per-agent subgraph each, +//! concurrent on its own lease, ONE DAG for the whole hive-wide op. +//! `stop`/`start` write the durable `wanted` intent via a head `SetWanted(w)` +//! node (holding the agent lease, so intent+reconcile is atomic per-agent); +//! `restart` writes no intent — it bounces the container and lets the tail +//! `Reconcile` converge to the agent's existing `wanted`. //! //! ```text //! rebuild(a): MetaSync(a) → Prebuild(a) → StopForUpdate(a) → Swap(a) →(ok) PostSwap(a) →(any) Reconcile(a) @@ -26,113 +24,77 @@ //! reparent(moves): Reparent(moves) [no rebuild — topology.json is read live] //! ``` //! +//! Nodes are **named, not counted**: a template declares a node and holds the +//! handle it gets back, so an edge says which node it waits on instead of +//! computing where that node landed. There is no submit-time cycle validation +//! left to do — `hive_jobq`'s builder inserts in declaration order and rejects +//! a reference to a node declared later, so every edge points backwards and a +//! cycle is unrepresentable rather than merely rejected. +//! //! For the dynamic power-op shapes (`stop` / `start` / `restart`, built from //! live online/offline state), see `submit.rs`. -use anyhow::{Result, bail}; +use hive_jobq::TerminalState; -use hive_jobq::{DepWhen, TerminalState}; +use super::model::{DagSpec, NodeKind, PermPayload, Source}; +use super::{Handle, Job}; -use super::model::{DagSpec, Dep, NodeKind, NodeSpec, PermPayload, Source}; - -/// After-ok edge on the previous node — the common chain link. Shared with -/// the async power-op builders in `submit.rs` (which assemble per-agent -/// chains dynamically from live container state). -pub(crate) fn after_ok(on: u64) -> Vec { - vec![Dep { - on, - when: DepWhen::AFTER_OK, - }] -} - -/// `AfterOk` edges onto every one of a DAG's **group-roots** — the success -/// branch of a per-outcome tail pair, and the aggregator the failure branch -/// keys off. +/// Declare one node carrying `kind`, with the resources that kind needs. /// -/// Group-roots are the right granularity, not "every node": a root's state *is* -/// its subtree's roll-up, so edging the roots covers every descendant while -/// keeping the dep list small and stable as subtrees grow. Because every edge is -/// `AFTER_OK`, this node runs only if *all* of them succeeded — and is ruled out -/// ([`TerminalState::Skipped`]) the moment one doesn't, which is precisely the -/// signal [`on_elimination_of`] waits for. -pub(crate) fn after_ok_all(ons: &[u64]) -> Vec { - ons.iter() - .map(|&on| Dep { - on, - when: DepWhen::AFTER_OK, - }) - .collect() -} - -/// `AFTER_ANY` edges onto every group-root — "wait for all of these to finish, -/// however they went". Ordering only; it accepts any outcome except the DAG -/// being dropped. -pub(crate) fn after_any_all(ons: &[u64]) -> Vec { - ons.iter() - .map(|&on| Dep { - on, - when: DepWhen::AFTER_ANY, - }) - .collect() -} - -/// A single edge satisfied only when `on` was **ruled out** by its own edges. +/// The resource declaration is [`NodeKind::resource_deps`] applied at the +/// construction site — a build slot for nix-heavy kinds, the agent lease for +/// container-affecting ones, the global meta window for meta-mutating ones. A +/// node that declares a resource an ancestor already holds re-enters that +/// grant rather than taking a fresh unit, so declaring costs nothing. /// -/// Dependency edges are conjunctive, so "any one of these several nodes failed" -/// cannot be written directly. This is the composition that expresses it: point -/// the success branch at every root with [`after_ok_all`], then hang the failure -/// branch off *that* node's elimination. Exactly one of the pair ever runs. -/// -/// Note it accepts `Skipped` and **not** `Cancelled`: if the whole DAG was -/// dropped before it started, the success branch is marked `Cancelled` directly -/// and this branch is ruled out too — a job nobody ran reports nothing. -pub(crate) fn on_elimination_of(on: u64) -> Vec { - vec![Dep { - on, - when: DepWhen::of(&[TerminalState::Skipped]), - }] -} - -/// A single edge satisfied only by the listed outcomes of `on` — for the -/// one-tail-per-outcome shape an approval DAG uses. -pub(crate) fn on_outcome(on: u64, outcomes: &[TerminalState]) -> Vec { - vec![Dep { - on, - when: DepWhen::of(outcomes), - }] +/// The returned handle is where edges and grouping are declared, and is `Copy` +/// — naming a node as a dependency does not consume the ability to name it +/// again. +pub(crate) fn node(b: &Job, kind: NodeKind) -> Handle<'_> { + // Read the resources off the kind before handing it over — the payload is + // moved into the node, not cloned for it. + let resources = kind.resource_deps(); + let mut handle = b.node(kind); + for (name, count) in resources { + handle = handle.needs_units(name, count); + } + handle } /// The `Rebuilt`-reporting tail pair for a rebuild-shaped DAG: the success node /// gated on every group-root in `roots`, and the failure node gated on *its* -/// elimination. `base` is the spec index the pair starts at. +/// elimination. /// /// Exactly one runs on a DAG that executed, and neither runs on one the operator -/// dropped — see [`on_elimination_of`]. -fn emit_rebuilt_tails(agent: &str, roots: &[u64], base: u64) -> Vec { +/// dropped — see [`hive_jobq::NodeRef::on_elimination_of`]. +fn emit_rebuilt_tails(b: &Job, agent: &str, roots: &[Handle<'_>]) { + let ok = roots.iter().fold( + node( + b, + NodeKind::EmitRebuilt { + agent: agent.to_owned(), + ok: true, + }, + ), + hive_jobq::NodeRef::after_ok, + ); // The failure branch needs *both*: the ok branch being ruled out (that is the // "something went wrong" signal) **and** every root actually finished. The // second half is easy to forget and gets the ordering wrong without it — a // failed `Prebuild` eliminates the ok branch immediately, while the recovery // `Reconcile` is still bringing the container back up, so reporting straight // off the elimination would announce the failure mid-recovery. - let mut on_fail = after_any_all(roots); - on_fail.extend(on_elimination_of(base)); - vec![ - node( - NodeKind::EmitRebuilt { - agent: agent.to_owned(), - ok: true, - }, - after_ok_all(roots), - ), + let _failed = roots.iter().fold( node( + b, NodeKind::EmitRebuilt { agent: agent.to_owned(), ok: false, }, - on_fail, - ), - ] + ) + .on_elimination_of(ok), + hive_jobq::NodeRef::after_any, + ); } /// The approval-resolving tails for an approval-carrying DAG: one per outcome of @@ -140,48 +102,20 @@ fn emit_rebuilt_tails(agent: &str, roots: &[u64], base: u64) -> Vec { /// /// The `Cancelled` node is what keeps a dropped approval DAG from dangling its /// row forever — its edge is the only one [`super::JobQueue::cancel`] spares. -fn resolve_approval_tails(approval_id: i64, root: u64) -> Vec { - [ +fn resolve_approval_tails(b: &Job, approval_id: i64, root: Handle<'_>) { + for outcome in [ TerminalState::Done, TerminalState::Failed, TerminalState::Cancelled, - ] - .into_iter() - .map(|outcome| { - node( + ] { + let _ = node( + b, NodeKind::ResolveApproval { approval_id, outcome, }, - on_outcome(root, &[outcome]), ) - }) - .collect() -} - -/// Build one **top-level (group-root)** node — `parent = None`. `kind` carries -/// the agent it targets ([`NodeKind`] is the payload directly). Shared with -/// `submit.rs`'s dynamic power-op builders. A root owns whatever resource it -/// declares for its whole subtree; its descendants borrow it (agent-lease / -/// build-slot continuity). Ordering vs other nodes is `deps`; grouping is -/// `parent`. -pub(crate) fn node(kind: NodeKind, deps: Vec) -> NodeSpec { - NodeSpec { - kind, - deps, - parent: None, - } -} - -/// Build a **child** node whose structural parent is spec-index `parent`. The -/// child runs once its parent reaches `Finishing` (the parent gate), so it must -/// NOT `deps` on `parent` (dep-scope validation rejects a dep on one's own -/// parent). `deps` here order the child against its *siblings* only. -pub(crate) fn child(parent: u64, kind: NodeKind, deps: Vec) -> NodeSpec { - NodeSpec { - kind, - deps, - parent: Some(parent), + .on_outcome(root, &[outcome]); } } @@ -198,28 +132,51 @@ pub(crate) struct RebuildOpts { pub graceful: bool, } -/// The rebuild node subtree (nested, three group roots). `base` is the spec -/// index of the first node (`MetaSync`). Structure: -/// - `MetaSync` (base+0, **root**): the meta-repo preamble (dir prep, agent -/// sync, optional relock). Owns the global `MetaWindow` — and *only* for its -/// own short duration, which is why it is a sibling root rather than -/// `Prebuild`'s parent: a resource is held across the holder's whole subtree, -/// so parenting the build under it would extend a hive-global window over -/// every rebuild's nix build. -/// - `Prebuild` (base+1, **root**): `AfterOk` `MetaSync`. Owns the build slot -/// for the whole mechanical subtree below it. Lease-exempt — the nix build -/// overlaps other DAGs on the same agent. -/// - the **stop root** (base+2, child of `Prebuild`): owns the agent lease and -/// runs once `Prebuild` reaches `Finishing` (parent gate). Non-graceful that -/// is `StopForUpdate` itself; graceful it is `Signal`, with `Drain` and then +/// The group-roots a [`rebuild_nodes`] subgraph exposes to its caller: what a +/// tail node edges onto, and what a follow-up node waits for. +/// +/// Only the roots — a root's state *is* its subtree's roll-up, so these three +/// cover every node in the subgraph without the caller knowing its shape. +#[derive(Debug, Clone, Copy)] +pub(crate) struct RebuildRoots<'a> { + /// The meta-repo preamble. + pub meta_sync: Handle<'a>, + /// The build root — its roll-up carries the whole + /// `StopForUpdate` → `Swap` → `PostSwap` subtree. + pub prebuild: Handle<'a>, + /// The recovery/convergence tail root. + pub reconcile: Handle<'a>, +} + +impl<'a> RebuildRoots<'a> { + /// The three roots as a slice, for edging a tail onto all of them. + fn all(self) -> [Handle<'a>; 3] { + [self.meta_sync, self.prebuild, self.reconcile] + } +} + +/// The rebuild node subtree (nested, three group roots). `after`, when given, is +/// the node this subgraph chains behind. Structure: +/// - `MetaSync` (**root**): the meta-repo preamble (dir prep, agent sync, +/// optional relock). Owns the global `MetaWindow` — and *only* for its own +/// short duration, which is why it is a sibling root rather than `Prebuild`'s +/// parent: a resource is held across the holder's whole subtree, so parenting +/// the build under it would extend a hive-global window over every rebuild's +/// nix build. +/// - `Prebuild` (**root**): `AfterOk` `MetaSync`. Owns the build slot for the +/// whole mechanical subtree below it. Lease-exempt — the nix build overlaps +/// other DAGs on the same agent. +/// - the **stop root** (child of `Prebuild`): owns the agent lease and runs once +/// `Prebuild` reaches `Finishing` (parent gate). Non-graceful that is +/// `StopForUpdate` itself; graceful it is `Signal`, with `Drain` and then /// `StopForUpdate` as its children so the lease stays continuous across the /// whole stop — siblings would each take the lease separately and leave a gap /// another DAG could claim the agent in, mid-bounce. /// - `Swap` (child of `StopForUpdate`): borrows the agent lease from its /// ancestors and the build slot from `Prebuild` — both continuous. -/// - `PostSwap` (child of `StopForUpdate`): the swap's Ok-only -/// bookkeeping tail (rev marker, forge/matrix sync, kick, rescan), `AfterOk` -/// its sibling `Swap`. +/// - `PostSwap` (child of `StopForUpdate`): the swap's Ok-only bookkeeping tail +/// (rev marker, forge/matrix sync, kick, rescan), `AfterOk` its sibling +/// `Swap`. /// - `Reconcile` (**last, root**): `AfterAny` `Prebuild`, which rolls up /// terminal only once its whole mechanical subtree (SFU→Swap→PostSwap) has /// settled — so `Reconcile` runs after the swap regardless of outcome, and as @@ -228,56 +185,47 @@ pub(crate) struct RebuildOpts { /// cancel-cascades `Prebuild`, i.e. terminal, so the tail still runs). It /// takes a fresh lease; the tiny gap is harmless — `Reconcile` converges to /// the persisted `wanted` idempotently. -pub(crate) fn rebuild_nodes(agent: &str, opts: RebuildOpts, base: u64) -> Vec { +pub(crate) fn rebuild_nodes<'a>( + b: &'a Job, + agent: &str, + opts: RebuildOpts, + after: Option>, +) -> RebuildRoots<'a> { let a = || agent.to_owned(); let RebuildOpts { relock, graceful } = opts; - let mut nodes = vec![ - node( - NodeKind::MetaSync { agent: a(), relock }, - if base == 0 { - Vec::new() - } else { - after_ok(base - 1) - }, - ), - node(NodeKind::Prebuild { agent: a() }, after_ok(base)), - ]; + + let mut meta_sync = node(b, NodeKind::MetaSync { agent: a(), relock }); + if let Some(after) = after { + meta_sync = meta_sync.after_ok(after); + } + let prebuild = node(b, NodeKind::Prebuild { agent: a() }).after_ok(meta_sync); + // The stop root hangs off `Prebuild` and owns the agent lease for - // everything below it. - let stop_root = base + 2; - if graceful { - nodes.push(child(base + 1, NodeKind::Signal { agent: a() }, Vec::new())); + // everything below it. `StopForUpdate` parents the swap pair either way. + let stop_for_update = if graceful { + let signal = node(b, NodeKind::Signal { agent: a() }).part_of(prebuild); // `Drain` is a *child* of `Signal`, so the parent gate already orders // it — a child must not dep on its own parent (dep-scope). - nodes.push(child(stop_root, NodeKind::Drain { agent: a() }, Vec::new())); - nodes.push(child( - stop_root, - NodeKind::StopForUpdate { agent: a() }, - after_ok(stop_root + 1), - )); + let drain = node(b, NodeKind::Drain { agent: a() }).part_of(signal); + node(b, NodeKind::StopForUpdate { agent: a() }) + .part_of(signal) + .after_ok(drain) } else { - nodes.push(child( - base + 1, - NodeKind::StopForUpdate { agent: a() }, - Vec::new(), - )); + node(b, NodeKind::StopForUpdate { agent: a() }).part_of(prebuild) + }; + + let swap = node(b, NodeKind::Swap { agent: a() }).part_of(stop_for_update); + let _post_swap = node(b, NodeKind::PostSwap { agent: a() }) + .part_of(stop_for_update) + .after_ok(swap); + + let reconcile = node(b, NodeKind::Reconcile { agent: a() }).after_any(prebuild); + + RebuildRoots { + meta_sync, + prebuild, + reconcile, } - // Index of `StopForUpdate`, which parents the swap pair either way. - let sfu = if graceful { stop_root + 2 } else { stop_root }; - nodes.push(child(sfu, NodeKind::Swap { agent: a() }, Vec::new())); - nodes.push(child( - sfu, - NodeKind::PostSwap { agent: a() }, - after_ok(sfu + 1), - )); - nodes.push(node( - NodeKind::Reconcile { agent: a() }, - vec![Dep { - on: base + 1, - when: DepWhen::AFTER_ANY, - }], - )); - nodes } /// The rebuild subgraph a [`NodeKind::DeployApply`] grows into its own DAG once @@ -288,7 +236,7 @@ pub(crate) fn rebuild_nodes(agent: &str, opts: RebuildOpts, base: u64) -> Vec Vec Vec { - let mut nodes = rebuild_nodes( +pub(crate) fn deploy_rebuild_nodes(agent: &str, approval_id: i64) -> Job { + let b = Job::new(); + let roots = rebuild_nodes( + &b, agent, RebuildOpts { relock: false, graceful: false, }, - 0, + None, ); - let reconcile = reconcile_index(&nodes, 0); - nodes.push(node( + let _finalize = node( + &b, NodeKind::FinalizeDeploy { agent: agent.to_owned(), approval_id, }, - vec![ - Dep { - on: 1, - when: DepWhen::AFTER_OK, - }, - Dep { - on: reconcile, - when: DepWhen::AFTER_OK, - }, - ], - )); - nodes -} - -/// Spec index of the `Reconcile` root a [`rebuild_nodes`] subgraph ends on, -/// for callers that gate a tail on it. Read off the emitted list rather than -/// hard-coded, because the subgraph's length depends on [`RebuildOpts`]. -fn reconcile_index(rebuild: &[NodeSpec], base: u64) -> u64 { - base + u64::try_from(rebuild.len()).unwrap_or(0).saturating_sub(1) + ) + .after_ok(roots.prebuild) + .after_ok(roots.reconcile); + b } /// One uniform rebuild shape — no `was_running` branch. `StopForUpdate` @@ -352,41 +287,41 @@ fn reconcile_index(rebuild: &[NodeSpec], base: u64) -> u64 { /// node. Edging `Reconcile` alone would not do: it is `AfterAny` `Prebuild`, so /// it reaches `Done` even after a failed swap and the tail would report success. pub fn rebuild(agent: &str, source: Source, reason: String, relock: bool) -> DagSpec { - let mut nodes = rebuild_nodes( + let job = Job::new(); + let roots = rebuild_nodes( + &job, agent, RebuildOpts { relock, graceful: false, }, - 0, + None, ); - let reconcile = reconcile_index(&nodes, 0); - let tail_base = u64::try_from(nodes.len()).unwrap_or(0); - nodes.extend(emit_rebuilt_tails(agent, &[0, 1, reconcile], tail_base)); + emit_rebuilt_tails(&job, agent, &roots.all()); DagSpec { source, reason, - nodes, + job, } } /// Approval-driven deploy (`MergeConfigPr`) as a phase subtree rather than the /// single opaque node it used to be. Structure: -/// - `DeployWindow` (0, **root**): the resource holder — global meta window, -/// agent lease, build slot — held across every child below. No work of its -/// own; it reaches `Finishing` immediately and the children run inside it. -/// - `MergeVerify` (1, child): drift-gate + fetch + eval-verify. Mutates -/// nothing, so a failure here cancel-cascades its siblings with the forge and -/// the applied repo exactly as they were. -/// - `DeployApply` (2, child, `AfterOk` `MergeVerify`): the irreversible half — +/// - `DeployWindow` (**root**): the resource holder — global meta window, agent +/// lease, build slot — held across every child below. No work of its own; it +/// reaches `Finishing` immediately and the children run inside it. +/// - `MergeVerify` (child): drift-gate + fetch + eval-verify. Mutates nothing, +/// so a failure here cancel-cascades its siblings with the forge and the +/// applied repo exactly as they were. +/// - `DeployApply` (child, `AfterOk` `MergeVerify`): the irreversible half — /// ff-merge + `prepare_deploy`. It doesn't rebuild inline; it grows /// [`deploy_rebuild_nodes`] into this DAG as its own children, so the build /// and the closing `FinalizeDeploy` are real nodes under the same window. -/// - `DeployTail` (3, child, `AfterAny` `DeployApply`): the compensation + +/// - `DeployTail` (child, `AfterAny` `DeployApply`): the compensation + /// bookkeeping tail — rollback when a merge landed unfinalized, forge tag /// mirror, PR failure comment (see [`NodeKind::DeployTail`]). /// -/// - `ResolveApproval` (4, **root**, `AfterAny` `DeployWindow`): resolves the +/// - `ResolveApproval` (**root**, `AfterAny` `DeployWindow`): resolves the /// approval row. A root rather than another child, so it isn't inside the /// window's resource subtree — it runs once the window has released the meta /// window, lease and build slot. One edge suffices here: `DeployWindow` is the @@ -397,48 +332,47 @@ pub fn rebuild(agent: &str, source: Source, reason: String, relock: bool) -> Dag /// leaves `flake.lock` staged-uncommitted for the build's whole duration. pub fn approval_deploy(agent: &str, approval_id: i64, reason: String) -> DagSpec { let a = || agent.to_owned(); + let job = Job::new(); + + let window = node( + &job, + NodeKind::DeployWindow { + agent: a(), + approval_id, + }, + ); + let verify = node( + &job, + NodeKind::MergeVerify { + agent: a(), + approval_id, + }, + ) + .part_of(window); + let apply = node( + &job, + NodeKind::DeployApply { + agent: a(), + approval_id, + }, + ) + .part_of(window) + .after_ok(verify); + let _tail = node( + &job, + NodeKind::DeployTail { + agent: a(), + approval_id, + }, + ) + .part_of(window) + .after_any(apply); + + resolve_approval_tails(&job, approval_id, window); DagSpec { source: Source::Approval, reason, - nodes: vec![ - node( - NodeKind::DeployWindow { - agent: a(), - approval_id, - }, - Vec::new(), - ), - child( - 0, - NodeKind::MergeVerify { - agent: a(), - approval_id, - }, - Vec::new(), - ), - child( - 0, - NodeKind::DeployApply { - agent: a(), - approval_id, - }, - after_ok(1), - ), - child( - 0, - NodeKind::DeployTail { - agent: a(), - approval_id, - }, - vec![Dep { - on: 2, - when: DepWhen::AFTER_ANY, - }], - ), - ] - .into_iter() - .chain(resolve_approval_tails(approval_id, 0)) - .collect(), + job, } } @@ -449,15 +383,17 @@ pub fn approval_deploy(agent: &str, approval_id: i64, reason: String) -> DagSpec /// in the queue tests); production paths no longer emit a bare reconcile. #[cfg(test)] pub fn reconcile_only(agent: &str, source: Source, reason: String) -> DagSpec { + let job = Job::new(); + let _reconcile = node( + &job, + NodeKind::Reconcile { + agent: agent.to_owned(), + }, + ); DagSpec { source, reason, - nodes: vec![node( - NodeKind::Reconcile { - agent: agent.to_owned(), - }, - Vec::new(), - )], + job, } } @@ -473,53 +409,56 @@ pub fn reconcile_only(agent: &str, source: Source, reason: String) -> DagSpec { /// `AfterAny` onto `Provision` — the DAG's only other group-root, so its roll-up /// already carries the whole cascade. pub fn spawn(agent: &str, approval_id: i64, reason: String) -> DagSpec { + let a = || agent.to_owned(); + let job = Job::new(); + + let provision = node(&job, NodeKind::Provision { agent: a() }); + let create = node(&job, NodeKind::Create { agent: a() }).part_of(provision); + let dropin = node(&job, NodeKind::WriteDropin { agent: a() }).part_of(create); + let _reconcile = node(&job, NodeKind::Reconcile { agent: a() }) + .part_of(create) + .after_ok(dropin); + + resolve_approval_tails(&job, approval_id, provision); DagSpec { source: Source::Approval, reason, - nodes: { - let a = || agent.to_owned(); - vec![ - node(NodeKind::Provision { agent: a() }, Vec::new()), - child(0, NodeKind::Create { agent: a() }, Vec::new()), - child(1, NodeKind::WriteDropin { agent: a() }, Vec::new()), - child(1, NodeKind::Reconcile { agent: a() }, after_ok(2)), - ] - .into_iter() - .chain(resolve_approval_tails(approval_id, 0)) - .collect() - }, + job, } } /// Perm change: commit the JSON file(s), then the rebuild subgraph so /// the updated `HIVE_TOOL_GROUPS` / `HIVE_CAPABILITIES` env var takes -/// effect in the container. Group-roots are `WritePermFile`(0) plus the rebuild -/// subgraph's `MetaSync`(1) / `Prebuild`(2) / `Reconcile`(6), so the -/// `EmitRebuilt` tail edges all four. +/// effect in the container. Group-roots are `WritePermFile` plus the rebuild +/// subgraph's `MetaSync` / `Prebuild` / `Reconcile`, so the `EmitRebuilt` tail +/// edges all four. pub fn perm_change(agent: &str, source: Source, reason: String, payload: PermPayload) -> DagSpec { - let mut nodes = vec![node( + let job = Job::new(); + let write = node( + &job, NodeKind::WritePermFile { agent: agent.to_owned(), payload, }, - Vec::new(), - )]; - let rebuild = rebuild_nodes( + ); + let roots = rebuild_nodes( + &job, agent, RebuildOpts { relock: true, graceful: false, }, - 1, + Some(write), + ); + emit_rebuilt_tails( + &job, + agent, + &[write, roots.meta_sync, roots.prebuild, roots.reconcile], ); - let reconcile = reconcile_index(&rebuild, 1); - nodes.extend(rebuild); - let tail_base = u64::try_from(nodes.len()).unwrap_or(0); - nodes.extend(emit_rebuilt_tails(agent, &[0, 1, 2, reconcile], tail_base)); DagSpec { source, reason, - nodes, + job, } } @@ -539,25 +478,26 @@ pub fn meta_update( reason: String, approval_id: Option, ) -> DagSpec { - let mut nodes = vec![node( + let job = Job::new(); + let lock = node( + &job, NodeKind::MetaLock { sweep: false, fanout: None, inputs, }, - Vec::new(), - )]; + ); // The bump itself has no side effect, so an operator-driven one ends at the // `MetaLock`; an approval-driven one still has its row to resolve and gets the // per-outcome tails edged onto that single group-root — whose roll-up covers // the rebuild subgraphs `MetaLock` grows into itself. if let Some(approval_id) = approval_id { - nodes.extend(resolve_approval_tails(approval_id, 0)); + resolve_approval_tails(&job, approval_id, lock); } DagSpec { source, reason, - nodes, + job, } } @@ -575,10 +515,12 @@ pub fn reparent( source: Source, reason: String, ) -> DagSpec { + let job = Job::new(); + let _reparent = node(&job, NodeKind::Reparent { moves }); DagSpec { source, reason, - nodes: vec![node(NodeKind::Reparent { moves }, Vec::new())], + job, } } @@ -586,45 +528,3 @@ pub fn reparent( // as ONE `Boot` DAG (a sweep `MetaLock` root that grows rebuild subgraphs // in-DAG, plus a `Reconcile` root per drifted agent) — no anchor node and no // per-agent child DAGs. - -/// Validate a spec before it enters the queue: node ids are dense -/// (index = id), deps + parents reference existing *earlier* nodes, and the -/// dep graph is acyclic (petgraph `toposort`). Rejecting cycles here fixes the -/// old queue's documented "circular dep silently deadlocks forever" caveat. -pub fn validate(spec: &DagSpec) -> Result<()> { - if spec.nodes.is_empty() { - bail!("dag spec {:?} has no nodes", spec.reason); - } - let n = spec.nodes.len(); - let mut graph = petgraph::graph::DiGraph::::new(); - let idx: Vec<_> = (0..n) - .map(|i| graph.add_node(u32::try_from(i).unwrap_or(u32::MAX))) - .collect(); - for (i, node) in spec.nodes.iter().enumerate() { - // A `parent` must index an earlier node — `insert_group` resolves it to - // an already-inserted `NodeId`, so a forward/out-of-bounds parent would - // otherwise panic there. - if let Some(p) = node.parent - && usize::try_from(p).is_ok_and(|p| p >= i) - { - bail!( - "dag spec {:?} node {i} has invalid parent {p} (must be an earlier node)", - spec.reason - ); - } - for dep in &node.deps { - let Some(&dep_idx) = usize::try_from(dep.on).ok().and_then(|i| idx.get(i)) else { - bail!( - "dag spec {:?} node {i} depends on unknown node {}", - spec.reason, - dep.on - ); - }; - graph.add_edge(dep_idx, idx[i], ()); - } - } - if petgraph::algo::toposort(&graph, None).is_err() { - bail!("dag spec {:?} contains a dependency cycle", spec.reason); - } - Ok(()) -} diff --git a/hive-c0re/src/job_queue/tests.rs b/hive-c0re/src/job_queue/tests.rs index 49d55bd4..a4e3a66b 100644 --- a/hive-c0re/src/job_queue/tests.rs +++ b/hive-c0re/src/job_queue/tests.rs @@ -6,9 +6,7 @@ //! scheduler's async loop is a thin claim/complete pump over the same //! methods exercised here. -use hive_jobq::DepWhen; - -use super::model::{Dep, NodeKind, NodeSpec}; +use super::model::NodeKind; use super::*; fn submit(q: &JobQueue, spec: DagSpec) -> u64 { @@ -141,71 +139,22 @@ fn resubmit_while_running_is_new_dag() { assert_eq!(q.snapshot().len(), 2); } -// ---- cycle rejection ---- - -#[test] -fn cyclic_dag_is_rejected_at_submit() { - let q = JobQueue::new(1); - let mut spec = rebuild("agent-a", "cyclic"); - // 0 → 1 → 0 cycle. - spec.nodes = vec![ - NodeSpec { - kind: NodeKind::StopForUpdate { - agent: "agent-a".to_owned(), - }, - deps: vec![Dep { - on: 1, - when: DepWhen::AFTER_OK, - }], - parent: None, - }, - NodeSpec { - kind: NodeKind::Reconcile { - agent: "agent-a".to_owned(), - }, - deps: vec![Dep { - on: 0, - when: DepWhen::AFTER_OK, - }], - parent: None, - }, - ]; - assert!(q.submit(spec).is_err(), "cyclic spec must be refused"); - assert!(q.snapshot().is_empty()); -} - -#[test] -fn unknown_dep_is_rejected_at_submit() { - let q = JobQueue::new(1); - let mut spec = rebuild("agent-a", "bad dep"); - spec.nodes = vec![NodeSpec { - kind: NodeKind::Reconcile { - agent: "agent-a".to_owned(), - }, - deps: vec![Dep { - on: 9, - when: DepWhen::AFTER_OK, - }], - parent: None, - }]; - assert!(q.submit(spec).is_err()); -} - -#[test] -fn invalid_parent_is_rejected_at_submit() { - let q = JobQueue::new(1); - let mut spec = rebuild("agent-a", "bad parent"); - // A forward/out-of-bounds parent index must be refused at validate, not - // panic in `insert_group`. - spec.nodes = vec![NodeSpec { - kind: NodeKind::Reconcile { - agent: "agent-a".to_owned(), - }, - deps: Vec::new(), - parent: Some(3), - }]; - assert!(q.submit(spec).is_err()); -} +// ---- malformed specs: no longer expressible ---- +// +// Three tests lived here — a dependency cycle, a dependency on a node that +// does not exist, and an out-of-range parent index — each asserting that +// `submit` refused the spec. All three built their spec by hand out of +// positional indices, which is exactly the representation that made those +// shapes possible: an index can name a node that isn't there, or one that +// comes later. +// +// A job is now declared against handles that only exist for nodes already +// declared, so there is no index to put out of range, and every edge points +// backwards — a cycle needs a forward edge. The guard those tests covered was +// deleted along with the failure mode. What remains — a handle used against a +// builder that never issued it — is `hive_jobq`'s to reject, and its builder +// tests cover it (`a_forward_edge_is_rejected_by_name`, +// `a_forward_parent_is_rejected_by_name`, `graph_rejection_surfaces_as_is`). // ---- dependency order within a DAG ---- @@ -243,20 +192,24 @@ fn rebuild_chain_claims_in_dep_order() { #[test] fn graceful_rebuild_chain_drains_before_stopping() { let q = JobQueue::new(1); - let spec = DagSpec { - source: Source::AutoUpdate, - reason: "sweep".to_owned(), - - nodes: templates::rebuild_nodes( - "agent-a", - templates::RebuildOpts { - relock: true, - graceful: true, - }, - 0, - ), - }; - let id = submit(&q, spec); + let job = Job::new(); + templates::rebuild_nodes( + &job, + "agent-a", + templates::RebuildOpts { + relock: true, + graceful: true, + }, + None, + ); + let id = submit( + &q, + DagSpec { + source: Source::AutoUpdate, + reason: "sweep".to_owned(), + job, + }, + ); for expected in [ "meta_sync", "prebuild", @@ -284,17 +237,34 @@ fn graceful_rebuild_chain_drains_before_stopping() { /// drain window, so `StopForUpdate` still hangs straight off `Prebuild`. #[test] fn non_graceful_rebuild_has_no_signal_or_drain() { - let kinds: Vec = templates::rebuild_nodes( + // Read the shape off the queue rather than out of a node list: a declared + // job keeps its nodes to itself and inserts them, so what it built is + // observable where it matters — in what the scheduler runs. + let q = JobQueue::new(1); + let job = Job::new(); + templates::rebuild_nodes( + &job, "agent-a", templates::RebuildOpts { relock: true, graceful: false, }, - 0, - ) - .iter() - .map(|n| n.kind.as_str().to_owned()) - .collect(); + None, + ); + let id = submit( + &q, + DagSpec { + source: Source::Manual, + reason: "manual".to_owned(), + job, + }, + ); + let mut kinds = Vec::new(); + for _ in 0..6 { + let c = claim_one(&q); + kinds.push(c.kind.as_str().to_owned()); + q.complete_node(c.node_id, Ok(())); + } assert_eq!( kinds, vec![ @@ -306,6 +276,8 @@ fn non_graceful_rebuild_has_no_signal_or_drain() { "reconcile" ] ); + // Settled after exactly those six — nothing else was declared. + assert_eq!(state_of(&q, id), State::Done); } /// A cleanly-finished DAG leaves the snapshot even though its not-taken @@ -720,19 +692,19 @@ fn append_subgraph_roots_on_emitter_and_rebases_local_deps() { // subgraph per stale agent into its OWN DAG. Each subgraph is rooted on // the emitter and its LOCAL 0-based deps are rebased onto the DAG. let q = JobQueue::new(4); + let job = Job::new(); + let _lock = templates::node( + &job, + NodeKind::MetaLock { + sweep: true, + fanout: None, + inputs: Vec::new(), + }, + ); let spec = DagSpec { source: Source::AutoUpdate, reason: "sweep".to_owned(), - - nodes: vec![NodeSpec { - kind: NodeKind::MetaLock { - sweep: true, - fanout: None, - inputs: Vec::new(), - }, - deps: Vec::new(), - parent: None, - }], + job, }; let id = submit(&q, spec); let emitter = claim_one(&q); @@ -742,18 +714,21 @@ fn append_subgraph_roots_on_emitter_and_rebases_local_deps() { // StopForUpdate → Swap → Reconcile, local 0-based deps. `graceful` must // match the sweep arm of `run_meta_lock` or this stops tracking production. let subgraph = |agent: &str| { + let job = Job::new(); templates::rebuild_nodes( + &job, agent, templates::RebuildOpts { relock: true, graceful: true, }, - 0, - ) + None, + ); + job }; // Must append BEFORE completing the emitter (the documented contract). - q.append_subgraph(id, &subgraph("a"), emitter.node_id); - q.append_subgraph(id, &subgraph("b"), emitter.node_id); + q.append_subgraph(id, subgraph("a"), emitter.node_id); + q.append_subgraph(id, subgraph("b"), emitter.node_id); q.complete_node(emitter.node_id, Ok(())); // Still ONE DAG; both subgraph roots become ready once the emitter is // Done (rooted on it), each on its own agent lease. Their `MetaSync` heads @@ -845,18 +820,17 @@ fn meta_update_grows_cascade_in_dag() { // Simulate the executor growing the cascade in-DAG (`relock = false` — a // cascade child must not re-lock and revert the parent's bump). for agent in ["alice", "bob"] { - q.append_subgraph( - id, - &templates::rebuild_nodes( - agent, - templates::RebuildOpts { - relock: false, - graceful: false, - }, - 0, - ), - meta_lock.node_id, + let job = Job::new(); + templates::rebuild_nodes( + &job, + agent, + templates::RebuildOpts { + relock: false, + graceful: false, + }, + None, ); + q.append_subgraph(id, job, meta_lock.node_id); } q.complete_node(meta_lock.node_id, Ok(())); // Still ONE DAG — no child DAGs — and both cascade rebuild subgraphs root @@ -1119,15 +1093,21 @@ fn cancelled_power_op_runs_no_compensating_node() { ), ]; for (name, writes_intent, spec) in cases { + let q = JobQueue::new(1); + let id = submit(&q, spec); + // Read the intent head off the submitted DAG rather than out of + // the spec: a declared job holds its own nodes and inserts them. assert_eq!( - spec.nodes + q.snapshot() .iter() - .any(|n| matches!(n.kind, NodeKind::SetWanted { .. })), + .find(|d| d.id == id) + .expect("submitted dag") + .nodes + .iter() + .any(|n| n.kind == "set_wanted"), writes_intent, "{name} intent head (graceful={graceful}, running={running})" ); - let q = JobQueue::new(1); - let id = submit(&q, spec); assert!(q.cancel(id), "cancelled while queued"); assert_eq!(state_of(&q, id), State::Cancelled); assert!( @@ -1272,7 +1252,7 @@ fn deploy_apply_grows_rebuild_subgraph_and_finalizes_after_it() { // gate immediately and letting the deploy "finish" before it had built. let grown = q.append_subgraph( id, - &templates::deploy_rebuild_nodes("agent-a", 11), + templates::deploy_rebuild_nodes("agent-a", 11), apply.node_id, ); assert!(!grown.is_empty(), "subgraph grafted onto the apply node"); @@ -1330,7 +1310,7 @@ fn deploy_dag_skips_finalize_but_still_tails_a_failed_graft() { let apply = claim_one(&q); q.append_subgraph( id, - &templates::deploy_rebuild_nodes("agent-a", 13), + templates::deploy_rebuild_nodes("agent-a", 13), apply.node_id, ); q.complete_node(apply.node_id, Ok(())); diff --git a/hive-c0re/src/workers/auto_update.rs b/hive-c0re/src/workers/auto_update.rs index 599d6796..0cb146d8 100644 --- a/hive-c0re/src/workers/auto_update.rs +++ b/hive-c0re/src/workers/auto_update.rs @@ -334,7 +334,7 @@ fn submit_boot_tree( n_deferred: usize, n_skipped: usize, ) { - use crate::job_queue::{DagSpec, NodeKind, NodeSpec, Source}; + use crate::job_queue::{DagSpec, Job, NodeKind, Source, templates}; // Fully-quiet boot (nothing stale, nothing drifted) submits nothing. if !any_stale && drifted.is_empty() { @@ -348,32 +348,27 @@ fn submit_boot_tree( n_skipped, ); - let mut nodes: Vec = Vec::new(); + let job = Job::new(); // Sweep whenever ANY marker is stale — even when every stale agent is // wanted-offline: the hyperhive lock bump must land now so their later // start-upgrade rebuilds build against it. No stale agents ⇒ no MetaLock // ⇒ no meta commit on a no-change boot. The `fanout` list rides the // MetaLock into `run_meta_lock`, which appends the rebuild subgraphs. if any_stale { - nodes.push(NodeSpec { - kind: NodeKind::MetaLock { + let _ = templates::node( + &job, + NodeKind::MetaLock { sweep: true, fanout: Some(fanout), // A sweep bumps `hyperhive` alone (`lock_update_hyperhive`), // so it names no inputs. inputs: Vec::new(), }, - deps: Vec::new(), - parent: None, - }); + ); } // One boot Reconcile per drifted agent — independent roots. for name in drifted { - nodes.push(NodeSpec { - kind: NodeKind::Reconcile { agent: name }, - deps: Vec::new(), - parent: None, - }); + let _ = templates::node(&job, NodeKind::Reconcile { agent: name }); } let spec = DagSpec { @@ -384,7 +379,7 @@ fn submit_boot_tree( // Rebuilding when the sweep will grow rebuild subgraphs (per-agent // crash-watch suppression during their Swap, applied at claim time); // a reconcile-only boot needs no transient. - nodes, + job, }; if let Err(e) = coord.job_queue.submit(spec) { tracing::warn!(error = ?e, "boot: sweep DAG submit failed"); diff --git a/hive-jobq/src/builder.rs b/hive-jobq/src/builder.rs index 884389da..e0a55abe 100644 --- a/hive-jobq/src/builder.rs +++ b/hive-jobq/src/builder.rs @@ -112,6 +112,13 @@ impl JobBuilder { Self::default() } + /// Whether nothing has been declared yet — for a caller deciding whether an + /// insertion is worth taking a lock for. + #[must_use] + pub fn is_empty(&self) -> bool { + self.nodes.borrow().is_empty() + } + /// Add a node carrying `payload`, with no edges, resources, or parent yet. /// /// The returned handle is where those are declared; it is [`Copy`], so it @@ -275,6 +282,14 @@ impl NodeRef<'_, N, R> { self.edge(on.into(), DepWhen::AFTER_ANY) } + /// Run only on the listed outcomes of `on` — the general form of + /// [`NodeRef::after_ok`] / [`NodeRef::after_any`], for the + /// one-node-per-outcome shape a job with a row to resolve uses. + #[must_use] + pub fn on_outcome(self, on: impl Into, outcomes: &[TerminalState]) -> Self { + self.edge(on.into(), DepWhen::of(outcomes)) + } + /// Run only if `on` was **ruled out** — i.e. its own `after_ok` edges did /// not hold. ///