hyperhive/hive-c0re/src/job_queue/exec.rs
atlas d3d73b5ffb refactor(#2815): derive the transient pill from the running node
The dashboard pill was declared once per DAG at submit time, so a rebuild
reported `rebuilding` for its entire life — through the prebuild, the
stop, the swap, the tail and the reconcile. It named the intent of the
request, not what was happening.

It is now read off the nodes actually running. A node lights a pill when
it is `Running` and declares the agent's own resource. Declaring is the
test, not targeting: `Prebuild` and `MetaSync` name an agent but are
lease-exempt on purpose (the container keeps serving), so they must not
light one. It is also not the lease *owner* — `resource_state()` answers
"who holds the slot", which is a different question from "what is
running", and a descendant that borrows an ancestor's grant never
appears in that map.

`TransientKind` is gone entirely rather than being re-derived. The label
is the node's own wire tag (`NodeKind::as_str`) — the same vocabulary
`NodeView.kind` already ships, so a pill and a DAG node name an operation
identically and there is no second taxonomy to keep in step. Work with no
node behind it (destroy, migration) supplies its own literal.

`DagSpec::transient`, `Claim::transient`, `DagMeta::transient` and
`NodeKind::Dag`'s `transient` field all go with it.

## the safety half, which is deliberately not the display half

`crash_watch::is_deliberate_stop` used to match a `TransientKind` to
decide whether a vanished container was intentional or a crash. That made
a pill's display vocabulary decide an alerting question, so renaming or
adding a label would silently move the alerting boundary.

`TransientState` now carries two independent fields: `label` (rendered,
nothing branches on it) and `deliberate_stop` (read only by the crash
watcher). The producer sets the second, because the producer is the only
thing that knows — it is not recoverable from the first.

For queue work that value is `NodeKind::takes_container_down()`, and it
is emphatically not "holds a lease": `Create` and `Start` hold the
agent's lease exactly like `Stop` does, and a container dying *while
starting* is a real crash that must keep reporting as one. The default is
`false` on purpose — a wrong `false` costs a spurious crash event, a
wrong `true` swallows a real crash silently.

## known cost, accepted on the issue

A restart no longer reads `restarting`. No `NodeKind` is unique to a
restart — `restart_chain` reuses `Signal` / `StopForUpdate` / `Drain` /
`Reconcile` — because "restart" is a property of the DAG's shape, not of
any node. A restart now reads `signal` / `stop_for_update`, then the
agent returns.

`Start` / `Stop` / `PostSwap` run inside a lease-holding ancestor and
re-declare nothing, so they light no pill and the agent reads idle for
those windows. Closing that is the resources-where-constructed work
(#2818), not this change.

Checked with clippy (`--all-targets -D warnings`), `cargo test -p
hive-c0re -p hive-jobq` (321 + 40 passed) and `nix fmt`.
2026-08-01 16:06:06 +02:00

684 lines
31 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;
use hive_jobq::TerminalState;
use super::model::{NodeKind, NodeSpec};
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(Debug, Default)]
pub struct NodeOutput {
/// Whole per-agent *subgraphs* to append into *this same* DAG at
/// runtime — the single in-DAG-growth channel. Each inner
/// `Vec<NodeSpec>` is one independent subgraph whose `deps` are local
/// (0-based within that subgraph); the scheduler appends each via
/// [`super::JobQueue::append_subgraph`], which rebases the deps onto the DAG's
/// node-id space and roots the subgraph 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<Vec<NodeSpec>>,
}
/// 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 takes the agent's lifecycle lease (see
/// `NodeKind::needs_lease`) 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 (`NodeKind::needs_meta_window`, 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 (`NodeKind::needs_meta_window`) 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| {
super::templates::rebuild_nodes(
agent,
super::templates::RebuildOpts {
relock: true,
graceful: true,
},
0,
)
})
.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| {
super::templates::rebuild_nodes(
agent,
super::templates::RebuildOpts {
relock: false,
graceful: false,
},
0,
)
})
.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| vec![vec![super::templates::node(kind, Vec::new())]];
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 (`NodeKind::needs_meta_window`): 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 (`NodeKind::needs_meta_window`), 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
}