hyperhive/hive-c0re/src/job_queue/templates.rs
atlas 15a9d5b652 feat(#2454): drain agents before stopping them in the boot sweep
A host restart brings hive-c0re up and the startup sweep rebuilds every
stale agent. Until now that stop was mechanical: `StopForUpdate` hung
straight off `Prebuild`, so an agent that was mid-turn when the host went
down had its turn cut off rather than finished.

The sweep now builds the same `Signal` -> `Drain` -> `StopForUpdate` chain
a graceful `hivectl restart` already uses, reusing the existing nodes and
`GRACEFUL_STOP_TIMEOUT` unchanged. Cost is bounded: the per-agent drains
overlap, so the sweep waits one timeout in total rather than one per agent.

`Signal` parents the rest of the stop instead of sitting beside it. All
three of `Signal` / `Drain` / `StopForUpdate` declare the agent lease, and
a resource is held across its holder's whole subtree — as siblings each
would take the lease separately, leaving a window between them for another
DAG to claim the agent mid-bounce.

Scope is the boot sweep alone: a manual rebuild, a meta-update cascade
child and a deploy all still stop mechanically, and a test pins that shape.

`rebuild_nodes` takes a `RebuildOpts` struct rather than a second
positional `bool`, which two adjacent flags would have made easy to swap at
a call site. Its callers no longer hard-code the subgraph's length either:
the `EmitRebuilt` tails and `FinalizeDeploy` used literal indices that
silently encoded "this builder emits exactly six nodes with `Reconcile`
last", which a variable-length subgraph turns into a wrong-node edge rather
than a compile error. They read the index off the emitted list now.
2026-07-27 20:16:56 +02:00

632 lines
26 KiB
Rust

//! DAG shape builders — every operation as a template over the shared
//! node primitives — plus submit-time cycle validation (petgraph is
//! confined to this validation; the runtime store stays the plain
//! `Vec<Node>` + `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)
//! rebuild(a) graceful: … Prebuild(a) → Signal(a) → Drain(a) → StopForUpdate(a) → … [boot sweep only]
//! spawn(a): Provision(a) → Create(a) → WriteDropin(a) → Reconcile(a) [wanted=Up at approve]
//! perm-change(a): WritePermFile(a) → «rebuild subgraph»
//! meta-update(inp): MetaLock(inp) →«in-DAG rebuild subgraph per affected a»
//! reparent(moves): Reparent(moves) [no rebuild — topology.json is read live]
//! ```
//!
//! For the dynamic power-op shapes (`stop` / `start` / `restart`, built from
//! live online/offline state), see `submit.rs`.
use anyhow::{Result, bail};
use hive_jobq::{DepWhen, TerminalState};
use super::model::{DagSpec, Dep, NodeKind, NodeSpec, PermPayload, Source};
use crate::coordinator::TransientKind;
/// 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<Dep> {
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.
///
/// 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<Dep> {
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<Dep> {
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.
///
/// 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<Dep> {
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<Dep> {
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. `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 [`on_elimination_of`].
fn emit_rebuilt_tails(agent: &str, roots: &[u64], base: u64) -> Vec<NodeSpec> {
// The failure branch needs *both*: the ok branch being ruled out (that is the
// "something went wrong" signal) **and** every root actually finished. The
// second half is easy to forget and gets the ordering wrong without it — a
// failed `Prebuild` eliminates the ok branch immediately, while the recovery
// `Reconcile` is still bringing the container back up, so reporting straight
// off the elimination would announce the failure mid-recovery.
let mut on_fail = after_any_all(roots);
on_fail.extend(on_elimination_of(base));
vec![
node(
NodeKind::EmitRebuilt {
agent: agent.to_owned(),
ok: true,
},
after_ok_all(roots),
),
node(
NodeKind::EmitRebuilt {
agent: agent.to_owned(),
ok: false,
},
on_fail,
),
]
}
/// The approval-resolving tails for an approval-carrying DAG: one per outcome of
/// the DAG's single group-root `root`, each accepting only its own.
///
/// The `Cancelled` node is what keeps a dropped approval DAG from dangling its
/// row forever — its edge is the only one [`super::JobQueue::cancel`] spares.
fn resolve_approval_tails(approval_id: i64, root: u64) -> Vec<NodeSpec> {
[
TerminalState::Done,
TerminalState::Failed,
TerminalState::Cancelled,
]
.into_iter()
.map(|outcome| {
node(
NodeKind::ResolveApproval {
approval_id,
outcome,
},
on_outcome(root, &[outcome]),
)
})
.collect()
}
/// Build one **top-level (group-root)** node — `parent = None`. `kind` carries
/// the agent it targets ([`NodeKind`] is the payload directly). Shared with
/// `submit.rs`'s dynamic power-op builders. A root owns whatever resource it
/// declares for its whole subtree; its descendants borrow it (agent-lease /
/// build-slot continuity). Ordering vs other nodes is `deps`; grouping is
/// `parent`.
pub(crate) fn node(kind: NodeKind, deps: Vec<Dep>) -> 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<Dep>) -> NodeSpec {
NodeSpec {
kind,
deps,
parent: Some(parent),
}
}
/// Knobs for [`rebuild_nodes`]. A struct rather than two positional `bool`s so
/// a call site cannot silently swap them.
#[derive(Debug, Clone, Copy)]
pub(crate) struct RebuildOpts {
/// Re-lock the meta flake inside `MetaSync`.
pub relock: bool,
/// Give the agent its `Signal` → `Drain` window to finish the turn in
/// flight before the container is stopped, instead of stopping it
/// outright. Costs up to one `GRACEFUL_STOP_TIMEOUT` per subgraph, and
/// those overlap across agents.
pub graceful: bool,
}
/// The rebuild node subtree (nested, three group roots). `base` is the spec
/// index of the first node (`MetaSync`). Structure:
/// - `MetaSync` (base+0, **root**): the meta-repo preamble (dir prep, agent
/// sync, optional relock). Owns the global `MetaWindow` — and *only* for its
/// own short duration, which is why it is a sibling root rather than
/// `Prebuild`'s parent: a resource is held across the holder's whole subtree,
/// so parenting the build under it would extend a hive-global window over
/// every rebuild's nix build.
/// - `Prebuild` (base+1, **root**): `AfterOk` `MetaSync`. Owns the build slot
/// for the whole mechanical subtree below it. Lease-exempt — the nix build
/// overlaps other DAGs on the same agent.
/// - the **stop root** (base+2, child of `Prebuild`): owns the agent lease and
/// runs once `Prebuild` reaches `Finishing` (parent gate). Non-graceful that
/// is `StopForUpdate` itself; graceful it is `Signal`, with `Drain` and then
/// `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`.
/// - `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
/// a top-level root it survives the cancel-cascade of a failed `Prebuild`
/// (recovery-start invariant, which also covers a failed `MetaSync`: that
/// cancel-cascades `Prebuild`, i.e. terminal, so the tail still runs). It
/// takes a fresh lease; the tiny gap is harmless — `Reconcile` converges to
/// the persisted `wanted` idempotently.
pub(crate) fn rebuild_nodes(agent: &str, opts: RebuildOpts, base: u64) -> Vec<NodeSpec> {
let a = || agent.to_owned();
let RebuildOpts { relock, graceful } = opts;
let mut nodes = vec![
node(
NodeKind::MetaSync { agent: a(), relock },
if base == 0 {
Vec::new()
} else {
after_ok(base - 1)
},
),
node(NodeKind::Prebuild { agent: a() }, after_ok(base)),
];
// The stop root hangs off `Prebuild` and owns the agent lease for
// everything below it.
let stop_root = base + 2;
if graceful {
nodes.push(child(base + 1, NodeKind::Signal { agent: a() }, Vec::new()));
// `Drain` is a *child* of `Signal`, so the parent gate already orders
// it — a child must not dep on its own parent (dep-scope).
nodes.push(child(stop_root, NodeKind::Drain { agent: a() }, Vec::new()));
nodes.push(child(
stop_root,
NodeKind::StopForUpdate { agent: a() },
after_ok(stop_root + 1),
));
} else {
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
/// the merge has landed and `prepare_deploy` has staged the lock, plus the
/// [`NodeKind::FinalizeDeploy`] that closes the window behind it.
///
/// `relock = false` is the whole reason this composes: `prepare_deploy` already
/// 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** 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*
/// cancelled child rolls the parent up `Failed`). That's the old
/// `build_result`.
/// - `AfterOk` `Reconcile` — the old call passed `deferred_start = false` on
/// purpose: the container had to come back up *before* the deploy was
/// finalized. `Reconcile` alone would not do, being `AfterAny` — it reaches
/// `Done` even after a failed `Swap`.
///
/// Appended, not submitted: the roots below become children of the emitting
/// `DeployApply` (see [`super::JobQueue::append_subgraph`]), which puts them
/// 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) -> Vec<NodeSpec> {
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(),
},
vec![
Dep {
on: 1,
when: DepWhen::AFTER_OK,
},
Dep {
on: reconcile,
when: DepWhen::AFTER_OK,
},
],
));
nodes
}
/// Spec index of the `Reconcile` root a [`rebuild_nodes`] subgraph ends on,
/// for callers that gate a tail on it. Read off the emitted list rather than
/// hard-coded, because the subgraph's length depends on [`RebuildOpts`].
fn reconcile_index(rebuild: &[NodeSpec], base: u64) -> u64 {
base + u64::try_from(rebuild.len()).unwrap_or(0).saturating_sub(1)
}
/// One uniform rebuild shape — no `was_running` branch. `StopForUpdate`
/// noops when already down; the tail `Reconcile` auto-noops the start
/// when `wanted = Offline` (a rebuild of a deliberately-stopped agent
/// leaves it stopped). `relock = false` only for meta-update cascade
/// children.
///
/// Closed by an [`NodeKind::EmitRebuilt`] tail edged onto all three group-roots
/// (`MetaSync`, `Prebuild`, `Reconcile`) — `Prebuild`'s roll-up carries the
/// 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 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,
approval_id: None,
inputs: Vec::new(),
transient: Some(TransientKind::Rebuilding),
nodes,
}
}
/// Approval-driven deploy (`MergeConfigPr`) as a phase subtree rather than the
/// single opaque node it used to be. Structure:
/// - `DeployWindow` (0, **root**): the resource holder — global meta window,
/// agent lease, build slot — held across every child below. No work of its
/// own; it reaches `Finishing` immediately and the children run inside it.
/// - `MergeVerify` (1, child): drift-gate + fetch + eval-verify. Mutates
/// nothing, so a failure here cancel-cascades its siblings with the forge and
/// the applied repo exactly as they were.
/// - `DeployApply` (2, child, `AfterOk` `MergeVerify`): the irreversible half —
/// ff-merge + `prepare_deploy`. It doesn't rebuild inline; it grows
/// [`deploy_rebuild_nodes`] into this DAG as its own children, so the build
/// and the closing `FinalizeDeploy` are real nodes under the same window.
/// - `DeployTail` (3, child, `AfterAny` `DeployApply`): the compensation +
/// bookkeeping tail — rollback when a merge landed unfinalized, forge tag
/// mirror, PR failure comment (see [`NodeKind::DeployTail`]).
///
/// - `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
/// DAG's only other group-root, so its roll-up already *is* the whole
/// pipeline's outcome.
///
/// 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 a = || agent.to_owned();
DagSpec {
source: Source::Approval,
reason,
approval_id: Some(approval_id),
inputs: Vec::new(),
transient: Some(TransientKind::Rebuilding),
nodes: vec![
node(NodeKind::DeployWindow { agent: a() }, Vec::new()),
child(0, NodeKind::MergeVerify { agent: a() }, Vec::new()),
child(0, NodeKind::DeployApply { agent: a() }, after_ok(1)),
child(
0,
NodeKind::DeployTail { agent: a() },
vec![Dep {
on: 2,
when: DepWhen::AFTER_ANY,
}],
),
]
.into_iter()
.chain(resolve_approval_tails(approval_id, 0))
.collect(),
}
}
/// A single `Reconcile` node that converges observed power state to the
/// persisted intent — `wanted` is untouched (no `SetWanted`), unlike the
/// operator `start`/`stop` templates. Test-only helper now (used to build
/// 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,
transient: Option<TransientKind>,
) -> DagSpec {
DagSpec {
source,
reason,
approval_id: None,
inputs: Vec::new(),
transient,
nodes: vec![node(
NodeKind::Reconcile {
agent: agent.to_owned(),
},
Vec::new(),
)],
}
}
/// First-deploy spawn (approval-driven): `Provision` (proposed/applied
/// repos, state subvolume, meta registration) then `Create`
/// (`nixos-container create`), drop-in write, then `Reconcile` starts
/// the container (`wanted = Up` written at approve time). All-or-nothing:
/// `Provision` (lease-exempt, precedes the container) is the group root;
/// `Create` (child) owns the agent lease; `WriteDropin` + `Reconcile`
/// (children of `Create`) borrow it. A failure cancel-cascades the rest —
/// unlike rebuild there's no recovery-reconcile (nothing to converge if the
/// 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 {
DagSpec {
source: Source::Approval,
reason,
approval_id: Some(approval_id),
inputs: Vec::new(),
transient: Some(TransientKind::Spawning),
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`(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,
approval_id: None,
inputs: Vec::new(),
transient: Some(TransientKind::Rebuilding),
nodes,
}
}
/// Meta-input lock bump. The `MetaLock` executor grows one rebuild subgraph
/// per affected agent into *this same* DAG on completion (via
/// `append_subgraph`) — appended *after* the bump lands so their prebuilds
/// run against the post-bump lock, and a failed bump appends nothing
/// (replacing the old fan-out-child-DAGs dance).
/// `transient = Rebuilding` because those appended subgraphs are rebuilds:
/// it's applied per-agent at claim time (the `MetaLock` head needs no lease,
/// 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(
inputs: Vec<String>,
source: Source,
reason: String,
approval_id: Option<i64>,
) -> DagSpec {
let mut nodes = vec![node(
NodeKind::MetaLock {
sweep: false,
fanout: None,
},
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,
approval_id,
inputs,
transient: Some(TransientKind::Rebuilding),
nodes,
}
}
/// Topology move(s) as a single-node DAG. `moves` is `(child, new_parent)`
/// pairs — len 1 for `set-parent`, len N for `set-parent-bulk`, applied
/// uniformly by the one [`NodeKind::Reparent`] node (which holds the global
/// meta window for its duration, same precedent as [`NodeKind::WritePermFile`]).
/// No rebuild subgraph: `topology.json` is read live by every consumer
/// (dashboard tree, `<parent>`/`<children>` sentinel routing, permission
/// 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(
moves: Vec<(hive_types::Ident, Option<hive_types::Ident>)>,
source: Source,
reason: String,
) -> DagSpec {
DagSpec {
source,
reason,
approval_id: None,
inputs: Vec::new(),
transient: None,
nodes: vec![node(NodeKind::Reparent { moves }, Vec::new())],
}
}
// The boot is assembled inline in `workers/auto_update.rs::submit_boot_tree`
// 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::<u32, ()>::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(())
}