Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
0ab6b764be | ||
|
|
600bc051e1 | ||
|
|
be2dfa8cd3 | ||
|
|
2294cd4516 | ||
|
|
a78280feed | ||
|
|
8834161fb9 | ||
|
|
50172d1716 | ||
|
|
456847eaa1 | ||
|
|
a5c321a1a0 |
16 changed files with 1359 additions and 1029 deletions
1
Cargo.lock
generated
1
Cargo.lock
generated
|
|
@ -1620,6 +1620,7 @@ dependencies = [
|
|||
"forgejo-api",
|
||||
"hive-core-agent-sock",
|
||||
"hive-host-sock",
|
||||
"hive-jobq",
|
||||
"hive-priv-sock",
|
||||
"hive-sh4re",
|
||||
"hive-types",
|
||||
|
|
|
|||
|
|
@ -53,6 +53,7 @@ clap_complete = "4"
|
|||
indicatif = "0.18"
|
||||
hive-sh4re = { path = "hive-sh4re" }
|
||||
hive-agent-sock = { path = "hive-agent-sock" }
|
||||
hive-jobq = { path = "hive-jobq" }
|
||||
hive-core-agent-sock = { path = "hive-core-agent-sock" }
|
||||
hive-claude = { path = "hive-claude" }
|
||||
hive-host-sock = { path = "hive-host-sock" }
|
||||
|
|
|
|||
|
|
@ -33,6 +33,7 @@ indicatif.workspace = true
|
|||
hive-core-agent-sock.workspace = true
|
||||
hive-sh4re.workspace = true
|
||||
hive-host-sock.workspace = true
|
||||
hive-jobq.workspace = true
|
||||
hive-priv-sock.workspace = true
|
||||
hive-types.workspace = true
|
||||
libc.workspace = true
|
||||
|
|
|
|||
|
|
@ -369,7 +369,8 @@ impl Drop for MetaUpdateGuard {
|
|||
}
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone, Copy)]
|
||||
#[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Serialize)]
|
||||
#[serde(rename_all = "snake_case")]
|
||||
pub enum TransientKind {
|
||||
/// `lifecycle::spawn` is running (nixos-container create + update + start).
|
||||
Spawning,
|
||||
|
|
|
|||
|
|
@ -127,8 +127,10 @@ pub(super) async fn post_rebuild_queue_cancel(
|
|||
State(state): State<AppState>,
|
||||
AxumPath(id): AxumPath<u64>,
|
||||
) -> Response {
|
||||
let cancelled = state.coord.job_queue.cancel(id);
|
||||
if cancelled {
|
||||
if let Some(terminal) = state.coord.job_queue.cancel(id) {
|
||||
// Fire the DAG's inline terminal hook (power-op intent revert / approval
|
||||
// resolution) off the cancel roll-up, then surface the flip live.
|
||||
crate::job_queue::exec::run_terminal_hook(&state.coord, &terminal).await;
|
||||
state.coord.emit_rebuild_queue_snapshot();
|
||||
axum::Json(serde_json::json!({"cancelled": true})).into_response()
|
||||
} else {
|
||||
|
|
|
|||
|
|
@ -10,8 +10,8 @@ use std::sync::Arc;
|
|||
|
||||
use anyhow::{Context as _, Result};
|
||||
|
||||
use super::model::{NodeKind, NodeSpec, State, Template};
|
||||
use super::{Claim, TerminalDag};
|
||||
use super::Claim;
|
||||
use super::model::{NodeKind, NodeSpec, State};
|
||||
use crate::coordinator::Coordinator;
|
||||
use crate::power::{ReconcileAction, reconcile_action};
|
||||
|
||||
|
|
@ -82,24 +82,86 @@ pub(super) async fn run_node(coord: &Arc<Coordinator>, claim: &Claim) -> Result<
|
|||
node_id: claim.node_id,
|
||||
};
|
||||
match &claim.kind {
|
||||
NodeKind::Prebuild { relock } => run_prebuild(coord, claim, &ctx, *relock).await,
|
||||
NodeKind::Swap => run_swap(coord, claim, &ctx).await,
|
||||
NodeKind::PostSwap => run_post_swap(coord, claim, &ctx).await,
|
||||
NodeKind::Provision => run_provision(coord, claim, &ctx).await,
|
||||
NodeKind::Create => run_create(claim, &ctx).await,
|
||||
NodeKind::Prebuild { relock, .. } => run_prebuild(coord, claim, &ctx, *relock).await,
|
||||
NodeKind::Swap { .. } => run_swap(coord, claim, &ctx).await,
|
||||
NodeKind::PostSwap { .. } => run_post_swap(coord, claim, &ctx).await,
|
||||
NodeKind::Provision { .. } => run_provision(coord, claim, &ctx).await,
|
||||
NodeKind::Create { .. } => run_create(claim, &ctx).await,
|
||||
NodeKind::MetaLock { sweep, fanout } => {
|
||||
run_meta_lock(coord, claim, &ctx, *sweep, fanout.clone()).await
|
||||
}
|
||||
NodeKind::Reconcile => run_reconcile(coord, claim).await,
|
||||
NodeKind::Start => run_start(coord, claim, &ctx).await,
|
||||
NodeKind::Stop => run_stop(coord, claim, &ctx).await,
|
||||
NodeKind::StopForUpdate => run_stop_for_update(coord, claim, &ctx).await,
|
||||
NodeKind::Signal => Ok(run_signal(coord, claim, &ctx)),
|
||||
NodeKind::Drain => run_drain(coord, claim, &ctx).await,
|
||||
NodeKind::WriteDropin => run_write_dropin(coord, claim).await,
|
||||
NodeKind::WritePermFile => run_write_perm_file(coord, claim, &ctx).await,
|
||||
NodeKind::ApprovalDeploy => run_approval_deploy(coord, claim).await,
|
||||
NodeKind::SetWanted { up } => run_set_wanted(coord, claim, *up),
|
||||
NodeKind::Reconcile { .. } => run_reconcile(coord, claim).await,
|
||||
NodeKind::Start { .. } => run_start(coord, claim, &ctx).await,
|
||||
NodeKind::Stop { .. } => run_stop(coord, claim, &ctx).await,
|
||||
NodeKind::StopForUpdate { .. } => run_stop_for_update(coord, claim, &ctx).await,
|
||||
NodeKind::Signal { .. } => Ok(run_signal(coord, claim, &ctx)),
|
||||
NodeKind::Drain { .. } => run_drain(coord, claim, &ctx).await,
|
||||
NodeKind::WriteDropin { .. } => run_write_dropin(coord, claim).await,
|
||||
NodeKind::WritePermFile { .. } => run_write_perm_file(coord, claim, &ctx).await,
|
||||
NodeKind::ApprovalDeploy { .. } => run_approval_deploy(coord, claim).await,
|
||||
NodeKind::SetWanted { up, .. } => run_set_wanted(coord, claim, *up),
|
||||
// Pure grouping container — no work; completing it lets it reach
|
||||
// `Finishing` so its child template nodes start. The DAG's terminal
|
||||
// hook fires (inline, via `run_terminal_hook`) when the container itself
|
||||
// rolls up terminal — not as a scheduled node.
|
||||
NodeKind::Dag { .. } => Ok(NodeOutput::default()),
|
||||
}
|
||||
}
|
||||
|
||||
/// Run a settled DAG's inline terminal hook, dispatched off its rolled-up
|
||||
/// summary — the container-terminal replacement for the old per-DAG hook node.
|
||||
/// Always best-effort: a hook failure is logged inside, never surfaced.
|
||||
pub(crate) async fn run_terminal_hook(coord: &Arc<Coordinator>, terminal: &super::TerminalDag) {
|
||||
match super::terminal_hook(terminal.template, terminal.approval_id) {
|
||||
Some(super::HookKind::ResolveApproval) => {
|
||||
crate::actions::resolve_approval_dag(coord, terminal).await;
|
||||
}
|
||||
Some(super::HookKind::EmitRebuilt) => emit_rebuilt(coord, terminal),
|
||||
Some(super::HookKind::RevertIntent) => revert_intent(coord, terminal).await,
|
||||
None => {}
|
||||
}
|
||||
}
|
||||
|
||||
/// Rebuild / perm-change hook: emit one `Rebuilt` manager event per targeted
|
||||
/// agent — `ok` on `Done`, `!ok` on `Failed`, none on cancel.
|
||||
fn emit_rebuilt(coord: &Arc<Coordinator>, terminal: &super::TerminalDag) {
|
||||
for agent in &terminal.agents {
|
||||
match terminal.state {
|
||||
State::Done => coord.notify_manager(&hive_sh4re::HelperEvent::Rebuilt {
|
||||
agent: agent.clone(),
|
||||
ok: true,
|
||||
note: None,
|
||||
sha: None,
|
||||
tag: None,
|
||||
}),
|
||||
State::Failed => coord.notify_manager(&hive_sh4re::HelperEvent::Rebuilt {
|
||||
agent: agent.clone(),
|
||||
ok: false,
|
||||
note: terminal.error.clone(),
|
||||
sha: None,
|
||||
tag: None,
|
||||
}),
|
||||
_ => {}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// Power-op hook: on a *cancelled* DAG, revert each targeted agent's `wanted`
|
||||
/// intent to its observed state — the operator's cancel means "don't do it", so
|
||||
/// the intent snaps back instead of the flip executing as a surprise side effect
|
||||
/// of some later reconcile. Noop on any non-cancelled outcome.
|
||||
async fn revert_intent(coord: &Arc<Coordinator>, terminal: &super::TerminalDag) {
|
||||
if terminal.state != State::Cancelled {
|
||||
return;
|
||||
}
|
||||
for agent in &terminal.agents {
|
||||
let running = crate::lifecycle::is_running(agent).await;
|
||||
if let Err(e) = coord
|
||||
.power
|
||||
.set(agent, crate::power::Wanted::from_running(running))
|
||||
{
|
||||
tracing::warn!(%agent, error = ?e, "agent_power: cancel revert failed");
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
|
|
@ -336,14 +398,17 @@ async fn run_reconcile(coord: &Arc<Coordinator>, claim: &Claim) -> Result<NodeOu
|
|||
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. The old
|
||||
// `append_node` inherited the emitter's agent implicitly; `append_subgraph`
|
||||
// carries it on the `NodeSpec`, so stamp `claim.agent` explicitly (same
|
||||
// effect, one in-DAG-growth channel instead of two).
|
||||
let sub = |kind| vec![vec![super::templates::node(name, kind, Vec::new())]];
|
||||
// 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),
|
||||
ReconcileAction::Stop => sub(NodeKind::Stop),
|
||||
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()
|
||||
|
|
@ -475,26 +540,30 @@ async fn run_write_perm_file(
|
|||
) -> 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");
|
||||
};
|
||||
ctx.step("writing + committing perm file");
|
||||
// Deploy-window gate: 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).
|
||||
let _window = crate::meta::exclusive().await;
|
||||
match &claim.perm_payload {
|
||||
Some(PermPayload::ToolGroups { groups }) => {
|
||||
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();
|
||||
}
|
||||
Some(PermPayload::Capabilities { caps }) => {
|
||||
PermPayload::Capabilities { caps } => {
|
||||
crate::meta::commit_capabilities(name, caps)
|
||||
.await
|
||||
.with_context(|| format!("commit capabilities for {name}"))?;
|
||||
coord.emit_capabilities_snapshot();
|
||||
}
|
||||
Some(PermPayload::Combined { groups, caps }) => {
|
||||
PermPayload::Combined { groups, caps } => {
|
||||
crate::meta::commit_perms(name, groups.as_deref(), caps.as_deref())
|
||||
.await
|
||||
.with_context(|| format!("commit perms for {name}"))?;
|
||||
|
|
@ -505,10 +574,6 @@ async fn run_write_perm_file(
|
|||
coord.emit_capabilities_snapshot();
|
||||
}
|
||||
}
|
||||
None => anyhow::bail!(
|
||||
"perm_change dag {} for {name} is missing perm_payload",
|
||||
claim.dag_id
|
||||
),
|
||||
}
|
||||
Ok(NodeOutput::default())
|
||||
}
|
||||
|
|
@ -531,71 +596,6 @@ async fn run_approval_deploy(coord: &Arc<Coordinator>, claim: &Claim) -> Result<
|
|||
.map(|()| NodeOutput::default())
|
||||
}
|
||||
|
||||
/// Terminal-roll-up hook, fired exactly once per DAG (node completion
|
||||
/// and cancel paths alike — the queue buffers roll-ups and the
|
||||
/// scheduler drains them). Three concerns:
|
||||
/// - approval DAGs resolve their approval row (except the opaque
|
||||
/// deploy pipeline, which resolves inside its node — unless it was
|
||||
/// cancelled while still queued and the node never ran);
|
||||
/// - non-approval rebuild-shaped DAGs emit exactly one `Rebuilt`
|
||||
/// manager event: ok on `Done`, !ok on `Failed`, none on cancel;
|
||||
/// - a cancelled power-op DAG reverts the `wanted` intent its submit
|
||||
/// wrote: the operator's cancel means "don't do it", so intent
|
||||
/// snaps back to the observed state instead of the flip executing
|
||||
/// as a surprise side effect of some later reconcile.
|
||||
pub(super) async fn on_dag_terminal(coord: &Arc<Coordinator>, terminal: &TerminalDag) {
|
||||
if terminal.state == State::Cancelled
|
||||
&& matches!(
|
||||
terminal.template,
|
||||
Template::Start
|
||||
| Template::Stop
|
||||
| Template::GracefulStop
|
||||
| Template::Restart
|
||||
| Template::GracefulRestart
|
||||
)
|
||||
{
|
||||
// Revert each targeted agent's power intent to its observed state —
|
||||
// the operator's cancel means "don't do it". Single-agent power-op
|
||||
// DAGs have one agent here.
|
||||
for agent in &terminal.agents {
|
||||
let running = crate::lifecycle::is_running(agent).await;
|
||||
if let Err(e) = coord
|
||||
.power
|
||||
.set(agent, crate::power::Wanted::from_running(running))
|
||||
{
|
||||
tracing::warn!(%agent, error = ?e, "agent_power: cancel revert failed");
|
||||
}
|
||||
}
|
||||
}
|
||||
if terminal.approval_id.is_some() {
|
||||
crate::actions::resolve_approval_dag(coord, terminal).await;
|
||||
return;
|
||||
}
|
||||
if matches!(terminal.template, Template::Rebuild | Template::PermChange) {
|
||||
// Rebuild / PermChange are single-agent; emit one `Rebuilt` per
|
||||
// targeted agent (exactly one today).
|
||||
for agent in &terminal.agents {
|
||||
match terminal.state {
|
||||
State::Done => coord.notify_manager(&hive_sh4re::HelperEvent::Rebuilt {
|
||||
agent: agent.clone(),
|
||||
ok: true,
|
||||
note: None,
|
||||
sha: None,
|
||||
tag: None,
|
||||
}),
|
||||
State::Failed => coord.notify_manager(&hive_sh4re::HelperEvent::Rebuilt {
|
||||
agent: agent.clone(),
|
||||
ok: false,
|
||||
note: terminal.error.clone(),
|
||||
sha: None,
|
||||
tag: None,
|
||||
}),
|
||||
_ => {}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/// 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;
|
||||
|
|
|
|||
File diff suppressed because it is too large
Load diff
|
|
@ -11,9 +11,11 @@
|
|||
//! DAG can span agents). See `docs/coordinator.md::Job queue` for the
|
||||
//! full design.
|
||||
|
||||
pub use hive_sh4re::jobs::{DagView, NodeId, NodeView, PermPayload, Source, State, Template};
|
||||
pub use hive_sh4re::jobs::{DagView, NodeId, PermPayload, Source, State, Template};
|
||||
use serde::Serialize;
|
||||
|
||||
use crate::coordinator::TransientKind;
|
||||
|
||||
/// When a dependency edge is considered satisfied.
|
||||
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
|
||||
#[serde(rename_all = "snake_case")]
|
||||
|
|
@ -54,12 +56,12 @@ pub enum NodeKind {
|
|||
/// skipped when the container is already down — it only exists to
|
||||
/// shrink the swap's downtime, which a stopped agent doesn't need
|
||||
/// (the sync + dir prep still run; `Swap` builds inline).
|
||||
Prebuild { relock: bool },
|
||||
Prebuild { agent: String, relock: bool },
|
||||
/// `nixos-container update` profile-swap (requires the container
|
||||
/// stopped). Re-applies nspawn flags + resource limits first —
|
||||
/// rebuild is the reconcile verb. The post-rebuild bookkeeping tail
|
||||
/// lives in the sibling `PostSwap` node.
|
||||
Swap,
|
||||
Swap { agent: String },
|
||||
/// The post-`Swap` bookkeeping tail as a first-class node: rev marker,
|
||||
/// forge + matrix sync, manager kick, container rescan, meta-inputs
|
||||
/// snapshot. Split out of `Swap` for dashboard visibility + retry
|
||||
|
|
@ -69,16 +71,16 @@ pub enum NodeKind {
|
|||
/// recovery still runs. Store/forge/matrix work only — no nix build, so
|
||||
/// build-slot-exempt; the agent lease taken at `Swap` is held across the
|
||||
/// whole chain until `Reconcile` settles, so it's not re-declared here.
|
||||
PostSwap,
|
||||
PostSwap { agent: String },
|
||||
/// First-spawn pre-create provisioning: proposed/applied repos,
|
||||
/// state subvolume, and meta registration (`sync_agents`). Runs
|
||||
/// ahead of `Create` so the `nixos-container create --flake
|
||||
/// meta#<name>` ref resolves. Store/meta-only — no container yet —
|
||||
/// so it's lease- and build-slot-exempt like `Prebuild`.
|
||||
Provision,
|
||||
Provision { agent: String },
|
||||
/// First-spawn `nixos-container create` proper. Assumes the
|
||||
/// upstream `Provision` node already registered the agent in meta.
|
||||
Create,
|
||||
Create { agent: String },
|
||||
/// Meta flake lock bump. `sweep = false`: `meta::lock_update`
|
||||
/// (commit fused, under `META_LOCK`) with the DAG's `inputs`;
|
||||
/// `sweep = true`: `meta::lock_update_hyperhive`, *non-fatal* (a
|
||||
|
|
@ -95,36 +97,37 @@ pub enum NodeKind {
|
|||
/// `Offline` & up, else noop). The mechanical work is not done in
|
||||
/// this node — it fans a child [`NodeKind::Start`] / [`NodeKind::Stop`]
|
||||
/// DAG out at runtime so the sub-step is a first-class DAG node.
|
||||
Reconcile,
|
||||
Reconcile { agent: String },
|
||||
/// Mechanical container start: the start preamble (runtime dir +
|
||||
/// drop-ins), `start_with_fallback`, MCP listener registration, and
|
||||
/// the manager kick. Fanned out by a [`NodeKind::Reconcile`] that
|
||||
/// observed `wanted = Up` and the container down.
|
||||
Start,
|
||||
Start { agent: String },
|
||||
/// Mechanical container stop: `nixos-container` kill, MCP listener
|
||||
/// unregister, and the `Killed` manager notify. Fanned out by a
|
||||
/// [`NodeKind::Reconcile`] that observed `wanted = Offline` and up.
|
||||
Stop,
|
||||
Stop { agent: String },
|
||||
/// Mechanical `nixos-container stop` for the profile swap. Never
|
||||
/// touches `wanted`. Noop if already stopped.
|
||||
StopForUpdate,
|
||||
StopForUpdate { agent: String },
|
||||
/// Set the graceful-stop fence + kick the harness so it runs one
|
||||
/// stop-checkpoint turn.
|
||||
Signal,
|
||||
Signal { agent: String },
|
||||
/// Await the harness clearing the fence, bounded by
|
||||
/// `GRACEFUL_STOP_TIMEOUT`. Resolves ok either way — the
|
||||
/// downstream `Reconcile` performs the actual stop.
|
||||
Drain,
|
||||
Drain { agent: String },
|
||||
/// `set_nspawn_flags` + `set_resource_limits` + daemon-reload.
|
||||
WriteDropin,
|
||||
/// Commit `tool-groups.json` / `capabilities.json` per the DAG's
|
||||
/// `perm_payload` (commit fused under `META_LOCK`).
|
||||
WritePermFile,
|
||||
WriteDropin { agent: String },
|
||||
/// Commit `tool-groups.json` / `capabilities.json` per its `payload`
|
||||
/// (commit fused under `META_LOCK`). The payload rides this node — the only
|
||||
/// consumer — rather than the generic DAG container.
|
||||
WritePermFile { agent: String, payload: PermPayload },
|
||||
/// Opaque approval deploy pipeline (`MergeConfigPr`): the two-phase
|
||||
/// prepare/finalize/abort meta deploy stays inside `actions.rs` in v1 —
|
||||
/// deliberately not
|
||||
/// modeled as scheduler nodes (see the design doc §9).
|
||||
ApprovalDeploy,
|
||||
ApprovalDeploy { agent: String },
|
||||
/// Write the agent's durable power intent (`wanted = Up` when `up`, else
|
||||
/// `Offline`) as a first-class DAG node, at the head of a power-op
|
||||
/// template so the downstream `Reconcile` reads it. Replaces the old
|
||||
|
|
@ -138,7 +141,24 @@ pub enum NodeKind {
|
|||
/// the DAG. (In `stale_start` the lease is thus held across the head
|
||||
/// `Prebuild`, but that's a no-op there — the agent is down, so prebuild
|
||||
/// is skipped.)
|
||||
SetWanted { up: bool },
|
||||
SetWanted { agent: String, up: bool },
|
||||
/// The **DAG container** node: one per submitted DAG, carrying the group's
|
||||
/// domain metadata. Every template node hangs *under* it (its subtree), so
|
||||
/// the container's `NodeId` **is** the DAG id, its rolled-up state **is** the
|
||||
/// DAG state, and it reaching terminal **is** the completion signal that
|
||||
/// fires the DAG's inline hook (approval-resolve / rebuilt-emit /
|
||||
/// intent-revert, dispatched off `template`). Pure grouping — lease- and
|
||||
/// build-slot-exempt; the executor instant-completes it (`Done`) so it
|
||||
/// reaches `Finishing` and its children start.
|
||||
Dag {
|
||||
template: Template,
|
||||
source: Source,
|
||||
reason: String,
|
||||
transient: Option<TransientKind>,
|
||||
approval_id: Option<i64>,
|
||||
inputs: Vec<String>,
|
||||
created_at: i64,
|
||||
},
|
||||
}
|
||||
|
||||
impl NodeKind {
|
||||
|
|
@ -146,21 +166,47 @@ impl NodeKind {
|
|||
pub fn as_str(&self) -> &'static str {
|
||||
match self {
|
||||
NodeKind::Prebuild { .. } => "prebuild",
|
||||
NodeKind::Swap => "swap",
|
||||
NodeKind::PostSwap => "post_swap",
|
||||
NodeKind::Provision => "provision",
|
||||
NodeKind::Create => "create",
|
||||
NodeKind::Swap { .. } => "swap",
|
||||
NodeKind::PostSwap { .. } => "post_swap",
|
||||
NodeKind::Provision { .. } => "provision",
|
||||
NodeKind::Create { .. } => "create",
|
||||
NodeKind::MetaLock { .. } => "meta_lock",
|
||||
NodeKind::Reconcile => "reconcile",
|
||||
NodeKind::Start => "start",
|
||||
NodeKind::Stop => "stop",
|
||||
NodeKind::StopForUpdate => "stop_for_update",
|
||||
NodeKind::Signal => "signal",
|
||||
NodeKind::Drain => "drain",
|
||||
NodeKind::WriteDropin => "write_dropin",
|
||||
NodeKind::WritePermFile => "write_perm_file",
|
||||
NodeKind::ApprovalDeploy => "approval_deploy",
|
||||
NodeKind::Reconcile { .. } => "reconcile",
|
||||
NodeKind::Start { .. } => "start",
|
||||
NodeKind::Stop { .. } => "stop",
|
||||
NodeKind::StopForUpdate { .. } => "stop_for_update",
|
||||
NodeKind::Signal { .. } => "signal",
|
||||
NodeKind::Drain { .. } => "drain",
|
||||
NodeKind::WriteDropin { .. } => "write_dropin",
|
||||
NodeKind::WritePermFile { .. } => "write_perm_file",
|
||||
NodeKind::ApprovalDeploy { .. } => "approval_deploy",
|
||||
NodeKind::SetWanted { .. } => "set_wanted",
|
||||
NodeKind::Dag { .. } => "dag",
|
||||
}
|
||||
}
|
||||
|
||||
/// The agent this node targets, or `""` for agentless kinds
|
||||
/// ([`NodeKind::MetaLock`] on the `hyperhive` pseudo-agent, and the
|
||||
/// [`NodeKind::Dag`] container).
|
||||
#[must_use]
|
||||
pub fn agent(&self) -> &str {
|
||||
match self {
|
||||
NodeKind::Prebuild { agent, .. }
|
||||
| NodeKind::Swap { agent }
|
||||
| NodeKind::PostSwap { agent }
|
||||
| NodeKind::Provision { agent }
|
||||
| NodeKind::Create { agent }
|
||||
| NodeKind::Reconcile { agent }
|
||||
| NodeKind::Start { agent }
|
||||
| NodeKind::Stop { agent }
|
||||
| NodeKind::StopForUpdate { agent }
|
||||
| NodeKind::Signal { agent }
|
||||
| NodeKind::Drain { agent }
|
||||
| NodeKind::WriteDropin { agent }
|
||||
| NodeKind::WritePermFile { agent, .. }
|
||||
| NodeKind::ApprovalDeploy { agent }
|
||||
| NodeKind::SetWanted { agent, .. } => agent,
|
||||
NodeKind::MetaLock { .. } | NodeKind::Dag { .. } => "",
|
||||
}
|
||||
}
|
||||
|
||||
|
|
@ -170,10 +216,10 @@ impl NodeKind {
|
|||
matches!(
|
||||
self,
|
||||
NodeKind::Prebuild { .. }
|
||||
| NodeKind::Swap
|
||||
| NodeKind::Create
|
||||
| NodeKind::Swap { .. }
|
||||
| NodeKind::Create { .. }
|
||||
| NodeKind::MetaLock { .. }
|
||||
| NodeKind::ApprovalDeploy
|
||||
| NodeKind::ApprovalDeploy { .. }
|
||||
)
|
||||
}
|
||||
|
||||
|
|
@ -188,58 +234,44 @@ impl NodeKind {
|
|||
pub fn needs_lease(&self) -> bool {
|
||||
matches!(
|
||||
self,
|
||||
NodeKind::Swap
|
||||
| NodeKind::Create
|
||||
| NodeKind::Reconcile
|
||||
| NodeKind::StopForUpdate
|
||||
| NodeKind::Signal
|
||||
| NodeKind::Drain
|
||||
| NodeKind::WriteDropin
|
||||
| NodeKind::ApprovalDeploy
|
||||
NodeKind::Swap { .. }
|
||||
| NodeKind::Create { .. }
|
||||
| NodeKind::Reconcile { .. }
|
||||
| NodeKind::StopForUpdate { .. }
|
||||
| NodeKind::Signal { .. }
|
||||
| NodeKind::Drain { .. }
|
||||
| NodeKind::WriteDropin { .. }
|
||||
| NodeKind::ApprovalDeploy { .. }
|
||||
| NodeKind::SetWanted { .. }
|
||||
)
|
||||
}
|
||||
}
|
||||
|
||||
/// One schedulable unit inside a DAG.
|
||||
#[derive(Debug, Clone)]
|
||||
pub struct Node {
|
||||
pub id: NodeId,
|
||||
/// The agent this node's work targets. Per-node so a single DAG can
|
||||
/// span agents (e.g. a hive-wide restart); the lifecycle lease is
|
||||
/// acquired against *this* agent (still globally exclusive per agent
|
||||
/// across all DAGs). `"hyperhive"` for meta-level nodes.
|
||||
pub agent: String,
|
||||
pub kind: NodeKind,
|
||||
pub deps: Vec<Dep>,
|
||||
pub state: State,
|
||||
/// Live sub-label while `Running` (kept for parity with the old
|
||||
/// per-entry `step`).
|
||||
pub step: Option<String>,
|
||||
/// Row id of the `build_logs` entry this node opened (`Prebuild` /
|
||||
/// `Swap` / `ApprovalDeploy`), for the dashboard's live-stream link.
|
||||
pub build_log_id: Option<i64>,
|
||||
pub started_at: Option<i64>,
|
||||
pub finished_at: Option<i64>,
|
||||
/// Populated when `state == Failed` (truncated by the queue).
|
||||
pub error: Option<String>,
|
||||
}
|
||||
|
||||
/// Submit-time spec for one node.
|
||||
#[derive(Debug, Clone)]
|
||||
pub struct NodeSpec {
|
||||
/// The agent this node targets (see [`Node::agent`]). Built by the
|
||||
/// `templates.rs` `node` helper, which stamps the template's agent
|
||||
/// onto every node.
|
||||
pub agent: String,
|
||||
/// The node's payload — [`NodeKind`] is the queue's payload type directly,
|
||||
/// and each variant carries the agent it targets (a DAG can span agents;
|
||||
/// the queue derives per-agent leasing from [`NodeKind::agent`]).
|
||||
pub kind: NodeKind,
|
||||
pub deps: Vec<Dep>,
|
||||
/// The **structural parent** axis — the spec-local index of this node's
|
||||
/// group parent, or `None` for a top-level (group-root) node. Independent
|
||||
/// of `deps`: `deps` order execution, `parent` groups nodes into a subtree
|
||||
/// whose resource the whole subtree borrows (the agent lease is owned by a
|
||||
/// group root and re-entered by its descendants for continuity). A child
|
||||
/// runs once its parent reaches `Finishing` (the parent gate), so a child
|
||||
/// never `deps` on its own parent (that would deadlock — dep-scope
|
||||
/// validation rejects it).
|
||||
pub parent: Option<u64>,
|
||||
}
|
||||
|
||||
/// Submit-time spec for a whole DAG. Built by `templates.rs`; validated
|
||||
/// (cycle rejection) by `JobQueue::submit`. No DAG-level `agent` — every
|
||||
/// node carries its own (a DAG can span agents), and the queue derives
|
||||
/// per-agent leasing from [`NodeSpec::agent`].
|
||||
/// per-agent leasing from [`NodeKind::agent`]. Type-specific payloads
|
||||
/// (`PermChange`'s file payload) ride the node that consumes them
|
||||
/// ([`NodeKind::WritePermFile`]), not this generic spec.
|
||||
#[derive(Debug, Clone)]
|
||||
pub struct DagSpec {
|
||||
pub template: Template,
|
||||
|
|
@ -251,147 +283,8 @@ pub struct DagSpec {
|
|||
/// `MetaUpdate`-only: the inputs to bump (also part of the dedup
|
||||
/// key for that template). Display copy lives on the DAG.
|
||||
pub inputs: Vec<String>,
|
||||
/// `PermChange`-only payload.
|
||||
pub perm_payload: Option<PermPayload>,
|
||||
/// Dashboard transient pill (and crash-watch suppression) held for
|
||||
/// the lease window — from lease acquisition to DAG terminal.
|
||||
pub transient: Option<crate::coordinator::TransientKind>,
|
||||
pub nodes: Vec<NodeSpec>,
|
||||
}
|
||||
|
||||
/// A live DAG in the queue. No DAG-level `agent`: agent is per-[`Node`],
|
||||
/// so a DAG can span agents. Per-agent leasing is derived from the
|
||||
/// nodes' agents.
|
||||
#[derive(Debug, Clone)]
|
||||
pub struct Dag {
|
||||
pub id: u64,
|
||||
pub template: Template,
|
||||
pub source: Source,
|
||||
pub reason: String,
|
||||
pub approval_id: Option<i64>,
|
||||
pub inputs: Vec<String>,
|
||||
pub perm_payload: Option<PermPayload>,
|
||||
pub transient: Option<crate::coordinator::TransientKind>,
|
||||
pub created_at: i64,
|
||||
pub nodes: Vec<Node>,
|
||||
/// Terminal roll-up already reported to the scheduler's hooks
|
||||
/// (approval resolution, transient release). Internal bookkeeping,
|
||||
/// never serialized.
|
||||
pub terminal_reported: bool,
|
||||
}
|
||||
|
||||
impl Dag {
|
||||
/// Roll-up state: `Failed` if any node failed; else `Running` if
|
||||
/// any running; else `Queued` if any queued; else `Cancelled` if
|
||||
/// any cancelled; else `Done`.
|
||||
pub fn rollup(&self) -> State {
|
||||
let mut any_cancelled = false;
|
||||
let mut any_queued = false;
|
||||
let mut any_running = false;
|
||||
for n in &self.nodes {
|
||||
match n.state {
|
||||
State::Failed => return State::Failed,
|
||||
State::Running => any_running = true,
|
||||
State::Queued => any_queued = true,
|
||||
State::Cancelled => any_cancelled = true,
|
||||
State::Done => {}
|
||||
}
|
||||
}
|
||||
if any_running {
|
||||
State::Running
|
||||
} else if any_queued {
|
||||
State::Queued
|
||||
} else if any_cancelled {
|
||||
State::Cancelled
|
||||
} else {
|
||||
State::Done
|
||||
}
|
||||
}
|
||||
|
||||
/// True when every node is terminal.
|
||||
pub fn is_terminal(&self) -> bool {
|
||||
self.nodes.iter().all(|n| n.state.is_terminal())
|
||||
}
|
||||
|
||||
/// True when no live (non-terminal) node of this DAG still targets
|
||||
/// `agent` — i.e. that agent's subgraph within the DAG has settled.
|
||||
/// Used to release an agent's lifecycle lease the moment its own
|
||||
/// work is done, rather than waiting for the whole DAG to terminate.
|
||||
/// Vacuously true for an agent the DAG has no node for; callers gate
|
||||
/// on actually holding that agent's lease first.
|
||||
pub fn agent_subgraph_terminal(&self, agent: &str) -> bool {
|
||||
self.nodes
|
||||
.iter()
|
||||
.filter(|n| n.agent == agent)
|
||||
.all(|n| n.state.is_terminal())
|
||||
}
|
||||
|
||||
/// First failed node's error, for the roll-up `error` field.
|
||||
pub fn first_error(&self) -> Option<&str> {
|
||||
self.nodes
|
||||
.iter()
|
||||
.find(|n| n.state == State::Failed)
|
||||
.and_then(|n| n.error.as_deref())
|
||||
}
|
||||
|
||||
pub fn node(&self, id: NodeId) -> Option<&Node> {
|
||||
self.nodes.iter().find(|n| n.id == id)
|
||||
}
|
||||
|
||||
pub fn node_mut(&mut self, id: NodeId) -> Option<&mut Node> {
|
||||
self.nodes.iter_mut().find(|n| n.id == id)
|
||||
}
|
||||
|
||||
/// Distinct agents this DAG's nodes target, in first-seen order.
|
||||
/// Used for per-agent lease release and the terminal cancel-revert —
|
||||
/// a single-agent DAG yields one, a multi-agent DAG yields several.
|
||||
pub fn agents(&self) -> Vec<String> {
|
||||
let mut seen: Vec<String> = Vec::new();
|
||||
for n in &self.nodes {
|
||||
if !seen.iter().any(|a| a == &n.agent) {
|
||||
seen.push(n.agent.clone());
|
||||
}
|
||||
}
|
||||
seen
|
||||
}
|
||||
}
|
||||
|
||||
impl Dag {
|
||||
pub fn view(&self) -> DagView {
|
||||
let started_at = self.nodes.iter().filter_map(|n| n.started_at).min();
|
||||
let finished_at = if self.is_terminal() {
|
||||
self.nodes.iter().filter_map(|n| n.finished_at).max()
|
||||
} else {
|
||||
None
|
||||
};
|
||||
DagView {
|
||||
id: self.id,
|
||||
kind: self.template,
|
||||
state: self.rollup(),
|
||||
source: self.source,
|
||||
reason: self.reason.clone(),
|
||||
enqueued_at: self.created_at,
|
||||
started_at,
|
||||
finished_at,
|
||||
inputs: self.inputs.clone(),
|
||||
approval_id: self.approval_id,
|
||||
perm_payload: self.perm_payload.clone(),
|
||||
nodes: self
|
||||
.nodes
|
||||
.iter()
|
||||
.map(|n| NodeView {
|
||||
id: n.id,
|
||||
agent: n.agent.clone(),
|
||||
kind: n.kind.as_str().to_owned(),
|
||||
deps: n.deps.iter().map(|d| d.on).collect(),
|
||||
state: n.state,
|
||||
step: n.step.clone(),
|
||||
build_log_id: n.build_log_id,
|
||||
started_at: n.started_at,
|
||||
finished_at: n.finished_at,
|
||||
error: n.error.clone(),
|
||||
})
|
||||
.collect(),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
|
|||
51
hive-c0re/src/job_queue/resource.rs
Normal file
51
hive-c0re/src/job_queue/resource.rs
Normal file
|
|
@ -0,0 +1,51 @@
|
|||
//! The concrete resource type the rebuild queue schedules over — the bridge
|
||||
//! from hive-c0re's [`NodeKind`] onto the domain-agnostic `hive-jobq` crate.
|
||||
//! `hive-jobq` is generic over a resource type `R: Clone + Eq + Hash` and a node
|
||||
//! payload `N`; here `R` is [`Resource`] and `N` is [`NodeKind`] directly (each
|
||||
//! variant carries the agent it targets).
|
||||
|
||||
use hive_jobq::Dep;
|
||||
|
||||
use super::model::NodeKind;
|
||||
|
||||
/// The two resource classes the queue gates concurrency on, as the crate's
|
||||
/// generic resource type `R`.
|
||||
#[derive(Debug, Clone, PartialEq, Eq, Hash)]
|
||||
pub enum Resource {
|
||||
/// One of the `buildSlots` permits, held by a nix-heavy node for its
|
||||
/// duration. Capacity is `services.hyperhive.c0re.buildSlots` (default 1),
|
||||
/// set on the [`hive_jobq::resources::ResourceTable`] at construction.
|
||||
BuildSlot,
|
||||
/// The per-agent lifecycle lease — globally exclusive per agent across all
|
||||
/// DAGs (unconfigured, so the crate's default capacity 1 applies). Held by
|
||||
/// a DAG's first container-affecting node for that agent and re-entered by
|
||||
/// the rest of that agent's subtree via the crate's recursive lock, so two
|
||||
/// DAGs never interleave container ops on one agent.
|
||||
Agent(String),
|
||||
}
|
||||
|
||||
impl NodeKind {
|
||||
/// The [`Dep::Resource`] edges this node must acquire to run, derived from
|
||||
/// its kind + agent: a build slot for nix-heavy kinds
|
||||
/// ([`NodeKind::needs_build_slot`]) and the agent lease for
|
||||
/// container-affecting kinds ([`NodeKind::needs_lease`]). Lease-exempt
|
||||
/// container ops (`Start` / `Stop`, fanned out by a lease-holding
|
||||
/// `Reconcile`) hold no lease of their own — they re-enter the ancestor's
|
||||
/// `Agent` lock through the crate's recursive re-entrancy.
|
||||
pub fn resource_deps(&self) -> Vec<Dep<Resource>> {
|
||||
let mut deps = Vec::new();
|
||||
if self.needs_build_slot() {
|
||||
deps.push(Dep::Resource {
|
||||
name: Resource::BuildSlot,
|
||||
count: 1,
|
||||
});
|
||||
}
|
||||
if self.needs_lease() {
|
||||
deps.push(Dep::Resource {
|
||||
name: Resource::Agent(self.agent().to_owned()),
|
||||
count: 1,
|
||||
});
|
||||
}
|
||||
deps
|
||||
}
|
||||
}
|
||||
|
|
@ -1,18 +1,25 @@
|
|||
//! The single scheduler task that drives all DAGs: claim every ready
|
||||
//! node (as many as the build slots / leases allow), spawn one
|
||||
//! executor task per claim, and on any completion re-evaluate.
|
||||
//! Concurrency comes from the build-slot count, not multiple workers.
|
||||
//! The single scheduler task that drives all DAGs: claim every ready node (as
|
||||
//! many as the build slots / leases allow), spawn one executor task per claim,
|
||||
//! and on any completion re-evaluate. Concurrency comes from the build-slot
|
||||
//! count, not multiple workers.
|
||||
//!
|
||||
//! Also owns the per-DAG transient guard (dashboard pill + crash-watch
|
||||
//! suppression) that the sync queue core can't hold itself — created when a
|
||||
//! DAG acquires its agent lease, dropped when the DAG settles terminal.
|
||||
//! Owns the per-DAG transient guard (dashboard pill + crash-watch suppression)
|
||||
//! that the sync queue core can't hold itself. The guard set is *reconciled*
|
||||
//! from live lease ownership ([`super::JobQueue::held_transients`]) each loop:
|
||||
//! a `(dag, agent)` pill exists for exactly as long as that agent's lease is
|
||||
//! held, so it appears when the agent's owner node starts and disappears when
|
||||
//! its subgraph settles — one pill per agent a DAG touches.
|
||||
//!
|
||||
//! In-DAG growth (a `MetaLock` growing rebuild subgraphs after the lock
|
||||
//! bump, a `Reconcile` fanning its `Start`/`Stop`) flows through
|
||||
//! `NodeOutput.append_subgraph`, applied before the emitting node
|
||||
//! completes — see `handle_completion`.
|
||||
//! Per-DAG terminal work (approval resolution, `Rebuilt`, cancelled-power-op
|
||||
//! intent revert) is not drained here: it runs as the DAG's focused terminal
|
||||
//! node (`ResolveApproval` / `EmitRebuilt` / `RevertIntent`), dispatched through
|
||||
//! `exec::run_node` like any other node once the DAG settles.
|
||||
//!
|
||||
//! In-DAG growth (a `MetaLock` growing rebuild subgraphs, a `Reconcile` fanning
|
||||
//! its `Start`/`Stop`) flows through `NodeOutput.append_subgraph`, applied
|
||||
//! before the emitting node completes — see `handle_completion`.
|
||||
|
||||
use std::collections::HashMap;
|
||||
use std::collections::{HashMap, HashSet};
|
||||
use std::sync::Arc;
|
||||
|
||||
use super::Claim;
|
||||
|
|
@ -26,38 +33,24 @@ struct NodeDone {
|
|||
|
||||
/// Scheduler loop. Spawned once at hive-c0re startup from `main.rs`.
|
||||
///
|
||||
/// Shutdown semantics: subscribes to `coord.shutdown_rx()`. On a true
|
||||
/// signal the loop exits immediately; already-running node tasks ride
|
||||
/// the runtime down with the process, and pending `Queued` DAGs are
|
||||
/// dropped — desired state is re-derived on next boot (boot sweep +
|
||||
/// reconcile), so the in-memory queue is deliberately not durable.
|
||||
/// Shutdown semantics: subscribes to `coord.shutdown_rx()`. On a true signal
|
||||
/// the loop exits immediately; already-running node tasks ride the runtime down
|
||||
/// with the process, and pending `Queued` DAGs are dropped — desired state is
|
||||
/// re-derived on next boot (boot sweep + reconcile), so the in-memory queue is
|
||||
/// deliberately not durable.
|
||||
pub async fn run_worker(coord: Arc<Coordinator>) {
|
||||
let mut shutdown = coord.shutdown_rx();
|
||||
let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel::<NodeDone>();
|
||||
// (DAG id, agent) → transient guard held for that agent's lease
|
||||
// window. Keyed per-agent so a multi-agent DAG shows one transient
|
||||
// pill per agent it touches.
|
||||
// (DAG id, agent) → transient guard held for that agent's lease window.
|
||||
let mut transients: HashMap<(u64, String), crate::coordinator::TransientGuard> = HashMap::new();
|
||||
loop {
|
||||
// Terminal roll-ups can appear without a node completion —
|
||||
// the cancel surfaces settle DAGs directly and wake this loop
|
||||
// via notify — so drain on every iteration, not just inside
|
||||
// handle_completion.
|
||||
process_terminals(&coord, &mut transients).await;
|
||||
reconcile_transients(&coord, &mut transients);
|
||||
let claims = coord.job_queue.claim_ready();
|
||||
if !claims.is_empty() {
|
||||
for claim in claims {
|
||||
if claim.lease_acquired
|
||||
&& let Some(kind) = claim.transient
|
||||
{
|
||||
transients.insert(
|
||||
(claim.dag_id, claim.agent.clone()),
|
||||
coord.transient_guard(&claim.agent, kind),
|
||||
);
|
||||
}
|
||||
tracing::info!(
|
||||
dag = claim.dag_id,
|
||||
node = claim.node_id,
|
||||
node = claim.node_id.get(),
|
||||
kind = claim.kind.as_str(),
|
||||
agent = %claim.agent,
|
||||
template = claim.template.as_str(),
|
||||
|
|
@ -71,6 +64,8 @@ pub async fn run_worker(coord: Arc<Coordinator>) {
|
|||
let _ = tx.send(NodeDone { claim, result });
|
||||
});
|
||||
}
|
||||
// Newly-started owner nodes now hold their leases — surface the pills.
|
||||
reconcile_transients(&coord, &mut transients);
|
||||
coord.emit_rebuild_queue_snapshot();
|
||||
continue;
|
||||
}
|
||||
|
|
@ -83,82 +78,85 @@ pub async fn run_worker(coord: Arc<Coordinator>) {
|
|||
}
|
||||
}
|
||||
Some(done) = rx.recv() => {
|
||||
handle_completion(&coord, &mut transients, done).await;
|
||||
handle_completion(&coord, done);
|
||||
}
|
||||
() = coord.job_queue.notify.notified() => {}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
async fn handle_completion(
|
||||
coord: &Arc<Coordinator>,
|
||||
transients: &mut HashMap<(u64, String), crate::coordinator::TransientGuard>,
|
||||
done: NodeDone,
|
||||
) {
|
||||
fn handle_completion(coord: &Arc<Coordinator>, done: NodeDone) {
|
||||
let NodeDone { claim, result } = done;
|
||||
match result {
|
||||
Ok(output) => {
|
||||
tracing::info!(
|
||||
dag = claim.dag_id,
|
||||
node = claim.node_id,
|
||||
node = claim.node_id.get(),
|
||||
"job_queue: node done"
|
||||
);
|
||||
// Append any in-DAG subgraphs BEFORE completing this node, so
|
||||
// completing it doesn't roll the DAG terminal while the appended
|
||||
// work is still pending — that keeps the lease-window transient
|
||||
// held across it. Each subgraph is independent, rooted on this
|
||||
// node (`AfterOk`), so it becomes ready the instant this one
|
||||
// settles `Done` just below. Covers both the multi-node case (a
|
||||
// `MetaLock` growing per-agent rebuild subgraphs) and the
|
||||
// single-node case (a `Reconcile` planner's `Start` / `Stop`).
|
||||
for subgraph in output.append_subgraph {
|
||||
// work is still pending. Each subgraph roots on this node
|
||||
// (`AfterOk`), so it becomes ready the instant this one settles
|
||||
// `Done` just below — covers both the multi-node case (a `MetaLock`
|
||||
// growing per-agent rebuild subgraphs) and the single-node case (a
|
||||
// `Reconcile` planner's `Start` / `Stop`).
|
||||
for subgraph in &output.append_subgraph {
|
||||
coord
|
||||
.job_queue
|
||||
.append_subgraph(claim.dag_id, subgraph, claim.node_id);
|
||||
}
|
||||
coord
|
||||
let terminal = coord
|
||||
.job_queue
|
||||
.complete_node(claim.dag_id, claim.node_id, Ok(()));
|
||||
fire_terminal_hook(coord, terminal);
|
||||
}
|
||||
Err(e) => {
|
||||
let msg = format!("{e:#}");
|
||||
tracing::warn!(
|
||||
dag = claim.dag_id,
|
||||
node = claim.node_id,
|
||||
node = claim.node_id.get(),
|
||||
kind = claim.kind.as_str(),
|
||||
agent = %claim.agent,
|
||||
error = %msg,
|
||||
"job_queue: node failed"
|
||||
);
|
||||
coord
|
||||
let terminal = coord
|
||||
.job_queue
|
||||
.complete_node(claim.dag_id, claim.node_id, Err(msg));
|
||||
fire_terminal_hook(coord, terminal);
|
||||
}
|
||||
}
|
||||
process_terminals(coord, transients).await;
|
||||
// The next loop iteration re-reconciles the transient pills against the
|
||||
// post-completion lease state (a settled subgraph drops its pill).
|
||||
coord.emit_rebuild_queue_snapshot();
|
||||
}
|
||||
|
||||
/// Drain per-agent lease releases and buffered terminal roll-ups.
|
||||
///
|
||||
/// Per-agent first: an agent's subgraph within a DAG went terminal (its
|
||||
/// lease was freed in `settle`), so drop that agent's `(dag, agent)`
|
||||
/// transient pill now — ahead of whole-DAG terminal for a multi-agent
|
||||
/// DAG. Then the whole-DAG terminals: drop any remaining transient the
|
||||
/// DAG still held and run the terminal hook (approval resolution,
|
||||
/// `Rebuilt` events, cancelled-power-op intent revert).
|
||||
async fn process_terminals(
|
||||
/// Fire a settled DAG's inline terminal hook (approval-resolve / rebuilt-emit /
|
||||
/// intent-revert) off the container-terminal summary `complete_node` returned —
|
||||
/// spawned so the async hook doesn't block the scheduler loop.
|
||||
fn fire_terminal_hook(coord: &Arc<Coordinator>, terminal: Option<super::TerminalDag>) {
|
||||
let Some(terminal) = terminal else {
|
||||
return;
|
||||
};
|
||||
let coord = Arc::clone(coord);
|
||||
tokio::spawn(async move {
|
||||
exec::run_terminal_hook(&coord, &terminal).await;
|
||||
});
|
||||
}
|
||||
|
||||
/// Reconcile the transient-guard set against live lease ownership: drop pills
|
||||
/// whose lease is no longer held, create one for each newly-held `(dag, agent)`.
|
||||
fn reconcile_transients(
|
||||
coord: &Arc<Coordinator>,
|
||||
transients: &mut HashMap<(u64, String), crate::coordinator::TransientGuard>,
|
||||
) {
|
||||
for rel in coord.job_queue.drain_agent_releases() {
|
||||
transients.remove(&(rel.dag_id, rel.agent));
|
||||
}
|
||||
for terminal in coord.job_queue.drain_terminal() {
|
||||
// Drop any per-agent transient guard the DAG still held (the
|
||||
// per-agent pass above already dropped the ones whose subgraphs
|
||||
// settled early).
|
||||
transients.retain(|(dag_id, _), _| *dag_id != terminal.dag_id);
|
||||
exec::on_dag_terminal(coord, &terminal).await;
|
||||
let held = coord.job_queue.held_transients();
|
||||
let keys: HashSet<(u64, String)> = held.iter().map(|(d, a, _)| (*d, a.clone())).collect();
|
||||
transients.retain(|k, _| keys.contains(k));
|
||||
for (dag_id, agent, kind) in held {
|
||||
transients
|
||||
.entry((dag_id, agent.clone()))
|
||||
.or_insert_with(|| coord.transient_guard(&agent, kind));
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -25,7 +25,7 @@
|
|||
use std::sync::Arc;
|
||||
|
||||
use super::model::{DagSpec, Dep, NodeKind, NodeSpec, Template};
|
||||
use super::templates::{after_ok, node, rebuild_nodes};
|
||||
use super::templates::{after_ok, child, node, rebuild_nodes};
|
||||
use super::{Source, templates};
|
||||
use crate::coordinator::{Coordinator, TransientKind};
|
||||
use crate::lifecycle;
|
||||
|
|
@ -60,13 +60,23 @@ pub fn rebuild(coord: &Arc<Coordinator>, agent: &str, source: Source, reason: St
|
|||
/// stays even for a down agent so a race-up between the state read and exec
|
||||
/// is still stopped in-DAG.
|
||||
fn stop_chain(agent: &str, graceful: bool, running: bool) -> Vec<NodeSpec> {
|
||||
let mut n = vec![node(agent, NodeKind::SetWanted { up: false }, Vec::new())];
|
||||
// `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).
|
||||
let a = || agent.to_owned();
|
||||
let mut n = vec![node(
|
||||
NodeKind::SetWanted {
|
||||
agent: a(),
|
||||
up: false,
|
||||
},
|
||||
Vec::new(),
|
||||
)];
|
||||
if graceful && running {
|
||||
n.push(node(agent, NodeKind::Signal, after_ok(0)));
|
||||
n.push(node(agent, NodeKind::Drain, after_ok(1)));
|
||||
n.push(node(agent, NodeKind::Reconcile, after_ok(2)));
|
||||
n.push(child(0, NodeKind::Signal { agent: a() }, Vec::new()));
|
||||
n.push(child(0, NodeKind::Drain { agent: a() }, after_ok(1)));
|
||||
n.push(child(0, NodeKind::Reconcile { agent: a() }, after_ok(2)));
|
||||
} else {
|
||||
n.push(node(agent, NodeKind::Reconcile, after_ok(0)));
|
||||
n.push(child(0, NodeKind::Reconcile { agent: a() }, Vec::new()));
|
||||
}
|
||||
n
|
||||
}
|
||||
|
|
@ -76,13 +86,26 @@ fn stop_chain(agent: &str, graceful: bool, running: bool) -> Vec<NodeSpec> {
|
|||
/// current derivations), otherwise a plain `Reconcile` (which starts a down
|
||||
/// agent and noops an already-running one).
|
||||
fn start_chain(agent: &str, running: bool, stale: bool) -> Vec<NodeSpec> {
|
||||
let mut n = vec![node(agent, NodeKind::SetWanted { up: true }, Vec::new())];
|
||||
let mut n = vec![node(
|
||||
NodeKind::SetWanted {
|
||||
agent: agent.to_owned(),
|
||||
up: true,
|
||||
},
|
||||
Vec::new(),
|
||||
)];
|
||||
if !running && stale {
|
||||
// Rebuild subgraph rooted at the SetWanted head (base = 1, so
|
||||
// `Prebuild` deps `after_ok(0)` = the head).
|
||||
// Rebuild subtree after the SetWanted head (base = 1, so the rebuild's
|
||||
// `Prebuild` root deps `after_ok(0)` = the head). `Prebuild` +
|
||||
// `Reconcile` are their own group roots (top-level, per `rebuild_nodes`).
|
||||
n.extend(rebuild_nodes(agent, true, 1));
|
||||
} else {
|
||||
n.push(node(agent, NodeKind::Reconcile, after_ok(0)));
|
||||
n.push(child(
|
||||
0,
|
||||
NodeKind::Reconcile {
|
||||
agent: agent.to_owned(),
|
||||
},
|
||||
Vec::new(),
|
||||
));
|
||||
}
|
||||
n
|
||||
}
|
||||
|
|
@ -97,22 +120,39 @@ fn start_chain(agent: &str, running: bool, stale: bool) -> Vec<NodeSpec> {
|
|||
/// converges to intent — a stopped (`wanted = Off`) agent stays stopped,
|
||||
/// a crashed (`wanted = Up`) agent comes back up.
|
||||
fn restart_chain(agent: &str, graceful: bool, running: bool) -> Vec<NodeSpec> {
|
||||
let a = || agent.to_owned();
|
||||
if !running {
|
||||
// Nothing to bounce — a lone Reconcile converges to intent.
|
||||
return vec![node(agent, NodeKind::Reconcile, Vec::new())];
|
||||
return vec![node(NodeKind::Reconcile { agent: a() }, Vec::new())];
|
||||
}
|
||||
// Running: mechanical stop then Reconcile. The first stop node is the
|
||||
// subgraph root (no SetWanted head) and acquires the agent lease.
|
||||
let mut n = Vec::new();
|
||||
if graceful {
|
||||
n.push(node(agent, NodeKind::Signal, Vec::new()));
|
||||
n.push(node(agent, NodeKind::Drain, after_ok(0)));
|
||||
n.push(node(agent, NodeKind::StopForUpdate, after_ok(1)));
|
||||
// Running: mechanical stop then Reconcile. The first stop node is the group
|
||||
// ROOT (no SetWanted head) and owns the agent lease; the rest are its
|
||||
// children (borrow the lease, dep-ordered), so the bounce holds one
|
||||
// continuous lease and `Reconcile` cancel-cascades if a stop step fails.
|
||||
let mut n = vec![if graceful {
|
||||
node(NodeKind::Signal { agent: a() }, Vec::new())
|
||||
} else {
|
||||
n.push(node(agent, NodeKind::StopForUpdate, Vec::new()));
|
||||
node(NodeKind::StopForUpdate { agent: a() }, Vec::new())
|
||||
}];
|
||||
if graceful {
|
||||
n.push(child(0, NodeKind::Drain { agent: a() }, Vec::new()));
|
||||
n.push(child(
|
||||
0,
|
||||
NodeKind::StopForUpdate { agent: a() },
|
||||
after_ok(1),
|
||||
));
|
||||
}
|
||||
let stop_idx = u32::try_from(n.len() - 1).unwrap_or(0);
|
||||
n.push(node(agent, NodeKind::Reconcile, after_ok(stop_idx)));
|
||||
// `Reconcile` gates on the last mechanical step. When the only step is the
|
||||
// root itself (non-graceful, `StopForUpdate` == index 0), the parent gate
|
||||
// already orders `Reconcile` after it — a child must NOT dep on its own
|
||||
// parent (dep-scope). So the sibling dep is added only for a graceful
|
||||
// bounce, where the last step is a sibling child.
|
||||
let deps = if n.len() > 1 {
|
||||
after_ok(u64::try_from(n.len() - 1).unwrap_or(0))
|
||||
} else {
|
||||
Vec::new()
|
||||
};
|
||||
n.push(child(0, NodeKind::Reconcile { agent: a() }, deps));
|
||||
n
|
||||
}
|
||||
|
||||
|
|
@ -124,7 +164,7 @@ fn restart_chain(agent: &str, graceful: bool, running: bool) -> Vec<NodeSpec> {
|
|||
fn concat_subgraphs(chains: Vec<Vec<NodeSpec>>) -> Vec<NodeSpec> {
|
||||
let mut out: Vec<NodeSpec> = Vec::new();
|
||||
for chain in chains {
|
||||
let base = u32::try_from(out.len()).unwrap_or(u32::MAX);
|
||||
let base = u64::try_from(out.len()).unwrap_or(u64::MAX);
|
||||
for spec in chain {
|
||||
let deps = spec
|
||||
.deps
|
||||
|
|
@ -135,9 +175,12 @@ fn concat_subgraphs(chains: Vec<Vec<NodeSpec>>) -> Vec<NodeSpec> {
|
|||
})
|
||||
.collect();
|
||||
out.push(NodeSpec {
|
||||
agent: spec.agent,
|
||||
kind: spec.kind,
|
||||
deps,
|
||||
// Rebase the structural parent by the same offset (a subgraph
|
||||
// root keeps `parent = None`, so the per-agent groups stay
|
||||
// independent + concurrent).
|
||||
parent: spec.parent.map(|p| base + p),
|
||||
});
|
||||
}
|
||||
}
|
||||
|
|
@ -158,7 +201,6 @@ fn power_dag(
|
|||
reason,
|
||||
approval_id: None,
|
||||
inputs: Vec::new(),
|
||||
perm_payload: None,
|
||||
transient: Some(transient),
|
||||
nodes,
|
||||
}
|
||||
|
|
|
|||
|
|
@ -35,52 +35,78 @@ 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: u32) -> Vec<Dep> {
|
||||
pub(crate) fn after_ok(on: u64) -> Vec<Dep> {
|
||||
vec![Dep {
|
||||
on,
|
||||
when: DepWhen::AfterOk,
|
||||
}]
|
||||
}
|
||||
|
||||
/// Build one node targeting `agent`. The single place a node's agent is
|
||||
/// stamped. Shared with `submit.rs`'s dynamic power-op builders.
|
||||
pub(crate) fn node(agent: &str, kind: NodeKind, deps: Vec<Dep>) -> NodeSpec {
|
||||
/// 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 {
|
||||
agent: agent.to_owned(),
|
||||
kind,
|
||||
deps,
|
||||
parent: None,
|
||||
}
|
||||
}
|
||||
|
||||
/// The rebuild node chain. `PostSwap` carries the swap's Ok-only
|
||||
/// bookkeeping tail (rev marker, forge/matrix sync, kick, rescan) and deps
|
||||
/// `Swap` with `AfterOk`. `Reconcile` then deps on `PostSwap` with
|
||||
/// `AfterAny`: it must run even when the swap failed, so a previously-up
|
||||
/// agent comes back on its old config (today's recovery-start). On swap
|
||||
/// failure the `AfterOk` `PostSwap` is cancel-cascaded to a terminal state,
|
||||
/// which still satisfies `Reconcile`'s `AfterAny` edge — the only `AfterAny`
|
||||
/// edge in v1. Pointing `Reconcile` at `PostSwap` (not `Swap`) also
|
||||
/// serializes the tail ahead of the reconcile, so there's no double
|
||||
/// rescan/kick race.
|
||||
pub(crate) fn rebuild_nodes(agent: &str, relock: bool, base: u32) -> Vec<NodeSpec> {
|
||||
/// 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),
|
||||
}
|
||||
}
|
||||
|
||||
/// The rebuild node subtree (nested, two group roots). `base` is the spec index
|
||||
/// of the first node (`Prebuild`). Structure:
|
||||
/// - `Prebuild` (base+0, **root**): owns the build slot for the whole subtree.
|
||||
/// Lease-exempt — the nix build overlaps other DAGs on the same agent.
|
||||
/// - `StopForUpdate` (base+1, child of `Prebuild`): owns the agent lease. Runs
|
||||
/// once `Prebuild` reaches `Finishing` (parent gate).
|
||||
/// - `Swap` (base+2, child of `StopForUpdate`): borrows the agent lease from its
|
||||
/// parent and the build slot from grand-ancestor `Prebuild` — both continuous.
|
||||
/// - `PostSwap` (base+3, child of `StopForUpdate`): the swap's Ok-only
|
||||
/// bookkeeping tail (rev marker, forge/matrix sync, kick, rescan), `AfterOk`
|
||||
/// its sibling `Swap`.
|
||||
/// - `Reconcile` (base+4, **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). It takes a fresh lease; the tiny gap is
|
||||
/// harmless — `Reconcile` converges to the persisted `wanted` idempotently.
|
||||
pub(crate) fn rebuild_nodes(agent: &str, relock: bool, base: u64) -> Vec<NodeSpec> {
|
||||
let a = || agent.to_owned();
|
||||
vec![
|
||||
node(
|
||||
agent,
|
||||
NodeKind::Prebuild { relock },
|
||||
NodeKind::Prebuild { agent: a(), relock },
|
||||
if base == 0 {
|
||||
Vec::new()
|
||||
} else {
|
||||
after_ok(base - 1)
|
||||
},
|
||||
),
|
||||
node(agent, NodeKind::StopForUpdate, after_ok(base)),
|
||||
node(agent, NodeKind::Swap, after_ok(base + 1)),
|
||||
node(agent, NodeKind::PostSwap, after_ok(base + 2)),
|
||||
child(base, NodeKind::StopForUpdate { agent: a() }, Vec::new()),
|
||||
child(base + 1, NodeKind::Swap { agent: a() }, Vec::new()),
|
||||
child(
|
||||
base + 1,
|
||||
NodeKind::PostSwap { agent: a() },
|
||||
after_ok(base + 2),
|
||||
),
|
||||
node(
|
||||
agent,
|
||||
NodeKind::Reconcile,
|
||||
NodeKind::Reconcile { agent: a() },
|
||||
vec![Dep {
|
||||
on: base + 3,
|
||||
on: base,
|
||||
when: DepWhen::AfterAny,
|
||||
}],
|
||||
),
|
||||
|
|
@ -99,7 +125,6 @@ pub fn rebuild(agent: &str, source: Source, reason: String, relock: bool) -> Dag
|
|||
reason,
|
||||
approval_id: None,
|
||||
inputs: Vec::new(),
|
||||
perm_payload: None,
|
||||
transient: Some(TransientKind::Rebuilding),
|
||||
nodes: rebuild_nodes(agent, relock, 0),
|
||||
}
|
||||
|
|
@ -115,9 +140,13 @@ pub fn approval_deploy(agent: &str, approval_id: i64, reason: String) -> DagSpec
|
|||
reason,
|
||||
approval_id: Some(approval_id),
|
||||
inputs: Vec::new(),
|
||||
perm_payload: None,
|
||||
transient: Some(TransientKind::Rebuilding),
|
||||
nodes: vec![node(agent, NodeKind::ApprovalDeploy, Vec::new())],
|
||||
nodes: vec![node(
|
||||
NodeKind::ApprovalDeploy {
|
||||
agent: agent.to_owned(),
|
||||
},
|
||||
Vec::new(),
|
||||
)],
|
||||
}
|
||||
}
|
||||
|
||||
|
|
@ -140,16 +169,25 @@ pub fn reconcile_only(
|
|||
reason,
|
||||
approval_id: None,
|
||||
inputs: Vec::new(),
|
||||
perm_payload: None,
|
||||
transient,
|
||||
nodes: vec![node(agent, NodeKind::Reconcile, Vec::new())],
|
||||
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).
|
||||
/// 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).
|
||||
pub fn spawn(agent: &str, approval_id: i64, reason: String) -> DagSpec {
|
||||
DagSpec {
|
||||
template: Template::Spawn,
|
||||
|
|
@ -157,14 +195,16 @@ pub fn spawn(agent: &str, approval_id: i64, reason: String) -> DagSpec {
|
|||
reason,
|
||||
approval_id: Some(approval_id),
|
||||
inputs: Vec::new(),
|
||||
perm_payload: None,
|
||||
transient: Some(TransientKind::Spawning),
|
||||
nodes: vec![
|
||||
node(agent, NodeKind::Provision, Vec::new()),
|
||||
node(agent, NodeKind::Create, after_ok(0)),
|
||||
node(agent, NodeKind::WriteDropin, after_ok(1)),
|
||||
node(agent, NodeKind::Reconcile, after_ok(2)),
|
||||
],
|
||||
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)),
|
||||
]
|
||||
},
|
||||
}
|
||||
}
|
||||
|
||||
|
|
@ -172,7 +212,13 @@ pub fn spawn(agent: &str, approval_id: i64, reason: String) -> DagSpec {
|
|||
/// the updated `HIVE_TOOL_GROUPS` / `HIVE_CAPABILITIES` env var takes
|
||||
/// effect in the container.
|
||||
pub fn perm_change(agent: &str, source: Source, reason: String, payload: PermPayload) -> DagSpec {
|
||||
let mut nodes = vec![node(agent, NodeKind::WritePermFile, Vec::new())];
|
||||
let mut nodes = vec![node(
|
||||
NodeKind::WritePermFile {
|
||||
agent: agent.to_owned(),
|
||||
payload,
|
||||
},
|
||||
Vec::new(),
|
||||
)];
|
||||
nodes.extend(rebuild_nodes(agent, true, 1));
|
||||
DagSpec {
|
||||
template: Template::PermChange,
|
||||
|
|
@ -180,7 +226,6 @@ pub fn perm_change(agent: &str, source: Source, reason: String, payload: PermPay
|
|||
reason,
|
||||
approval_id: None,
|
||||
inputs: Vec::new(),
|
||||
perm_payload: Some(payload),
|
||||
transient: Some(TransientKind::Rebuilding),
|
||||
nodes,
|
||||
}
|
||||
|
|
@ -208,10 +253,8 @@ pub fn meta_update(
|
|||
reason,
|
||||
approval_id,
|
||||
inputs,
|
||||
perm_payload: None,
|
||||
transient: Some(TransientKind::Rebuilding),
|
||||
nodes: vec![node(
|
||||
"hyperhive",
|
||||
NodeKind::MetaLock {
|
||||
sweep: false,
|
||||
fanout: None,
|
||||
|
|
@ -227,9 +270,9 @@ pub fn meta_update(
|
|||
// per-agent child DAGs.
|
||||
|
||||
/// Validate a spec before it enters the queue: node ids are dense
|
||||
/// (index = id), deps reference existing nodes, and the dep graph is
|
||||
/// acyclic (petgraph `toposort`). Rejecting cycles here fixes the old
|
||||
/// queue's documented "circular dep silently deadlocks forever" caveat.
|
||||
/// (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.template);
|
||||
|
|
@ -240,8 +283,19 @@ pub fn validate(spec: &DagSpec) -> Result<()> {
|
|||
.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.template
|
||||
);
|
||||
}
|
||||
for dep in &node.deps {
|
||||
let Some(&dep_idx) = idx.get(dep.on as usize) else {
|
||||
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.template,
|
||||
|
|
|
|||
|
|
@ -109,20 +109,24 @@ fn cyclic_dag_is_rejected_at_submit() {
|
|||
// 0 → 1 → 0 cycle.
|
||||
spec.nodes = vec![
|
||||
NodeSpec {
|
||||
agent: "agent-a".to_owned(),
|
||||
kind: NodeKind::StopForUpdate,
|
||||
kind: NodeKind::StopForUpdate {
|
||||
agent: "agent-a".to_owned(),
|
||||
},
|
||||
deps: vec![Dep {
|
||||
on: 1,
|
||||
when: DepWhen::AfterOk,
|
||||
}],
|
||||
parent: None,
|
||||
},
|
||||
NodeSpec {
|
||||
agent: "agent-a".to_owned(),
|
||||
kind: NodeKind::Reconcile,
|
||||
kind: NodeKind::Reconcile {
|
||||
agent: "agent-a".to_owned(),
|
||||
},
|
||||
deps: vec![Dep {
|
||||
on: 0,
|
||||
when: DepWhen::AfterOk,
|
||||
}],
|
||||
parent: None,
|
||||
},
|
||||
];
|
||||
assert!(q.submit(spec).is_err(), "cyclic spec must be refused");
|
||||
|
|
@ -134,12 +138,30 @@ fn unknown_dep_is_rejected_at_submit() {
|
|||
let q = JobQueue::new(1);
|
||||
let mut spec = rebuild("agent-a", "bad dep");
|
||||
spec.nodes = vec![NodeSpec {
|
||||
agent: "agent-a".to_owned(),
|
||||
kind: NodeKind::Reconcile,
|
||||
kind: NodeKind::Reconcile {
|
||||
agent: "agent-a".to_owned(),
|
||||
},
|
||||
deps: vec![Dep {
|
||||
on: 9,
|
||||
when: DepWhen::AfterOk,
|
||||
}],
|
||||
parent: None,
|
||||
}];
|
||||
assert!(q.submit(spec).is_err());
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn invalid_parent_is_rejected_at_submit() {
|
||||
let q = JobQueue::new(1);
|
||||
let mut spec = rebuild("agent-a", "bad parent");
|
||||
// A forward/out-of-bounds parent index must be refused at validate, not
|
||||
// panic in `insert_group`.
|
||||
spec.nodes = vec![NodeSpec {
|
||||
kind: NodeKind::Reconcile {
|
||||
agent: "agent-a".to_owned(),
|
||||
},
|
||||
deps: Vec::new(),
|
||||
parent: Some(3),
|
||||
}];
|
||||
assert!(q.submit(spec).is_err());
|
||||
}
|
||||
|
|
@ -180,13 +202,16 @@ fn build_slot_serializes_nix_heavy_nodes() {
|
|||
assert_eq!(first.dag_id, a);
|
||||
assert_eq!(first.kind.as_str(), "prebuild");
|
||||
q.complete_node(a, first.node_id, Ok(()));
|
||||
// With the slot free again, FIFO gives... a's StopForUpdate is
|
||||
// slot-free (lease) and b's Prebuild takes the slot — both run.
|
||||
// Uniform hold: agent-a keeps the build slot across its whole build chain
|
||||
// (Swap re-enters it), so a's StopForUpdate (lease, slot-free) runs but b's
|
||||
// Prebuild must wait for a's slot-needers (through Swap) to finish.
|
||||
let claims = q.claim_ready();
|
||||
let kinds: Vec<(u64, &str)> = claims.iter().map(|c| (c.dag_id, c.kind.as_str())).collect();
|
||||
assert!(kinds.contains(&(a, "stop_for_update")));
|
||||
assert!(kinds.contains(&(b, "prebuild")));
|
||||
assert_eq!(claims.len(), 2);
|
||||
assert_eq!(kinds, vec![(a, "stop_for_update")]);
|
||||
assert!(
|
||||
!kinds.iter().any(|&(d, _)| d == b),
|
||||
"b's build waits — slot held across a's chain"
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
|
|
@ -208,9 +233,27 @@ fn fifo_fairness_for_the_slot() {
|
|||
let first = claim_one(&q);
|
||||
assert_eq!(first.dag_id, a, "submit order wins the slot");
|
||||
q.complete_node(a, first.node_id, Ok(()));
|
||||
let next: Vec<u64> = q.claim_ready().iter().map(|cl| cl.dag_id).collect();
|
||||
assert!(next.contains(&b), "b's prebuild before c's");
|
||||
assert!(!next.contains(&c));
|
||||
// Uniform hold: the slot stays with agent-a until its Swap (the last
|
||||
// slot-needer) completes. Drive a's chain; the moment its slot frees,
|
||||
// submit order (b before c) wins it.
|
||||
let mut freed_to = None;
|
||||
for _ in 0..6 {
|
||||
let claims = q.claim_ready();
|
||||
if let Some(nb) = claims.iter().find(|cl| cl.dag_id == b || cl.dag_id == c) {
|
||||
freed_to = Some(nb.dag_id);
|
||||
break;
|
||||
}
|
||||
for cl in claims {
|
||||
if cl.dag_id == a {
|
||||
q.complete_node(a, cl.node_id, Ok(()));
|
||||
}
|
||||
}
|
||||
}
|
||||
assert_eq!(
|
||||
freed_to,
|
||||
Some(b),
|
||||
"b's prebuild wins the freed slot before c's"
|
||||
);
|
||||
}
|
||||
|
||||
// ---- per-agent lease ----
|
||||
|
|
@ -234,17 +277,19 @@ fn lease_serializes_two_lifecycle_dags_for_same_agent() {
|
|||
let first = claim_one(&q);
|
||||
assert_eq!(first.dag_id, restart);
|
||||
assert_eq!(first.kind.as_str(), "stop_for_update");
|
||||
assert!(first.lease_acquired);
|
||||
q.complete_node(restart, first.node_id, Ok(()));
|
||||
// Same DAG keeps the lease through the tail Reconcile.
|
||||
// Same DAG keeps the lease through the tail Reconcile (re-entered from the
|
||||
// dep graph — no fresh acquire), since stop's Reconcile can't re-enter it.
|
||||
let second = claim_one(&q);
|
||||
assert_eq!(second.dag_id, restart);
|
||||
assert_eq!(second.kind.as_str(), "reconcile");
|
||||
assert!(!second.lease_acquired, "lease already held by this DAG");
|
||||
q.complete_node(restart, second.node_id, Ok(()));
|
||||
// Restart terminal → lease released → stop's Reconcile runs.
|
||||
// Restart's work is terminal → its lease releases, so stop's now-unblocked
|
||||
// Reconcile becomes ready (restart's inline hook fired off the returned
|
||||
// summary — no terminal-hook node).
|
||||
let third = claim_one(&q);
|
||||
assert_eq!(third.dag_id, stop);
|
||||
assert_eq!(third.kind.as_str(), "reconcile");
|
||||
q.complete_node(stop, third.node_id, Ok(()));
|
||||
assert_eq!(state_of(&q, restart), State::Done);
|
||||
assert_eq!(state_of(&q, stop), State::Done);
|
||||
|
|
@ -288,8 +333,15 @@ fn lease_exempt_prebuild_overlaps_other_dag_on_same_agent() {
|
|||
.expect("reconcile claim")
|
||||
.clone();
|
||||
q.complete_node(stop, reconcile.node_id, Ok(()));
|
||||
let next = claim_one(&q);
|
||||
assert_eq!(next.kind.as_str(), "stop_for_update");
|
||||
// stop's Reconcile done → its lease frees, so rebuild's StopForUpdate
|
||||
// unblocks. (stop's DAG rolls up terminal; its inline hook fires off the
|
||||
// returned summary — no terminal-hook node in the claim set.)
|
||||
let after = q.claim_ready();
|
||||
let sfu = after
|
||||
.iter()
|
||||
.find(|c| c.kind.as_str() == "stop_for_update")
|
||||
.expect("rebuild StopForUpdate unblocked once the lease frees");
|
||||
assert_eq!(sfu.agent, "agent-a");
|
||||
}
|
||||
|
||||
#[test]
|
||||
|
|
@ -315,16 +367,16 @@ fn multi_agent_restart_is_one_dag_with_concurrent_per_agent_subgraphs() {
|
|||
// lease (no contention across distinct agents), all inside the single DAG.
|
||||
let claims = q.claim_ready();
|
||||
assert!(claims.iter().all(|c| c.dag_id == id));
|
||||
let mut heads: Vec<(&str, &str, bool)> = claims
|
||||
let mut heads: Vec<(&str, &str)> = claims
|
||||
.iter()
|
||||
.map(|c| (c.agent.as_str(), c.kind.as_str(), c.lease_acquired))
|
||||
.map(|c| (c.agent.as_str(), c.kind.as_str()))
|
||||
.collect();
|
||||
heads.sort_unstable();
|
||||
assert_eq!(
|
||||
heads,
|
||||
vec![
|
||||
("agent-a", "stop_for_update", true),
|
||||
("agent-b", "stop_for_update", true),
|
||||
("agent-a", "stop_for_update"),
|
||||
("agent-b", "stop_for_update"),
|
||||
],
|
||||
"both per-agent subgraphs start concurrently, each acquiring its own lease"
|
||||
);
|
||||
|
|
@ -387,17 +439,14 @@ fn multi_agent_stop_is_one_dag_with_concurrent_per_agent_subgraphs() {
|
|||
assert_eq!(q.snapshot().len(), 1);
|
||||
let claims = q.claim_ready();
|
||||
assert!(claims.iter().all(|c| c.dag_id == id));
|
||||
let mut heads: Vec<(&str, &str, bool)> = claims
|
||||
let mut heads: Vec<(&str, &str)> = claims
|
||||
.iter()
|
||||
.map(|c| (c.agent.as_str(), c.kind.as_str(), c.lease_acquired))
|
||||
.map(|c| (c.agent.as_str(), c.kind.as_str()))
|
||||
.collect();
|
||||
heads.sort_unstable();
|
||||
assert_eq!(
|
||||
heads,
|
||||
vec![
|
||||
("agent-a", "set_wanted", true),
|
||||
("agent-b", "set_wanted", true),
|
||||
],
|
||||
vec![("agent-a", "set_wanted"), ("agent-b", "set_wanted")],
|
||||
"both per-agent stop subgraphs start concurrently, each on its own lease"
|
||||
);
|
||||
}
|
||||
|
|
@ -512,15 +561,14 @@ fn append_subgraph_roots_on_emitter_and_rebases_local_deps() {
|
|||
reason: "sweep".to_owned(),
|
||||
approval_id: None,
|
||||
inputs: Vec::new(),
|
||||
perm_payload: None,
|
||||
transient: None,
|
||||
nodes: vec![NodeSpec {
|
||||
agent: "hyperhive".to_owned(),
|
||||
kind: NodeKind::MetaLock {
|
||||
sweep: true,
|
||||
fanout: None,
|
||||
},
|
||||
deps: Vec::new(),
|
||||
parent: None,
|
||||
}],
|
||||
};
|
||||
let id = submit(&q, spec);
|
||||
|
|
@ -532,8 +580,8 @@ fn append_subgraph_roots_on_emitter_and_rebases_local_deps() {
|
|||
// tracks any drift in that builder's root-first (`base = 0`) shape.
|
||||
let subgraph = |agent: &str| templates::rebuild_nodes(agent, true, 0);
|
||||
// Must append BEFORE completing the emitter (the documented contract).
|
||||
q.append_subgraph(id, subgraph("a"), emitter.node_id);
|
||||
q.append_subgraph(id, subgraph("b"), emitter.node_id);
|
||||
q.append_subgraph(id, &subgraph("a"), emitter.node_id);
|
||||
q.append_subgraph(id, &subgraph("b"), emitter.node_id);
|
||||
q.complete_node(id, emitter.node_id, Ok(()));
|
||||
// Still ONE DAG; both subgraph roots become ready once the emitter is
|
||||
// Done (rooted on it), each on its own agent lease.
|
||||
|
|
@ -580,7 +628,7 @@ fn meta_update_carries_rebuilding_transient_and_grows_cascade_in_dag() {
|
|||
for agent in ["alice", "bob"] {
|
||||
q.append_subgraph(
|
||||
id,
|
||||
templates::rebuild_nodes(agent, false, 0),
|
||||
&templates::rebuild_nodes(agent, false, 0),
|
||||
meta_lock.node_id,
|
||||
);
|
||||
}
|
||||
|
|
@ -726,7 +774,10 @@ fn failed_reconcile_marks_dag_failed() {
|
|||
fn cancel_clears_queued_dag() {
|
||||
let q = JobQueue::new(1);
|
||||
let id = submit(&q, rebuild("agent-a", "r"));
|
||||
assert!(q.cancel(id));
|
||||
// Cancel returns the terminal summary (state `Cancelled`) — the inline hook
|
||||
// fires off it at the caller; there's no terminal-hook node to claim.
|
||||
let terminal = q.cancel(id).expect("cancelled");
|
||||
assert_eq!(terminal.state, State::Cancelled);
|
||||
assert_eq!(state_of(&q, id), State::Cancelled);
|
||||
assert!(q.claim_ready().is_empty());
|
||||
}
|
||||
|
|
@ -736,29 +787,32 @@ fn cancel_refuses_running_dag() {
|
|||
let q = JobQueue::new(1);
|
||||
let id = submit(&q, rebuild("agent-a", "r"));
|
||||
let _ = claim_one(&q);
|
||||
assert!(!q.cancel(id));
|
||||
assert!(q.cancel(id).is_none());
|
||||
assert_eq!(state_of(&q, id), State::Running);
|
||||
}
|
||||
|
||||
// ---- terminal reporting + lease release ----
|
||||
|
||||
#[test]
|
||||
fn terminal_dag_reported_exactly_once_and_lease_released() {
|
||||
fn dag_settles_terminal_and_releases_lease_after_work() {
|
||||
let q = JobQueue::new(1);
|
||||
let id = submit(&q, restart_online(&["agent-a"], false, "r"));
|
||||
// restart = StopForUpdate → Reconcile; not terminal until the last
|
||||
// node completes.
|
||||
// restart = StopForUpdate → Reconcile.
|
||||
let stop = claim_one(&q);
|
||||
assert_eq!(stop.kind.as_str(), "stop_for_update");
|
||||
q.complete_node(id, stop.node_id, Ok(()));
|
||||
assert!(q.drain_terminal().is_empty(), "dag not terminal yet");
|
||||
let rec = claim_one(&q);
|
||||
q.complete_node(id, rec.node_id, Ok(()));
|
||||
let reports = q.drain_terminal();
|
||||
assert_eq!(reports.len(), 1);
|
||||
assert_eq!(reports[0].dag_id, id);
|
||||
assert_eq!(reports[0].state, State::Done);
|
||||
assert!(q.drain_terminal().is_empty(), "reported exactly once");
|
||||
// Lease released: a new DAG for the agent can claim immediately.
|
||||
assert_eq!(rec.kind.as_str(), "reconcile");
|
||||
// Completing the last work node rolls the container up terminal and returns
|
||||
// the summary the inline hook consumes — there is no terminal-hook node.
|
||||
let summary = q
|
||||
.complete_node(id, rec.node_id, Ok(()))
|
||||
.expect("terminal summary");
|
||||
assert_eq!(summary.state, State::Done);
|
||||
assert!(q.claim_ready().is_empty(), "no terminal-hook node to claim");
|
||||
assert_eq!(state_of(&q, id), State::Done);
|
||||
// Lease released when the work chain settled: a new DAG for the agent claims
|
||||
// immediately.
|
||||
let next = submit(
|
||||
&q,
|
||||
templates::reconcile_only(
|
||||
|
|
@ -771,31 +825,33 @@ fn terminal_dag_reported_exactly_once_and_lease_released() {
|
|||
);
|
||||
let c = claim_one(&q);
|
||||
assert_eq!(c.dag_id, next);
|
||||
assert!(c.lease_acquired);
|
||||
}
|
||||
|
||||
/// A DAG cancelled while fully queued must still surface a terminal
|
||||
/// roll-up for the scheduler's hooks — otherwise a queued approval
|
||||
/// DAG cancelled by the operator would dangle its approval forever.
|
||||
#[test]
|
||||
fn cancelled_dag_reports_terminal_once() {
|
||||
fn cancelled_dag_finalizes_with_terminal_rollup() {
|
||||
let q = JobQueue::new(1);
|
||||
let id = submit(
|
||||
&q,
|
||||
templates::approval_deploy("agent-a", 7, "approval #7".to_owned()),
|
||||
);
|
||||
assert!(q.cancel(id));
|
||||
let reports = q.drain_terminal();
|
||||
assert_eq!(reports.len(), 1);
|
||||
assert_eq!(reports[0].dag_id, id);
|
||||
assert_eq!(reports[0].state, State::Cancelled);
|
||||
assert_eq!(reports[0].approval_id, Some(7));
|
||||
// Never re-reported by later activity.
|
||||
// Cancel rolls the DAG up terminal and returns its summary — the inline hook
|
||||
// (approval resolution) runs off it at the caller. Cancelled + approval id 7.
|
||||
let summary = q.cancel(id).expect("cancelled");
|
||||
assert_eq!(summary.state, State::Cancelled);
|
||||
assert_eq!(summary.approval_id, Some(7));
|
||||
// The cancelled DAG's summary stays available (until history-trimmed) and
|
||||
// unrelated later activity doesn't disturb it.
|
||||
let other = submit(&q, rebuild("agent-b", "r"));
|
||||
let c = claim_one(&q);
|
||||
assert_eq!(c.dag_id, other);
|
||||
q.complete_node(other, c.node_id, Err("boom".to_owned()));
|
||||
assert!(q.drain_terminal().iter().all(|t| t.dag_id != id));
|
||||
assert_eq!(
|
||||
q.terminal_summary(id).map(|t| t.state),
|
||||
Some(State::Cancelled)
|
||||
);
|
||||
}
|
||||
|
||||
// ---- steps, build logs, history ----
|
||||
|
|
@ -804,7 +860,10 @@ fn cancelled_dag_reports_terminal_once() {
|
|||
fn set_step_only_on_running_and_signals_change() {
|
||||
let q = JobQueue::new(1);
|
||||
let id = submit(&q, rebuild("agent-a", "r"));
|
||||
assert!(!q.set_step(id, 0, "too early"), "queued node refuses step");
|
||||
assert!(
|
||||
!q.set_step_running(id, "too early"),
|
||||
"no running node yet → refused"
|
||||
);
|
||||
let c = claim_one(&q);
|
||||
assert!(q.set_step(id, c.node_id, "nix build"));
|
||||
assert!(
|
||||
|
|
@ -823,7 +882,10 @@ fn set_step_only_on_running_and_signals_change() {
|
|||
fn set_build_log_id_links_running_node() {
|
||||
let q = JobQueue::new(1);
|
||||
let id = submit(&q, rebuild("agent-a", "r"));
|
||||
assert!(!q.set_build_log_id(id, 0, 41), "queued node refuses log id");
|
||||
assert!(
|
||||
!q.set_build_log_id_running(id, 41),
|
||||
"no running node yet → refused"
|
||||
);
|
||||
let c = claim_one(&q);
|
||||
assert!(q.set_build_log_id(id, c.node_id, 42));
|
||||
assert!(q.set_build_log_id_running(id, 43));
|
||||
|
|
@ -848,6 +910,8 @@ fn history_evicts_old_terminals_per_template() {
|
|||
),
|
||||
);
|
||||
let c = claim_one(&q);
|
||||
// Completing the single work node rolls the container up terminal (its
|
||||
// inline hook fires off the returned summary — no terminal-hook node).
|
||||
q.complete_node(id, c.node_id, Ok(()));
|
||||
}
|
||||
// Fresh terminals are inside the grace window: nothing evicts yet,
|
||||
|
|
@ -859,8 +923,7 @@ fn history_evicts_old_terminals_per_template() {
|
|||
"grace window protects fresh terminals"
|
||||
);
|
||||
// Past the grace window the per-template cap applies.
|
||||
q.trim_ignoring_grace();
|
||||
assert_eq!(q.snapshot().len(), 5, "per-template history cap");
|
||||
assert_eq!(q.snapshot_no_grace().len(), 5, "per-template history cap");
|
||||
assert_eq!(q.live_count(), 0);
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -325,20 +325,20 @@ fn submit_boot_tree(
|
|||
// MetaLock into `run_meta_lock`, which appends the rebuild subgraphs.
|
||||
if any_stale {
|
||||
nodes.push(NodeSpec {
|
||||
agent: "hyperhive".to_owned(),
|
||||
kind: NodeKind::MetaLock {
|
||||
sweep: true,
|
||||
fanout: Some(fanout),
|
||||
},
|
||||
deps: Vec::new(),
|
||||
parent: None,
|
||||
});
|
||||
}
|
||||
// One boot Reconcile per drifted agent — independent roots.
|
||||
for name in drifted {
|
||||
nodes.push(NodeSpec {
|
||||
agent: name,
|
||||
kind: NodeKind::Reconcile,
|
||||
kind: NodeKind::Reconcile { agent: name },
|
||||
deps: Vec::new(),
|
||||
parent: None,
|
||||
});
|
||||
}
|
||||
|
||||
|
|
@ -348,7 +348,6 @@ fn submit_boot_tree(
|
|||
reason,
|
||||
approval_id: None,
|
||||
inputs: Vec::new(),
|
||||
perm_payload: None,
|
||||
// Rebuilding when the sweep will grow rebuild subgraphs (per-agent
|
||||
// crash-watch suppression during their Swap, applied at claim time);
|
||||
// a reconcile-only boot needs no transient.
|
||||
|
|
|
|||
|
|
@ -134,8 +134,11 @@ pub enum PermPayload {
|
|||
},
|
||||
}
|
||||
|
||||
/// Node id, unique within its DAG.
|
||||
pub type NodeId = u32;
|
||||
/// Node id. Carries the scheduler crate's globally-monotonic node id
|
||||
/// (`hive_jobq::NodeId`) verbatim on the wire — unique across all DAGs, not
|
||||
/// just within one. Consumers treat it opaquely (grouping + dep matching),
|
||||
/// so the widening from the old dag-local `u32` is transparent.
|
||||
pub type NodeId = u64;
|
||||
|
||||
/// One node of a queued DAG, as serialized. Step labels, build-log
|
||||
/// links, errors, and timestamps are per-node; the DAG-level `state`
|
||||
|
|
@ -192,7 +195,5 @@ pub struct DagView {
|
|||
pub inputs: Vec<String>,
|
||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||||
pub approval_id: Option<i64>,
|
||||
#[serde(default, skip_serializing_if = "Option::is_none")]
|
||||
pub perm_payload: Option<PermPayload>,
|
||||
pub nodes: Vec<NodeView>,
|
||||
}
|
||||
|
|
|
|||
|
|
@ -325,7 +325,7 @@ mod tests {
|
|||
|
||||
use super::render_dag_line;
|
||||
|
||||
fn node(id: u32, agent: &str, kind: &str, state: State, step: Option<&str>) -> NodeView {
|
||||
fn node(id: u64, agent: &str, kind: &str, state: State, step: Option<&str>) -> NodeView {
|
||||
NodeView {
|
||||
id,
|
||||
agent: agent.to_owned(),
|
||||
|
|
@ -353,7 +353,6 @@ mod tests {
|
|||
finished_at: None,
|
||||
inputs: vec![],
|
||||
approval_id: None,
|
||||
perm_payload: None,
|
||||
nodes: vec![
|
||||
node(0, "alice", "prebuild", State::Done, None),
|
||||
node(1, "alice", "stop_for_update", State::Done, None),
|
||||
|
|
@ -392,7 +391,6 @@ mod tests {
|
|||
finished_at: Some(2),
|
||||
inputs: vec![],
|
||||
approval_id: None,
|
||||
perm_payload: None,
|
||||
nodes: vec![failed],
|
||||
};
|
||||
let line = render_dag_line(&dag);
|
||||
|
|
|
|||
Loading…
Reference in a new issue