diff --git a/hive-jobq/src/guard.rs b/hive-jobq/src/guard.rs new file mode 100644 index 00000000..fb150625 --- /dev/null +++ b/hive-jobq/src/guard.rs @@ -0,0 +1,130 @@ +//! RAII guard objects over [`ResourceTable`] — owning resource grants. +//! +//! A running node acquires its resources through [`SharedResources::acquire`], +//! which hands back a [`ResourceGuard`] owning those units. Dropping the guard +//! releases exactly what it acquired, so a node's resources are freed when its +//! grant goes out of scope — there is no explicit release call to forget. +//! +//! Re-entrancy (a sub-node reusing a resource its ancestor group already holds) +//! is not expressed here: the scheduler tracks it with a single borrow slot per +//! `(holder, resource)` and never re-acquires, so these guards are always owning. +//! +//! Single-owner by design: the scheduler drives one settle loop, so the shared +//! table is `Rc>` (single-threaded interior mutability), not +//! `Arc>` — there is no cross-thread contention to guard against. + +use std::cell::RefCell; +use std::rc::Rc; + +use std::hash::Hash; + +use crate::resources::ResourceTable; + +/// A [`ResourceTable`] shared between the scheduler and the live guards that +/// release back into it on drop. Cheap to clone — an `Rc` refcount bump. +#[derive(Debug, Clone)] +pub struct SharedResources(Rc>>); + +impl Default for SharedResources { + fn default() -> Self { + Self::new(ResourceTable::new()) + } +} + +impl SharedResources { + /// Wrap an existing table so guards can release into it. + #[must_use] + pub fn new(table: ResourceTable) -> Self { + Self(Rc::new(RefCell::new(table))) + } + + /// Atomically acquire every requested `(name, count)` or none of them. + /// + /// Returns an owning [`ResourceGuard`] (releases on drop) when the whole + /// request fits in what is available right now; returns `None` and leaves + /// the table completely untouched otherwise. Duplicate names are summed and + /// an over-capacity request can never succeed — same all-or-nothing + /// semantics as [`ResourceTable::try_acquire_all`]. + #[must_use] + pub fn acquire(&self, reqs: Vec<(R, u32)>) -> Option> { + if self.0.borrow_mut().try_acquire_all(&reqs) { + Some(ResourceGuard { + table: self.clone(), + reqs, + }) + } else { + None + } + } + + /// Observe the underlying table — test-only (the scheduler's tests assert + /// on resource availability). Gated `#[cfg(test)]` so it is compiled out of + /// the shipped crate: no consumer can reach the raw table through it. + #[cfg(test)] + pub fn with(&self, f: impl FnOnce(&ResourceTable) -> T) -> T { + f(&self.0.borrow()) + } +} + +/// An RAII grant of resources: dropping it releases exactly the units it +/// acquired back into the shared table. +#[derive(Debug)] +pub struct ResourceGuard { + table: SharedResources, + reqs: Vec<(R, u32)>, +} + +impl ResourceGuard { + /// The `(name, count)` units this guard releases on drop. + #[must_use] + pub fn held(&self) -> &[(R, u32)] { + &self.reqs + } +} + +impl Drop for ResourceGuard { + fn drop(&mut self) { + self.table.0.borrow_mut().release_all(&self.reqs); + } +} + +#[cfg(test)] +mod tests { + use super::*; + + fn res(name: &str) -> String { + name.to_owned() + } + + fn shared_with(slots: u32) -> SharedResources { + let mut t = ResourceTable::new(); + t.set_capacity(res("build-slot"), slots); + SharedResources::new(t) + } + + #[test] + fn owning_guard_releases_on_drop() { + let sr = shared_with(2); + let slot = res("build-slot"); + { + let g = sr.acquire(vec![(slot.clone(), 2)]).expect("fits"); + assert_eq!(g.held(), &[(slot.clone(), 2)]); + // Both units held → any further acquire fails. + assert!(sr.acquire(vec![(slot.clone(), 1)]).is_none()); + } // guard dropped here → its units are released + // Full capacity is available again. + assert!(sr.acquire(vec![(slot.clone(), 2)]).is_some()); + } + + #[test] + fn acquire_returns_none_and_leaves_table_untouched_when_it_does_not_fit() { + let sr = shared_with(1); + let slot = res("build-slot"); + let held = sr.acquire(vec![(slot.clone(), 1)]).expect("first fits"); + assert!(sr.acquire(vec![(slot.clone(), 1)]).is_none()); + // The failed acquire took nothing extra: dropping the one real grant + // frees exactly one unit, so a single-unit acquire then fits. + drop(held); + assert!(sr.acquire(vec![(slot.clone(), 1)]).is_some()); + } +} diff --git a/hive-jobq/src/lib.rs b/hive-jobq/src/lib.rs index ee669f1e..29970de4 100644 --- a/hive-jobq/src/lib.rs +++ b/hive-jobq/src/lib.rs @@ -7,11 +7,11 @@ //! inserts a self-contained **node group** and returns its id; the scheduler //! runs a continuous loop, starting every node whose [`Dep`]s are satisfied: //! -//! - **Resource** deps are named counting semaphores ([`ResourceName`]): -//! `build-slot` (capacity N), `agent/` (capacity 1), or any name -//! (capacity 1, created on use). A node acquires *all* its resource deps -//! atomically at start (all-or-nothing) — no hold-and-wait, so no deadlock -//! and no cycle detection needed. +//! - **Resource** deps are named counting semaphores over a caller-chosen +//! type `R` (a `String` or an enum): `build-slot` (cap N), `agent/` +//! (cap 1), or any name (cap 1, created on use). A node acquires *all* its +//! resource deps atomically at start (all-or-nothing) — no hold-and-wait, +//! so no deadlock and no cycle detection needed. //! - **Node** deps wait on a node/group per [`DepWhen`]: `AfterOk` needs //! success (a failed dep cancels the dependent), `AfterAny` only terminal. //! @@ -26,10 +26,12 @@ //! node kind. Resources are held by the acquiring node and released on //! completion via guard objects, recursive within a group. //! -//! The guards and the scheduler loop are follow-ups; the named-counter resource -//! machinery lives in [`resources`], the base this data model builds on. +//! The [`scheduler`] settle loop drives execution; the resource machinery +//! lives in [`resources`] and the RAII lock guards over it in `guard`. +pub(crate) mod guard; pub mod resources; +pub mod scheduler; /// Opaque, stable, monotonic node identifier. /// @@ -47,16 +49,6 @@ pub mod resources; )] pub struct NodeId(pub(crate) u64); -/// A named counting semaphore. -/// -/// Examples: `build-slot` (capacity configured to the number of build slots), -/// `agent/` (capacity 1 — the per-agent lifecycle lock), or any other -/// name, which is assumed to have capacity 1 and is created on first use. -#[derive( - Debug, Clone, PartialEq, Eq, PartialOrd, Ord, Hash, serde::Serialize, serde::Deserialize, -)] -pub struct ResourceName(pub String); - /// When a [`Dep::Node`] edge is satisfied — the strong/weak distinction the /// current queue carries as `DepWhen`, load-bearing for failure safety. #[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Serialize, serde::Deserialize)] @@ -87,7 +79,7 @@ impl DepWhen { /// edge it names is satisfied (per its [`DepWhen`]) *and* every [`Dep::Resource`] /// it names can be acquired (all of them, atomically). #[derive(Debug, Clone, PartialEq, Eq, serde::Serialize, serde::Deserialize)] -pub enum Dep { +pub enum Dep { /// Depend on another node (or a group, by its group node's id). Whether a /// *failed* dependency satisfies the edge is decided by `when`: `AfterOk` /// requires success (and cancels this node if the dep fails), `AfterAny` @@ -103,7 +95,7 @@ pub enum Dep { /// released when the node completes. Resource { /// The resource to acquire. - name: ResourceName, + name: R, /// How many units to hold (usually 1). count: u32, }, @@ -141,7 +133,7 @@ impl State { /// payload; the caller supplies `N` (its own node kind) and a runner to execute /// a claimed node. #[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] -pub struct Node { +pub struct Node { /// Stable identity, assigned on insert. pub id: NodeId, /// The group this node belongs to, if any. `None` for a top-level group @@ -150,7 +142,7 @@ pub struct Node { /// Caller-defined payload (the node's kind / work description). pub payload: N, /// What must hold before this node runs (other nodes + resources). - pub deps: Vec, + pub deps: Vec>, /// Lifecycle state. pub state: State, } @@ -186,9 +178,15 @@ pub enum GraphError { /// walks this graph filling open slots. Completed groups are retained (no /// pruning in v1). #[derive(Debug, serde::Serialize, serde::Deserialize)] -#[serde(try_from = "GraphData")] -pub struct Graph { - nodes: Vec>, +#[serde( + try_from = "GraphData", + bound( + serialize = "N: serde::Serialize, R: serde::Serialize", + deserialize = "N: serde::Deserialize<'de>, R: serde::Deserialize<'de>" + ) +)] +pub struct Graph { + nodes: Vec>, next_id: u64, } @@ -196,15 +194,16 @@ pub struct Graph { // below — which runs [`Graph::validate`], so a loaded graph can never carry a // dangling id reference (Serialize does not validate; Deserialize always does). #[derive(serde::Deserialize)] -struct GraphData { - nodes: Vec>, +#[serde(bound(deserialize = "N: serde::Deserialize<'de>, R: serde::Deserialize<'de>"))] +struct GraphData { + nodes: Vec>, next_id: u64, } -impl TryFrom> for Graph { +impl TryFrom> for Graph { type Error = GraphError; - fn try_from(data: GraphData) -> Result { + fn try_from(data: GraphData) -> Result { let graph = Graph { nodes: data.nodes, next_id: data.next_id, @@ -216,13 +215,13 @@ impl TryFrom> for Graph { // A `derive(Default)` would wrongly require `N: Default` (an empty graph holds // no payload); an empty `Vec>` needs no such bound, so impl it directly. -impl Default for Graph { +impl Default for Graph { fn default() -> Self { Self::new() } } -impl Graph { +impl Graph { /// An empty graph. #[must_use] pub fn new() -> Self { @@ -253,7 +252,7 @@ impl Graph { pub fn insert( &mut self, payload: N, - deps: Vec, + deps: Vec>, parent: Option, ) -> Result { if let Some(parent_id) = parent @@ -281,15 +280,33 @@ impl Graph { /// Borrow a node by id. #[must_use] - pub fn node(&self, id: NodeId) -> Option<&Node> { + pub fn node(&self, id: NodeId) -> Option<&Node> { self.nodes.iter().find(|n| n.id == id) } /// The direct children of a group node (nodes whose `parent` is `id`). - pub fn children(&self, id: NodeId) -> impl Iterator> { + pub fn children(&self, id: NodeId) -> impl Iterator> { self.nodes.iter().filter(move |n| n.parent == Some(id)) } + /// Every node in the graph, in insertion order. The scheduler iterates + /// this to find runnable pending nodes. + pub fn nodes(&self) -> impl Iterator> { + self.nodes.iter() + } + + /// Set a node's lifecycle state, returning `false` for an unknown id. The + /// scheduler drives every state transition — nothing else mutates state, + /// which is what keeps the resource guards + terminality in sync. + pub(crate) fn set_state(&mut self, id: NodeId, state: State) -> bool { + if let Some(node) = self.nodes.iter_mut().find(|n| n.id == id) { + node.state = state; + true + } else { + false + } + } + /// A group is terminal once the group node itself is terminal *and* every /// node inside it (recursively) is terminal. The node's own state matters: /// a group node still `Pending`/`Running` is not terminal even with no @@ -345,7 +362,7 @@ mod tests { #[test] fn insert_mints_stable_monotonic_ids() { - let mut g: Graph<&str> = Graph::new(); + let mut g: Graph<&str, String> = Graph::new(); let a = g.insert("sweep", vec![], None).unwrap(); let b = g .insert( @@ -364,24 +381,19 @@ mod tests { assert_eq!(g.node(a).unwrap().parent, None); } - fn set_state(g: &mut Graph, id: NodeId, state: State) { - let idx = g.nodes.iter().position(|n| n.id == id).unwrap(); - g.nodes[idx].state = state; - } - #[test] fn group_terminal_requires_the_group_node_and_all_children_terminal() { - let mut g: Graph<&str> = Graph::new(); + let mut g: Graph<&str, String> = Graph::new(); let group = g.insert("group", vec![], None).unwrap(); let child = g.insert("child", vec![], Some(group)).unwrap(); // Both pending → not terminal. assert!(!g.group_terminal(group)); // Child done, but the group node itself is still pending → NOT terminal: // the group node's own state is load-bearing, not just its children. - set_state(&mut g, child, State::Done); + g.set_state(child, State::Done); assert!(!g.group_terminal(group)); // Group node terminal too → the whole group is terminal. - set_state(&mut g, group, State::Done); + g.set_state(group, State::Done); assert!(g.group_terminal(group)); } @@ -389,12 +401,12 @@ mod tests { fn empty_running_group_is_not_terminal() { // A running node with no children yet may still append some, so it must // not read as terminal just because its child set is currently empty. - let mut g: Graph<&str> = Graph::new(); + let mut g: Graph<&str, String> = Graph::new(); let group = g.insert("group", vec![], None).unwrap(); - set_state(&mut g, group, State::Running); + g.set_state(group, State::Running); assert!(!g.group_terminal(group)); // Once it finishes (having grown no children), it is terminal. - set_state(&mut g, group, State::Done); + g.set_state(group, State::Done); assert!(g.group_terminal(group)); } @@ -424,7 +436,7 @@ mod tests { #[test] fn insert_rejects_unknown_parent() { - let mut g: Graph<&str> = Graph::new(); + let mut g: Graph<&str, String> = Graph::new(); let bogus = NodeId(7); assert_eq!( g.insert("x", vec![], Some(bogus)).unwrap_err(), @@ -434,7 +446,7 @@ mod tests { #[test] fn insert_rejects_unknown_dep() { - let mut g: Graph<&str> = Graph::new(); + let mut g: Graph<&str, String> = Graph::new(); let bogus = NodeId(42); let deps = vec![Dep::Node { id: bogus, @@ -448,7 +460,7 @@ mod tests { #[test] fn valid_graph_round_trips_through_serde() { - let mut g: Graph = Graph::new(); + let mut g: Graph = Graph::new(); let a = g.insert("a".to_owned(), vec![], None).unwrap(); g.insert( "b".to_owned(), @@ -460,7 +472,7 @@ mod tests { ) .unwrap(); let json = serde_json::to_string(&g).unwrap(); - let back: Graph = serde_json::from_str(&json).unwrap(); + let back: Graph = serde_json::from_str(&json).unwrap(); assert!(back.validate().is_ok()); assert_eq!(back.node(a).unwrap().payload, "a"); } @@ -469,7 +481,7 @@ mod tests { fn deserialize_rejects_a_dangling_dependency() { // Build a graph whose only node depends on a non-existent id, serialize // it (Serialize does not validate), and confirm deserialize rejects it. - let bad = Graph:: { + let bad = Graph:: { nodes: vec![Node { id: NodeId(0), parent: None, @@ -483,13 +495,13 @@ mod tests { next_id: 1, }; let json = serde_json::to_string(&bad).unwrap(); - let err = serde_json::from_str::>(&json).unwrap_err(); + let err = serde_json::from_str::>(&json).unwrap_err(); assert!(err.to_string().contains("unknown node")); } #[test] fn validate_rejects_next_id_that_would_remint() { - let bad = Graph::<&str> { + let bad = Graph::<&str, String> { nodes: vec![Node { id: NodeId(5), parent: None, diff --git a/hive-jobq/src/resources.rs b/hive-jobq/src/resources.rs index bf215fe7..52198366 100644 --- a/hive-jobq/src/resources.rs +++ b/hive-jobq/src/resources.rs @@ -20,8 +20,8 @@ //! //! [`Dep::Resource`]: crate::Dep::Resource -use crate::ResourceName; use std::collections::HashMap; +use std::hash::Hash; /// A set of named counting semaphores. /// @@ -30,22 +30,22 @@ use std::collections::HashMap; /// units atomically via [`ResourceTable::try_acquire_all`] / /// [`ResourceTable::release_all`]. #[derive(Debug, Clone)] -pub struct ResourceTable { - /// Configured capacities, keyed by name. Missing ⇒ `default_capacity`. - capacities: HashMap, - /// Units currently held, keyed by name. Missing ⇒ 0. - held: HashMap, - /// Capacity assumed for a name with no configured entry. +pub struct ResourceTable { + /// Configured capacities, keyed by resource. Missing ⇒ `default_capacity`. + capacities: HashMap, + /// Units currently held, keyed by resource. Missing ⇒ 0. + held: HashMap, + /// Capacity assumed for a resource with no configured entry. default_capacity: u32, } -impl Default for ResourceTable { +impl Default for ResourceTable { fn default() -> Self { Self::new() } } -impl ResourceTable { +impl ResourceTable { /// An empty table whose unconfigured names default to capacity 1. #[must_use] pub fn new() -> Self { @@ -61,14 +61,14 @@ impl ResourceTable { /// Overwrites any previous capacity for that name. Lowering capacity below /// the currently-held count is allowed — the table simply reports zero /// available until enough is released; it never rejects a config change. - pub fn set_capacity(&mut self, name: ResourceName, capacity: u32) { + pub fn set_capacity(&mut self, name: R, capacity: u32) { self.capacities.insert(name, capacity); } /// The capacity of a name — its configured value, or the default (1) if it /// was never configured. #[must_use] - pub fn capacity(&self, name: &ResourceName) -> u32 { + pub fn capacity(&self, name: &R) -> u32 { self.capacities .get(name) .copied() @@ -77,13 +77,13 @@ impl ResourceTable { /// Units of a name currently held (0 if none). #[must_use] - pub fn held(&self, name: &ResourceName) -> u32 { + pub fn held(&self, name: &R) -> u32 { self.held.get(name).copied().unwrap_or(0) } /// Units of a name available to acquire right now (`capacity - held`). #[must_use] - pub fn available(&self, name: &ResourceName) -> u32 { + pub fn available(&self, name: &R) -> u32 { self.capacity(name).saturating_sub(self.held(name)) } @@ -95,7 +95,7 @@ impl ResourceTable { /// `false` and leaves the table completely untouched. A request for more /// units than a name's capacity can therefore never succeed — the scheduler /// should reject such a node at insert time so it does not wait forever. - pub fn try_acquire_all(&mut self, reqs: &[(ResourceName, u32)]) -> bool { + pub(crate) fn try_acquire_all(&mut self, reqs: &[(R, u32)]) -> bool { let wanted = aggregate(reqs); // All-or-nothing: bail before mutating if any request cannot be met. for (name, &count) in &wanted { @@ -113,7 +113,7 @@ impl ResourceTable { /// /// Duplicate names are summed. Releasing more than is held saturates at zero /// rather than underflowing, so a double release is harmless. - pub fn release_all(&mut self, reqs: &[(ResourceName, u32)]) { + pub(crate) fn release_all(&mut self, reqs: &[(R, u32)]) { for (name, count) in aggregate(reqs) { if let Some(h) = self.held.get_mut(name) { *h = h.saturating_sub(count); @@ -123,8 +123,8 @@ impl ResourceTable { } /// Sum a request list into per-name totals so duplicate names are one entry. -fn aggregate(reqs: &[(ResourceName, u32)]) -> HashMap<&ResourceName, u32> { - let mut wanted: HashMap<&ResourceName, u32> = HashMap::new(); +fn aggregate(reqs: &[(R, u32)]) -> HashMap<&R, u32> { + let mut wanted: HashMap<&R, u32> = HashMap::new(); for (name, count) in reqs { *wanted.entry(name).or_insert(0) += *count; } @@ -135,8 +135,8 @@ fn aggregate(reqs: &[(ResourceName, u32)]) -> HashMap<&ResourceName, u32> { mod tests { use super::*; - fn res(name: &str) -> ResourceName { - ResourceName(name.to_owned()) + fn res(name: &str) -> String { + name.to_owned() } #[test] diff --git a/hive-jobq/src/scheduler.rs b/hive-jobq/src/scheduler.rs new file mode 100644 index 00000000..b6fc8883 --- /dev/null +++ b/hive-jobq/src/scheduler.rs @@ -0,0 +1,425 @@ +//! The settle loop — drives a [`Graph`] to completion over the resource pool. +//! +//! [`Scheduler::settle`] claims every currently-runnable pending node (its +//! [`Dep::Node`] edges satisfied *and* all its [`Dep::Resource`] units acquired +//! atomically), marks it `Running`, holds its resource guards, and returns the +//! newly-started ids for the caller's runner to execute. The runner reports each +//! node's result back with [`Scheduler::complete`]; a running node may grow its +//! own sub-group first via [`Scheduler::append`]. Concurrency is emergent from +//! resource capacity — there is no separate active-node cap. +//! +//! A resource is held for the acquiring node's *entire subtree* lifetime: the +//! owned guard is released only when that node and every descendant is terminal +//! (`group_terminal`), not when the node's own work finishes. Single-owner and +//! synchronous — the caller drives `settle` / `complete`; no async or locking +//! lives here (that's the runner's job, one layer up). +//! +//! Recursive-lock re-entrancy (a sub-node reusing an ancestor group's lock) and +//! the eager `AfterOk` failure cascade are layered on top of this owned core. + +use std::collections::HashMap; +use std::hash::Hash; + +use crate::guard::{ResourceGuard, SharedResources}; +use crate::resources::ResourceTable; +use crate::{Dep, DepWhen, Graph, GraphError, NodeId, State}; + +/// The result of a node's own execution, reported to [`Scheduler::complete`]. +/// +/// `Cancelled` is not an outcome a runner reports — it is scheduler-driven (an +/// `AfterOk` dependency failed), so a runner only ever says `Done` or `Failed`. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum Outcome { + /// The node's work succeeded. + Done, + /// The node's work failed. + Failed, +} + +/// Drives a [`Graph`] over a shared resource pool: claim runnable nodes, hold +/// their resources for the subtree's lifetime, release on subtree-terminal. +pub struct Scheduler { + graph: Graph, + resources: SharedResources, + /// Owned resource guards, keyed by the node that acquired them. Dropped + /// (releasing the units) when that node's whole subtree is terminal. + owned: HashMap>>, + /// The single re-entrancy slot per `(ancestor-holder, resource)`: the id of + /// the descendant currently *borrowing* that ancestor's lock. Present ⇒ the + /// slot is taken, so no other descendant may re-enter the same lock until + /// the borrower's subtree is terminal — "only one node at a time within". + borrow_slots: HashMap<(NodeId, R), NodeId>, +} + +impl Scheduler { + /// A scheduler over `graph` with `resources` as the capacity pool. + #[must_use] + pub fn new(graph: Graph, resources: ResourceTable) -> Self { + Self { + graph, + resources: SharedResources::new(resources), + owned: HashMap::new(), + borrow_slots: HashMap::new(), + } + } + + /// The graph, for inspection (state, hierarchy, UI rendering). + #[must_use] + pub fn graph(&self) -> &Graph { + &self.graph + } + + /// Append a node — e.g. a running node growing its own sub-group. Delegates + /// to [`Graph::insert`]; call [`Scheduler::settle`] afterwards to start it + /// once it is runnable. + /// + /// # Errors + /// Propagates [`GraphError`] for a dangling parent or dependency id. + pub fn append( + &mut self, + payload: N, + deps: Vec>, + parent: Option, + ) -> Result { + self.graph.insert(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 owned guards held, and its id + /// returned for the runner to execute. A single pass suffices — a node + /// started here is `Running`, not terminal, so it cannot satisfy another + /// node's dependency in the same pass; it only consumes resources. + #[must_use] + pub fn settle(&mut self) -> Vec { + let pending: Vec = self + .graph + .nodes() + .filter(|n| n.state == State::Pending) + .map(|n| n.id) + .collect(); + let mut started = Vec::new(); + for id in pending { + if self.node_deps_satisfied(id) && self.try_start(id) { + started.push(id); + } + } + started + } + + /// Try to start node `id`: classify each resource dep as *owned* (no + /// ancestor holds it → acquire real units) or *borrowed* (an ancestor group + /// already holds it → re-enter, gated by the one re-entrancy slot), then + /// take everything atomically or nothing. Returns whether it started. + fn try_start(&mut self, id: NodeId) -> bool { + let mut owned_reqs = Vec::new(); + let mut borrows = Vec::new(); + for (name, count) in self.resource_reqs(id) { + if let Some(ancestor) = self.ancestor_owning(id, &name) { + // Re-entrant reuse: allowed only if the slot is free. + if self.borrow_slots.contains_key(&(ancestor, name.clone())) { + return false; + } + borrows.push((ancestor, name)); + } else { + owned_reqs.push((name, count)); + } + } + // Owned units are all-or-nothing; borrow slots were all confirmed free + // above, so this is the only fallible step. Nothing mutated until here. + let Some(guard) = self.resources.acquire(owned_reqs) else { + return false; + }; + self.owned.entry(id).or_default().push(guard); + for slot in borrows { + self.borrow_slots.insert(slot, id); + } + self.graph.set_state(id, State::Running); + true + } + + /// The nearest ancestor of `id` that *owns* (holds real units of) `name`, + /// or `None` if no ancestor holds it (⇒ `id` must own-acquire it itself). + fn ancestor_owning(&self, id: NodeId, name: &R) -> Option { + let mut cursor = self.graph.node(id)?.parent; + while let Some(ancestor) = cursor { + if self.node_owns(ancestor, name) { + return Some(ancestor); + } + cursor = self.graph.node(ancestor)?.parent; + } + None + } + + /// Whether node `holder` holds an owned guard covering resource `name`. + fn node_owns(&self, holder: NodeId, name: &R) -> bool { + self.owned.get(&holder).is_some_and(|guards| { + guards + .iter() + .any(|g| g.held().iter().any(|(n, _)| n == name)) + }) + } + + /// Report a running node's own execution result. Sets its state, then + /// releases the owned guards of every node whose whole subtree has become + /// terminal — a parent keeps its lock until its last descendant finishes. + /// Call [`Scheduler::settle`] again afterwards to start newly-unblocked work. + pub fn complete(&mut self, id: NodeId, outcome: Outcome) { + let state = match outcome { + Outcome::Done => State::Done, + Outcome::Failed => State::Failed, + }; + self.graph.set_state(id, state); + if outcome == Outcome::Failed { + self.cascade_cancel(id); + } + self.release_settled_subtrees(); + } + + /// Eagerly cancel the transitive `AfterOk` dependents of a just-failed node: + /// they can never run (a strong dependency failed), so mark them `Cancelled` + /// now — before they could claim resources. A dependent is always still + /// `Pending` here (a `Running` node's `AfterOk` deps were `Done` when it + /// started, and `Done` is terminal), so no resources need releasing. + fn cascade_cancel(&mut self, failed: NodeId) { + let mut stack = vec![failed]; + while let Some(dep) = stack.pop() { + let dependents: Vec = self + .graph + .nodes() + .filter(|n| { + n.state == State::Pending + && n.deps.iter().any( + |d| matches!(d, Dep::Node { id, when: DepWhen::AfterOk } if *id == dep), + ) + }) + .map(|n| n.id) + .collect(); + for d in dependents { + self.graph.set_state(d, State::Cancelled); + stack.push(d); + } + } + } + + /// Release everything whose subtree has become terminal: drop the owned + /// guards of any holder (→ frees its units) and free any re-entrancy slot + /// held by a borrower — both are held for the whole subtree lifetime. + fn release_settled_subtrees(&mut self) { + let holders: Vec = self.owned.keys().copied().collect(); + for holder in holders { + if self.graph.group_terminal(holder) { + self.owned.remove(&holder); // drops guards → releases the units + } + } + self.borrow_slots + .retain(|_, borrower| !self.graph.group_terminal(*borrower)); + } + + /// Whether every [`Dep::Node`] edge of `id` is satisfied. `Dep::Resource` + /// edges are handled by the atomic acquire in [`Scheduler::settle`], not here. + fn node_deps_satisfied(&self, id: NodeId) -> bool { + let Some(node) = self.graph.node(id) else { + return false; + }; + node.deps.iter().all(|dep| match dep { + Dep::Resource { .. } => true, + Dep::Node { id, when } => self.dep_node_satisfied(*id, *when), + }) + } + + /// Whether a node/group dependency `id` satisfies edge kind `when`. A group + /// is depended on as a whole: `AfterAny` needs its subtree terminal (any + /// outcome), `AfterOk` needs its whole subtree to have succeeded. + fn dep_node_satisfied(&self, id: NodeId, when: DepWhen) -> bool { + match when { + DepWhen::AfterAny => self.graph.group_terminal(id), + DepWhen::AfterOk => self.subtree_all_done(id), + } + } + + /// Whether `id` and every descendant reached [`State::Done`] — the success + /// condition for an `AfterOk` edge onto a (possibly group) node. + fn subtree_all_done(&self, id: NodeId) -> bool { + let Some(node) = self.graph.node(id) else { + return false; + }; + node.state == State::Done && self.graph.children(id).all(|c| self.subtree_all_done(c.id)) + } + + /// The `(name, count)` resource units `id` must hold to run. + fn resource_reqs(&self, id: NodeId) -> Vec<(R, u32)> { + let Some(node) = self.graph.node(id) else { + return Vec::new(); + }; + node.deps + .iter() + .filter_map(|dep| match dep { + Dep::Resource { name, count } => Some((name.clone(), *count)), + Dep::Node { .. } => None, + }) + .collect() + } +} + +#[cfg(test)] +mod tests { + use super::*; + + fn res(name: &str) -> String { + name.to_owned() + } + + /// A graph + a resource table with `build-slot` set to `slots`. + fn scheduler_with_slots(slots: u32) -> Scheduler<&'static str, String> { + let mut table = ResourceTable::new(); + table.set_capacity(res("build-slot"), slots); + Scheduler::new(Graph::new(), table) + } + + fn slot_dep() -> Vec> { + vec![Dep::Resource { + name: res("build-slot"), + count: 1, + }] + } + + fn resource_dep(name: &str) -> Vec> { + vec![Dep::Resource { + name: res(name), + count: 1, + }] + } + + #[test] + fn resource_node_starts_then_releases_on_complete() { + let mut s = scheduler_with_slots(1); + let n = s.append("build", slot_dep(), None).expect("insert"); + // settle claims it (a slot is free) and marks it Running. + assert_eq!(s.settle(), vec![n]); + assert_eq!(s.graph().node(n).unwrap().state, State::Running); + // Slot is held. + assert!(s.resources.with(|t| t.available(&res("build-slot")) == 0)); + // Completing it releases the slot (subtree is just this node). + s.complete(n, Outcome::Done); + assert_eq!(s.graph().node(n).unwrap().state, State::Done); + assert!(s.resources.with(|t| t.available(&res("build-slot")) == 1)); + } + + #[test] + fn build_slot_cap_limits_concurrency_and_release_unblocks() { + let mut s = scheduler_with_slots(2); + let a = s.append("a", slot_dep(), None).expect("a"); + let b = s.append("b", slot_dep(), None).expect("b"); + let c = s.append("c", slot_dep(), None).expect("c"); + // cap 2 → a + b start, c blocks on the exhausted slot. + let started = s.settle(); + assert_eq!(started, vec![a, b]); + assert_eq!(s.graph().node(c).unwrap().state, State::Pending); + // a finishes → its slot frees → c can now start. + s.complete(a, Outcome::Done); + assert_eq!(s.settle(), vec![c]); + assert_eq!(s.graph().node(c).unwrap().state, State::Running); + } + + #[test] + fn parent_holds_resource_until_child_subtree_done() { + let mut s = scheduler_with_slots(1); + // Parent grabs the single build-slot and runs. + let parent = s.append("parent", slot_dep(), None).expect("parent"); + assert_eq!(s.settle(), vec![parent]); + // Parent grows a child (no resource dep of its own) and finishes its + // OWN work — but its subtree is not terminal, so it keeps the slot. + let child = s.append("child", vec![], Some(parent)).expect("child"); + s.complete(parent, Outcome::Done); + assert!( + s.resources.with(|t| t.available(&res("build-slot")) == 0), + "parent must keep its lock while a child is still pending/running" + ); + // The child starts and completes → now the whole subtree is terminal → + // the parent's slot is released exactly once. + assert_eq!(s.settle(), vec![child]); + s.complete(child, Outcome::Done); + assert!(s.resources.with(|t| t.available(&res("build-slot")) == 1)); + } + + #[test] + fn recursive_lock_serializes_re_entrant_descendants() { + // `agent/foo` is unconfigured → default capacity 1. + let mut s = Scheduler::new(Graph::new(), ResourceTable::new()); + let agent = res("agent/foo"); + // A group node owns agent/foo and runs. + let group = s + .append("group", resource_dep("agent/foo"), None) + .expect("group"); + assert_eq!(s.settle(), vec![group]); + assert!(s.resources.with(|t| t.available(&agent) == 0)); + // Two sub-nodes each need agent/foo → they re-enter the group's lock, + // but only ONE at a time (the single re-entrancy slot). + let c1 = s + .append("c1", resource_dep("agent/foo"), Some(group)) + .expect("c1"); + let c2 = s + .append("c2", resource_dep("agent/foo"), Some(group)) + .expect("c2"); + let started = s.settle(); + assert_eq!( + started, + vec![c1], + "only one descendant may borrow at a time" + ); + assert_eq!(s.graph().node(c2).unwrap().state, State::Pending); + // The lock was NOT re-acquired — still just the group's one unit held. + assert!(s.resources.with(|t| t.available(&agent) == 0)); + // c1 finishes → its borrow slot frees → c2 can now re-enter. + s.complete(c1, Outcome::Done); + assert_eq!(s.settle(), vec![c2]); + assert_eq!(s.graph().node(c2).unwrap().state, State::Running); + // Still no double-acquire; the group's single unit is the only hold. + assert!(s.resources.with(|t| t.available(&agent) == 0)); + } + + #[test] + fn failed_after_ok_dep_cancels_dependents_but_after_any_still_runs() { + let mut s: Scheduler<&str, String> = Scheduler::new(Graph::new(), ResourceTable::new()); + let root = s.append("root", vec![], None).expect("root"); + let strong1 = s + .append( + "strong1", + vec![Dep::Node { + id: root, + when: DepWhen::AfterOk, + }], + None, + ) + .expect("strong1"); + let strong2 = s + .append( + "strong2", + vec![Dep::Node { + id: strong1, + when: DepWhen::AfterOk, + }], + None, + ) + .expect("strong2"); + let weak = s + .append( + "weak", + vec![Dep::Node { + id: root, + when: DepWhen::AfterAny, + }], + None, + ) + .expect("weak"); + assert_eq!(s.settle(), vec![root]); + s.complete(root, Outcome::Failed); + // The AfterOk chain strong1→strong2 is eagerly cancelled (a strong dep + // failed)… + assert_eq!(s.graph().node(strong1).unwrap().state, State::Cancelled); + assert_eq!(s.graph().node(strong2).unwrap().state, State::Cancelled); + // …but the AfterAny dependent still runs — it converges regardless. + assert_eq!(s.settle(), vec![weak]); + } +}