From 379c9bb5706406015ea8c8b42a7f9fa6f97fea82 Mon Sep 17 00:00:00 2001 From: atlas Date: Mon, 3 Aug 2026 12:53:52 +0200 Subject: [PATCH] jobq: the core alias is a JobBuilder, not a Job MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit A JobBuilder holds pending nodes that are not in the graph yet — it is the thing you declare into. Naming the alias Job claimed it was the work itself, and the name propagated into every parameter derived from it (job: super::Job in run_node read as if it carried the DAG). Prose uses meaning the job *queue* are left alone: main.rs's "Job-queue scheduler" comment and the docs/coordinator.md reference. 315 tests pass unchanged. --- hive-c0re/src/job_queue/exec.rs | 6 ++--- hive-c0re/src/job_queue/mod.rs | 8 +++---- hive-c0re/src/job_queue/submit.rs | 18 +++++++------- hive-c0re/src/job_queue/templates.rs | 36 ++++++++++++++-------------- hive-c0re/src/job_queue/tests.rs | 14 +++++------ hive-c0re/src/workers/auto_update.rs | 2 +- 6 files changed, 42 insertions(+), 42 deletions(-) diff --git a/hive-c0re/src/job_queue/exec.rs b/hive-c0re/src/job_queue/exec.rs index 403f11f6..a0348d4e 100644 --- a/hive-c0re/src/job_queue/exec.rs +++ b/hive-c0re/src/job_queue/exec.rs @@ -36,7 +36,7 @@ pub const GRACEFUL_STOP_TIMEOUT: std::time::Duration = std::time::Duration::from /// /// ⚠️ Taken **by value and handed back**, not by reference. A `JobBuilder` is /// `RefCell`-backed: owned it is `Send`, but `&JobBuilder` is not (a shared ref -/// is `Send` only if the referent is `Sync`, and `RefCell` never is). A `&Job` +/// is `Send` only if the referent is `Sync`, and `RefCell` never is). A `&JobBuilder` /// parameter would be live across every `.await` in this fn and make the whole /// future non-`Send`, which the scheduler's `tokio::spawn` rejects. So the /// growth executors below return *what to grow* and the declaration happens @@ -48,10 +48,10 @@ pub const GRACEFUL_STOP_TIMEOUT: std::time::Duration = std::time::Duration::from /// themselves. Nothing here needs a claim to exist as a type. pub(super) async fn run_node( coord: &Arc, - job: super::Job, + job: super::JobBuilder, id: NodeId, kind: &NodeKind, -) -> (super::Job, Result<()>) { +) -> (super::JobBuilder, Result<()>) { // The agent this node targets rides the payload — empty for the agentless // container kinds (`MetaLock`, `Dag`), which never read it. let agent = kind.agent(); diff --git a/hive-c0re/src/job_queue/mod.rs b/hive-c0re/src/job_queue/mod.rs index 12b7c08e..a58923d8 100644 --- a/hive-c0re/src/job_queue/mod.rs +++ b/hive-c0re/src/job_queue/mod.rs @@ -25,7 +25,7 @@ //! The queue is runtime-only (no persistence): an empty graph on boot; desired //! state is re-derived by the reconcile sweep. A single scheduler task //! ([`scheduler::run_worker`]) drives it; concurrency comes from the build-slot -//! capacity, not multiple workers. Design: `docs/coordinator.md::Job queue`. +//! capacity, not multiple workers. Design: `docs/coordinator.md::JobBuilder queue`. pub mod exec; pub mod model; @@ -53,7 +53,7 @@ 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; +pub type JobBuilder = hive_jobq::JobBuilder; /// A handle to one node a template declared — where its edges, grouping and /// resources are declared. `Copy`; naming a node as a dependency does not @@ -156,7 +156,7 @@ fn outcome_of(result: Result<(), String>) -> Outcome { /// Propagates a crate graph-insert error (malformed dep/parent / dep-scope). fn insert_group( inner: &mut Sched, - declare: impl FnOnce(&Job), + declare: impl FnOnce(&JobBuilder), group_parent: Option, ) -> anyhow::Result<()> { inner @@ -213,7 +213,7 @@ impl JobQueue { &self, source: Source, reason: String, - declare: impl FnOnce(&Job), + declare: impl FnOnce(&JobBuilder), ) -> anyhow::Result { let mut inner = self.lock(); let container = inner diff --git a/hive-c0re/src/job_queue/submit.rs b/hive-c0re/src/job_queue/submit.rs index 359a49c0..08241e35 100644 --- a/hive-c0re/src/job_queue/submit.rs +++ b/hive-c0re/src/job_queue/submit.rs @@ -6,7 +6,7 @@ //! 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 -//! (`Job::node` + `templates::rebuild_nodes`), all declaring into ONE job +//! (`JobBuilder::node` + `templates::rebuild_nodes`), all declaring into ONE job //! (independent per-agent roots, concurrent on their own leases). //! //! Dynamic shape rule: `stop`/`start` carry a head `SetWanted(w)` (durable @@ -27,7 +27,7 @@ use std::sync::Arc; use super::model::NodeKind; use super::resource::Resource; use super::templates::rebuild_nodes; -use super::{Job, Source, templates}; +use super::{JobBuilder, Source, templates}; use crate::coordinator::Coordinator; use crate::lifecycle; @@ -35,7 +35,7 @@ fn submit_and_emit( coord: &Arc, source: Source, reason: String, - declare: impl FnOnce(&Job), + declare: impl FnOnce(&JobBuilder), ) -> u64 { let id = coord .job_queue @@ -67,7 +67,7 @@ pub fn rebuild(coord: &Arc, agent: &str, source: Source, reason: St /// 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(b: &JobBuilder, agent: &str, graceful: bool, running: bool) { // `SetWanted` is the group root and owns the agent lease; the mechanical // steps are its children (borrow the lease, run once it reaches `Finishing`, // dep-ordered among themselves). @@ -107,7 +107,7 @@ fn stop_chain(b: &Job, agent: &str, graceful: bool, running: bool) { /// 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) { +fn start_chain(b: &JobBuilder, agent: &str, running: bool, stale: bool) { let wanted = b .node(NodeKind::SetWanted { agent: agent.to_owned(), @@ -138,7 +138,7 @@ 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(b: &JobBuilder, agent: &str, graceful: bool, running: bool) { let a = || agent.to_owned(); if !running { // Nothing to bounce — a lone Reconcile converges to intent. @@ -198,7 +198,7 @@ fn restart_chain(b: &Job, agent: &str, graceful: bool, running: bool) { // subgraph's indices onto another's used to be a function. /// Declare the stop DAG from explicit `(agent, running)` targets. -pub(crate) fn stop_nodes(b: &Job, targets: &[(String, bool)], graceful: bool) { +pub(crate) fn stop_nodes(b: &JobBuilder, targets: &[(String, bool)], graceful: bool) { for (agent, running) in targets { stop_chain(b, agent, graceful, *running); } @@ -210,14 +210,14 @@ pub(crate) fn stop_nodes(b: &Job, targets: &[(String, bool)], graceful: bool) { /// running under its lease, so a down+stale agent that grew a rebuild subgraph /// reports `rebuilding` during its swap and `starting` at its reconcile, /// without the DAG having to guess one label covering every target. -pub(crate) fn start_nodes(b: &Job, targets: &[(String, bool, bool)]) { +pub(crate) fn start_nodes(b: &JobBuilder, targets: &[(String, bool, bool)]) { for (agent, running, stale) in targets { start_chain(b, agent, *running, *stale); } } /// Declare the restart DAG from explicit `(agent, running)` targets. -pub(crate) fn restart_nodes(b: &Job, targets: &[(String, bool)], graceful: bool) { +pub(crate) fn restart_nodes(b: &JobBuilder, targets: &[(String, bool)], graceful: bool) { for (agent, running) in targets { restart_chain(b, agent, graceful, *running); } diff --git a/hive-c0re/src/job_queue/templates.rs b/hive-c0re/src/job_queue/templates.rs index 01b35b60..a7198835 100644 --- a/hive-c0re/src/job_queue/templates.rs +++ b/hive-c0re/src/job_queue/templates.rs @@ -16,19 +16,19 @@ //! ``` //! //! 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 +//! [`JobBuilder::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 ([`rebuild_nodes`]) over `Job::node`. +//! primitives this module exports ([`rebuild_nodes`]) over `JobBuilder::node`. use hive_jobq::TerminalState; use super::model::{NodeKind, PermPayload}; use super::resource::Resource; -use super::{Handle, Job}; +use super::{Handle, JobBuilder}; /// 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* @@ -36,7 +36,7 @@ use super::{Handle, Job}; /// /// 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<'_>]) { +fn emit_rebuilt_tails(b: &JobBuilder, agent: &str, roots: &[Handle<'_>]) { let ok = roots.iter().fold( b.node(NodeKind::EmitRebuilt { agent: agent.to_owned(), @@ -65,7 +65,7 @@ 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<'_>) { +fn resolve_approval_tails(b: &JobBuilder, approval_id: i64, root: Handle<'_>) { for outcome in [ TerminalState::Done, TerminalState::Failed, @@ -90,7 +90,7 @@ fn resolve_approval_tails(b: &Job, approval_id: i64, root: Handle<'_>) { /// /// Same reason as [`fanned_out_mechanical`] for living here: this was the /// second construction site declaring nodes inline in an executor. -pub(crate) fn grown_rebuilds(b: &Job, agents: &[String], relock: bool) { +pub(crate) fn grown_rebuilds(b: &JobBuilder, agents: &[String], relock: bool) { for agent in agents { rebuild_nodes(b, agent, relock, None); } @@ -99,7 +99,7 @@ pub(crate) fn grown_rebuilds(b: &Job, agents: &[String], relock: bool) { /// As [`grown_rebuilds`], but each agent gets its `Signal` → `Drain` window /// before being stopped. The boot sweep's flavour: it stops agents that were /// mid-turn when the host came up, so they drain rather than being cut off. -pub(crate) fn grown_graceful_rebuilds(b: &Job, agents: &[String], relock: bool) { +pub(crate) fn grown_graceful_rebuilds(b: &JobBuilder, agents: &[String], relock: bool) { for agent in agents { graceful_rebuild_nodes(b, agent, relock, None); } @@ -117,7 +117,7 @@ pub(crate) fn grown_graceful_rebuilds(b: &Job, agents: &[String], relock: bool) /// declaration does: this is the one construction site that was hiding in an /// executor, which meant the only test of it had to re-declare the same two /// calls itself and would have kept passing if the executor changed. -pub(crate) fn fanned_out_mechanical(b: &Job, kind: NodeKind) { +pub(crate) fn fanned_out_mechanical(b: &JobBuilder, kind: NodeKind) { let lease = Resource::Agent(kind.agent().to_owned()); let _ = b.node(kind).needs(lease); } @@ -176,7 +176,7 @@ impl<'a> RebuildRoots<'a> { /// takes a fresh lease; the tiny gap is harmless — `Reconcile` converges to /// the persisted `wanted` idempotently. fn rebuild_subtree<'a>( - b: &'a Job, + b: &'a JobBuilder, agent: &str, relock: bool, graceful: bool, @@ -250,7 +250,7 @@ fn rebuild_subtree<'a>( /// when given, is the node this subgraph chains behind. See /// [`rebuild_subtree`] for the structure. pub(crate) fn rebuild_nodes<'a>( - b: &'a Job, + b: &'a JobBuilder, agent: &str, relock: bool, after: Option>, @@ -266,7 +266,7 @@ pub(crate) fn rebuild_nodes<'a>( /// *prepend* nodes — it **re-parents** the stop root, so a caller cannot /// declare it without being handed the internals. Only the boot sweep wants it. pub(crate) fn graceful_rebuild_nodes<'a>( - b: &'a Job, + b: &'a JobBuilder, agent: &str, relock: bool, after: Option>, @@ -298,7 +298,7 @@ pub(crate) fn graceful_rebuild_nodes<'a>( /// `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(b: &Job, agent: &str, approval_id: i64) { +pub(crate) fn deploy_rebuild_nodes(b: &JobBuilder, agent: &str, approval_id: i64) { let roots = rebuild_nodes(b, agent, false, None); let _finalize = b .node(NodeKind::FinalizeDeploy { @@ -321,7 +321,7 @@ pub(crate) fn deploy_rebuild_nodes(b: &Job, agent: &str, approval_id: i64) { /// 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(b: &Job, agent: &str, relock: bool) { +pub fn rebuild(b: &JobBuilder, agent: &str, relock: bool) { let roots = rebuild_nodes(b, agent, relock, None); emit_rebuilt_tails(b, agent, &roots.all()); } @@ -351,7 +351,7 @@ pub fn rebuild(b: &Job, agent: &str, relock: bool) { /// /// 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(b: &Job, agent: &str, approval_id: i64) { +pub fn approval_deploy(b: &JobBuilder, agent: &str, approval_id: i64) { let a = || agent.to_owned(); // The window is the widest holder in the tree: it brackets a nix // build (`BuildSlot`), takes the container down across the swap @@ -402,7 +402,7 @@ pub fn approval_deploy(b: &Job, agent: &str, approval_id: i64) { /// 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(b: &Job, agent: &str, approval_id: i64) { +pub fn spawn(b: &JobBuilder, agent: &str, approval_id: i64) { let a = || agent.to_owned(); let provision = b .node(NodeKind::Provision { agent: a() }) @@ -430,7 +430,7 @@ pub fn spawn(b: &Job, agent: &str, approval_id: i64) { /// 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(b: &Job, agent: &str, payload: PermPayload) { +pub fn perm_change(b: &JobBuilder, agent: &str, payload: PermPayload) { let write = b .node(NodeKind::WritePermFile { agent: agent.to_owned(), @@ -455,7 +455,7 @@ pub fn perm_change(b: &Job, agent: &str, payload: PermPayload) { /// so the "hyperhive" pseudo-agent gets no pill), giving each cascade agent /// crash-watch suppression during its `Swap` — the property the old child /// `Rebuild` DAGs carried via their own transient. -pub fn meta_update(b: &Job, inputs: Vec, approval_id: Option) { +pub fn meta_update(b: &JobBuilder, inputs: Vec, approval_id: Option) { let lock = b .node(NodeKind::MetaLock { sweep: false, @@ -483,7 +483,7 @@ pub fn meta_update(b: &Job, inputs: Vec, approval_id: Option) { /// checks), so a parent move needs no container rebuild to take effect. /// No transient pill either — the node is agentless (no lease to hang one /// off of) and near-instant. No tail node: the write is the whole effect. -pub fn reparent(b: &Job, moves: Vec<(hive_types::Ident, Option)>) { +pub fn reparent(b: &JobBuilder, moves: Vec<(hive_types::Ident, Option)>) { let _reparent = b .node(NodeKind::Reparent { moves }) .needs(Resource::MetaWindow); diff --git a/hive-c0re/src/job_queue/tests.rs b/hive-c0re/src/job_queue/tests.rs index 94feb164..f9a062d4 100644 --- a/hive-c0re/src/job_queue/tests.rs +++ b/hive-c0re/src/job_queue/tests.rs @@ -20,7 +20,7 @@ use super::*; /// Submit a declared shape with the metadata every mechanics test uses. /// `Source::Manual` because none of these exercise provenance — the tests that /// do name their own source at the call site. -fn submit(q: &JobQueue, reason: &str, declare: impl FnOnce(&Job)) -> u64 { +fn submit(q: &JobQueue, reason: &str, declare: impl FnOnce(&JobBuilder)) -> u64 { q.submit(Source::Manual, reason.to_owned(), declare) .expect("valid shape") } @@ -29,7 +29,7 @@ fn ident(s: &str) -> hive_types::Ident { hive_types::Ident::parse(s).expect("valid test ident") } -fn rebuild(b: &Job, agent: &str) { +fn rebuild(b: &JobBuilder, agent: &str) { templates::rebuild(b, agent, true); } @@ -37,14 +37,14 @@ fn rebuild(b: &Job, agent: &str) { /// 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_nodes`). -fn restart_online(b: &Job, agents: &[&str], graceful: bool) { +fn restart_online(b: &JobBuilder, agents: &[&str], graceful: bool) { let targets: Vec<(String, bool)> = agents.iter().map(|a| ((*a).to_owned(), true)).collect(); submit::restart_nodes(b, &targets, graceful); } /// Stop shape with every agent treated as **running** — the online shape /// (`SetWanted → [Signal→Drain→](graceful) Reconcile`). -fn stop_online(b: &Job, agents: &[&str], graceful: bool) { +fn stop_online(b: &JobBuilder, agents: &[&str], graceful: bool) { let targets: Vec<(String, bool)> = agents.iter().map(|a| ((*a).to_owned(), true)).collect(); submit::stop_nodes(b, &targets, graceful); } @@ -426,7 +426,7 @@ fn rebuild_chain_is_declared_serial() { fn graceful_rebuild_chain_drains_before_stopping() { let q = JobQueue::new(1); let id = q - .submit(Source::AutoUpdate, "sweep".to_owned(), |b: &Job| { + .submit(Source::AutoUpdate, "sweep".to_owned(), |b: &JobBuilder| { templates::graceful_rebuild_nodes(b, "agent-a", true, None); }) .expect("valid shape"); @@ -737,7 +737,7 @@ fn boot_sweep_nodes_declare_their_own_resources() { // compile; only an exhaustive caller list would have caught it. let q = JobQueue::new(4); let id = q - .submit(Source::AutoUpdate, "boot".to_owned(), |b: &Job| { + .submit(Source::AutoUpdate, "boot".to_owned(), |b: &JobBuilder| { crate::workers::auto_update::boot_nodes( b, true, @@ -954,7 +954,7 @@ fn a_meta_lock_grows_one_rebuild_subgraph_per_agent() { let q = JobQueue::new(4); let agents = vec!["alice".to_owned(), "bob".to_owned()]; let id = q - .submit(Source::AutoUpdate, "sweep".to_owned(), |b: &Job| { + .submit(Source::AutoUpdate, "sweep".to_owned(), |b: &JobBuilder| { templates::grown_graceful_rebuilds(b, &agents, true); }) .expect("valid shape"); diff --git a/hive-c0re/src/workers/auto_update.rs b/hive-c0re/src/workers/auto_update.rs index c9b705d1..9c52e25c 100644 --- a/hive-c0re/src/workers/auto_update.rs +++ b/hive-c0re/src/workers/auto_update.rs @@ -327,7 +327,7 @@ pub async fn run(coord: Arc) -> Result<()> { /// kind-derived resources were removed, which drops the agent lease a boot /// reconcile needs to not race another DAG's container ops. pub(crate) fn boot_nodes( - b: &crate::job_queue::Job, + b: &crate::job_queue::JobBuilder, any_stale: bool, fanout: Vec, drifted: Vec,