jobq: complete is the completion, and failure grows nothing
Two review findings from the previous round, re-checked against the actual tree rather than against my notes. `Scheduler::complete` was still a public wrapper whose entire body was `self.finish(id, outcome)`. Its docstring argued the split was not a redirect because both completion forms shared `finish` -- but sharing a private helper is not a reason for two public names. `finish`'s body now lives in `complete`, and `complete_growing` calls it. Same sharing, one name, no redirect. Growth on a failed node is now dropped by `complete_growing` instead of by the host loop. Failure cancel-cascades to every pending child of the completing node, and grown work is inserted as its children, so anything appended here is Skipped by the next statement -- the insert is not wrong, it is provably pointless. That is a consequence of this crate's cascade rule, so this crate should be the one enforcing it; a host that has to remember it can forget it. Behaviour is unchanged: hive-c0re already dropped growth before calling, and now no longer has to. `Scheduler::new_job` is left alone but documented for what it is: the hole in `JobBuilder::new`'s pub(crate) wall, with no non-test caller since claim_next mints a builder per running node. Closing it is a venue question rather than a rename, so it stays for now.
This commit is contained in:
parent
ab53f6710d
commit
8459cc66bd
4 changed files with 53 additions and 52 deletions
|
|
@ -252,22 +252,6 @@ impl JobQueue {
|
||||||
self.lock().graph().root_of(node).map(NodeId::get)
|
self.lock().graph().root_of(node).map(NodeId::get)
|
||||||
}
|
}
|
||||||
|
|
||||||
/// A builder for a node to declare more work into while it runs.
|
|
||||||
///
|
|
||||||
/// Handed to [`exec::run_node`] and returned to the crate's completion.
|
|
||||||
/// Only `hive_jobq` can construct one, which is why this goes through the
|
|
||||||
/// scheduler rather than `Job::default()`.
|
|
||||||
///
|
|
||||||
/// The completion wrappers that used to live beside this (`complete_node`,
|
|
||||||
/// `complete_node_growing`) are **gone**: a node is completed inside the
|
|
||||||
/// future [`hive_jobq::scheduler::Scheduler::claim_next`] hands back, so
|
|
||||||
/// this layer has nothing left to wrap. The tests keep their own extension
|
|
||||||
/// trait for driving completions by hand.
|
|
||||||
#[must_use]
|
|
||||||
pub fn new_job(&self) -> Job {
|
|
||||||
self.lock().new_job()
|
|
||||||
}
|
|
||||||
|
|
||||||
/// Cancel a DAG that hasn't started yet: every work node is still `Pending`,
|
/// Cancel a DAG that hasn't started yet: every work node is still `Pending`,
|
||||||
/// so each is cancelled. `false` once any work node is running or terminal —
|
/// so each is cancelled. `false` once any work node is running or terminal —
|
||||||
/// an in-flight nix build isn't interruptible.
|
/// an in-flight nix build isn't interruptible.
|
||||||
|
|
|
||||||
|
|
@ -101,17 +101,14 @@ pub async fn run_worker(coord: Arc<Coordinator>) {
|
||||||
"job_queue: node failed"
|
"job_queue: node failed"
|
||||||
),
|
),
|
||||||
}
|
}
|
||||||
let outcome = super::outcome_of(result.map_err(|e| format!("{e:#}")));
|
// Growth on a failed node is dropped by `complete_growing`,
|
||||||
// Growth is dropped on failure: a node that declared
|
// not here: failure cancel-cascades inside jobq, so that
|
||||||
// follow-up work and *then* failed does not want it run —
|
// rule is the crate's to enforce and this loop does not get
|
||||||
// failure cancel-cascades, so inserting it would only add
|
// to forget it.
|
||||||
// nodes to immediately cancel.
|
(
|
||||||
let grown = if matches!(outcome, hive_jobq::scheduler::Outcome::Failed(_)) {
|
grown,
|
||||||
coord.job_queue.new_job()
|
super::outcome_of(result.map_err(|e| format!("{e:#}"))),
|
||||||
} else {
|
)
|
||||||
grown
|
|
||||||
};
|
|
||||||
(grown, outcome)
|
|
||||||
}
|
}
|
||||||
})
|
})
|
||||||
};
|
};
|
||||||
|
|
|
||||||
|
|
@ -122,9 +122,23 @@ impl ClaimReady for JobQueue {
|
||||||
trait CompleteNode {
|
trait CompleteNode {
|
||||||
fn complete_node(&self, node_id: NodeId, result: Result<(), String>);
|
fn complete_node(&self, node_id: NodeId, result: Result<(), String>);
|
||||||
fn complete_node_growing(&self, node_id: NodeId, result: Result<(), String>, grown: Job);
|
fn complete_node_growing(&self, node_id: NodeId, result: Result<(), String>, grown: Job);
|
||||||
|
fn new_job(&self) -> Job;
|
||||||
}
|
}
|
||||||
|
|
||||||
impl CompleteNode for JobQueue {
|
impl CompleteNode for JobQueue {
|
||||||
|
/// Mint a builder to declare growth into.
|
||||||
|
///
|
||||||
|
/// Test-only for the same reason as the rest of this trait: production
|
||||||
|
/// never mints one, because `claim_next` hands each running node its
|
||||||
|
/// builder and takes it back. That leaves `Scheduler::new_job` with no
|
||||||
|
/// non-test caller either — see the note on that fn.
|
||||||
|
fn new_job(&self) -> Job {
|
||||||
|
self.sched()
|
||||||
|
.lock()
|
||||||
|
.expect("job_queue mutex poisoned")
|
||||||
|
.new_job()
|
||||||
|
}
|
||||||
|
|
||||||
fn complete_node(&self, node_id: NodeId, result: Result<(), String>) {
|
fn complete_node(&self, node_id: NodeId, result: Result<(), String>) {
|
||||||
self.sched()
|
self.sched()
|
||||||
.lock()
|
.lock()
|
||||||
|
|
|
||||||
|
|
@ -324,7 +324,18 @@ impl<N, R: Clone + Eq + Hash> Scheduler<N, R> {
|
||||||
/// propagates up the parent chain. Call [`Scheduler::settle`] again afterwards
|
/// propagates up the parent chain. Call [`Scheduler::settle`] again afterwards
|
||||||
/// to start newly-unblocked work.
|
/// to start newly-unblocked work.
|
||||||
pub fn complete(&mut self, id: NodeId, outcome: Outcome) {
|
pub fn complete(&mut self, id: NodeId, outcome: Outcome) {
|
||||||
self.finish(id, outcome);
|
match outcome {
|
||||||
|
Outcome::Failed(error) => {
|
||||||
|
// Record the reason before the terminal transition so it's set
|
||||||
|
// by the time `set_state` stamps `finished_at`.
|
||||||
|
self.graph.set_error(id, error);
|
||||||
|
self.graph.set_state(id, State::Failed);
|
||||||
|
self.cascade_cancel(id);
|
||||||
|
}
|
||||||
|
Outcome::Done => self.settle_terminal(id),
|
||||||
|
}
|
||||||
|
self.roll_up_ancestors(id);
|
||||||
|
self.release_ready();
|
||||||
}
|
}
|
||||||
|
|
||||||
/// A fresh builder for a **running** node to declare more work into.
|
/// A fresh builder for a **running** node to declare more work into.
|
||||||
|
|
@ -335,9 +346,13 @@ impl<N, R: Clone + Eq + Hash> Scheduler<N, R> {
|
||||||
/// [`NodeId`]s only at insert), so it can be filled in freely and handed
|
/// [`NodeId`]s only at insert), so it can be filled in freely and handed
|
||||||
/// back to [`Scheduler::complete_growing`], which inserts it under the lock.
|
/// back to [`Scheduler::complete_growing`], which inserts it under the lock.
|
||||||
///
|
///
|
||||||
/// This is the only way to get one — [`JobBuilder::new`] is `pub(crate)` and
|
/// ⚠️ **This has no non-test caller left, and it is the hole in
|
||||||
/// there is no `Default` impl — so a caller can declare work but never
|
/// [`JobBuilder::new`]'s `pub(crate)` wall** — it hands out exactly the
|
||||||
/// insert it itself.
|
/// builder that fn withholds. [`Scheduler::claim_next`] mints one per
|
||||||
|
/// running node itself, so production never asks. Kept only so the host's
|
||||||
|
/// graph-growth tests can still declare work by hand; the fix is a venue
|
||||||
|
/// question (move those tests here vs. a closure-form completion), not a
|
||||||
|
/// rename.
|
||||||
#[must_use]
|
#[must_use]
|
||||||
pub fn new_job(&self) -> JobBuilder<N, R> {
|
pub fn new_job(&self) -> JobBuilder<N, R> {
|
||||||
JobBuilder::new()
|
JobBuilder::new()
|
||||||
|
|
@ -353,6 +368,13 @@ impl<N, R: Clone + Eq + Hash> Scheduler<N, R> {
|
||||||
/// outright, which is the overwhelmingly common case (most nodes grow no
|
/// outright, which is the overwhelmingly common case (most nodes grow no
|
||||||
/// work at all).
|
/// work at all).
|
||||||
///
|
///
|
||||||
|
/// **A failed node grows nothing**, whatever it declared. Failure
|
||||||
|
/// cancel-cascades to every pending child of `id`, so work inserted here
|
||||||
|
/// would be `Skipped` by the very next statement — the insert is not wrong,
|
||||||
|
/// it is provably pointless. This lives here rather than in the caller
|
||||||
|
/// because it is a consequence of *this crate's* cascade rule; a host that
|
||||||
|
/// had to remember it could forget it.
|
||||||
|
///
|
||||||
/// # Errors
|
/// # Errors
|
||||||
/// [`BuildError`] if `grown` is malformed — **and the node is still
|
/// [`BuildError`] if `grown` is malformed — **and the node is still
|
||||||
/// completed**. Its own work already happened; refusing to complete it
|
/// completed**. Its own work already happened; refusing to complete it
|
||||||
|
|
@ -371,7 +393,10 @@ impl<N, R: Clone + Eq + Hash> Scheduler<N, R> {
|
||||||
// dangling `parent` edge rather than being rejected. The host used to
|
// dangling `parent` edge rather than being rejected. The host used to
|
||||||
// carry this guard itself, as a lookup before a separate append call;
|
// carry this guard itself, as a lookup before a separate append call;
|
||||||
// it belongs here, where the graph is and where it cannot be skipped.
|
// it belongs here, where the graph is and where it cannot be skipped.
|
||||||
let grew = if grown.is_empty() || self.graph.node(id).is_none() {
|
let grew = if grown.is_empty()
|
||||||
|
|| matches!(outcome, Outcome::Failed(_))
|
||||||
|
|| self.graph.node(id).is_none()
|
||||||
|
{
|
||||||
Ok(())
|
Ok(())
|
||||||
} else {
|
} else {
|
||||||
let graph = &mut self.graph;
|
let graph = &mut self.graph;
|
||||||
|
|
@ -381,29 +406,10 @@ impl<N, R: Clone + Eq + Hash> Scheduler<N, R> {
|
||||||
})
|
})
|
||||||
.map(|_ids| ())
|
.map(|_ids| ())
|
||||||
};
|
};
|
||||||
self.finish(id, outcome);
|
self.complete(id, outcome);
|
||||||
grew
|
grew
|
||||||
}
|
}
|
||||||
|
|
||||||
/// The completion half, shared by [`Scheduler::complete`] and
|
|
||||||
/// [`Scheduler::complete_growing`] so neither is a redirect through the
|
|
||||||
/// other: the growing form must insert *before* this runs, and the plain
|
|
||||||
/// form must not pay for an empty job.
|
|
||||||
fn finish(&mut self, id: NodeId, outcome: Outcome) {
|
|
||||||
match outcome {
|
|
||||||
Outcome::Failed(error) => {
|
|
||||||
// Record the reason before the terminal transition so it's set
|
|
||||||
// by the time `set_state` stamps `finished_at`.
|
|
||||||
self.graph.set_error(id, error);
|
|
||||||
self.graph.set_state(id, State::Failed);
|
|
||||||
self.cascade_cancel(id);
|
|
||||||
}
|
|
||||||
Outcome::Done => self.settle_terminal(id),
|
|
||||||
}
|
|
||||||
self.roll_up_ancestors(id);
|
|
||||||
self.release_ready();
|
|
||||||
}
|
|
||||||
|
|
||||||
/// Whether every direct child of `id` is terminal.
|
/// Whether every direct child of `id` is terminal.
|
||||||
fn all_children_terminal(&self, id: NodeId) -> bool {
|
fn all_children_terminal(&self, id: NodeId) -> bool {
|
||||||
self.graph
|
self.graph
|
||||||
|
|
|
||||||
Loading…
Reference in a new issue