diff --git a/CLAUDE.md b/CLAUDE.md index a47d22b4..d853b55c 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -57,8 +57,7 @@ hand-maintained per-file tree drifts out of sync with the code. - **`hive-jobq/`** — persistent job-DAG scheduler, extracted from hive-c0re's in-tree `job_queue` as a domain-agnostic library. One persistent graph for the whole system (not a DAG per job); enqueuing - inserts a self-contained sub-DAG and returns the ids of the nodes the job - asked for, in the order it named them. Generic over + inserts a self-contained sub-DAG and returns its node ids. Generic over the node payload `N` and the resource name `R`; resource deps are named counting semaphores acquired all-or-nothing at node start. `hive-c0re`'s remaining `job_queue/` module is the c0re-specific layer *over* this diff --git a/Cargo.lock b/Cargo.lock index 8ac44a4c..6223c7ac 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1758,7 +1758,6 @@ dependencies = [ "serde", "serde_json", "thiserror 2.0.18", - "uuid", ] [[package]] diff --git a/Cargo.toml b/Cargo.toml index eb79b132..9111c620 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -64,7 +64,6 @@ hive-sock-client = { path = "hive-sock-client" } hive-types = { path = "hive-types" } thiserror = "2" tower-http = { version = "0.7", features = ["fs"] } -uuid = { version = "1", features = ["v4"] } rmcp = { version = "2", default-features = false, features = [ "server", "macros", diff --git a/docs/coordinator.md b/docs/coordinator.md index dc87418d..af66f329 100644 --- a/docs/coordinator.md +++ b/docs/coordinator.md @@ -23,15 +23,9 @@ dashboard group; the **node** is the unit of scheduling / execution / build-log. Deps are intra-DAG edges only (`AfterOk` by default: the dep must succeed, a failed/cancelled dep cancels the dependent — cancel-downstream). Cross-DAG ordering comes from the per-agent lease, -never from edges between DAGs. - -A DAG is **declared, not described**: a template builds it through -`hive_jobq::JobBuilder`, naming each node it depends on via the handle -`b.node(kind)` handed back, and the builder inserts the nodes itself. A handle -only exists for a node already declared, so every edge points backwards and a -cycle cannot be written down — there is no submit-time validation pass, because -there is no malformed spec to reject. (The old queue had a petgraph `toposort` -here, guarding against the positional indices that used to express edges.) +never from edges between DAGs. Submit-time validation (petgraph `toposort`) +rejects cyclic specs outright, fixing the old queue's "circular dep silently +deadlocks" caveat. ### Node inventory (primitives) diff --git a/hive-c0re/src/job_queue/exec.rs b/hive-c0re/src/job_queue/exec.rs index d6f44a37..7dcdf4dd 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, Declare}; +use super::Claim; use hive_jobq::TerminalState; -use super::model::NodeKind; +use super::model::{NodeKind, NodeSpec}; use crate::coordinator::Coordinator; use crate::power::{ReconcileAction, reconcile_action}; @@ -27,32 +27,22 @@ pub const GRACEFUL_STOP_TIMEOUT: std::time::Duration = std::time::Duration::from /// Extra signal an executor hands back to the scheduler alongside /// success. -#[derive(Default)] +#[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 [`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, -} - -impl std::fmt::Debug for NodeOutput { - fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { - // The subgraphs are closures — how many were emitted is the only thing - // there is to say about them before the queue runs them. - f.debug_struct("NodeOutput") - .field("append_subgraph", &self.append_subgraph.len()) - .finish() - } + /// 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>, } /// Build-log sink for one claimed node. @@ -352,18 +342,14 @@ async fn run_meta_lock( .unwrap_or_default() .iter() .map(|agent| { - let agent = agent.clone(); - Box::new(move |b: &super::Job| { - super::templates::rebuild_nodes( - b, - &agent, - super::templates::RebuildOpts { - relock: true, - graceful: true, - }, - None, - ); - }) as Declare + super::templates::rebuild_nodes( + agent, + super::templates::RebuildOpts { + relock: true, + graceful: true, + }, + 0, + ) }) .collect(); return Ok(NodeOutput { append_subgraph }); @@ -385,18 +371,14 @@ async fn run_meta_lock( let append_subgraph = cascade .iter() .map(|agent| { - let agent = agent.clone(); - Box::new(move |b: &super::Job| { - super::templates::rebuild_nodes( - b, - &agent, - super::templates::RebuildOpts { - relock: false, - graceful: false, - }, - None, - ); - }) as Declare + super::templates::rebuild_nodes( + agent, + super::templates::RebuildOpts { + relock: false, + graceful: false, + }, + 0, + ) }) .collect(); Ok(NodeOutput { append_subgraph }) @@ -416,11 +398,7 @@ 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 9b9db781..bbe90925 100644 --- a/hive-c0re/src/job_queue/mod.rs +++ b/hive-c0re/src/job_queue/mod.rs @@ -47,29 +47,9 @@ use hive_jobq::{Dep, Graph, NodeId}; use tokio::sync::Notify; pub use hive_jobq::TerminalState; -pub use model::{DagSpec, DagView, NodeKind, PermPayload, Source, State}; +pub use model::{DagSpec, DagView, NodeKind, NodeSpec, 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 a -/// borrowed one; only `hive_jobq` can make or insert it. -pub type Job = hive_jobq::JobBuilder; - -/// A job's shape as a **recipe**: given a builder, declare the nodes. -/// -/// What a template returns and what an executor hands back, because neither -/// can build a job itself — `hive_jobq` creates the builder inside its own -/// insertion call and never lets one out. So the transferable thing is the -/// declaring closure, and the queue runs it at the moment it inserts. -/// -/// `Send` because an executor's output crosses the scheduler's task boundary. -pub type Declare = Box; - -/// 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. @@ -163,36 +143,52 @@ impl Default for JobQueue { } } -/// 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. +/// 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). /// /// # Errors /// Propagates a crate graph-insert error (malformed dep/parent / dep-scope). fn insert_group( inner: &mut QueueInner, - declare: impl FnOnce(&Job), + nodes: &[NodeSpec], 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, - // which `submit` inserts itself, and nothing downstream looks an - // individual step up by id. - Vec::new() - }) - .map_err(|e| anyhow::anyhow!("job_queue: graph insert failed: {e}"))?; - Ok(()) +) -> anyhow::Result> { + let mut ids: Vec = Vec::with_capacity(nodes.len()); + for ns in nodes { + let payload = ns.kind.clone(); + let mut deps = payload.resource_deps(); + 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); + inner.node_rt.insert(id, NodeRuntime::default()); + } + Ok(ids) } impl JobQueue { @@ -216,19 +212,16 @@ impl JobQueue { self.inner.lock().expect("job_queue mutex poisoned") } - /// Submit a DAG: insert a [`NodeKind::Dag`] **container node** carrying the - /// group's metadata, then insert the template's nodes as its subtree (their - /// roots re-parented to the container). Returns the container's id as the - /// DAG id — its rolled-up state is the DAG state. - /// - /// Takes the spec's recipe by generic, not as a boxed [`Declare`]: a spec - /// travels from the template that built it directly into this call, so - /// there is nothing to allocate for. + /// Submit a DAG. Validates the spec, inserts a [`NodeKind::Dag`] **container + /// node** carrying the group's metadata, then inserts the template's nodes as + /// its subtree (their roots re-parented to the container). Returns the + /// container's id as the DAG id — its rolled-up state is the DAG state. /// /// # Errors - /// Propagates a graph-insert error (dependencies that aren't - /// dependency-topological). - pub fn submit(&self, spec: DagSpec) -> anyhow::Result { + /// 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 @@ -243,7 +236,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.declare, Some(container))?; + insert_group(&mut inner, &spec.nodes, 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 @@ -262,11 +255,15 @@ impl JobQueue { /// gate — the children run once `dep_on` reaches `Finishing`. Because the /// emitting node stays `Finishing` until this appended subtree is terminal and /// the DAG's terminal node deps on the top root, roll-up keeps the DAG from - /// settling early with no explicit wiring. A no-op if the DAG is gone. - pub fn append_subgraph(&self, dag_id: u64, declare: Declare, dep_on: NodeId) { + /// 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() { + return Vec::new(); + } let mut inner = self.lock(); if inner.container(dag_id).is_none() { - return; + return Vec::new(); } // Insert the subgraph as a group rooted under the emitting node: the // subgraph's own root becomes a child of `dep_on`, its steps children of @@ -274,16 +271,20 @@ 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. - if let Err(e) = insert_group(&mut inner, declare, Some(dep_on)) { - tracing::error!( - dag = dag_id, - error = %e, - "job_queue: append_subgraph insert failed" - ); - return; - } + let ids = match insert_group(&mut inner, nodes, Some(dep_on)) { + Ok(ids) => ids, + Err(e) => { + tracing::error!( + dag = dag_id, + error = %e, + "job_queue: append_subgraph insert failed" + ); + return Vec::new(); + } + }; drop(inner); self.notify.notify_one(); + ids } /// Claim every currently-runnable node, acquiring its resources, and mark it @@ -449,14 +450,17 @@ impl JobQueue { .nodes() .filter(|n| matches!(n.state, State::Running)) .filter_map(|n| { - let agent = - n.payload - .resource_deps() - .into_iter() - .find_map(|(name, _)| match name { - Resource::Agent(a) => Some(a), - _ => None, - })?; + let agent = n + .payload + .resource_deps() + .into_iter() + .find_map(|d| match d { + Dep::Resource { + name: Resource::Agent(a), + .. + } => Some(a), + _ => None, + })?; Some(RunningTransient { agent, label: n.payload.as_str().to_owned(), @@ -674,6 +678,12 @@ 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 19556271..77ba0757 100644 --- a/hive-c0re/src/job_queue/model.rs +++ b/hive-c0re/src/job_queue/model.rs @@ -13,10 +13,18 @@ //! full design. use chrono::{DateTime, Utc}; -pub use hive_host_sock::jobs::{DagView, PermPayload, Source, State}; +pub use hive_host_sock::jobs::{DagView, NodeId, PermPayload, Source, State}; use serde::Serialize; -use hive_jobq::TerminalState; +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, +} /// The primitive operations — each kind maps to one executor fn in /// `exec.rs`, a thin wrapper over existing `lifecycle.rs` / `meta.rs` @@ -461,43 +469,35 @@ impl NodeKind { } } -/// 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, and no built job either: `declare` -/// is a *recipe* the queue runs against a builder `hive_jobq` owns, at the -/// moment it inserts. 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. -/// -/// Generic over the recipe rather than boxing it: a spec goes from the template -/// that returns it straight to the `submit` that consumes it, so the closure's -/// concrete type is known the whole way and needs neither an allocation nor a -/// `Send` bound. (The executor's `append_subgraph` is the case that *does* need -/// a boxed [`super::Declare`] — its recipes are collected into a `Vec` and -/// applied later, across a task boundary.) -pub struct DagSpec { +/// 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)] +pub struct DagSpec { pub source: Source, /// Free-form "why". pub reason: String, - /// Declares the DAG's nodes — their edges, grouping and resources — onto - /// the builder the queue hands it. - pub declare: F, -} - -impl std::fmt::Debug for DagSpec { - fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { - // The recipe is a closure; there is nothing to show of it, and its - // nodes do not exist until the queue runs it. - f.debug_struct("DagSpec") - .field("source", &self.source) - .field("reason", &self.reason) - .finish_non_exhaustive() - } + pub nodes: Vec, } diff --git a/hive-c0re/src/job_queue/resource.rs b/hive-c0re/src/job_queue/resource.rs index a3118137..0b46ab9d 100644 --- a/hive-c0re/src/job_queue/resource.rs +++ b/hive-c0re/src/job_queue/resource.rs @@ -4,6 +4,8 @@ //! payload `N`; here `R` is [`Resource`] and `N` is [`NodeKind`] directly (each //! variant carries the agent it targets). +use hive_jobq::Dep; + use super::model::NodeKind; /// The two resource classes the queue gates concurrency on, as the crate's @@ -34,7 +36,7 @@ pub enum Resource { } impl NodeKind { - /// The resources this node must acquire to run — `(name, units)` — derived from + /// The [`Dep::Resource`] edges this node must acquire to run, derived from /// its kind + agent: a build slot for nix-heavy kinds /// ([`NodeKind::needs_build_slot`]) and the agent lease for /// container-affecting kinds ([`NodeKind::needs_lease`]). Lease-exempt @@ -48,16 +50,25 @@ impl NodeKind { /// (`try_acquire_all`) — a node never holds one resource while waiting on /// another, so the multi-resource kinds (a `MetaLock` wants a build slot /// *and* the meta window) cannot deadlock against each other. - pub fn resource_deps(&self) -> Vec<(Resource, u32)> { + pub fn resource_deps(&self) -> Vec> { let mut deps = Vec::new(); if self.needs_build_slot() { - deps.push((Resource::BuildSlot, 1)); + deps.push(Dep::Resource { + name: Resource::BuildSlot, + count: 1, + }); } if self.needs_lease() { - deps.push((Resource::Agent(self.agent().to_owned()), 1)); + deps.push(Dep::Resource { + name: Resource::Agent(self.agent().to_owned()), + count: 1, + }); } if self.needs_meta_window() { - deps.push((Resource::MetaWindow, 1)); + deps.push(Dep::Resource { + name: Resource::MetaWindow, + count: 1, + }); } deps } diff --git a/hive-c0re/src/job_queue/scheduler.rs b/hive-c0re/src/job_queue/scheduler.rs index 3efdf913..105d5326 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 7aaf1878..43955ec9 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, rebuild_nodes}`), all declaring into ONE job -//! (independent per-agent roots, concurrent on their own leases). +//! (`templates::{node, after_ok, rebuild_nodes}`), and concatenate them into +//! ONE DAG (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,13 +24,13 @@ use std::sync::Arc; -use super::model::{DagSpec, NodeKind}; -use super::templates::{RebuildOpts, node, rebuild_nodes}; -use super::{Job, Source, templates}; +use super::model::{DagSpec, Dep, NodeKind, NodeSpec}; +use super::templates::{RebuildOpts, after_ok, child, node, rebuild_nodes}; +use super::{Source, templates}; use crate::coordinator::Coordinator; use crate::lifecycle; -fn submit_and_emit(coord: &Arc, spec: super::DagSpec) -> u64 { +fn submit_and_emit(coord: &Arc, spec: super::DagSpec) -> u64 { let id = coord .job_queue .submit(spec) @@ -51,75 +51,71 @@ 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 declares into the shared job it is handed, and -// names the nodes it depends on — so there is nothing to rebase. +// then hand it in. Each chain uses LOCAL (0-based) deps; `concat_subgraphs` +// rebases them into one DAG. /// 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(b: &Job, agent: &str, graceful: bool, running: bool) { +fn stop_chain(agent: &str, graceful: bool, running: bool) -> Vec { // `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 wanted = node( - b, + let mut n = vec![node( NodeKind::SetWanted { agent: a(), up: false, }, - ); - // Declaration order is dependency order: the quiesce steps come first so - // the `Reconcile` that waits on them can name them. + Vec::new(), + )]; if graceful && running { - 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); + 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))); } else { - let _ = node(b, NodeKind::Reconcile { agent: a() }).part_of(wanted); + n.push(child(0, NodeKind::Reconcile { agent: a() }, Vec::new())); } + 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(b: &Job, agent: &str, running: bool, stale: bool) { - let wanted = node( - b, +fn start_chain(agent: &str, running: bool, stale: bool) -> Vec { + let mut n = vec![node( NodeKind::SetWanted { agent: agent.to_owned(), up: true, }, - ); + Vec::new(), + )]; if !running && stale { - // Rebuild subtree chained behind the `SetWanted` head. `MetaSync`, + // Rebuild subtree after the SetWanted head (base = 1, so the rebuild's + // `MetaSync` root deps `after_ok(0)` = the head). `MetaSync`, // `Prebuild` + `Reconcile` are their own group roots (top-level, per // `rebuild_nodes`). - rebuild_nodes( - b, + n.extend(rebuild_nodes( agent, RebuildOpts { relock: true, graceful: false, }, - Some(wanted), - ); + 1, + )); } else { - let _ = node( - b, + n.push(child( + 0, NodeKind::Reconcile { agent: agent.to_owned(), }, - ) - .part_of(wanted); + Vec::new(), + )); } + n } /// One agent's **restart** subgraph. Restart NEVER rewrites `wanted` @@ -131,49 +127,82 @@ fn start_chain(b: &Job, agent: &str, running: bool, stale: bool) { /// 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(b: &Job, agent: &str, graceful: bool, running: bool) { +fn restart_chain(agent: &str, graceful: bool, running: bool) -> Vec { let a = || agent.to_owned(); if !running { // Nothing to bounce — a lone Reconcile converges to intent. - let _ = node(b, NodeKind::Reconcile { agent: a() }); - return; + return vec![node(NodeKind::Reconcile { agent: a() }, Vec::new())]; } // 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. - // - // `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 { - 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); + let mut n = vec![if graceful { + node(NodeKind::Signal { agent: a() }, Vec::new()) } else { - let stop = node(b, NodeKind::StopForUpdate { agent: a() }); - let _ = node(b, NodeKind::Reconcile { agent: a() }).part_of(stop); + node(NodeKind::StopForUpdate { agent: a() }, Vec::new()) + }]; + 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)) + } else { + Vec::new() + }; + n.push(child(0, NodeKind::Reconcile { agent: a() }, deps)); + n } -/// Wrap the per-agent subgraphs in a `DagSpec`. No tail node: a power op's +/// 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), + }); + } + } + out +} + +/// Wrap assembled power-op `nodes` in a `DagSpec`. No tail node: a power op's /// effect is its nodes (`SetWanted` + `Reconcile`), with nothing left to do once /// they settle. -/// -/// 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, declare: F) -> DagSpec { +fn power_dag(source: Source, reason: String, nodes: Vec) -> DagSpec { DagSpec { source, reason, - declare, + nodes, } } @@ -188,17 +217,12 @@ pub(crate) fn stop_spec( graceful: bool, source: Source, reason: String, -) -> DagSpec> { - let targets = targets.to_vec(); - power_dag( - source, - reason, - Box::new(move |b: &Job| { - for (agent, running) in targets { - stop_chain(b, &agent, graceful, running); - } - }), - ) +) -> DagSpec { + let chains = targets + .iter() + .map(|(agent, running)| stop_chain(agent, graceful, *running)) + .collect(); + power_dag(source, reason, concat_subgraphs(chains)) } /// Assemble the start DAG from explicit `(agent, running, stale)` targets. @@ -211,17 +235,12 @@ pub(crate) fn start_spec( targets: &[(String, bool, bool)], source: Source, reason: String, -) -> DagSpec> { - let targets = targets.to_vec(); - power_dag( - source, - reason, - Box::new(move |b: &Job| { - for (agent, running, stale) in targets { - start_chain(b, &agent, running, stale); - } - }), - ) +) -> DagSpec { + let chains = targets + .iter() + .map(|(agent, running, stale)| start_chain(agent, *running, *stale)) + .collect(); + power_dag(source, reason, concat_subgraphs(chains)) } /// Assemble the restart DAG from explicit `(agent, running)` targets. @@ -230,17 +249,12 @@ pub(crate) fn restart_spec( graceful: bool, source: Source, reason: String, -) -> DagSpec> { - let targets = targets.to_vec(); - power_dag( - source, - reason, - Box::new(move |b: &Job| { - for (agent, running) in targets { - restart_chain(b, &agent, graceful, running); - } - }), - ) +) -> DagSpec { + let chains = targets + .iter() + .map(|(agent, running)| restart_chain(agent, graceful, *running)) + .collect(); + power_dag(source, reason, concat_subgraphs(chains)) } /// 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 7d29b946..70c6d6a9 100644 --- a/hive-c0re/src/job_queue/templates.rs +++ b/hive-c0re/src/job_queue/templates.rs @@ -1,7 +1,21 @@ //! DAG shape builders — every operation as a template over the shared -//! node primitives. Pure (no I/O); each node carries its own `agent` (there is -//! no DAG-level agent), stamped by the [`node`] helper along with the -//! resources that node's kind needs. +//! node primitives — plus submit-time cycle validation (petgraph is +//! confined to this validation; the runtime store stays the plain +//! `Vec` + `deps`). +//! +//! 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`. //! //! ```text //! rebuild(a): MetaSync(a) → Prebuild(a) → StopForUpdate(a) → Swap(a) →(ok) PostSwap(a) →(any) Reconcile(a) @@ -12,76 +26,113 @@ //! reparent(moves): Reparent(moves) [no rebuild — topology.json is read live] //! ``` //! -//! Nodes are **named, not counted** — a template holds the handle -//! [`Job::node`] hands back, so an edge says which node it waits on. Why that -//! removes submit-time cycle validation: `docs/coordinator.md`. -//! -//! The hive-wide **power ops** (`stop` / `start` / `restart`) are NOT here: -//! their per-agent shape depends on live running state (an async -//! `lifecycle::is_running` read), so `submit.rs` assembles them out of the -//! primitives this module exports ([`node`], [`rebuild_nodes`]). +//! For the dynamic power-op shapes (`stop` / `start` / `restart`, built from +//! live online/offline state), see `submit.rs`. -use hive_jobq::TerminalState; +use anyhow::{Result, bail}; -use super::model::{DagSpec, NodeKind, PermPayload, Source}; -use super::{Declare, Handle, Job}; +use hive_jobq::{DepWhen, TerminalState}; -/// Declare one node carrying `kind`, with the resources that kind needs. +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. /// -/// 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. +/// 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 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 +/// 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 `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. +/// elimination. `base` is the spec index the pair starts at. /// /// Exactly one runs on a DAG that executed, and neither runs on one the operator -/// 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, - ); +/// dropped — see [`on_elimination_of`]. +fn emit_rebuilt_tails(agent: &str, roots: &[u64], base: u64) -> Vec { // 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 _failed = roots.iter().fold( + 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), + ), node( - b, NodeKind::EmitRebuilt { agent: agent.to_owned(), ok: false, }, - ) - .on_elimination_of(ok), - hive_jobq::NodeRef::after_any, - ); + on_fail, + ), + ] } /// The approval-resolving tails for an approval-carrying DAG: one per outcome of @@ -89,20 +140,48 @@ fn emit_rebuilt_tails(b: &Job, agent: &str, roots: &[Handle<'_>]) { /// /// 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(b: &Job, approval_id: i64, root: Handle<'_>) { - for outcome in [ +fn resolve_approval_tails(approval_id: i64, root: u64) -> Vec { + [ TerminalState::Done, TerminalState::Failed, TerminalState::Cancelled, - ] { - let _ = node( - b, + ] + .into_iter() + .map(|outcome| { + node( NodeKind::ResolveApproval { approval_id, outcome, }, + on_outcome(root, &[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), } } @@ -119,51 +198,28 @@ pub(crate) struct RebuildOpts { pub graceful: bool, } -/// 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 +/// 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 /// `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 @@ -172,47 +228,56 @@ impl<'a> RebuildRoots<'a> { /// 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<'a>( - b: &'a Job, - agent: &str, - opts: RebuildOpts, - after: Option>, -) -> RebuildRoots<'a> { +pub(crate) fn rebuild_nodes(agent: &str, opts: RebuildOpts, base: u64) -> Vec { let a = || agent.to_owned(); let RebuildOpts { relock, graceful } = opts; - - 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); - + 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)), + ]; // The stop root hangs off `Prebuild` and owns the agent lease for - // 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); + // everything below it. + let stop_root = base + 2; + if graceful { + nodes.push(child(base + 1, NodeKind::Signal { agent: a() }, Vec::new())); // `Drain` is a *child* of `Signal`, so the parent gate already orders // it — a child must not dep on its own parent (dep-scope). - let drain = node(b, NodeKind::Drain { agent: a() }).part_of(signal); - node(b, NodeKind::StopForUpdate { agent: a() }) - .part_of(signal) - .after_ok(drain) + 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), + )); } else { - 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, + nodes.push(child( + base + 1, + NodeKind::StopForUpdate { agent: a() }, + Vec::new(), + )); } + // 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 @@ -223,7 +288,7 @@ pub(crate) fn rebuild_nodes<'a>( /// relocked and staged `flake.lock`, so the appended `MetaSync` must do the dir /// prep + `sync_agents` *without* re-locking over it. /// -/// `FinalizeDeploy` waits on **two** roots, which together reproduce the gate +/// `FinalizeDeploy` waits on **two** siblings, which together reproduce the gate /// the old fused node had around its inline `rebuild_no_meta` call: /// - `AfterOk` `Prebuild` — a parent's state is its roll-up, so this is `Done` /// only once `StopForUpdate` → `Swap` → `PostSwap` all are (a failed *or* @@ -239,28 +304,40 @@ pub(crate) fn rebuild_nodes<'a>( /// inside the `DeployWindow`'s subtree — so the `MetaWindow` this subgraph's /// `MetaSync` and `FinalizeDeploy` declare is re-entered from the ancestor /// already holding it rather than deadlocking against it. -pub(crate) fn deploy_rebuild_nodes(agent: &str, approval_id: i64) -> Declare { - let agent = agent.to_owned(); - Box::new(move |b: &Job| { - let roots = rebuild_nodes( - b, - &agent, - RebuildOpts { - relock: false, - graceful: false, +pub(crate) fn deploy_rebuild_nodes(agent: &str, approval_id: i64) -> Vec { + let mut nodes = rebuild_nodes( + agent, + RebuildOpts { + relock: false, + graceful: false, + }, + 0, + ); + let reconcile = reconcile_index(&nodes, 0); + nodes.push(node( + NodeKind::FinalizeDeploy { + agent: agent.to_owned(), + approval_id, + }, + vec![ + Dep { + on: 1, + when: DepWhen::AFTER_OK, }, - None, - ); - let _finalize = node( - b, - NodeKind::FinalizeDeploy { - agent: agent.clone(), - approval_id, + Dep { + on: reconcile, + when: DepWhen::AFTER_OK, }, - ) - .after_ok(roots.prebuild) - .after_ok(roots.reconcile); - }) + ], + )); + 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) } /// One uniform rebuild shape — no `was_running` branch. `StopForUpdate` @@ -274,48 +351,42 @@ pub(crate) fn deploy_rebuild_nodes(agent: &str, approval_id: i64) -> Declare { /// whole `StopForUpdate`→`Swap`→`PostSwap` subtree, so those three cover every /// 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 agent = agent.to_owned(); +pub fn rebuild(agent: &str, source: Source, reason: String, relock: bool) -> DagSpec { + let mut nodes = rebuild_nodes( + agent, + RebuildOpts { + relock, + graceful: false, + }, + 0, + ); + 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)); DagSpec { source, reason, - declare: Box::new(move |b: &Job| { - let roots = rebuild_nodes( - b, - &agent, - RebuildOpts { - relock, - graceful: false, - }, - None, - ); - emit_rebuilt_tails(b, &agent, &roots.all()); - }), + nodes, } } /// Approval-driven deploy (`MergeConfigPr`) as a phase subtree rather than the /// single opaque node it used to be. Structure: -/// - `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 — +/// - `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 — /// 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` (child, `AfterAny` `DeployApply`): the compensation + +/// - `DeployTail` (3, child, `AfterAny` `DeployApply`): the compensation + /// bookkeeping tail — rollback when a merge landed unfinalized, forge tag /// mirror, PR failure comment (see [`NodeKind::DeployTail`]). /// -/// - `ResolveApproval` (**root**, `AfterAny` `DeployWindow`): resolves the +/// - `ResolveApproval` (4, **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 @@ -324,53 +395,50 @@ pub fn rebuild( /// /// The window still spans the container build, as it must: `prepare_deploy` /// leaves `flake.lock` staged-uncommitted for the build's whole duration. -pub fn approval_deploy( - agent: &str, - approval_id: i64, - reason: String, -) -> DagSpec> { - let agent = agent.to_owned(); +pub fn approval_deploy(agent: &str, approval_id: i64, reason: String) -> DagSpec { + let a = || agent.to_owned(); DagSpec { source: Source::Approval, reason, - declare: Box::new(move |b: &Job| { - let a = || agent.clone(); - let window = node( - b, + nodes: vec![ + node( NodeKind::DeployWindow { agent: a(), approval_id, }, - ); - let verify = node( - b, + Vec::new(), + ), + child( + 0, NodeKind::MergeVerify { agent: a(), approval_id, }, - ) - .part_of(window); - let apply = node( - b, + Vec::new(), + ), + child( + 0, NodeKind::DeployApply { agent: a(), approval_id, }, - ) - .part_of(window) - .after_ok(verify); - let _tail = node( - b, + after_ok(1), + ), + child( + 0, NodeKind::DeployTail { agent: a(), approval_id, }, - ) - .part_of(window) - .after_any(apply); - - resolve_approval_tails(b, approval_id, window); - }), + vec![Dep { + on: 2, + when: DepWhen::AFTER_ANY, + }], + ), + ] + .into_iter() + .chain(resolve_approval_tails(approval_id, 0)) + .collect(), } } @@ -380,18 +448,16 @@ pub fn approval_deploy( /// single-node lifecycle DAGs that exercise per-agent lease serialization /// 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 agent = agent.to_owned(); +pub fn reconcile_only(agent: &str, source: Source, reason: String) -> DagSpec { DagSpec { source, reason, - declare: Box::new(move |b: &Job| { - let _reconcile = node(b, NodeKind::Reconcile { agent }); - }), + nodes: vec![node( + NodeKind::Reconcile { + agent: agent.to_owned(), + }, + Vec::new(), + )], } } @@ -406,63 +472,54 @@ pub fn reconcile_only( /// container was never created). Closed by a `ResolveApproval` tail root edged /// `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 agent = agent.to_owned(); +pub fn spawn(agent: &str, approval_id: i64, reason: String) -> DagSpec { DagSpec { source: Source::Approval, reason, - declare: Box::new(move |b: &Job| { - let a = || agent.clone(); - let provision = node(b, NodeKind::Provision { agent: a() }); - let create = node(b, NodeKind::Create { agent: a() }).part_of(provision); - let dropin = node(b, NodeKind::WriteDropin { agent: a() }).part_of(create); - let _reconcile = node(b, NodeKind::Reconcile { agent: a() }) - .part_of(create) - .after_ok(dropin); - - resolve_approval_tails(b, approval_id, provision); - }), + 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() + }, } } /// 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` 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 agent = agent.to_owned(); +/// 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. +pub fn perm_change(agent: &str, source: Source, reason: String, payload: PermPayload) -> DagSpec { + let mut nodes = vec![node( + NodeKind::WritePermFile { + agent: agent.to_owned(), + payload, + }, + Vec::new(), + )]; + let rebuild = rebuild_nodes( + agent, + RebuildOpts { + relock: true, + graceful: false, + }, + 1, + ); + 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, - declare: Box::new(move |b: &Job| { - let write = node( - b, - NodeKind::WritePermFile { - agent: agent.clone(), - payload, - }, - ); - let roots = rebuild_nodes( - b, - &agent, - RebuildOpts { - relock: true, - graceful: false, - }, - Some(write), - ); - emit_rebuilt_tails( - b, - &agent, - &[write, roots.meta_sync, roots.prebuild, roots.reconcile], - ); - }), + nodes, } } @@ -481,28 +538,26 @@ pub fn meta_update( source: Source, reason: String, approval_id: Option, -) -> DagSpec> { +) -> DagSpec { + let mut nodes = vec![node( + 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)); + } DagSpec { source, reason, - declare: Box::new(move |b: &Job| { - let lock = node( - b, - NodeKind::MetaLock { - sweep: false, - fanout: None, - inputs, - }, - ); - // 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 { - resolve_approval_tails(b, approval_id, lock); - } - }), + nodes, } } @@ -519,13 +574,11 @@ pub fn reparent( moves: Vec<(hive_types::Ident, Option)>, source: Source, reason: String, -) -> DagSpec> { +) -> DagSpec { DagSpec { source, reason, - declare: Box::new(move |b: &Job| { - let _reparent = node(b, NodeKind::Reparent { moves }); - }), + nodes: vec![node(NodeKind::Reparent { moves }, Vec::new())], } } @@ -533,3 +586,45 @@ 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 8b026b01..49d55bd4 100644 --- a/hive-c0re/src/job_queue/tests.rs +++ b/hive-c0re/src/job_queue/tests.rs @@ -6,32 +6,20 @@ //! scheduler's async loop is a thin claim/complete pump over the same //! methods exercised here. -use super::model::NodeKind; +use hive_jobq::DepWhen; + +use super::model::{Dep, NodeKind, NodeSpec}; use super::*; -fn submit(q: &JobQueue, spec: DagSpec) -> u64 { +fn submit(q: &JobQueue, spec: DagSpec) -> u64 { q.submit(spec).expect("valid spec") } -/// Erase a spec's recipe to the boxed [`Declare`] so specs of *different* -/// shapes can share one type — e.g. a table of `(name, spec)` cases. -/// -/// Production never needs this: each submit path builds one spec and hands it -/// straight to `submit`, so the concrete closure type is known end to end. A -/// test table is the case where several shapes must be one type. -fn erase(spec: DagSpec) -> DagSpec { - DagSpec { - source: spec.source, - reason: spec.reason, - declare: Box::new(spec.declare), - } -} - fn ident(s: &str) -> hive_types::Ident { hive_types::Ident::parse(s).expect("valid test ident") } -fn rebuild(agent: &str, reason: &str) -> DagSpec> { +fn rebuild(agent: &str, reason: &str) -> DagSpec { templates::rebuild(agent, Source::Manual, reason.to_owned(), true) } @@ -39,22 +27,14 @@ fn rebuild(agent: &str, reason: &str) -> DagSpec> { /// shape (`[Signal→Drain→] StopForUpdate → Reconcile`, no `SetWanted` head) /// most queue-mechanics tests assume. Mirrors the pre-dynamic /// `templates::restart` (which is now the state-aware `submit::restart_spec`). -fn restart_online( - agents: &[&str], - graceful: bool, - reason: &str, -) -> DagSpec> { +fn restart_online(agents: &[&str], graceful: bool, reason: &str) -> DagSpec { let targets: Vec<(String, bool)> = agents.iter().map(|a| ((*a).to_owned(), true)).collect(); submit::restart_spec(&targets, graceful, Source::Manual, reason.to_owned()) } /// Stop DAG spec with every agent treated as **running** — the online shape /// (`SetWanted → [Signal→Drain→](graceful) Reconcile`). -fn stop_online( - agents: &[&str], - graceful: bool, - reason: &str, -) -> DagSpec> { +fn stop_online(agents: &[&str], graceful: bool, reason: &str) -> DagSpec { let targets: Vec<(String, bool)> = agents.iter().map(|a| ((*a).to_owned(), true)).collect(); submit::stop_spec(&targets, graceful, Source::Manual, reason.to_owned()) } @@ -161,22 +141,71 @@ fn resubmit_while_running_is_new_dag() { assert_eq!(q.snapshot().len(), 2); } -// ---- 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`). +// ---- 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()); +} // ---- dependency order within a DAG ---- @@ -214,24 +243,20 @@ fn rebuild_chain_claims_in_dep_order() { #[test] fn graceful_rebuild_chain_drains_before_stopping() { let q = JobQueue::new(1); - let id = submit( - &q, - DagSpec { - source: Source::AutoUpdate, - reason: "sweep".to_owned(), - declare: Box::new(|b: &Job| { - templates::rebuild_nodes( - b, - "agent-a", - templates::RebuildOpts { - relock: true, - graceful: true, - }, - None, - ); - }), - }, - ); + 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); for expected in [ "meta_sync", "prebuild", @@ -259,34 +284,17 @@ 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() { - // 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 id = submit( - &q, - DagSpec { - source: Source::Manual, - reason: "manual".to_owned(), - declare: Box::new(|b: &Job| { - templates::rebuild_nodes( - b, - "agent-a", - templates::RebuildOpts { - relock: true, - graceful: false, - }, - None, - ); - }), + let kinds: Vec = templates::rebuild_nodes( + "agent-a", + templates::RebuildOpts { + relock: true, + graceful: false, }, - ); - 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(())); - } + 0, + ) + .iter() + .map(|n| n.kind.as_str().to_owned()) + .collect(); assert_eq!( kinds, vec![ @@ -298,8 +306,6 @@ 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 @@ -717,16 +723,16 @@ fn append_subgraph_roots_on_emitter_and_rebases_local_deps() { let spec = DagSpec { source: Source::AutoUpdate, reason: "sweep".to_owned(), - declare: Box::new(|b: &Job| { - let _lock = templates::node( - b, - NodeKind::MetaLock { - sweep: true, - fanout: None, - inputs: Vec::new(), - }, - ); - }), + + nodes: vec![NodeSpec { + kind: NodeKind::MetaLock { + sweep: true, + fanout: None, + inputs: Vec::new(), + }, + deps: Vec::new(), + parent: None, + }], }; let id = submit(&q, spec); let emitter = claim_one(&q); @@ -735,23 +741,19 @@ fn append_subgraph_roots_on_emitter_and_rebases_local_deps() { // sweep MetaLock grows: root MetaSync → root Prebuild → Signal → Drain → // 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| -> Declare { - let agent = agent.to_owned(); - Box::new(move |b: &Job| { - templates::rebuild_nodes( - b, - &agent, - templates::RebuildOpts { - relock: true, - graceful: true, - }, - None, - ); - }) + let subgraph = |agent: &str| { + templates::rebuild_nodes( + agent, + templates::RebuildOpts { + relock: true, + graceful: true, + }, + 0, + ) }; // 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 @@ -843,18 +845,18 @@ 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"] { - let declare: Declare = Box::new(move |b: &Job| { - templates::rebuild_nodes( - b, + q.append_subgraph( + id, + &templates::rebuild_nodes( agent, templates::RebuildOpts { relock: false, graceful: false, }, - None, - ); - }); - q.append_subgraph(id, declare, meta_lock.node_id); + 0, + ), + meta_lock.node_id, + ); } q.complete_node(meta_lock.node_id, Ok(())); // Still ONE DAG — no child DAGs — and both cascade rebuild subgraphs root @@ -1095,55 +1097,37 @@ fn cancelled_power_op_runs_no_compensating_node() { for graceful in [false, true] { for running in [false, true] { let targets = vec![("agent-a".to_owned(), running)]; - // Erased to `DagSpec`: three different recipe types have to - // sit in one array. let cases = [ ( "restart", false, - erase(submit::restart_spec( - &targets, - graceful, - Source::Manual, - "bounce".to_owned(), - )), + submit::restart_spec(&targets, graceful, Source::Manual, "bounce".to_owned()), ), ( "stop", true, - erase(submit::stop_spec( - &targets, - graceful, - Source::Manual, - "stop".to_owned(), - )), + submit::stop_spec(&targets, graceful, Source::Manual, "stop".to_owned()), ), ( "start", true, - erase(submit::start_spec( + submit::start_spec( &[("agent-a".to_owned(), running, false)], Source::Manual, "start".to_owned(), - )), + ), ), ]; 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!( - q.snapshot() + spec.nodes .iter() - .find(|d| d.id == id) - .expect("submitted dag") - .nodes - .iter() - .any(|n| n.kind == "set_wanted"), + .any(|n| matches!(n.kind, NodeKind::SetWanted { .. })), 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!( @@ -1286,11 +1270,12 @@ fn deploy_apply_grows_rebuild_subgraph_and_finalizes_after_it() { // BEFORE the emitting node is completed. Completing first would settle the // apply node `Done` with nothing under it, opening the tail's `AfterAny` // gate immediately and letting the deploy "finish" before it had built. - q.append_subgraph( + 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"); q.complete_node(apply.node_id, Ok(())); // The grafted chain runs in rebuild order. `claim_one` asserts exactly one @@ -1345,7 +1330,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 7753c7f5..599d6796 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, Source, templates}; + use crate::job_queue::{DagSpec, NodeKind, NodeSpec, Source}; // Fully-quiet boot (nothing stale, nothing drifted) submits nothing. if !any_stale && drifted.is_empty() { @@ -348,29 +348,33 @@ fn submit_boot_tree( n_skipped, ); - let declare: crate::job_queue::Declare = Box::new(move |b| { - // 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 { - let _ = templates::node( - b, - NodeKind::MetaLock { - sweep: true, - fanout: Some(fanout), - // A sweep bumps `hyperhive` alone (`lock_update_hyperhive`), - // so it names no inputs. - inputs: Vec::new(), - }, - ); - } - // One boot Reconcile per drifted agent — independent roots. - for name in drifted { - let _ = templates::node(b, NodeKind::Reconcile { agent: name }); - } - }); + let mut nodes: Vec = Vec::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 { + 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 spec = DagSpec { // The sweep's own rebuild subgraphs emit their `Rebuilt` events as they @@ -380,7 +384,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. - declare, + nodes, }; if let Err(e) = coord.job_queue.submit(spec) { tracing::warn!(error = ?e, "boot: sweep DAG submit failed"); diff --git a/hive-jobq/Cargo.toml b/hive-jobq/Cargo.toml index 845aa46f..4fc70e0a 100644 --- a/hive-jobq/Cargo.toml +++ b/hive-jobq/Cargo.toml @@ -12,7 +12,6 @@ chrono = { workspace = true } enumflags2.workspace = true serde = { workspace = true } thiserror = { workspace = true } -uuid = { workspace = true } [dev-dependencies] serde_json = { workspace = true } diff --git a/hive-jobq/README.md b/hive-jobq/README.md index eac5c262..46bcad53 100644 --- a/hive-jobq/README.md +++ b/hive-jobq/README.md @@ -16,9 +16,8 @@ kinds, wires deps, and supplies a runner; the scheduler decides what can start. ## Model One **persistent graph** for the whole system, not a DAG per job. Enqueuing -inserts a self-contained sub-DAG and returns the ids of the nodes the job -*asked* for, in the order it named them; the scheduler runs a continuous loop, -starting every node whose deps are satisfied: +inserts a self-contained sub-DAG and returns the new node ids; the scheduler +runs a continuous loop, starting every node whose deps are satisfied: - **Resource deps** are named counting semaphores over a caller-chosen type `R` — e.g. `build-slot` (capacity N), `agent/` (capacity 1), or any diff --git a/hive-jobq/src/builder.rs b/hive-jobq/src/builder.rs deleted file mode 100644 index ed894cae..00000000 --- a/hive-jobq/src/builder.rs +++ /dev/null @@ -1,782 +0,0 @@ -//! Build a job's sub-DAG by **naming** nodes instead of counting them. -//! -//! [`JobBuilder::node`] hands back a handle, obtainable no other way; edges are -//! handle → handle. So a caller names the node it wants to depend on rather -//! than computing where that node landed, and there is no positional index to -//! get wrong. -//! -//! **An insertion API, not a spec factory.** A builder is only ever handed to a -//! closure by the single insertion entry point -//! ([`crate::scheduler::Scheduler::insert_job`]), which inserts the declared -//! nodes and returns the ids the job asked for. It cannot be constructed, held -//! or inserted from outside this crate, and there is no intermediate -//! node-description type to keep in sync with [`crate::Graph::insert`]'s signature — -//! so a job has no representation that can be passed around instead of being -//! inserted. -//! -//! **Payload-agnostic.** Generic over the same `N` and `R` as [`crate::Graph`]: the -//! builder knows nothing about what a node *does*, only how nodes relate. -//! -//! # Declaration order -//! -//! Nodes are inserted in declaration order, so a node must be declared *after* -//! everything it references. That is not a limitation the builder invents: node -//! ids are minted by the graph at insert time, so a forward edge has nothing to -//! point at. Declaring one is a [`BuildError::ForwardEdge`] rather than a -//! silent reorder — a builder that sorted for you would quietly accept a shape -//! the graph itself cannot express. - -use std::cell::RefCell; -use std::collections::HashMap; - -use crate::{Dep, DepWhen, GraphError, NodeId, TerminalState}; - -/// An opaque identity for a node **within the job being built**. -/// -/// Obtainable only from [`JobBuilder::node`], so it can only ever name a node -/// that exists. Deliberately not a [`NodeId`]: those are minted by the graph at -/// insert time and are meaningful system-wide, while this is a builder-local -/// label that stops existing once the job is inserted. -/// -/// It is a v4 uuid because it has to survive leaving its builder. A [`NodeRef`] -/// converts into a bare `NodeGuid`, dropping the borrow that tied it to the -/// builder — so a handle *can* be carried into a different job (captured by an -/// inner closure, say). Under a per-builder counter the two jobs' `0`s would be -/// equal and that handle would silently address an unrelated node; drawn at -/// random per node it simply isn't found, and the insert fails by name. -#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)] -pub struct NodeGuid(uuid::Uuid); - -/// Why a job could not be inserted. -#[derive(Debug, Clone, PartialEq, Eq, thiserror::Error)] -pub enum BuildError { - /// A node depends on one declared later. See the module docs on - /// declaration order. - #[error( - "node {node:?} depends on {dep:?}, which is declared later — a dependency \ - must be declared before the node that names it" - )] - ForwardEdge { - /// The node carrying the edge. - node: NodeGuid, - /// The not-yet-declared target. - dep: NodeGuid, - }, - /// A node is grouped under a parent declared later. Same cause as - /// [`BuildError::ForwardEdge`]. - #[error( - "node {node:?} is grouped under {parent:?}, which is declared later — a \ - parent must be declared before its children" - )] - ForwardParent { - /// The child node. - node: NodeGuid, - /// The not-yet-declared parent. - parent: NodeGuid, - }, - /// A handle named in the closure's return value belongs to a different - /// job. Same cause as a foreign handle on an edge: a [`NodeGuid`] outlives - /// the borrow that tied it to its builder, so one can be carried here. - #[error("handle {node:?} names no node in this job")] - UnknownNode { - /// The handle that resolved to nothing. - node: NodeGuid, - }, - /// An edge leaves the depender's own parent group — including an edge onto - /// its **own parent**, which would deadlock: a parent parks in - /// `Finishing` until its children settle, so a child waiting on it could - /// never run. The graph enforces the same rule - /// (`GraphError::DepOutsideParent`); this catches it before any insert. - #[error( - "node {node:?} depends on {dep:?}, which is outside its parent group — \ - an edge must stay within the depender's own group" - )] - DepOutsideGroup { - /// The node carrying the edge. - node: NodeGuid, - /// The out-of-group target. - dep: NodeGuid, - }, - /// An edge that no outcome can satisfy (an empty [`DepWhen`]), so the node - /// could never become runnable. Refused at declaration rather than - /// inserted — the graph only catches this when validating a *deserialized* - /// graph, so without this check it would insert and silently never run. - #[error("node {node:?} has an edge on {dep:?} that no outcome can satisfy")] - UnsatisfiableDep { - /// The node carrying the edge. - node: NodeGuid, - /// The target of the unsatisfiable edge. - dep: NodeGuid, - }, - /// The graph rejected an otherwise well-formed node — an out-of-group edge, - /// an unsatisfiable [`DepWhen`], and so on. - #[error(transparent)] - Graph(#[from] GraphError), -} - -/// Reject a job whose own declarations don't hold up — **before anything is -/// inserted**, so no failure can leave a partial job behind. -/// -/// This covers *every* rejection [`crate::Graph::insert`] can raise for a -/// builder-produced node, which is what makes the insert loop below infallible -/// in practice: -/// -/// | graph rejection | why it cannot reach the graph | -/// |---|---| -/// | `UnknownParent` / `UnknownDep` | a handle only exists if this job declared it, and the ids are minted here | -/// | `DepOutsideParent` | the grouping rule is re-checked here against the job's own parent chains | -/// -/// It also rejects an **empty [`DepWhen`]**, which the graph currently only -/// catches on the deserialize path (`Graph::validate`) — so a node that nothing -/// could ever satisfy is refused at declaration instead of being inserted and -/// silently never running. -/// -/// `root_parent` matters for the grouping rule: a node that declared no parent -/// hangs there, so what counts as "inside the group" depends on whether the job -/// is being attached under a container or at the top level. -/// -/// # Errors -/// [`BuildError::ForwardEdge`] / [`BuildError::ForwardParent`] for a reference -/// to a later node, [`BuildError::UnknownNode`] for a requested handle this job -/// never declared, [`BuildError::DepOutsideGroup`] for an edge leaving the -/// depender's parent group, [`BuildError::UnsatisfiableDep`] for an empty edge. -fn check_job_shape( - pending: &[Pending], - wanted: &[NodeGuid], - root_parent: Option, -) -> Result<(), BuildError> { - let mut declared: std::collections::HashSet = std::collections::HashSet::new(); - // Declared parent per node, for walking a job-internal parent chain. - let mut parent_of: HashMap> = HashMap::new(); - - for node in pending { - for (dep, when) in &node.deps { - if !declared.contains(dep) { - return Err(BuildError::ForwardEdge { - node: node.guid, - dep: *dep, - }); - } - if when.is_empty() { - return Err(BuildError::UnsatisfiableDep { - node: node.guid, - dep: *dep, - }); - } - if !dep_in_group(&parent_of, node.parent, *dep, root_parent) { - return Err(BuildError::DepOutsideGroup { - node: node.guid, - dep: *dep, - }); - } - } - if let Some(parent) = node.parent - && !declared.contains(&parent) - { - return Err(BuildError::ForwardParent { - node: node.guid, - parent, - }); - } - declared.insert(node.guid); - parent_of.insert(node.guid, node.parent); - } - for guid in wanted { - if !declared.contains(guid) { - return Err(BuildError::UnknownNode { node: *guid }); - } - } - Ok(()) -} - -/// The graph's grouping rule (`Graph::dep_target_in_group`) decided against the -/// job's own declarations, before any node exists. -/// -/// A depender whose declared parent is `Some(q)` may only name a **proper -/// descendant of `q`** — never `q` itself, which would deadlock (a parent parks -/// in `Finishing` until its children settle). A depender that declared no parent -/// hangs under `root_parent`: -/// - `root_parent = Some(_)` — every node in this job is somewhere under it, so -/// any job-internal target is in-group. -/// - `root_parent = None` — the depender is top-level, so the target must be -/// top-level too, i.e. it must also have declared no parent. -fn dep_in_group( - parent_of: &HashMap>, - node_parent: Option, - dep: NodeGuid, - root_parent: Option, -) -> bool { - match node_parent { - // Walk the target's declared chain looking for the depender's parent. - // Starts at the target's *parent*, so the target is never its own - // ancestor — matching `Graph::is_descendant`. - Some(group) => { - let mut cur = parent_of.get(&dep).copied().flatten(); - while let Some(p) = cur { - if p == group { - return true; - } - cur = parent_of.get(&p).copied().flatten(); - } - false - } - None if root_parent.is_some() => true, - None => parent_of.get(&dep).copied().flatten().is_none(), - } -} - -/// One node as the builder holds it: edges and parent still name *handles*, so -/// nothing here depends on ids the graph has not minted yet. -#[derive(Debug)] -struct Pending { - guid: NodeGuid, - payload: N, - deps: Vec<(NodeGuid, DepWhen)>, - resources: Vec<(R, u32)>, - parent: Option, -} - -/// Accumulates a job's nodes, handing back a handle for each. -/// -/// Sub-builders take `&JobBuilder` and **return their handles**, which is what -/// removes the base-offset parameter a positional API needs: a caller names the -/// node it wants rather than computing where that node landed. -#[derive(Debug)] -pub struct JobBuilder { - nodes: RefCell>>, -} - -// Hand-written rather than derived: `#[derive(Default)]` would demand -// `N: Default, R: Default`, which has nothing to do with an empty builder. -impl Default for JobBuilder { - fn default() -> Self { - Self { - nodes: RefCell::new(Vec::new()), - } - } -} - -impl JobBuilder { - /// A fresh, empty builder. - /// - /// **Crate-private, and that is the API.** A builder is only ever handed to - /// a closure by the single insertion entry point - /// ([`crate::scheduler::Scheduler::insert_job`]), which inserts the declared - /// nodes and returns the ids. Nothing job-shaped is constructible or - /// carryable outside this crate — otherwise it is a spec factory again, - /// just with a builder's name on it. - pub(crate) fn new() -> Self { - 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 - /// can be named as a dependency as many times as needed. - pub fn node(&self, payload: N) -> NodeRef<'_, N, R> { - let guid = NodeGuid(uuid::Uuid::new_v4()); - self.nodes.borrow_mut().push(Pending { - guid, - payload, - deps: Vec::new(), - resources: Vec::new(), - parent: None, - }); - NodeRef { - builder: self, - guid, - } - } - - /// Insert every declared node into `graph`, in declaration order, and - /// return the id each handle's node was minted as. - /// - /// `root_parent` is the group's attachment point: a node that declared no - /// [`NodeRef::part_of`] hangs there rather than at the top level. That is - /// what makes a job a *self-contained sub-DAG* — the whole thing goes under - /// one container node, or under the runtime node that emitted it, and the - /// job's own declarations stay relative. - /// - /// The returned map is how a caller finds out what its handles became: hold - /// the handles it cares about, then look them up here. Handles that were - /// never looked up cost nothing. - /// - /// # Errors - /// - /// [`BuildError::ForwardEdge`] / [`BuildError::ForwardParent`] if a node - /// references one declared after it, or [`BuildError::Graph`] if the graph - /// rejects a node (see [`crate::Graph::insert`]). - pub(crate) fn insert_with( - self, - root_parent: Option, - wanted: &[NodeGuid], - mut insert: impl FnMut(N, Vec>, Option) -> Result, - ) -> Result, BuildError> { - let pending = self.nodes.into_inner(); - check_job_shape(&pending, wanted, root_parent)?; - - let mut ids: HashMap = HashMap::new(); - for node in pending { - let parent = match node.parent { - None => root_parent, - Some(p) => Some(ids[&p]), - }; - let mut deps: Vec> = Vec::with_capacity(node.deps.len()); - for (on, when) in node.deps { - deps.push(Dep::Node { id: ids[&on], when }); - } - deps.extend( - node.resources - .into_iter() - .map(|(name, count)| Dep::Resource { name, count }), - ); - let id = insert(node.payload, deps, parent)?; - ids.insert(node.guid, id); - } - Ok(wanted.iter().map(|g| ids[g]).collect()) - } - - /// Apply `f` to the node named by `guid`. - fn with_node(&self, guid: NodeGuid, f: impl FnOnce(&mut Pending)) { - let mut nodes = self.nodes.borrow_mut(); - if let Some(n) = nodes.iter_mut().find(|n| n.guid == guid) { - f(n); - } - } -} - -/// A handle to one declared node, and the surface its edges are declared on. -/// -/// [`Copy`] because it is only a builder borrow plus a handle: naming a node as -/// a dependency must not consume the ability to name it again. -#[derive(Debug)] -pub struct NodeRef<'a, N, R> { - builder: &'a JobBuilder, - guid: NodeGuid, -} - -// Hand-written for the same reason as `Default` above: deriving would demand -// `N: Copy, R: Copy`, but a handle copies a borrow and an integer. -impl Clone for NodeRef<'_, N, R> { - fn clone(&self) -> Self { - *self - } -} -impl Copy for NodeRef<'_, N, R> {} - -impl From> for NodeGuid { - fn from(n: NodeRef<'_, N, R>) -> Self { - n.guid - } -} - -impl From<&NodeRef<'_, N, R>> for NodeGuid { - fn from(n: &NodeRef<'_, N, R>) -> Self { - n.guid - } -} - -impl NodeRef<'_, N, R> { - /// This node's handle, for callers that want to hold the identity without - /// the builder borrow (e.g. to look the id up after - /// [`JobBuilder::insert_into`]). - #[must_use] - pub fn guid(self) -> NodeGuid { - self.guid - } - - /// Run only after `on` succeeds — the common chain link. - /// - /// Chained against several nodes it means *all* of them succeeded, and this - /// node is ruled out ([`TerminalState::Skipped`]) the moment one doesn't — - /// which is the signal [`NodeRef::on_elimination_of`] waits for. - #[must_use] - pub fn after_ok(self, on: impl Into) -> Self { - self.edge(on.into(), DepWhen::AFTER_OK) - } - - /// Run once `on` is terminal, however it ended. The cleanup edge: a node - /// that must run whether or not the work it follows succeeded. - #[must_use] - pub fn after_any(self, on: impl Into) -> Self { - 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. - /// - /// This is how *"any one of several nodes failed"* gets expressed at all. - /// Dependency edges are conjunctive, so that condition cannot be written - /// directly; the composition that expresses it is a pair — point the - /// success branch at every root with [`NodeRef::after_ok`], then hang the - /// failure branch off *that* node's elimination. Exactly one of the pair - /// ever runs. - /// - /// It accepts [`TerminalState::Skipped`] and **not** - /// [`TerminalState::Cancelled`] on purpose: if the whole job was dropped - /// before it started, the success branch is `Cancelled` directly and this - /// branch is ruled out too — a job nobody ran reports nothing. - #[must_use] - pub fn on_elimination_of(self, on: impl Into) -> Self { - self.edge(on.into(), DepWhen::of(&[TerminalState::Skipped])) - } - - /// Hold one unit of the named resource while this node and its subtree run. - /// - /// Declared on every node that needs it, even when an ancestor already - /// holds the same resource — the ancestor's grant is re-used (a re-entrant - /// borrow), so declaring it costs nothing and makes the node's needs - /// readable where the node is written. - #[must_use] - pub fn needs(self, name: R) -> Self { - self.needs_units(name, 1) - } - - /// [`NodeRef::needs`] with an explicit unit count, for a resource whose - /// capacity is a budget rather than a mutex. - #[must_use] - pub fn needs_units(self, name: R, count: u32) -> Self { - self.builder - .with_node(self.guid, |n| n.resources.push((name, count))); - self - } - - /// Group this node under `parent`. - /// - /// The child runs once its parent reaches [`crate::State::Finishing`] (the - /// parent gate), so it must not also depend on that parent — dep-scope - /// validation rejects that, and it would deadlock. A group root owns - /// whatever resource it declares for its whole subtree. - #[must_use] - pub fn part_of(self, parent: impl Into) -> Self { - let parent = parent.into(); - self.builder - .with_node(self.guid, |n| n.parent = Some(parent)); - self - } - - /// Record one dependency edge and hand the handle back for chaining. - fn edge(self, on: NodeGuid, when: DepWhen) -> Self { - self.builder - .with_node(self.guid, |n| n.deps.push((on, when))); - self - } -} - -#[cfg(test)] -mod tests { - use super::BuildError; - use crate::resources::ResourceTable; - use crate::scheduler::Scheduler; - use crate::{Dep, DepWhen, Graph, NodeId}; - - /// A scheduler over a graph whose payload is a name and whose resources are - /// strings — the only way in, since insertion is a scheduler operation. - fn sched() -> Scheduler<&'static str, &'static str> { - Scheduler::new(Graph::new(), ResourceTable::new()) - } - - fn deps_of(g: &Scheduler<&'static str, &'static str>, id: NodeId) -> Vec> { - g.graph().node(id).expect("node present").deps.clone() - } - - /// The point of the handle layer: an edge declared against a *handle* comes - /// out addressing the id that node was actually minted as. - #[test] - fn edges_resolve_to_minted_ids() { - let mut g = sched(); - let ids = g - .insert_job(None, |b| { - let first = b.node("a"); - let second = b.node("b").after_ok(first); - vec![first.guid(), second.guid()] - }) - .expect("insert"); - let [first, second] = ids[..] else { - panic!("two ids back, in the order asked for") - }; - assert_eq!( - deps_of(&g, second), - vec![Dep::Node { - id: first, - when: DepWhen::AFTER_OK - }] - ); - assert!(deps_of(&g, first).is_empty()); - } - - #[test] - fn parent_resolves_to_a_minted_id() { - let mut g = sched(); - let ids = g - .insert_job(None, |b| { - let root = b.node("a"); - let child = b.node("b").part_of(root); - vec![root.guid(), child.guid()] - }) - .expect("insert"); - let [root, child] = ids[..] else { - panic!("two ids back") - }; - assert_eq!(g.graph().node(root).expect("root").parent, None); - assert_eq!(g.graph().node(child).expect("child").parent, Some(root)); - } - - /// A handle is `Copy`, so naming the same node as a dependency twice must - /// not consume it — the fan-out every composite job needs. - #[test] - fn one_handle_can_be_depended_on_twice() { - let mut g = sched(); - let ids = g - .insert_job(None, |b| { - let shared = b.node("a"); - let ok = b.node("b").after_ok(shared); - let any = b.node("c").after_any(shared); - vec![shared.guid(), ok.guid(), any.guid()] - }) - .expect("insert"); - let [shared, ok, any] = ids[..] else { - panic!("three ids back") - }; - assert_eq!( - deps_of(&g, ok), - vec![Dep::Node { - id: shared, - when: DepWhen::AFTER_OK - }] - ); - assert_eq!( - deps_of(&g, any), - vec![Dep::Node { - id: shared, - when: DepWhen::AFTER_ANY - }] - ); - } - - /// Resource deps ride along with the node deps, in one insert. - #[test] - fn resources_become_resource_deps() { - let mut g = sched(); - let ids = g - .insert_job(None, |b| { - vec![ - b.node("a") - .needs("agent/atlas") - .needs_units("build", 2) - .guid(), - ] - }) - .expect("insert"); - let [only] = ids[..] else { - panic!("one id back") - }; - assert_eq!( - deps_of(&g, only), - vec![ - Dep::Resource { - name: "agent/atlas", - count: 1 - }, - Dep::Resource { - name: "build", - count: 2 - }, - ] - ); - } - - /// A handle stays live across later `node()` calls, so it is possible to - /// aim an edge *backwards* in declaration order. The graph mints ids as it - /// inserts, so there is nothing for that edge to point at — say so instead - /// of quietly reordering. - #[test] - fn a_forward_edge_is_rejected_by_name() { - let mut g = sched(); - let mut named = None; - let err = g - .insert_job(None, |b| { - let first = b.node("a"); - let second = b.node("b"); - let _ = first.after_any(second); - named = Some((first.guid(), second.guid())); - // The insert fails, so nothing comes back to ask for — the - // handles under test travel out through `named` instead. - Vec::new() - }) - .expect_err("forward edge"); - let (first, second) = named.expect("declared"); - assert_eq!( - err, - BuildError::ForwardEdge { - node: first, - dep: second - } - ); - } - - #[test] - fn a_forward_parent_is_rejected_by_name() { - let mut g = sched(); - let mut named = None; - let err = g - .insert_job(None, |b| { - let child = b.node("a"); - let parent = b.node("b"); - let _ = child.part_of(parent); - named = Some((child.guid(), parent.guid())); - Vec::new() - }) - .expect_err("forward parent"); - let (child, parent) = named.expect("declared"); - assert_eq!( - err, - BuildError::ForwardParent { - node: child, - parent - } - ); - } - - /// A handle belongs to the job that issued it. Carrying one into a second - /// job — an inner closure capturing an outer handle, say — must not - /// silently address whichever node happens to sit at the same position. - /// - /// This is what a random [`NodeGuid`] buys: with a plain per-builder - /// counter both jobs' first handles would be equal, and the edge below - /// would resolve, wrongly, to the second job's own node. - #[test] - fn a_handle_from_another_job_is_not_silently_resolved() { - let mut g = sched(); - let mut foreign = None; - g.insert_job(None, |b| { - foreign = Some(b.node("first job").guid()); - Vec::new() - }) - .expect("first job inserts"); - let foreign = foreign.expect("declared"); - - let err = g - .insert_job(None, |b| { - let _ = b.node("second job").after_ok(foreign); - Vec::new() - }) - .expect_err("foreign handle"); - assert!( - matches!(err, BuildError::ForwardEdge { dep, .. } if dep == foreign), - "{err:?}" - ); - } - - /// The grouping rule is enforced **before the first insert**, so a job that - /// breaks it leaves the graph untouched rather than half-built. - /// - /// This used to assert `BuildError::Graph(_)` — i.e. that the graph's own - /// rejection surfaced through the builder. It does not reach the graph any - /// more: the same rule is now decided from the job's own declarations, and - /// the stronger property (nothing was inserted) is what this pins. - #[test] - fn an_out_of_group_dep_is_refused_before_anything_is_inserted() { - let mut g = sched(); - let err = g - .insert_job(None, |b| { - let root = b.node("root"); - // A child may not depend on its own parent: the parent gate - // already orders them, and the edge would deadlock. - let _ = b.node("child").part_of(root).after_ok(root); - Vec::new() - }) - .expect_err("out-of-group dep"); - assert!(matches!(err, BuildError::DepOutsideGroup { .. }), "{err:?}"); - assert_eq!(g.graph().nodes().count(), 0, "nothing may have landed"); - } - - /// An edge no outcome can satisfy is refused at declaration. The graph only - /// catches this when validating a deserialized graph, so without the - /// pre-pass such a node would insert cleanly and then never become - /// runnable — a silent hang rather than an error. - #[test] - fn an_edge_no_outcome_can_satisfy_is_refused() { - let mut g = sched(); - let err = g - .insert_job(None, |b| { - let first = b.node("a"); - let _ = b.node("b").on_outcome(first, &[]); - Vec::new() - }) - .expect_err("unsatisfiable edge"); - assert!( - matches!(err, BuildError::UnsatisfiableDep { .. }), - "{err:?}" - ); - assert_eq!(g.graph().nodes().count(), 0, "nothing may have landed"); - } - - /// A job's own roots hang under the group's attachment point, while a node - /// that named a parent inside the job keeps it. This is what lets a - /// template be written without knowing the container it will live under. - #[test] - fn root_parent_adopts_only_the_jobs_own_roots() { - let mut g = sched(); - let container = g.append("container", Vec::new(), None).expect("container"); - - let ids = g - .insert_job(Some(container), |b| { - let root = b.node("root"); - let child = b.node("child").part_of(root); - vec![root.guid(), child.guid()] - }) - .expect("insert"); - let [root, child] = ids[..] else { - panic!("two ids back") - }; - assert_eq!(g.graph().node(root).expect("root").parent, Some(container)); - assert_eq!(g.graph().node(child).expect("child").parent, Some(root)); - } - - /// A handle that names no node in *this* job is refused rather than - /// silently dropped from the returned ids — the return is positional, so a - /// short vector would misalign every id after it. - #[test] - fn asking_for_a_foreign_handle_is_an_error() { - let mut g = sched(); - let mut foreign = None; - g.insert_job(None, |b| { - foreign = Some(b.node("first job").guid()); - Vec::new() - }) - .expect("first job inserts"); - let foreign = foreign.expect("declared"); - - let err = g - .insert_job(None, |b| { - let _ = b.node("second job"); - vec![foreign] - }) - .expect_err("foreign handle asked for"); - assert_eq!(err, BuildError::UnknownNode { node: foreign }); - } - - #[test] - fn an_empty_builder_inserts_nothing() { - let mut g = sched(); - let ids = g.insert_job(None, |_| Vec::new()).expect("insert"); - assert!(ids.is_empty()); - assert_eq!(g.graph().nodes().count(), 0); - } -} diff --git a/hive-jobq/src/lib.rs b/hive-jobq/src/lib.rs index e6f61e54..d89a1763 100644 --- a/hive-jobq/src/lib.rs +++ b/hive-jobq/src/lib.rs @@ -26,12 +26,9 @@ //! re-entrant borrow, one branch at a time). Single-threaded — the scheduler //! owns the resource table and mutates it directly. See [`scheduler`]. -pub mod builder; pub mod resources; pub mod scheduler; -pub use builder::{BuildError, JobBuilder, NodeGuid, NodeRef}; - use chrono::{DateTime, Utc}; /// Opaque, stable, monotonic node identifier. diff --git a/hive-jobq/src/scheduler.rs b/hive-jobq/src/scheduler.rs index 2727f0fa..e61dd64d 100644 --- a/hive-jobq/src/scheduler.rs +++ b/hive-jobq/src/scheduler.rs @@ -31,7 +31,6 @@ use std::collections::HashMap; use std::hash::Hash; -use crate::builder::{BuildError, JobBuilder, NodeGuid}; use crate::resources::ResourceTable; use crate::{Dep, Graph, GraphError, NodeId, State, TerminalState}; @@ -100,41 +99,6 @@ impl Scheduler { self.graph.insert(payload, deps, parent) } - /// Insert a whole job under `root_parent`, returning the id each handle's - /// node was minted as. - /// - /// `declare` receives a fresh [`JobBuilder`], names the job's nodes on it, - /// and returns the handles whose ids it wants back — they come back in - /// that order. The builder never leaves this call. That is the whole - /// insertion API — a caller cannot construct a builder, hold one, or - /// insert one itself, so there is no way to end up with a job-shaped value - /// being passed around as a spec. - /// - /// The one insertion entry point: every node goes through - /// [`Scheduler::append`], so a caller never has to reach past the scheduler - /// at the graph underneath. Call [`Scheduler::settle`] afterwards to start - /// whatever became runnable. - /// - /// **Atomic in the job's own shape.** A forward edge, a forward parent, or - /// a request for a handle this job never declared is rejected *before* the - /// first node is inserted, so a malformed job leaves the graph untouched - /// rather than half-built. - /// - /// # Errors - /// Propagates [`BuildError`] — a forward reference in the job's own - /// declarations, a handle from a different job, or a graph rejection. - pub fn insert_job( - &mut self, - root_parent: Option, - declare: impl FnOnce(&JobBuilder) -> Vec, - ) -> Result, BuildError> { - let job = JobBuilder::new(); - let wanted = declare(&job); - job.insert_with(root_parent, &wanted, |payload, deps, parent| { - self.append(payload, deps, parent) - }) - } - /// Claim every currently-runnable pending node and start it: node-deps /// satisfied and all resource-deps acquired atomically (all-or-nothing). /// Each claimed node is marked `Running`, its acquired units recorded, and