hyperhive/hive-c0re/src/job_queue/exec.rs
atlas 58a9f218f2 job_queue: fix the boot sweep's lost declarations, drop the node wrapper
Two review findings on the resources-at-construction change.

argus: `workers::auto_update`'s boot sweep constructs nodes through
`templates::node` too, and it was not converted. With the kind-derived
declaration gone, its sweep `MetaLock` and its per-agent `Reconcile`
silently declared no resources at all — so a boot reconcile no longer
held the agent lease and could race another DAG's container ops, and the
sweep's meta commit could land inside another node's staged deploy
window. Nothing failed to compile: removing an implicit behaviour from a
helper is invisible at every call site that relied on it.

The declarations now live in a pure `boot_nodes`, split out of
`submit_boot_tree` so they can be exercised without a `Coordinator`.
That path is the only place job nodes are built outside `job_queue/`,
which is exactly why it had no coverage; `boot_sweep_nodes_declare_
their_own_resources` closes that, asserting against declared graph edges
rather than against the kind.

mara: `templates::node` is a redundant redirect now that it no longer
derives resources — deleted, and its 43 call sites use `Job::node`
directly. The reasoning it documented moved to the module docs of
`templates.rs` and `resource.rs`, which is where it stays true.
2026-08-02 16:29:06 +02:00

712 lines
32 KiB
Rust
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

//! Node executors — one async fn per [`NodeKind`], each a thin wrapper
//! over existing `lifecycle.rs` / `meta.rs` / `actions.rs` code. Node
//! executors keep their own internal error handling where it exists
//! today (cold-start fallback inside `Reconcile`, non-fatal boot-time
//! lock bump inside the sweep `MetaLock`, warn-only forge sync in the
//! `Swap` tail); DAG-level failure handling is cancel-downstream in
//! the queue.
use std::sync::Arc;
use anyhow::{Context as _, Result};
use super::{Claim, Declare};
use hive_jobq::TerminalState;
use super::model::NodeKind;
use super::resource::Resource;
use crate::coordinator::Coordinator;
use crate::power::{ReconcileAction, reconcile_action};
/// Max time `Drain` waits for the harness to run its stop-checkpoint
/// turn before falling back to the hard stop. Generous — a checkpoint
/// turn can take a while — but bounded so a wedged agent never blocks
/// the stop indefinitely. Drains hold no build slot, so a whole-hive
/// graceful stop overlaps every agent's drain instead of serialising
/// N × this timeout.
pub const GRACEFUL_STOP_TIMEOUT: std::time::Duration = std::time::Duration::from_mins(3);
/// Extra signal an executor hands back to the scheduler alongside
/// success.
#[derive(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<Declare>,
}
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()
}
}
/// Build-log sink for one claimed node.
struct Ctx<'a> {
coord: &'a Arc<Coordinator>,
dag_id: u64,
node_id: super::NodeId,
}
impl Ctx<'_> {
fn build_log(&self, log_id: i64) {
if self
.coord
.job_queue
.set_build_log_id(self.dag_id, self.node_id, log_id)
{
self.coord.emit_rebuild_queue_snapshot();
}
}
}
/// Run one claimed node to completion. Called from a task the
/// scheduler spawns per claim; the `Result` (stringified) becomes the
/// node's terminal state.
pub(super) async fn run_node(coord: &Arc<Coordinator>, claim: &Claim) -> Result<NodeOutput> {
let ctx = Ctx {
coord,
dag_id: claim.dag_id,
node_id: claim.node_id,
};
match &claim.kind {
NodeKind::MetaSync { relock, .. } => run_meta_sync(coord, claim, *relock).await,
NodeKind::Prebuild { .. } => run_prebuild(claim, &ctx).await,
NodeKind::Swap { .. } => run_swap(coord, claim, &ctx).await,
NodeKind::PostSwap { .. } => run_post_swap(coord, claim).await,
NodeKind::Provision { .. } => run_provision(coord, claim).await,
NodeKind::Create { .. } => run_create(claim).await,
NodeKind::MetaLock {
sweep,
fanout,
inputs,
} => run_meta_lock(coord, *sweep, fanout.clone(), inputs).await,
NodeKind::Reconcile { .. } => run_reconcile(coord, claim).await,
NodeKind::Start { .. } => run_start(coord, claim).await,
NodeKind::Stop { .. } => run_stop(coord, claim).await,
NodeKind::StopForUpdate { .. } => run_stop_for_update(coord, claim).await,
NodeKind::Signal { .. } => Ok(run_signal(coord, claim)),
NodeKind::Drain { .. } => run_drain(coord, claim).await,
NodeKind::WriteDropin { .. } => run_write_dropin(coord, claim).await,
NodeKind::WritePermFile { .. } => run_write_perm_file(coord, claim).await,
NodeKind::Reparent { .. } => run_reparent(coord, claim).await,
NodeKind::MergeVerify { approval_id, .. } => run_merge_verify(coord, *approval_id).await,
NodeKind::DeployApply { approval_id, .. } => {
run_deploy_apply(coord, claim, *approval_id).await
}
NodeKind::FinalizeDeploy { approval_id, .. } => {
run_finalize_deploy(coord, *approval_id).await
}
NodeKind::DeployTail { approval_id, .. } => {
run_deploy_tail(coord, claim, *approval_id).await
}
NodeKind::ResolveApproval {
approval_id,
outcome,
} => run_resolve_approval(coord, claim, *approval_id, *outcome).await,
NodeKind::EmitRebuilt { ok, .. } => Ok(run_emit_rebuilt(coord, claim, *ok)),
NodeKind::SetWanted { up, .. } => run_set_wanted(coord, claim, *up),
// The two nodes that carry no work of their own; completing either
// lets it reach `Finishing` so the nodes under it start.
// - `Dag`: pure grouping container. The DAG's terminal side effect, if
// any, is its own tail node in the graph.
// - `DeployWindow`: pure resource holder — the meta window, agent lease
// and build slot it declares stay held until its subtree settles.
NodeKind::Dag { .. } | NodeKind::DeployWindow { .. } => Ok(NodeOutput::default()),
}
}
/// Resolve the DAG's approval row the way this node's own `outcome` says.
///
/// Nothing is inspected: a template emits one of these per outcome, each edged to
/// accept only that one, so *which* node the scheduler let run already is the
/// answer. Best-effort — a resolution failure is logged inside
/// [`crate::actions::resolve_approval_dag`], never surfaced as a node failure,
/// since the work already happened and failing the tail would only misreport it.
async fn run_resolve_approval(
coord: &Arc<Coordinator>,
claim: &Claim,
approval_id: i64,
outcome: TerminalState,
) -> Result<NodeOutput> {
let reason = (outcome == TerminalState::Failed)
.then(|| coord.job_queue.first_error(claim.dag_id))
.flatten();
crate::actions::resolve_approval_dag(coord, approval_id, outcome, reason.as_deref()).await;
Ok(NodeOutput::default())
}
/// Emit this agent's `Rebuilt` manager event. `ok` is not computed — it is which
/// of the tail pair the graph let run. The failure note comes from the DAG's
/// first failing node, since the branch knows *that* it failed but not *why*.
fn run_emit_rebuilt(coord: &Arc<Coordinator>, claim: &Claim, ok: bool) -> NodeOutput {
coord.notify_manager(&hive_sh4re::HelperEvent::Rebuilt {
agent: claim.agent.clone(),
ok,
note: (!ok)
.then(|| coord.job_queue.first_error(claim.dag_id))
.flatten(),
sha: None,
tag: None,
});
NodeOutput::default()
}
/// Write the agent's durable power intent — the DAG-node form of the old
/// pre-submit `set_wanted` side effect. Store-only (no container touch), so
/// build-slot-exempt; but it declares the agent's lifecycle lease
/// (`Resource::Agent`) so the whole power-op DAG is atomic per-agent.
/// The downstream `Reconcile` reads the intent this writes. Unlike the old
/// warn-and-continue write, a failed write fails the node (cancel-downstream
/// cancels the `Reconcile`) rather than letting it converge to a stale
/// intent — that atomicity is the point of moving it into the DAG.
fn run_set_wanted(coord: &Arc<Coordinator>, claim: &Claim, up: bool) -> Result<NodeOutput> {
let wanted = if up {
crate::power::Wanted::Up
} else {
crate::power::Wanted::Offline
};
coord
.power
.set(&claim.agent, wanted)
.with_context(|| format!("set wanted={} for agent {}", wanted.as_str(), claim.agent))?;
Ok(NodeOutput::default())
}
/// The rebuild's meta preamble: runtime-dir prep, an idempotent meta
/// `sync_agents`, and the optional per-agent relock. Runs under the deploy
/// window (`Resource::MetaWindow`, held by the scheduler for this
/// node) so its commits can never land inside another node's staged
/// prepare→finalize window.
///
/// Deliberately a separate node from the [`run_prebuild`] it feeds: that
/// build takes minutes and only *reads* the store, so keeping the global
/// window off it is what lets rebuilds of different agents overlap.
async fn run_meta_sync(
coord: &Arc<Coordinator>,
claim: &Claim,
relock: bool,
) -> Result<NodeOutput> {
let name = &claim.agent;
// Runs while the agent is still up — the runtime dir and MCP listener
// already exist. Use the pure path accessor; no need to re-register the
// listener (event-driven: registered at start/create).
let agent_dir = crate::paths::agent_runtime_dir(name);
let hive = coord.hive_env();
let paths = Coordinator::agent_paths(name, agent_dir);
crate::lifecycle::prepare_rebuild_dirs(name, &paths).await?;
// Idempotent meta sync so a manual rebuild can also recover from a
// divergent meta repo; then bump just this agent's input. `relock =
// false` only for meta-update cascade children, where re-locking
// would revert the bump the cascade just committed.
let agents = crate::lifecycle::agents_for_meta_listing().await?;
crate::meta::sync_agents(&hive, &agents).await?;
if relock {
crate::meta::lock_update_for_rebuild(name).await?;
}
Ok(NodeOutput::default())
}
/// Out-of-band toplevel build while the container keeps serving: warm
/// `system.build.toplevel` so the later `Swap` hits cache and skips
/// straight to the profile-swap, against a meta repo the upstream
/// `MetaSync` node has already synced. The warm build is skipped when the
/// container is already down: its only purpose is to shrink the swap's
/// downtime window, so a stopped agent (no uptime to preserve) doesn't
/// pay the double eval — `Swap` builds inline instead.
async fn run_prebuild(claim: &Claim, ctx: &Ctx<'_>) -> Result<NodeOutput> {
let name = &claim.agent;
// Warm the toplevel build only when the container is up — the whole
// point of prebuild is to shrink the swap's downtime window. A
// stopped agent has no uptime to preserve, so skip the (expensive)
// eval and let the downstream `Swap` build inline.
if crate::lifecycle::is_running(name).await {
let flake_ref = format!("{}#{name}", crate::paths::meta_root().display());
crate::lifecycle::prebuild_toplevel(name, &flake_ref, &|log_id| ctx.build_log(log_id))
.await?;
}
Ok(NodeOutput::default())
}
/// Profile-swap: re-apply drop-ins (rebuild is the reconcile verb),
/// `nixos-container update`, then the post-rebuild bookkeeping tail
/// (rev marker, `Rebuilt` event, forge/matrix sync, kick, rescan).
/// The recovery-start on failure is NOT here — the DAG's tail
/// `Reconcile` runs after this node terminal ok *or* fail.
async fn run_swap(coord: &Arc<Coordinator>, claim: &Claim, ctx: &Ctx<'_>) -> Result<NodeOutput> {
let name = &claim.agent;
// Swap runs on an already-existing (stopped) container — runtime dir
// and listener were created earlier. Pure path accessor suffices.
let agent_dir = crate::paths::agent_runtime_dir(name);
let hive = coord.hive_env();
let paths = Coordinator::agent_paths(name, agent_dir);
let result = crate::lifecycle::swap_update(name, &hive, &paths, &|log_id| {
ctx.build_log(log_id);
})
.await;
// On success the Ok-only bookkeeping tail (rev marker, forge/matrix
// sync, kick, rescan, snapshot) runs in the sibling `PostSwap` node,
// which deps `AfterOk(Swap)`. On failure `PostSwap` is cancel-cascaded
// and the tail `Reconcile` (`AfterAny(PostSwap)`) handles recovery; here
// we only refresh the observed state so dashboards reflect the failed
// swap immediately. The `Rebuilt { ok: false }` manager event is emitted by
// the DAG's `EmitRebuilt` tail (any node may be the one that failed).
if result.is_err() {
coord.rescan_containers_and_emit().await;
}
result.map(|()| NodeOutput::default())
}
/// The post-`Swap` bookkeeping tail, split into its own node for dashboard
/// visibility + retry granularity. Deps `AfterOk(Swap)`, so reaching here
/// means the profile swap succeeded. Store/forge/matrix work only — no nix
/// build (build-slot-exempt); the agent lease taken at `Swap` is still held
/// (the whole chain up to `Reconcile` is one agent's subgraph).
async fn run_post_swap(coord: &Arc<Coordinator>, claim: &Claim) -> Result<NodeOutput> {
let name = &claim.agent;
if let Some(rev) = crate::auto_update::current_flake_rev(&coord.hyperhive_flake)
&& let Err(e) = std::fs::write(crate::paths::applied_rev_marker(name), rev)
{
tracing::warn!(%name, error = ?e, "write rev marker failed");
}
// The `Rebuilt` manager event is emitted exactly once per agent by the DAG's
// `EmitRebuilt` tail — emitting ok here and letting a failed tail `Reconcile`
// add a contradictory !ok would double-report the same rebuild.
// Full forge + matrix sync on every successful rebuild so the rebuild
// path is equivalent to the startup sweep: tokens, config-repo mirror,
// meta access all recover without a hive-c0re restart.
crate::forge::sync_agent(name, crate::forge::core_token().as_deref()).await;
crate::matrix::sync_agent_standalone(name).await;
// Wake the agent on its next turn so claude sees a "you were rebuilt"
// hint; rescan so dashboards drop the "needs update" chip; lock bump →
// meta-inputs re-render.
coord.kick_agent(name, "container rebuilt");
coord.rescan_containers_and_emit().await;
crate::dashboard::emit_meta_inputs_snapshot(coord);
Ok(NodeOutput::default())
}
/// First-spawn pre-create provisioning: proposed/applied repos, state
/// subvolume, and the meta `sync_agents` registration. Runs under the
/// deploy window (it declares `Resource::MetaWindow`) so its commit can't
/// land inside another node's staged deploy window.
async fn run_provision(coord: &Arc<Coordinator>, claim: &Claim) -> Result<NodeOutput> {
let name = &claim.agent;
let agent_dir = crate::paths::agent_runtime_dir(name);
let hive = coord.hive_env();
let paths = Coordinator::agent_paths(name, agent_dir);
crate::lifecycle::provision_container(name, &hive, &paths).await?;
Ok(NodeOutput::default())
}
/// `nixos-container create` proper — the upstream `Provision` node
/// already registered the agent in meta, so this only reads the store
/// (no deploy-window gate needed, mirroring `Prebuild`'s build). Runtime
/// dir creation and MCP listener registration are deferred to the tail
/// `Reconcile` (`converge_start_preamble` + `register_agent`) so this
/// node stays purely "create", not "create + start".
async fn run_create(claim: &Claim) -> Result<NodeOutput> {
crate::lifecycle::create_only(&claim.agent).await?;
Ok(NodeOutput::default())
}
/// Meta flake lock bump. Boot-sweep flavour is non-fatal (a failed
/// bump must not cancel the fan-out rebuilds — they proceed against
/// the current lock, exactly like today's sweep); the meta-update
/// flavour propagates errors, and a failed bump fans out nothing.
async fn run_meta_lock(
coord: &Arc<Coordinator>,
sweep: bool,
fanout: Option<Vec<String>>,
inputs: &[String],
) -> Result<NodeOutput> {
if sweep {
if let Err(e) = crate::meta::lock_update_hyperhive().await {
tracing::warn!(error = ?e, "startup sweep: meta lock_update_hyperhive failed");
}
// Grow one rebuild subgraph per stale agent into *this* boot DAG
// (rooted on this `MetaLock`, so they build against the post-bump
// lock), rather than fanning out child DAGs. `relock = true` — a
// boot sweep relocks per-agent like a manual rebuild.
//
// `graceful = true` here and nowhere else: a boot sweep stops agents
// that were already mid-turn when the host came up, so they get their
// drain window rather than being cut off. The per-agent drains overlap,
// so the sweep's cost ceiling is one `GRACEFUL_STOP_TIMEOUT` in total,
// not one per agent.
let append_subgraph = fanout
.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
})
.collect();
return Ok(NodeOutput { append_subgraph });
}
let _progress = coord.meta_update_guard();
crate::meta::lock_update(inputs).await?;
// Lock file changed — meta-inputs panel re-renders.
crate::dashboard::emit_meta_inputs_snapshot(coord);
let cascade = match fanout {
Some(list) => list,
None => meta_update_cascade_agents(inputs).await,
};
// Grow one rebuild subgraph per affected agent into *this* meta-update
// DAG (rooted on this `MetaLock`, so they build against the post-bump
// lock), rather than fanning out child DAGs. `relock = false` — the
// cascade children must NOT re-lock, which would revert the bump this
// node just committed (the property the old `fanout_specs` meta-update
// branch encoded).
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
})
.collect();
Ok(NodeOutput { append_subgraph })
}
/// Idempotent power-converge *planner*: compare `wanted` (durable
/// intent) against observed state and, when they diverge, fan the
/// mechanical `Start` / `Stop` out as a first-class node appended to
/// *this* DAG (a single-node `NodeOutput::append_subgraph` rooted on
/// this node). Does no container work itself — the sub-step becomes
/// visible in the DAG and the lease-window transient (or the sub-step's
/// own node-local guard) rides across it.
async fn run_reconcile(coord: &Arc<Coordinator>, claim: &Claim) -> Result<NodeOutput> {
let name = &claim.agent;
let running = crate::lifecycle::is_running(name).await;
let wanted = coord.power.get_or_seed(name, running)?;
// One node targeting this agent, rooted on this reconcile node. `NodeKind`
// carries the agent it targets, so stamp `claim.agent` into the fanned-out
// Start/Stop kind (one in-DAG-growth channel).
let sub = |kind: NodeKind| {
// `Start` / `Stop` declare the lease they run under. This node is their
// parent and holds it, so the declaration is a re-entrant borrow — no
// second unit, no deadlock. It exists so the requirement belongs to the
// node rather than to the fact that a `Reconcile` happens to fan it out.
let lease = Resource::Agent(kind.agent().to_owned());
vec![Box::new(move |b: &super::Job| {
let _ = b.node(kind).needs(lease);
}) as Declare]
};
let append_subgraph = match reconcile_action(wanted, running) {
ReconcileAction::Start => sub(NodeKind::Start {
agent: name.clone(),
}),
ReconcileAction::Stop => sub(NodeKind::Stop {
agent: name.clone(),
}),
ReconcileAction::Noop => {
tracing::debug!(%name, wanted = wanted.as_str(), running, "reconcile: noop");
Vec::new()
}
};
Ok(NodeOutput { append_subgraph })
}
/// Mechanical container start — the sub-step a `Reconcile` planner fans
/// out when it observes `wanted = Up` and the container down.
async fn run_start(coord: &Arc<Coordinator>, claim: &Claim) -> Result<NodeOutput> {
let name = &claim.agent;
// No node-local transient guard: the pill is derived from the running node
// set, and `Start` reports `Starting` via `NodeKind::transient_kind`. This
// used to take one "only when the DAG holds none", which was a second
// derivation covering the gap left by a DAG-level declaration that couldn't
// describe a sub-step.
// Run the typed start preamble: ensures the runtime dir exists and
// writes the nspawn/resource-limits drop-ins. The returned
// StartableAgent token is the only way to call start_with_fallback —
// omitting this becomes a compile error.
let agent_dir = crate::paths::agent_runtime_dir(name);
let hive = coord.hive_env();
let paths = Coordinator::agent_paths(name, agent_dir);
let token = crate::lifecycle::converge_start_preamble(name, &hive, &paths).await?;
crate::lifecycle::start_with_fallback(token).await?;
// Bind the MCP listener immediately after starting the container.
// The preamble created the runtime dir; the container is now coming
// up and will connect to this socket on its first turn. Event-driven
// (no background poll) — c0re owns the listener lifecycle, so
// register here rather than waiting for a sweep.
coord.register_agent(name)?;
coord.kick_agent(name, "container started");
coord.rescan_containers_and_emit().await;
Ok(NodeOutput::default())
}
/// Mechanical container stop — the sub-step a `Reconcile` planner fans
/// out when it observes `wanted = Offline` and the container up.
async fn run_stop(coord: &Arc<Coordinator>, claim: &Claim) -> Result<NodeOutput> {
let name = &claim.agent;
// See `run_start`: no node-local guard — `Stop` reports `Stopping` from its
// own kind now.
crate::lifecycle::kill(name).await?;
coord.unregister_agent(name);
coord.notify_manager(&hive_sh4re::HelperEvent::Killed {
agent: name.clone(),
});
coord.rescan_containers_and_emit().await;
Ok(NodeOutput::default())
}
/// Mechanical stop for the profile swap. Never *changes* `wanted`;
/// noop when already stopped.
async fn run_stop_for_update(coord: &Arc<Coordinator>, claim: &Claim) -> Result<NodeOutput> {
let name = &claim.agent;
if crate::lifecycle::is_running(name).await {
// Seed a missing agent_power row from the PRE-stop observation
// — the DAG's tail `Reconcile` observes only the mechanically
// stopped state and would otherwise seed a running-but-unknown
// agent as `Offline`, stranding it down after its own rebuild.
if let Err(e) = coord.power.get_or_seed(name, true) {
tracing::warn!(%name, error = ?e, "agent_power: pre-stop seed failed");
}
crate::lifecycle::kill(name).await?;
coord.rescan_containers_and_emit().await;
}
Ok(NodeOutput::default())
}
/// Set the graceful fence + kick so the harness sees it promptly and
/// runs its one stop-checkpoint turn.
///
/// Skipped entirely for a paused agent: its loop parks on the pause
/// marker without polling the broker, so it would never observe the
/// fence and the downstream drain would just burn
/// `GRACEFUL_STOP_TIMEOUT`. Safe because the harness tests the marker at
/// the top of its loop — a paused agent has no turn in flight, so there
/// is nothing to checkpoint.
fn run_signal(coord: &Arc<Coordinator>, claim: &Claim) -> NodeOutput {
if hive_types::Ident::parse(&claim.agent).is_ok_and(|a| Coordinator::is_paused(&a)) {
return NodeOutput::default();
}
coord.mark_graceful_stop(&claim.agent);
coord.kick_agent(&claim.agent, "graceful stop requested");
NodeOutput::default()
}
/// Await the harness clearing the fence (`GracefulStopComplete`) or
/// the timeout — either way the downstream `Reconcile` proceeds with
/// the actual stop.
async fn run_drain(coord: &Arc<Coordinator>, claim: &Claim) -> Result<NodeOutput> {
let name = &claim.agent;
let deadline = std::time::Instant::now() + GRACEFUL_STOP_TIMEOUT;
while coord.is_graceful_stop_pending(name) {
if std::time::Instant::now() >= deadline {
tracing::warn!(agent = %name, "graceful stop: drain timed out — hard-stopping");
break;
}
tokio::time::sleep(std::time::Duration::from_millis(500)).await;
}
coord.clear_graceful_stop(name);
Ok(NodeOutput::default())
}
/// `set_nspawn_flags` + `set_resource_limits` + daemon-reload.
async fn run_write_dropin(coord: &Arc<Coordinator>, claim: &Claim) -> Result<NodeOutput> {
let name = &claim.agent;
// write_dropins only needs the path value to build AgentPaths; the
// dir doesn't need to exist at this point (created by ensure_agent_runtime_dir
// on the upstream Prebuild/Start node).
let agent_dir = crate::paths::agent_runtime_dir(name);
let hive = coord.hive_env();
let paths = Coordinator::agent_paths(name, agent_dir);
crate::lifecycle::write_dropins(name, &hive, &paths).await?;
Ok(NodeOutput::default())
}
/// Write + commit the perm file(s) (fused under `META_LOCK` so the
/// working tree is never left dirty), then emit the P3RM1SS10NS-tab
/// snapshots so the dashboard reflects the new assignment.
async fn run_write_perm_file(coord: &Arc<Coordinator>, claim: &Claim) -> Result<NodeOutput> {
use super::model::PermPayload;
let name = &claim.agent;
// The perm file payload rides the node itself (the only consumer).
let NodeKind::WritePermFile { payload, .. } = &claim.kind else {
anyhow::bail!("run_write_perm_file on a non-WritePermFile node");
};
// Runs under the deploy window (it declares `Resource::MetaWindow`): a
// perm commit landing inside another node's staged prepare→finalize
// window would sweep the staged deploy lock into its commit (the
// commits are also path-limited in meta.rs — belt and braces).
match payload {
PermPayload::ToolGroups { groups } => {
crate::meta::commit_tool_groups(name, groups)
.await
.with_context(|| format!("commit tool-groups for {name}"))?;
coord.emit_tool_groups_snapshot();
}
PermPayload::Capabilities { caps } => {
crate::meta::commit_capabilities(name, caps)
.await
.with_context(|| format!("commit capabilities for {name}"))?;
coord.emit_capabilities_snapshot();
}
PermPayload::Combined { groups, caps } => {
crate::meta::commit_perms(name, groups.as_deref(), caps.as_deref())
.await
.with_context(|| format!("commit perms for {name}"))?;
if groups.is_some() {
coord.emit_tool_groups_snapshot();
}
if caps.is_some() {
coord.emit_capabilities_snapshot();
}
}
}
Ok(NodeOutput::default())
}
/// Apply the node's `(child, new_parent)` moves as one `META_LOCK`-fused
/// commit (`Coordinator::reparent_bulk_with_notify`, which already handles
/// both the single- and bulk-move case, sends the per-agent move
/// notifications, and rescans + diff-emits the container tree). Runs under
/// the deploy window (it declares `Resource::MetaWindow`), same reasoning as
/// `run_write_perm_file`: a topology commit landing inside another node's
/// staged deploy window would sweep the staged lock into its commit.
async fn run_reparent(coord: &Arc<Coordinator>, claim: &Claim) -> Result<NodeOutput> {
let NodeKind::Reparent { moves } = &claim.kind else {
anyhow::bail!("run_reparent on a non-Reparent node");
};
let refs: Vec<(&str, Option<&str>)> = moves
.iter()
.map(|(child, parent)| {
(
child.as_str(),
parent.as_ref().map(hive_types::Ident::as_str),
)
})
.collect();
coord
.reparent_bulk_with_notify(&refs)
.await
.map_err(|e| anyhow::anyhow!(e))?;
Ok(NodeOutput::default())
}
/// Deploy phase 1 — drift gate, fetch, eval-verify. Mutates nothing, so a
/// failure here cancel-cascades the rest of the subtree with the forge and the
/// applied repo exactly as they were.
async fn run_merge_verify(coord: &Arc<Coordinator>, approval_id: i64) -> Result<NodeOutput> {
crate::actions::run_deploy_merge_verify(coord, approval_id)
.await
.map(|()| NodeOutput::default())
}
/// Deploy phase 2 — the irreversible half: ff-merge, then phase 1 of the
/// two-phase meta deploy.
///
/// On success it grows the ordinary rebuild subgraph (plus its closing
/// `FinalizeDeploy`) into this DAG rooted on *this* node — which is what puts
/// the appended nodes inside the `DeployWindow`'s subtree, so the `MetaWindow`
/// their `MetaSync` declares is re-entered rather than deadlocked against the
/// ancestor already holding it. On failure nothing is appended and the tail
/// compensates, exactly as before.
async fn run_deploy_apply(
coord: &Arc<Coordinator>,
claim: &Claim,
approval_id: i64,
) -> Result<NodeOutput> {
crate::actions::run_deploy_apply(coord, approval_id).await?;
Ok(NodeOutput {
append_subgraph: vec![super::templates::deploy_rebuild_nodes(
claim.kind.agent(),
approval_id,
)],
})
}
/// Deploy phase 3 — close the staged-lock window once the appended rebuild has
/// come up clean: drop the rollback ref, plant the `deployed/<id>` tag, commit
/// the staged lock.
async fn run_finalize_deploy(coord: &Arc<Coordinator>, approval_id: i64) -> Result<NodeOutput> {
crate::actions::run_finalize_deploy(coord, approval_id)
.await
.map(|()| NodeOutput::default())
}
/// Deploy compensation + bookkeeping tail. `AfterAny` the apply node, so it
/// runs on every outcome; it is deliberately infallible (see
/// [`crate::actions::run_deploy_tail`]) — a failing tail must not flip an
/// otherwise-successful deploy's DAG state.
///
/// Takes the agent from the node payload so the tail can still compensate when
/// the approval row is gone (deny race, purge).
async fn run_deploy_tail(
coord: &Arc<Coordinator>,
claim: &Claim,
approval_id: i64,
) -> Result<NodeOutput> {
crate::actions::run_deploy_tail(coord, Some(claim.dag_id), claim.kind.agent(), approval_id)
.await;
Ok(NodeOutput::default())
}
/// Compute which agents a `nix flake update <inputs>` on the meta
/// flake affects — the fan-out set for `MetaUpdate` DAGs. Empty
/// `inputs` or any input under `hyperhive` → every container;
/// otherwise just the agents named by `agent-<name>` inputs.
/// Topology-sorted so parents rebuild before their children.
pub async fn meta_update_cascade_agents(inputs: &[String]) -> Vec<String> {
let touched_hyperhive = inputs
.iter()
.any(|i| i == "hyperhive" || i.starts_with("hyperhive/"));
let touched_agents: Vec<String> = inputs
.iter()
.filter_map(|i| i.strip_prefix("agent-"))
.map(|rest| rest.split('/').next().unwrap_or(rest).to_owned())
.collect();
let mut names = if touched_hyperhive || inputs.is_empty() {
crate::lifecycle::list()
.await
.unwrap_or_default()
.into_iter()
.filter_map(|c| {
c.strip_prefix(crate::lifecycle::AGENT_PREFIX)
.map(str::to_owned)
})
.collect()
} else {
touched_agents
};
let topo = crate::topology::read();
crate::auto_update::topology_sort(&mut names, &topo);
names
}