diff --git a/docs/coordinator.md b/docs/coordinator.md index b4e6b4c2..115d8ec2 100644 --- a/docs/coordinator.md +++ b/docs/coordinator.md @@ -106,7 +106,7 @@ subgraph each (independent roots, run concurrently on their own leases), not N separate DAGs. **These are built dynamically from each agent's live running state** (an -async `lifecycle::is_running` read), so they live in `job_queue/power.rs`, +async `lifecycle::is_running` read), so they live in `job_queue/submit.rs`, not the pure/sync `templates.rs`. Per-agent shape rule: `stop`/`start` carry a head `SetWanted` (intent) — `restart` does not; the tail `Reconcile` (convergence guarantee — cheap, noops when already converged) is ALWAYS @@ -161,12 +161,12 @@ Notable collapses: Per-agent power *intent* — `wanted: Up | Offline` — is durable as the `agent_power` table in the coordinator DB (`hive-c0re/src/stores/power.rs`). `container_view` remains the observed *status*; `Reconcile` nodes converge the -two. Setting `wanted` is never a queued node: the power layer -(`job_queue/power.rs`) writes the row synchronously, then inserts the DAG +two. Setting `wanted` is never a queued node: the submit layer +(`job_queue/submit.rs`) writes the row synchronously, then submits the DAG whose `Reconcile` reads the fresh value — rapid toggles are last-writer-wins. Power toggles never commit to the meta repo. Every operator power surface — dashboard buttons, the MCP tools, and `hivectl stop/start/restart/kill` — -rides the queue through that power layer, so intent, lease serialization, +rides the queue through that submit layer, so intent, lease serialization, and crash-watch suppression can't drift per surface; the only direct starts left are the root-agent bootstrap and infra containers (no lease, no harness). Cancelling a still-queued power DAG reverts `wanted` to the diff --git a/hive-c0re/src/actions.rs b/hive-c0re/src/actions.rs index 11c94c2f..7333f387 100644 --- a/hive-c0re/src/actions.rs +++ b/hive-c0re/src/actions.rs @@ -55,12 +55,13 @@ pub async fn approve(coord: Arc, id: i64) -> Result<()> { // nothing). let inputs: Vec = serde_json::from_str(&approval.commit_ref).unwrap_or_default(); - let inserted = coord.job_queue.insert_job(|b| { - crate::job_queue::templates::meta_update(b, inputs, Some(id)); - Vec::new() - }); - if let Err(e) = inserted { - return Err(e.context("insert meta-update dag")); + let submitted = coord.job_queue.submit( + crate::job_queue::Source::Approval, + format!("approval #{id} meta input update"), + |b| crate::job_queue::templates::meta_update(b, inputs, Some(id)), + ); + if let Err(e) = submitted { + return Err(e.context("submit meta-update dag")); } coord.emit_rebuild_queue_snapshot(); Ok(()) @@ -74,12 +75,13 @@ pub async fn approve(coord: Arc, id: i64) -> Result<()> { { tracing::warn!(agent = %approval.agent, error = ?e, "agent_power: seed on spawn failed"); } - let inserted = coord.job_queue.insert_job(|b| { - crate::job_queue::templates::spawn(b, approval.agent.as_str(), id); - Vec::new() - }); - if let Err(e) = inserted { - return Err(e.context("insert spawn dag")); + let submitted = coord.job_queue.submit( + crate::job_queue::Source::Approval, + format!("approval #{id} spawn"), + |b| crate::job_queue::templates::spawn(b, approval.agent.as_str(), id), + ); + if let Err(e) = submitted { + return Err(e.context("submit spawn dag")); } coord.emit_rebuild_queue_snapshot(); Ok(()) @@ -104,7 +106,12 @@ pub async fn approve(coord: Arc, id: i64) -> Result<()> { // `run_deploy_apply` (ff-merge, then grows the rebuild subgraph), // `run_finalize_deploy` (deploy tag + lock commit) and // `run_deploy_tail` (compensation + forge mirror). - enqueue_approval_rebuild(&coord, approval.agent.as_str(), id); + enqueue_approval_rebuild( + &coord, + approval.agent.as_str(), + id, + format!("approval #{id} merge config pr"), + ); Ok(()) } } @@ -114,12 +121,19 @@ pub async fn approve(coord: Arc, id: i64) -> Result<()> { /// dispatch arm — the work ends in a container rebuild routed through the /// queue. See [`crate::job_queue::templates::approval_deploy`] for the node /// shape; the executor dispatches each node to the `run_deploy_*` bodies below. -fn enqueue_approval_rebuild(coord: &Arc, agent: &str, approval_id: i64) { - if let Err(e) = coord.job_queue.insert_job(|b| { - crate::job_queue::templates::approval_deploy(b, agent, approval_id); - Vec::new() - }) { - tracing::error!(%agent, approval_id, error = ?e, "insert approval deploy dag failed"); +fn enqueue_approval_rebuild( + coord: &Arc, + agent: &str, + approval_id: i64, + reason: String, +) { + if let Err(e) = coord + .job_queue + .submit(crate::job_queue::Source::Approval, reason, |b| { + crate::job_queue::templates::approval_deploy(b, agent, approval_id); + }) + { + tracing::error!(%agent, approval_id, error = ?e, "submit approval deploy dag failed"); } coord.emit_rebuild_queue_snapshot(); } diff --git a/hive-c0re/src/dashboard/lifecycle_ops.rs b/hive-c0re/src/dashboard/lifecycle_ops.rs index e12723fa..87e7c2d0 100644 --- a/hive-c0re/src/dashboard/lifecycle_ops.rs +++ b/hive-c0re/src/dashboard/lifecycle_ops.rs @@ -1,12 +1,11 @@ //! Container lifecycle endpoints for the dashboard. //! //! Rebuild / restart / start / stop (hard + graceful) / update-all all -//! insert DAGs into the job queue — the power ops via -//! [`crate::job_queue::power`], the static shapes straight through -//! `JobQueue::insert` — so each shows a visible queued→running transient on -//! the dashboard; a direct sub-second start/stop only flashed the badge. -//! Start/stop also persist the agent's `wanted` power intent before -//! inserting; the DAG's `Reconcile` converges to it. Destroy delegates to +//! submit DAGs to the job queue (`job_queue::submit`), so each shows a +//! visible queued→running transient on the dashboard — a direct +//! sub-second start/stop only flashed the badge. Start/stop also +//! persist the agent's `wanted` power intent before submitting; the +//! DAG's `Reconcile` converges to it. Destroy delegates to //! `actions::destroy` (optionally purging). use axum::{ @@ -28,6 +27,7 @@ pub(super) struct GracefulParams { } use super::{AppState, Ident, error_response, guard_agent_name, strip_container_prefix}; +use crate::job_queue::{Source, submit}; use crate::{actions, lifecycle}; /// Queue a rebuild DAG for `name`. @@ -50,13 +50,12 @@ pub(super) async fn post_rebuild( if let Some(reject) = guard_agent_name(&state, &logical).await { return reject; } - if let Err(e) = state.coord.job_queue.insert_job(|b| { - crate::job_queue::templates::rebuild(b, &logical, true); - Vec::new() - }) { - tracing::error!(agent = %logical, error = ?e, "rebuild: insert failed"); - } - state.coord.emit_rebuild_queue_snapshot(); + submit::rebuild( + &state.coord, + &logical, + Source::Manual, + "manual via dashboard ↻ R3BU1LD button".to_owned(), + ); (StatusCode::OK, "ok").into_response() } @@ -93,12 +92,13 @@ pub(super) async fn post_kill( // timeout fallback to a hard stop). The agent's lifecycle // lease keeps it from racing an in-flight rebuild for the same // agent, and per-node progress surfaces on the queue snapshot. - if let Err(e) = - crate::job_queue::power::stop_many(&state.coord, std::slice::from_ref(&logical), true) - .await - { - tracing::error!(agent = %logical, error = ?e, "graceful stop: insert failed"); - } + submit::graceful_stop( + &state.coord, + &logical, + Source::Manual, + "manual via dashboard graceful stop".to_owned(), + ) + .await; return (StatusCode::OK, "ok").into_response(); } // Manager is stoppable from the dashboard like any other @@ -111,12 +111,13 @@ pub(super) async fn post_kill( // `socket_server.rs::Request::Kill` stays in place: a // manager calling Kill on its own container is self-suicide // mid-call, not a legitimate operator action. - if let Err(e) = - crate::job_queue::power::stop_many(&state.coord, std::slice::from_ref(&logical), false) - .await - { - tracing::error!(agent = %logical, error = ?e, "stop: insert failed"); - } + submit::stop( + &state.coord, + &logical, + Source::Manual, + "manual via dashboard stop".to_owned(), + ) + .await; (StatusCode::OK, "ok").into_response() } @@ -148,23 +149,22 @@ pub(super) async fn post_restart( return reject; } if params.graceful { - if let Err(e) = crate::job_queue::power::restart_many( + submit::graceful_restart( &state.coord, - std::slice::from_ref(&logical), - true, + &logical, + Source::Manual, + "manual via dashboard graceful restart".to_owned(), ) - .await - { - tracing::error!(agent = %logical, error = ?e, "graceful restart: insert failed"); - } + .await; return (StatusCode::OK, "ok").into_response(); } - if let Err(e) = - crate::job_queue::power::restart_many(&state.coord, std::slice::from_ref(&logical), false) - .await - { - tracing::error!(agent = %logical, error = ?e, "restart: insert failed"); - } + submit::restart( + &state.coord, + &logical, + Source::Manual, + "manual via dashboard ↺ R3START button".to_owned(), + ) + .await; (StatusCode::OK, "ok").into_response() } @@ -226,11 +226,13 @@ pub(super) async fn post_start( return (StatusCode::OK, "ok").into_response(); } } - if let Err(e) = - crate::job_queue::power::start_many(&state.coord, std::slice::from_ref(&logical)).await - { - tracing::error!(agent = %logical, error = ?e, "start: insert failed"); - } + submit::start( + &state.coord, + &logical, + Source::Manual, + "manual via dashboard start".to_owned(), + ) + .await; (StatusCode::OK, "ok").into_response() } @@ -409,14 +411,13 @@ pub(super) async fn post_update_all(State(state): State) -> Response { else { continue; }; - if let Err(e) = state.coord.job_queue.insert_job(|b| { - crate::job_queue::templates::rebuild(b, &logical, true); - Vec::new() - }) { - tracing::error!(agent = %logical, error = ?e, "update-all: insert failed"); - } + submit::rebuild( + &state.coord, + &logical, + Source::Manual, + "manual via dashboard 🌀 UPDATE ALL".to_owned(), + ); } - state.coord.emit_rebuild_queue_snapshot(); (StatusCode::OK, "ok").into_response() } diff --git a/hive-c0re/src/dashboard/meta_inputs.rs b/hive-c0re/src/dashboard/meta_inputs.rs index 3f6705d8..3efe6b74 100644 --- a/hive-c0re/src/dashboard/meta_inputs.rs +++ b/hive-c0re/src/dashboard/meta_inputs.rs @@ -208,18 +208,16 @@ pub(super) async fn post_meta_update( if inputs.is_empty() { return error_response("meta-update: no inputs selected"); } + let inputs_label = inputs.join(", "); // Cascade rebuild children fan out from the MetaLock node when the // lock bump lands — appended by the scheduler so they build against // the post-bump lock, and a failed bump simply fans out nothing. - state - .coord - .job_queue - .insert_job(|b| { - crate::job_queue::templates::meta_update(b, inputs, None); - Vec::new() - }) - .expect("template-declared shapes are acyclic"); - state.coord.emit_rebuild_queue_snapshot(); + crate::job_queue::submit::meta_update( + &state.coord, + inputs, + crate::job_queue::Source::Manual, + format!("meta-update via dashboard ({inputs_label})"), + ); (StatusCode::OK, "ok").into_response() } diff --git a/hive-c0re/src/dashboard/permissions.rs b/hive-c0re/src/dashboard/permissions.rs index d542fc8c..601ca400 100644 --- a/hive-c0re/src/dashboard/permissions.rs +++ b/hive-c0re/src/dashboard/permissions.rs @@ -155,21 +155,15 @@ pub(super) async fn post_tool_groups( // META_LOCK inside the WritePermFile node, so concurrent // batch-apply actions for different agents never race on the // shared tool-groups.json. - state - .coord - .job_queue - .insert_job(|b| { - crate::job_queue::templates::perm_change( - b, - &logical, - crate::job_queue::PermPayload::ToolGroups { - groups: body.groups.clone(), - }, - ); - Vec::new() - }) - .expect("template-declared shapes are acyclic"); - state.coord.emit_rebuild_queue_snapshot(); + crate::job_queue::submit::perm_change( + &state.coord, + &logical, + crate::job_queue::Source::Manual, + "tool-group change via permissions UI".to_owned(), + crate::job_queue::PermPayload::ToolGroups { + groups: body.groups.clone(), + }, + ); tracing::info!(agent = %logical, groups = ?body.groups, "operator: set tool-groups via dashboard"); Ok((StatusCode::OK, "ok").into_response()) } @@ -270,21 +264,15 @@ pub(super) async fn post_capabilities( // META_LOCK inside the WritePermFile node, so concurrent // batch-apply actions for different agents never race on the // shared capabilities.json. - state - .coord - .job_queue - .insert_job(|b| { - crate::job_queue::templates::perm_change( - b, - &logical, - crate::job_queue::PermPayload::Capabilities { - caps: body.caps.clone(), - }, - ); - Vec::new() - }) - .expect("template-declared shapes are acyclic"); - state.coord.emit_rebuild_queue_snapshot(); + crate::job_queue::submit::perm_change( + &state.coord, + &logical, + crate::job_queue::Source::Manual, + "capability change via dashboard".to_owned(), + crate::job_queue::PermPayload::Capabilities { + caps: body.caps.clone(), + }, + ); tracing::info!(agent = %logical, caps = ?body.caps, "operator: set capabilities via dashboard"); Ok((StatusCode::OK, "ok").into_response()) } @@ -371,19 +359,13 @@ pub(super) async fn post_permissions( } // Phase 2 — submit one combined PermChange DAG per affected agent. for (logical, groups, caps) in staged { - state - .coord - .job_queue - .insert_job(|b| { - crate::job_queue::templates::perm_change( - b, - &logical, - crate::job_queue::PermPayload::Combined { groups, caps }, - ); - Vec::new() - }) - .expect("template-declared shapes are acyclic"); - state.coord.emit_rebuild_queue_snapshot(); + crate::job_queue::submit::perm_change( + &state.coord, + &logical, + crate::job_queue::Source::Manual, + "batch permission change via permissions UI".to_owned(), + crate::job_queue::PermPayload::Combined { groups, caps }, + ); tracing::info!(agent = %logical, "operator: batch perm change via dashboard"); } Ok((StatusCode::OK, "ok").into_response()) diff --git a/hive-c0re/src/dashboard/topology.rs b/hive-c0re/src/dashboard/topology.rs index 37ec5045..8e8422fb 100644 --- a/hive-c0re/src/dashboard/topology.rs +++ b/hive-c0re/src/dashboard/topology.rs @@ -21,6 +21,7 @@ use utoipa::ToSchema; use problem_details::ProblemDetails; use super::{AppState, error_problem}; +use crate::job_queue::{Source, submit}; /// `POST /api/topology/set-parent` body. `child` is required. /// `new_parent` may be: @@ -93,15 +94,12 @@ pub(super) async fn post_set_parent( new_parent = ?new_parent, "operator: set-parent via dashboard" ); - state - .coord - .job_queue - .insert_job(|b| { - crate::job_queue::templates::reparent(b, vec![(child, new_parent)]); - Vec::new() - }) - .expect("template-declared shapes are acyclic"); - state.coord.emit_rebuild_queue_snapshot(); + submit::reparent( + &state.coord, + vec![(child, new_parent)], + Source::Manual, + "manual set-parent via dashboard".to_owned(), + ); Ok((StatusCode::OK, "ok").into_response()) } @@ -152,14 +150,11 @@ pub(super) async fn post_set_parent_bulk( .map_err(|e| error_problem(&e))?; let names: Vec<&str> = body.iter().map(|e| e.child.as_str()).collect(); tracing::info!(agents = ?names, "operator: set-parent-bulk via dashboard"); - state - .coord - .job_queue - .insert_job(|b| { - crate::job_queue::templates::reparent(b, moves); - Vec::new() - }) - .expect("template-declared shapes are acyclic"); - state.coord.emit_rebuild_queue_snapshot(); + submit::reparent( + &state.coord, + moves, + Source::Manual, + "manual set-parent-bulk via dashboard".to_owned(), + ); Ok((StatusCode::OK, "ok").into_response()) } diff --git a/hive-c0re/src/job_queue/exec.rs b/hive-c0re/src/job_queue/exec.rs index d8866327..fe0c4074 100644 --- a/hive-c0re/src/job_queue/exec.rs +++ b/hive-c0re/src/job_queue/exec.rs @@ -53,7 +53,7 @@ pub(super) async fn run_node( kind: &NodeKind, ) -> (super::JobBuilder, Result<()>) { // The agent this node targets rides the payload — empty for the agentless - // kinds (`MetaLock`, `Reparent`), which never read it. + // container kinds (`MetaLock`, `Dag`), which never read it. let agent = kind.agent(); // Every arm is `Result<()>`; the three that grow work declare into `builder` // *synchronously*, after their own awaits have finished. Borrowing `&builder` @@ -115,21 +115,26 @@ pub(super) async fn run_node( run_finalize_deploy(coord, *approval_id).await } NodeKind::DeployTail { approval_id, .. } => { - run_deploy_tail(coord, coord.job_queue.root_of(id), agent, *approval_id).await + run_deploy_tail(coord, coord.job_queue.dag_of(id), agent, *approval_id).await } NodeKind::ResolveApproval { approval_id, outcome, - } => run_resolve_approval(coord, coord.job_queue.root_of(id), *approval_id, *outcome).await, + } => run_resolve_approval(coord, coord.job_queue.dag_of(id), *approval_id, *outcome).await, NodeKind::EmitRebuilt { ok, .. } => { - run_emit_rebuilt(coord, agent, coord.job_queue.root_of(id), *ok).await; + run_emit_rebuilt(coord, agent, coord.job_queue.dag_of(id), *ok).await; Ok(()) } NodeKind::SetWanted { up, .. } => run_set_wanted(coord, agent, *up), - // Braces carry no work of their own; completing one lets it reach - // `Finishing` so the nodes under it start. What they declare stays held - // until their whole subtree settles. - NodeKind::DeployWindow { .. } | NodeKind::AgentWindow { .. } => Ok(()), + // The nodes that carry no work of their own; completing one 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` / `AgentWindow`: pure resource holders (braces) — + // what they declare stays held until their subtree settles. + NodeKind::Dag { .. } | NodeKind::DeployWindow { .. } | NodeKind::AgentWindow { .. } => { + Ok(()) + } }; (builder, result) } diff --git a/hive-c0re/src/job_queue/mod.rs b/hive-c0re/src/job_queue/mod.rs index 04ed00af..b9727493 100644 --- a/hive-c0re/src/job_queue/mod.rs +++ b/hive-c0re/src/job_queue/mod.rs @@ -10,16 +10,17 @@ //! carries the agent it targets ([`NodeKind::agent`]); the two resource //! classes are [`resource::Resource`] (`BuildSlot` node-held, `Agent` lease //! subtree-held), declared per node at its construction site; -//! - **a job has no container node.** A template declares its nodes and names -//! the roots it wants back; `insert_job` returns those ids. Grouping is the -//! parent axis (a root's rolled-up state *is* its subtree's), so membership is -//! a graph walk with no host-side side-tables. The lease is owned by a subtree -//! root and borrowed by its descendants (continuity); -//! - terminal work is an ordinary **tail node** +//! - a **DAG is a single container node** ([`NodeKind::Dag`], `parent = None`) +//! carrying the group's metadata, with the work nodes hung under it as +//! its subtree (the **parent axis** groups; `deps` order). So the container's +//! `NodeId` is the DAG id, its rolled-up state is the DAG state, and membership +//! is a graph walk — there are no host grouping side-tables. The lease is owned +//! by a subtree root and borrowed by its descendants (continuity); +//! - per-DAG terminal work is an ordinary **tail node** //! ([`NodeKind::ResolveApproval`] / [`NodeKind::EmitRebuilt`]) that the builder -//! appends in [`templates`], edged onto the job's group roots by the outcome it -//! reports. Templates emit one tail per outcome and the graph runs exactly one, -//! so nothing branches at runtime. +//! appends in [`templates`], edged onto the DAG's other group roots by the +//! outcome it reports. Templates emit one tail per outcome and the graph runs +//! exactly one, so nothing branches at runtime. //! //! The queue is runtime-only (no persistence): an empty graph on boot; desired //! state is re-derived by the reconcile sweep. A single scheduler task @@ -28,9 +29,9 @@ pub mod exec; pub mod model; -pub mod power; pub mod resource; pub mod scheduler; +pub mod submit; pub mod templates; #[cfg(test)] mod tests; @@ -45,7 +46,7 @@ use hive_jobq_wire::{GraphNode, GraphWire}; use tokio::sync::Notify; pub use hive_jobq::TerminalState; -pub use model::{NodeKind, PermPayload, State}; +pub use model::{NodeKind, PermPayload, Source, State}; use resource::Resource; /// A job under construction: `hive_jobq`'s builder over this queue's payload @@ -89,10 +90,12 @@ pub struct RunningTransient { /// The crate scheduler, specialised to this host's node + resource types. /// -/// A job is **just its nodes** — no container, no grouping side-tables. A -/// root's rolled-up state is its subtree's, so membership is a graph walk and -/// "which job is this node in" is [`JobQueue::root_of`]. One shared crate -/// [`Graph`] holds every job's nodes. +/// A **DAG is a single container node** ([`NodeKind::Dag`], `parent = None`) +/// whose subtree is the DAG's work — so the container's `NodeId` is the DAG id, +/// its rolled-up state is the DAG state, and there are no grouping side-tables: +/// membership + meta are graph queries ([`container`] + the `hive_jobq::Graph` +/// accessors, with the meta read straight off the container's payload). One +/// shared crate [`Graph`] holds every DAG. /// /// There is deliberately **no wrapper struct and no per-node side map**. The /// last map held the `build_logs` row id; that link now lives on the log row @@ -137,12 +140,35 @@ fn outcome_of(result: Result<(), String>) -> Outcome { } } -// `insert_group` lived here: a `group_parent`-taking insert whose only -// remaining caller was the DAG container, everything under it. Runtime growth -// never went through it — an executor declares into the builder `hive_jobq` -// hands it, which parents the new work under the emitting node by -// construction. With no container to be the other kind of parent, the -// distinction it existed to express is gone. +/// Insert a declared `job` into the shared graph, returning the inserted ids. +/// +/// A node that declared no parent hangs under `group_parent` — the DAG +/// container for a template, the emitting node for a runtime-appended +/// subgraph. Templates declare the parent axis + sibling ordering directly, so +/// there is no dep-on-root to drop and no lease to hoist: each node declares +/// its own resources, and the crate's borrow model keeps a resource continuous +/// across a subtree (a root owns it, descendants borrow it). Independent group +/// roots carry no cross-links, so a multi-agent DAG's per-agent subgraphs run +/// concurrently, each on its own lease. +/// +/// # Errors +/// Propagates a crate graph-insert error (malformed dep/parent / dep-scope). +fn insert_group( + inner: &mut Sched, + declare: impl FnOnce(&JobBuilder), + group_parent: Option, +) -> anyhow::Result<()> { + inner + .insert_job(group_parent, |b| { + declare(b); + // c0re names no handles: a DAG is addressed by its container node, + // which `submit` inserts itself, and nothing downstream looks an + // individual step up by id. + Vec::new() + }) + .map_err(|e| anyhow::anyhow!("job_queue: graph insert failed: {e}"))?; + Ok(()) +} impl JobQueue { #[must_use] @@ -162,34 +188,40 @@ impl JobQueue { self.sched.lock().expect("job_queue mutex poisoned") } - /// Insert a job's nodes into the shared graph, then wake the run loop. + /// Submit a DAG: insert a [`NodeKind::Dag`] **container node** carrying the + /// group's metadata, then insert the template's nodes as its subtree (their + /// roots re-parented to the container). Returns the container's id as the + /// DAG id — its rolled-up state is the DAG state. /// - /// Deliberately named for the [`hive_jobq`] primitive it wraps, because - /// that is nearly all it is. **The wrapper earns its place on the wake**: - /// the crate is sync and runtime-free — it holds no `Notify` at all — so - /// the channel the run loop parks on belongs to the host, and something has - /// to ping it. Left to call sites, an insert whose ping was forgotten would - /// leave a correct DAG sitting unscheduled until an unrelated event - /// happened along; nothing would fail, and no test in isolation would see - /// it. + /// The container is an ordinary node: it declares no resources, so the + /// scheduler claims it on the next pass, runs its (empty) logic and parks + /// it in `Finishing`, at which point its children become runnable. Nothing + /// here completes it by hand — a node with no work of its own still goes + /// the way every other node goes. /// - /// Returns exactly what the primitive returns: the ids of the nodes the - /// template named, in the order it named them. + /// `source` and `reason` are the container node's own payload — they are + /// arguments here rather than fields of a spec struct because that is all + /// they ever were. `declare` is the recipe, taken by generic and run + /// against a builder `hive_jobq` owns: it goes from the template straight + /// into this call, so there is nothing to allocate for. /// /// # Errors /// Propagates a graph-insert error (dependencies that aren't /// dependency-topological). - pub fn insert_job( + pub fn submit( &self, - declare: impl FnOnce(&JobBuilder) -> Vec, - ) -> anyhow::Result> { + source: Source, + reason: String, + declare: impl FnOnce(&JobBuilder), + ) -> anyhow::Result { let mut inner = self.lock(); - let named = inner - .insert_job(None, declare) - .map_err(|e| anyhow::anyhow!("job_queue: graph insert failed: {e}"))?; + let container = inner + .append(NodeKind::Dag { source, reason }, Vec::new(), None) + .map_err(|e| anyhow::anyhow!("job_queue: container insert failed: {e}"))?; + insert_group(&mut inner, declare, Some(container))?; drop(inner); self.notify.notify_one(); - Ok(named) + Ok(container.get()) } /// The scheduler itself, for `hive_jobq`'s run-loop seam @@ -203,16 +235,11 @@ impl JobQueue { &self.sched } - /// The id of the **group root** `node` belongs to, for log lines and the - /// dashboard. Derived from the graph rather than carried alongside the node - /// — the parent axis already knows it. - /// - /// Was `dag_of`, when a job's nodes hung under a container node that *was* - /// the group. Without it the parent chain ends at whichever root the - /// template declared, so this answers "which root owns this node", not - /// "which DAG is this in" — there is no longer such a thing. + /// The DAG container id owning `node`, for log lines and the dashboard. + /// Derived from the graph rather than carried alongside the node — the + /// parent axis already knows it. #[must_use] - pub fn root_of(&self, node: NodeId) -> Option { + pub fn dag_of(&self, node: NodeId) -> Option { self.lock().graph().root_of(node).map(NodeId::get) } @@ -397,9 +424,9 @@ fn find_node(sched: &Sched, id: u64) -> Option { /// root, plus the newest [`MAX_HISTORY_DAGS`] settled ones. /// /// Selected *structurally* — a root is a node with no parent. The typed -/// projection this replaced keyed on the since-removed container kind instead, -/// which made the visible set depend on one host node kind; nothing here knows -/// what a node means. +/// projection this replaced keyed on `NodeKind::Dag` instead, which made the +/// visible set depend on one host node kind; nothing here knows what a node +/// means. /// /// **This bound is load-bearing, not tidiness.** Nothing ever removes a node /// from the graph (bounded pruning is a Stage-C follow-up), so serving diff --git a/hive-c0re/src/job_queue/model.rs b/hive-c0re/src/job_queue/model.rs index 57cf0dca..f1f7936d 100644 --- a/hive-c0re/src/job_queue/model.rs +++ b/hive-c0re/src/job_queue/model.rs @@ -1,18 +1,18 @@ -//! Data model for the generic job-DAG queue: the node kinds — the primitive -//! operations — and what each one carries. The `State` / `PermPayload` wire -//! enums live in `hive_host_sock::jobs` (they travel on the host admin -//! socket) and are re-exported here for the queue's internal use. The graph -//! itself is served through `hive_jobq_wire`'s generic projection — there is -//! no second, typed view of it any more. +//! Data model for the generic job-DAG queue: node kinds (the primitive +//! operations), dependency edges, and the runtime `Dag` / `Node` store. +//! The `Source` / `State` / `PermPayload` wire enums live in +//! `hive_host_sock::jobs` (they travel on the host admin socket) and are +//! re-exported here for the queue's internal use. The graph itself is +//! served through `hive_jobq_wire`'s generic projection — there is no +//! second, typed view of it any more. //! -//! **One level, not two.** The node is the unit of everything: scheduling, -//! execution, build-log, cancel, and the dashboard group (a group root's -//! subtree *is* the group). A DAG used to be a second level above it, with -//! its own store and its own id; there is no container node any more, so a -//! job is exactly the nodes it declared. See `docs/coordinator.md::Job queue` -//! for the full design. +//! Two levels: the **DAG** is the unit of cancel / approval-resolution +//! and the dashboard group; the **node** is the unit of scheduling / +//! execution / build-log, and carries its own `agent` (a +//! DAG can span agents). See `docs/coordinator.md::Job queue` for the +//! full design. -pub use hive_host_sock::jobs::{PermPayload, State}; +pub use hive_host_sock::jobs::{PermPayload, Source, State}; use serde::Serialize; use hive_jobq::TerminalState; @@ -280,6 +280,18 @@ pub enum NodeKind { /// `Prebuild`, but that's a no-op there — the agent is down, so prebuild /// is skipped.) SetWanted { agent: String, up: bool }, + /// The **DAG container** node: one per submitted DAG, carrying the group's + /// domain metadata. Every node hangs *under* it (its subtree), so + /// the container's `NodeId` **is** the DAG id and its rolled-up state **is** + /// the DAG state. Pure grouping — lease- and + /// build-slot-exempt; the executor instant-completes it (`Done`) so it + /// reaches `Finishing` and its children start. + /// + /// No `created_at` here: the graph stamps [`hive_jobq::Node::created_at`] on + /// every node at insert, so the container already has one. A second copy in + /// the payload would be the same instant recorded twice, with only this + /// variant's version reachable to a generic viewer. + Dag { source: Source, reason: String }, } /// How a hive-c0re node describes itself to a generic graph viewer. @@ -350,13 +362,14 @@ impl NodeKind { NodeKind::ResolveApproval { .. } => "resolve_approval", NodeKind::EmitRebuilt { .. } => "emit_rebuilt", NodeKind::SetWanted { .. } => "set_wanted", + NodeKind::Dag { .. } => "dag", } } /// The agent this node targets, or `""` for agentless kinds /// ([`NodeKind::MetaLock`] on the `hyperhive` pseudo-agent, - /// [`NodeKind::Reparent`] which can span multiple agents, and - /// [`NodeKind::ResolveApproval`] which acts on an approval row). + /// [`NodeKind::Reparent`] which can span multiple agents, and the + /// [`NodeKind::Dag`] container). #[must_use] pub fn agent(&self) -> &str { match self { @@ -384,7 +397,8 @@ impl NodeKind { | NodeKind::SetWanted { agent, .. } => agent, NodeKind::MetaLock { .. } | NodeKind::Reparent { .. } - | NodeKind::ResolveApproval { .. } => "", + | NodeKind::ResolveApproval { .. } + | NodeKind::Dag { .. } => "", } } diff --git a/hive-c0re/src/job_queue/scheduler.rs b/hive-c0re/src/job_queue/scheduler.rs index f60f2028..3c19c529 100644 --- a/hive-c0re/src/job_queue/scheduler.rs +++ b/hive-c0re/src/job_queue/scheduler.rs @@ -93,7 +93,7 @@ pub async fn run_worker(coord: Arc) { let coord = node_coord; async move { tracing::info!( - dag = coord.job_queue.root_of(id).unwrap_or_default(), + dag = coord.job_queue.dag_of(id).unwrap_or_default(), node = id.get(), kind = kind.as_str(), agent = %kind.agent(), diff --git a/hive-c0re/src/job_queue/power.rs b/hive-c0re/src/job_queue/submit.rs similarity index 51% rename from hive-c0re/src/job_queue/power.rs rename to hive-c0re/src/job_queue/submit.rs index b1606a57..3ec3e11b 100644 --- a/hive-c0re/src/job_queue/power.rs +++ b/hive-c0re/src/job_queue/submit.rs @@ -1,27 +1,59 @@ -//! Power ops (`stop` / `start` / `restart`) — the DAG shapes whose per-agent -//! form depends on **live** container state, so they cannot be static -//! [`super::templates`] entries. +//! Request-level submit API — the surface the dashboard POST handlers, +//! the MCP socket handlers, and `hivectl` paths call. //! -//! The split is the purity line, not the subject matter: the `*_chain` / -//! `*_nodes` builders below are pure (they take `running` / `stale` as -//! parameters, which is what keeps them unit-testable without a container), -//! and only the `*_many` entry points do the async `lifecycle::is_running` -//! read that produces those parameters. +//! The **power ops** (`stop` / `start` / `restart`) are built here, not in +//! `templates.rs`: each agent's subgraph shape depends on its *live* running +//! state, which needs an async `lifecycle::is_running` read that a pure/sync +//! template can't do. So these fns are async — they read each agent's state, +//! assemble a per-agent subgraph out of the shared pure primitives +//! (`JobBuilder::node` + `templates::rebuild_nodes`), all declaring into ONE job +//! (independent per-agent roots, concurrent on their own leases). //! -//! Each entry point declares its group and inserts it. There is no metadata -//! parameter and no container node: attribution is not something every caller -//! has to invent, and a DAG is addressed by the nodes a template names. +//! Dynamic shape rule: `stop`/`start` carry a head `SetWanted(w)` (durable +//! intent write) — `restart` does NOT (it bounces the container but leaves +//! `wanted` untouched, so a deliberately-stopped agent isn't forced up). The +//! tail `Reconcile` (the convergence guarantee — cheap, noops when already +//! converged) is ALWAYS present; only the *mechanical* nodes +//! (`Signal`/`Drain`/`StopForUpdate`) are state-conditional (skipped for a +//! down agent — nothing to quiesce/stop). Keeping `Reconcile` in every shape +//! closes the TOCTOU window: if an agent flips state between the `is_running` +//! read and node execution, the tail `Reconcile` still converges it in-DAG, +//! with `StopForUpdate`-noop as the backstop — no reliance on an external +//! reconcile sweep. Every helper emits a fresh queue snapshot so the +//! dashboard shows the new DAG immediately. use std::sync::Arc; -use hive_jobq::NodeId; - +use super::model::NodeKind; use super::resource::Resource; use super::templates::rebuild_nodes; -use super::{JobBuilder, NodeKind}; +use super::{JobBuilder, Source, templates}; use crate::coordinator::Coordinator; use crate::lifecycle; +fn submit_and_emit( + coord: &Arc, + source: Source, + reason: String, + declare: impl FnOnce(&JobBuilder), +) -> u64 { + let id = coord + .job_queue + .submit(source, reason, declare) + .expect("template-declared shapes are acyclic"); + coord.emit_rebuild_queue_snapshot(); + id +} + +/// Manual/approval-independent rebuild (always relocks the agent's +/// meta input — the meta-update cascade grows its own rebuild subgraphs +/// in-DAG instead of going through this surface). +pub fn rebuild(coord: &Arc, agent: &str, source: Source, reason: String) -> u64 { + submit_and_emit(coord, source, reason, |builder| { + templates::rebuild(builder, agent, true); + }) +} + // ---- dynamic power-op DAG assembly ---------------------------------------- // // The pure per-agent chain builders below take `running` (and `stale`) @@ -35,16 +67,7 @@ use crate::lifecycle; /// actually running (nothing to drain on a down container). The `Reconcile` /// stays even for a down agent so a race-up between the state read and exec /// is still stopped in-DAG. -/// Returns the group root's guid, which is what the caller names so -/// `insert_job` hands its id back — that id is how `hivectl` polls this agent's -/// progress. A chain that returned nothing would insert correctly and leave the -/// caller with nothing to wait on. -fn stop_chain( - builder: &JobBuilder, - agent: &str, - graceful: bool, - running: bool, -) -> hive_jobq::NodeGuid { +fn stop_chain(builder: &JobBuilder, agent: &str, graceful: bool, running: bool) { // `SetWanted` is the group root and owns the agent lease; the mechanical // steps are its children (borrow the lease, run once it reaches `Finishing`, // dep-ordered among themselves). @@ -73,26 +96,13 @@ fn stop_chain( .needs(Resource::Agent(a())) .part_of(wanted); } - wanted.guid() } /// One agent's **start** subgraph. `SetWanted(Up)` head; a down + stale-rev /// agent gets the rebuild subgraph (its tail `Reconcile` starts it on /// current derivations), otherwise a plain `Reconcile` (which starts a down /// agent and noops an already-running one). -/// -/// Returns **every** group root — see [`stop_chain`] for why they are named. -/// -/// ⚠️ More than one in the stale branch: `rebuild_nodes` chains its roots -/// *behind* `SetWanted` with `after_ok`, it does **not** nest them under it. So -/// `SetWanted` rolls up only itself, and naming it alone would report the whole -/// start finished while the rebuild was still running. -fn start_chain( - builder: &JobBuilder, - agent: &str, - running: bool, - stale: bool, -) -> Vec { +fn start_chain(builder: &JobBuilder, agent: &str, running: bool, stale: bool) { let wanted = builder .node(NodeKind::SetWanted { agent: agent.to_owned(), @@ -100,16 +110,10 @@ fn start_chain( }) .needs(Resource::Agent(agent.to_owned())); if !running && stale { - // Rebuild subtree chained behind the `SetWanted` head. `MetaSync`, the - // `AgentWindow` brace and `Reconcile` are their own group roots - // (top-level, per `rebuild_nodes`) — hence all four names. - let roots = rebuild_nodes(builder, agent, true, Some(wanted)); - vec![ - wanted.guid(), - roots.meta_sync.guid(), - roots.agent_window.guid(), - roots.reconcile.guid(), - ] + // Rebuild subtree chained behind the `SetWanted` head. `MetaSync`, + // `Prebuild` + `Reconcile` are their own group roots (top-level, per + // `rebuild_nodes`). + rebuild_nodes(builder, agent, true, Some(wanted)); } else { let _ = builder .node(NodeKind::Reconcile { @@ -117,7 +121,6 @@ fn start_chain( }) .needs(Resource::Agent(agent.to_owned())) .part_of(wanted); - vec![wanted.guid()] } } @@ -130,22 +133,14 @@ fn start_chain( /// before `Reconcile`; a down agent gets just `Reconcile`, which /// converges to intent — a stopped (`wanted = Off`) agent stays stopped, /// a crashed (`wanted = Up`) agent comes back up. -/// -/// Returns the group root's guid — see [`stop_chain`]. -fn restart_chain( - builder: &JobBuilder, - agent: &str, - graceful: bool, - running: bool, -) -> hive_jobq::NodeGuid { +fn restart_chain(builder: &JobBuilder, agent: &str, graceful: bool, running: bool) { let a = || agent.to_owned(); if !running { - // Nothing to bounce — a lone Reconcile converges to intent, and is - // itself the root. - return builder + // Nothing to bounce — a lone Reconcile converges to intent. + let _ = builder .node(NodeKind::Reconcile { agent: a() }) - .needs(Resource::Agent(a())) - .guid(); + .needs(Resource::Agent(a())); + return; } // 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 @@ -181,7 +176,6 @@ fn restart_chain( .needs(Resource::Agent(a())) .part_of(signal) .after_ok(stop); - signal.guid() } else { let stop = builder .node(NodeKind::StopForUpdate { agent: a() }) @@ -190,7 +184,6 @@ fn restart_chain( .node(NodeKind::Reconcile { agent: a() }) .needs(Resource::Agent(a())) .part_of(stop); - stop.guid() } } @@ -208,15 +201,10 @@ fn restart_chain( // subgraph's indices onto another's used to be a function. /// Declare the stop DAG from explicit `(agent, running)` targets. -pub(crate) fn stop_nodes( - builder: &JobBuilder, - targets: &[(String, bool)], - graceful: bool, -) -> Vec { - targets - .iter() - .map(|(agent, running)| stop_chain(builder, agent, graceful, *running)) - .collect() +pub(crate) fn stop_nodes(builder: &JobBuilder, targets: &[(String, bool)], graceful: bool) { + for (agent, running) in targets { + stop_chain(builder, agent, graceful, *running); + } } /// Assemble the start DAG from explicit `(agent, running, stale)` targets. @@ -225,68 +213,76 @@ pub(crate) fn stop_nodes( /// running under its lease, so a down+stale agent that grew a rebuild subgraph /// reports `rebuilding` during its swap and `starting` at its reconcile, /// without the DAG having to guess one label covering every target. -pub(crate) fn start_nodes( - builder: &JobBuilder, - targets: &[(String, bool, bool)], -) -> Vec { - targets - .iter() - .flat_map(|(agent, running, stale)| start_chain(builder, agent, *running, *stale)) - .collect() +pub(crate) fn start_nodes(builder: &JobBuilder, targets: &[(String, bool, bool)]) { + for (agent, running, stale) in targets { + start_chain(builder, agent, *running, *stale); + } } /// Declare the restart DAG from explicit `(agent, running)` targets. -pub(crate) fn restart_nodes( - builder: &JobBuilder, - targets: &[(String, bool)], - graceful: bool, -) -> Vec { - targets - .iter() - .map(|(agent, running)| restart_chain(builder, agent, graceful, *running)) - .collect() +pub(crate) fn restart_nodes(builder: &JobBuilder, targets: &[(String, bool)], graceful: bool) { + for (agent, running) in targets { + restart_chain(builder, agent, graceful, *running); + } } -// ---- entry points --------------------------------------------------------- -// -// Each reads the live state its shape depends on, then declares + inserts. +/// Restart a single agent. Thin wrapper over [`restart_many`]. +pub async fn restart(coord: &Arc, agent: &str, source: Source, reason: String) -> u64 { + restart_many(coord, &[agent.to_owned()], false, source, reason).await +} -/// Restart `agents` in a **single** DAG — one per-agent subgraph each, built -/// from live running state and run concurrently on their own leases. A running -/// agent gets the stop→reconcile chain (`graceful` prepends signal→drain); a -/// down agent gets a lone `Reconcile`. Restart never writes `wanted`, so the +/// Graceful restart of a single agent (signal → drain → stop → reconcile, +/// when running). Thin wrapper over [`restart_many`] with `graceful = true`. +pub async fn graceful_restart( + coord: &Arc, + agent: &str, + source: Source, + reason: String, +) -> u64 { + restart_many(coord, &[agent.to_owned()], true, source, reason).await +} + +/// Restart `agents` (one or many) in a **single** DAG — one per-agent +/// subgraph each, built dynamically from live running state and run +/// concurrently on their own leases. A running agent gets the stop→reconcile +/// chain (`graceful` prepends signal→drain); a down agent gets just a lone +/// `Reconcile` (nothing to stop). Restart never writes `wanted`, so the /// tail `Reconcile` converges each agent to its EXISTING intent — a -/// deliberately-stopped agent stays down. -/// -/// # Errors -/// Propagates a graph-insert error. +/// deliberately-stopped agent stays down. The whole hive-wide +/// `hivectl restart` is one DAG. pub async fn restart_many( coord: &Arc, agents: &[String], graceful: bool, -) -> anyhow::Result> { + source: Source, + reason: String, +) -> u64 { let mut targets = Vec::with_capacity(agents.len()); for agent in agents { targets.push((agent.clone(), lifecycle::is_running(agent).await)); } - let ids = coord - .job_queue - .insert_job(|b| restart_nodes(b, &targets, graceful))?; - coord.emit_rebuild_queue_snapshot(); - Ok(ids) + submit_and_emit(coord, source, reason, |builder| { + restart_nodes(builder, &targets, graceful); + }) } -/// Start `agents` in a **single** DAG. A down agent gets -/// `SetWanted(Up) → Reconcile` (or, rev stale, a rebuild-then-start so it comes -/// up on current derivations); an already-running agent gets the same shape -/// with the reconcile noop'ing. -/// -/// # Errors -/// Propagates a graph-insert error. +/// Start a single agent. Thin wrapper over [`start_many`]. +pub async fn start(coord: &Arc, agent: &str, source: Source, reason: String) -> u64 { + start_many(coord, &[agent.to_owned()], source, reason).await +} + +/// Start `agents` (one or many) in a **single** DAG — one per-agent subgraph +/// each, built dynamically from live state and run concurrently on their own +/// leases. A down agent gets `SetWanted(Up) → Reconcile` (or, rev stale, a +/// rebuild-then-start so it comes up on current derivations); an already- +/// running agent gets `SetWanted(Up) → Reconcile` (the reconcile noops). The +/// whole hive-wide `hivectl start` is one DAG. pub async fn start_many( coord: &Arc, agents: &[String], -) -> anyhow::Result> { + source: Source, + reason: String, +) -> u64 { let current = crate::auto_update::current_flake_rev(&coord.hyperhive_flake); let mut targets = Vec::with_capacity(agents.len()); for agent in agents { @@ -300,30 +296,88 @@ pub async fn start_many( } targets.push((agent.clone(), running, stale)); } - let ids = coord.job_queue.insert_job(|b| start_nodes(b, &targets))?; - coord.emit_rebuild_queue_snapshot(); - Ok(ids) + submit_and_emit(coord, source, reason, |builder| { + start_nodes(builder, &targets); + }) } -/// Stop `agents` in a **single** DAG. A running agent gets -/// `SetWanted(Off) → [Signal → Drain →](graceful) Reconcile`; a down agent -/// skips the pointless quiesce but keeps the `Reconcile` as the race-up -/// backstop. -/// -/// # Errors -/// Propagates a graph-insert error. +/// Hard stop a single agent. Thin wrapper over [`stop_many`]. +pub async fn stop(coord: &Arc, agent: &str, source: Source, reason: String) -> u64 { + stop_many(coord, &[agent.to_owned()], false, source, reason).await +} + +/// Graceful stop of a single agent (signal → drain → reconcile, when +/// running). Thin wrapper over [`stop_many`] with `graceful = true`. +pub async fn graceful_stop( + coord: &Arc, + agent: &str, + source: Source, + reason: String, +) -> u64 { + stop_many(coord, &[agent.to_owned()], true, source, reason).await +} + +/// Stop `agents` (one or many) in a **single** DAG — one per-agent subgraph +/// each, built dynamically from live state and run concurrently on their own +/// leases. A running agent gets `SetWanted(Off) → [Signal → Drain →](graceful) +/// Reconcile`; a down agent gets just `SetWanted(Off) → Reconcile` (skips the +/// pointless quiesce, keeps the Reconcile as the race-up backstop). The whole +/// hive-wide `hivectl stop` is one DAG. pub async fn stop_many( coord: &Arc, agents: &[String], graceful: bool, -) -> anyhow::Result> { + source: Source, + reason: String, +) -> u64 { let mut targets = Vec::with_capacity(agents.len()); for agent in agents { targets.push((agent.clone(), lifecycle::is_running(agent).await)); } - let ids = coord - .job_queue - .insert_job(|b| stop_nodes(b, &targets, graceful))?; - coord.emit_rebuild_queue_snapshot(); - Ok(ids) + submit_and_emit(coord, source, reason, |builder| { + stop_nodes(builder, &targets, graceful); + }) +} + +/// Perm change: commit the JSON file(s) then rebuild. +pub fn perm_change( + coord: &Arc, + agent: &str, + source: Source, + reason: String, + payload: super::PermPayload, +) -> u64 { + submit_and_emit(coord, source, reason, |builder| { + templates::perm_change(builder, agent, payload); + }) +} + +/// Meta-input lock bump; cascade rebuilds fan out on completion. +pub fn meta_update( + coord: &Arc, + inputs: Vec, + source: Source, + reason: String, +) -> u64 { + submit_and_emit(coord, source, reason, |builder| { + templates::meta_update(builder, inputs, None); + }) +} + +/// Topology move(s) as a queue DAG. `moves` is `(child, new_parent)` pairs — +/// one entry for `set-parent`, N for `set-parent-bulk`. Fire-and-forget like +/// everything else in this module: submits and returns a DAG id +/// immediately, the caller learns the outcome async (dashboard job view / +/// `hivectl`'s `QueueDag` poll). Wired from `server.rs`'s `HostRequest:: +/// SetParent` (hivectl) and `dashboard/topology.rs`'s `set-parent`/ +/// `set-parent-bulk` handlers. +pub fn reparent( + coord: &Arc, + moves: Vec<(hive_types::Ident, Option)>, + source: Source, + reason: String, +) -> u64 { + submit_and_emit(coord, source, reason, |builder| { + templates::reparent(builder, moves); + }) } diff --git a/hive-c0re/src/job_queue/templates.rs b/hive-c0re/src/job_queue/templates.rs index b90827b4..536093fc 100644 --- a/hive-c0re/src/job_queue/templates.rs +++ b/hive-c0re/src/job_queue/templates.rs @@ -17,8 +17,8 @@ //! //! The hive-wide **power ops** (`stop` / `start` / `restart`) are NOT here: //! their per-agent shape depends on live running state (an async -//! `lifecycle::is_running` read), so [`super::power`] assembles them out of -//! the primitives this module exports ([`rebuild_nodes`]). +//! `lifecycle::is_running` read), so `submit.rs` assembles them out of the +//! primitives this module exports ([`rebuild_nodes`]). use hive_jobq::TerminalState; @@ -317,22 +317,13 @@ pub(crate) fn deploy_rebuild_nodes(builder: &JobBuilder, agent: &str, approval_i /// children. /// /// Closed by an [`NodeKind::EmitRebuilt`] tail edged onto all three group-roots -/// (`MetaSync`, the `AgentWindow` brace, `Reconcile`) — the brace's roll-up -/// carries the whole `StopForUpdate`→`Swap`→`RebuildBookkeeping` subtree, so -/// those three cover every node. Edging `Reconcile` alone would not do: it is -/// `AfterAny` the brace, so it reaches `Done` even after a failed swap and the -/// tail would report success. -/// -/// Returns those three roots, so a caller that needs to wait on the rebuild can -/// name them. -pub fn rebuild(builder: &JobBuilder, agent: &str, relock: bool) -> Vec { +/// (`MetaSync`, `Prebuild`, `Reconcile`) — `Prebuild`'s roll-up carries the +/// whole `StopForUpdate`→`Swap`→`RebuildBookkeeping` subtree, so those three cover every +/// node. Edging `Reconcile` alone would not do: it is `AfterAny` `Prebuild`, so +/// it reaches `Done` even after a failed swap and the tail would report success. +pub fn rebuild(builder: &JobBuilder, agent: &str, relock: bool) { let roots = rebuild_nodes(builder, agent, relock, None); emit_rebuilt_tails(builder, agent, &roots.all()); - vec![ - roots.meta_sync.guid(), - roots.agent_window.guid(), - roots.reconcile.guid(), - ] } /// Approval-driven deploy (`MergeConfigPr`) as a phase subtree rather than the @@ -492,16 +483,10 @@ pub fn meta_update(builder: &JobBuilder, inputs: Vec, approval_id: Optio /// checks), so a parent move needs no container rebuild to take effect. /// No transient pill either — the node is agentless (no lease to hang one /// off of) and near-instant. No tail node: the write is the whole effect. -/// -/// Returns the single node's guid so a caller can wait on it. -pub fn reparent( - builder: &JobBuilder, - moves: Vec<(hive_types::Ident, Option)>, -) -> hive_jobq::NodeGuid { - builder +pub fn reparent(builder: &JobBuilder, moves: Vec<(hive_types::Ident, Option)>) { + let _reparent = builder .node(NodeKind::Reparent { moves }) - .needs(Resource::MetaWindow) - .guid() + .needs(Resource::MetaWindow); } // The boot is assembled inline in `workers/auto_update.rs::submit_boot_tree` diff --git a/hive-c0re/src/job_queue/tests.rs b/hive-c0re/src/job_queue/tests.rs index 41e45018..4bc199a9 100644 --- a/hive-c0re/src/job_queue/tests.rs +++ b/hive-c0re/src/job_queue/tests.rs @@ -4,7 +4,7 @@ //! that graph (wire projection, history retention, error truncation). //! //! **Nothing here runs a node.** Everything a template declares is in the graph -//! the moment `insert` returns, so the assertions read it there. Whether the +//! the moment `submit` returns, so the assertions read it there. Whether the //! scheduler then honours those declarations — cascade, roll-up, grant //! borrow/release, fairness, the `Finishing` gate — is `hive_jobq`'s property //! and is tested in `hive_jobq`, against its own primitives rather than through @@ -17,64 +17,36 @@ use super::model::NodeKind; use super::*; -/// Insert a declared job, naming nothing — the shape assertions read the whole -/// graph. A test that needs a handle calls `q.insert` directly and names the -/// node it cares about. -fn insert(q: &JobQueue, declare: impl FnOnce(&JobBuilder)) { - q.insert_job(|b| { - declare(b); - Vec::new() - }) - .expect("valid shape"); -} - -/// Insert a declared job and hand back the ids of the nodes it **named**, in -/// the order it named them. -/// -/// This is the handle that replaced the DAG id: there is no container to point -/// at any more, so a test that needs to cancel a job or read its state names -/// the roots it cares about — exactly what production does with the ids -/// [`JobQueue::insert_job`] returns. -/// Handed back as raw `u64`, the same form production passes to `cancel` and -/// `node_subtrees` — a `NodeId` cannot be fabricated, so the read surface takes -/// raw ids and searches for them. -fn insert_named( - q: &JobQueue, - declare: impl FnOnce(&JobBuilder) -> Vec, -) -> Vec { - q.insert_job(declare) +/// Submit a declared shape with the metadata every mechanics test uses. +/// `Source::Manual` because none of these exercise provenance — the tests that +/// do name their own source at the call site. +fn submit(q: &JobQueue, reason: &str, declare: impl FnOnce(&JobBuilder)) -> u64 { + q.submit(Source::Manual, reason.to_owned(), declare) .expect("valid shape") - .into_iter() - .map(NodeId::get) - .collect() } fn ident(s: &str) -> hive_types::Ident { hive_types::Ident::parse(s).expect("valid test ident") } -fn rebuild(builder: &JobBuilder, agent: &str) -> Vec { - templates::rebuild(builder, agent, true) +fn rebuild(builder: &JobBuilder, agent: &str) { + templates::rebuild(builder, agent, true); } /// Restart shape with every agent treated as **running** — the online /// shape (`[Signal→Drain→] StopForUpdate → Reconcile`, no `SetWanted` head) /// most queue-mechanics tests assume. Mirrors the pre-dynamic -/// `templates::restart` (which is now the state-aware `power::restart_nodes`). -fn restart_online( - builder: &JobBuilder, - agents: &[&str], - graceful: bool, -) -> Vec { +/// `templates::restart` (which is now the state-aware `submit::restart_nodes`). +fn restart_online(builder: &JobBuilder, agents: &[&str], graceful: bool) { let targets: Vec<(String, bool)> = agents.iter().map(|a| ((*a).to_owned(), true)).collect(); - power::restart_nodes(builder, &targets, graceful) + submit::restart_nodes(builder, &targets, graceful); } /// Stop shape with every agent treated as **running** — the online shape /// (`SetWanted → [Signal→Drain→](graceful) Reconcile`). -fn stop_online(builder: &JobBuilder, agents: &[&str], graceful: bool) -> Vec { +fn stop_online(builder: &JobBuilder, agents: &[&str], graceful: bool) { let targets: Vec<(String, bool)> = agents.iter().map(|a| ((*a).to_owned(), true)).collect(); - power::stop_nodes(builder, &targets, graceful) + submit::stop_nodes(builder, &targets, graceful); } // `Claimed` / `ClaimReady` / `CompleteNode` lived here: a claim snapshot type @@ -90,8 +62,8 @@ fn stop_online(builder: &JobBuilder, agents: &[&str], graceful: bool) -> Vec, /// Kinds this node declared a node-dep on, in declaration order, each with /// the outcome set that satisfies it. @@ -124,35 +96,33 @@ fn when_tag(when: hive_jobq::DepWhen) -> String { .join("|") } -/// Every work node in the graph, in insertion order, as its declared shape. +/// Every work node under `dag`, in insertion order, as its declared shape. /// /// **This is what the template tests are actually about.** A template's output -/// is fully determined the moment `insert` returns: the kinds, the parent +/// is fully determined the moment `submit` returns: the kinds, the parent /// nesting and the dep edges are all sitting in the graph. Reading them here /// keeps the assertion on c0re's own product. Whether the scheduler then /// *honours* those edges — runs a chain serially, holds a grant across a /// subtree — is `hive_jobq`'s property and is tested in `hive_jobq`. -fn declared_shape(q: &JobQueue) -> Vec { - declared_shape_filtered(q, &|_| true) +fn declared_shape(q: &JobQueue, dag: u64) -> Vec { + declared_shape_filtered(q, dag, &|_| true) } -/// Every node in the graph, since each test inserts into a fresh [`JobQueue`]. -/// -/// This used to take a DAG id and filter by `root_of(n) == Some(container)`. -/// With no container node there is no per-DAG root to filter on — and nothing -/// to exclude either, because every node in the graph is now real work. Tests -/// that insert more than one job name a node per job and assert on the ids -/// [`JobQueue::insert`] hands back. -fn declared_shape_filtered(q: &JobQueue, keep: &dyn Fn(&NodeKind) -> bool) -> Vec { +fn declared_shape_filtered( + q: &JobQueue, + dag: u64, + keep: &dyn Fn(&NodeKind) -> bool, +) -> Vec { let sched = q.sched().lock().expect("job_queue mutex poisoned"); let graph = sched.graph(); + let root = graph.resolve_id(dag).expect("dag id is a real node id"); let kind_of = |id: NodeId| graph.node(id).map(|n| n.payload.as_str()); graph .nodes() - .filter(|n| keep(&n.payload)) + .filter(|n| n.id != root && graph.root_of(n.id) == Some(root) && keep(&n.payload)) .map(|n| Declared { kind: n.payload.as_str(), - parent: n.parent.and_then(kind_of), + parent: n.parent.filter(|p| *p != root).and_then(kind_of), after: n .deps .iter() @@ -165,66 +135,78 @@ fn declared_shape_filtered(q: &JobQueue, keep: &dyn Fn(&NodeKind) -> bool) -> Ve .collect() } -/// The id of the one node of `kind`, for the resource assertions. +/// The id of the one node of `kind` under `dag`, for the resource assertions. /// /// Panics unless there is exactly one — every caller is about a shape where the /// kind is unique, so two would mean the assertion had quietly stopped being /// about the node the test names. -fn node_of(q: &JobQueue, kind: &str) -> hive_jobq::NodeId { +fn node_of(q: &JobQueue, dag: u64, kind: &str) -> hive_jobq::NodeId { let sched = q.sched().lock().expect("job_queue mutex poisoned"); let graph = sched.graph(); + let root = graph.resolve_id(dag).expect("dag id is a real node id"); let mut found: Vec<_> = graph .nodes() - .filter(|n| n.payload.as_str() == kind) + .filter(|n| n.id != root && graph.root_of(n.id) == Some(root) && n.payload.as_str() == kind) .map(|n| n.id) .collect(); assert_eq!( found.len(), 1, - "expected exactly one {kind} node in the graph" + "expected exactly one {kind} node in the dag" ); found.pop().expect("checked above") } -/// The payload of the one node of `kind`, for assertions about what a node -/// *carries* rather than how it is wired. -fn payload_of(q: &JobQueue, kind: &str) -> NodeKind { - let id = node_of(q, kind); +/// The payload of the one node of `kind` under `dag`, for assertions about what +/// a node *carries* rather than how it is wired. +fn payload_of(q: &JobQueue, dag: u64, kind: &str) -> NodeKind { + let id = node_of(q, dag, kind); let sched = q.sched().lock().expect("job_queue mutex poisoned"); sched.graph().node(id).expect("node exists").payload.clone() } -/// Kinds of every node still `Pending` — the nodes that could yet run. -/// Stronger than asking the scheduler what is *ready right now*: a node +/// Kinds of every node under `dag` still `Pending` — the nodes that could yet +/// run. Stronger than asking the scheduler what is *ready right now*: a node /// blocked on a dep is not ready but is very much still alive. -fn pending_kinds(q: &JobQueue) -> Vec<&'static str> { - pending_kinds_filtered(q, &|_| true) +fn pending_kinds(q: &JobQueue, dag: u64) -> Vec<&'static str> { + pending_kinds_filtered(q, dag, &|_| true) } /// [`pending_kinds`] restricted to the nodes whose payload names `agent`. -fn pending_kinds_for(q: &JobQueue, agent: &str) -> Vec<&'static str> { - pending_kinds_filtered(q, &|kind: &NodeKind| kind.agent() == agent) +fn pending_kinds_for(q: &JobQueue, dag: u64, agent: &str) -> Vec<&'static str> { + pending_kinds_filtered(q, dag, &|kind: &NodeKind| kind.agent() == agent) } -/// The payloads of every node still `Pending`, for the cases where *which* of a -/// family of same-kind nodes survived is the assertion — a template emits one -/// tail per outcome and they differ only in what they carry. -fn pending_payloads(q: &JobQueue) -> Vec { +/// The payloads of every node under `dag` still `Pending`, for the cases where +/// *which* of a family of same-kind nodes survived is the assertion — a +/// template emits one tail per outcome and they differ only in what they carry. +fn pending_payloads(q: &JobQueue, dag: u64) -> Vec { let sched = q.sched().lock().expect("job_queue mutex poisoned"); let graph = sched.graph(); + let root = graph.resolve_id(dag).expect("dag id is a real node id"); graph .nodes() - .filter(|n| n.state == State::Pending) + .filter(|n| n.id != root && graph.root_of(n.id) == Some(root) && n.state == State::Pending) .map(|n| n.payload.clone()) .collect() } -fn pending_kinds_filtered(q: &JobQueue, keep: &dyn Fn(&NodeKind) -> bool) -> Vec<&'static str> { +fn pending_kinds_filtered( + q: &JobQueue, + dag: u64, + keep: &dyn Fn(&NodeKind) -> bool, +) -> Vec<&'static str> { let sched = q.sched().lock().expect("job_queue mutex poisoned"); let graph = sched.graph(); + let root = graph.resolve_id(dag).expect("dag id is a real node id"); graph .nodes() - .filter(|n| n.state == State::Pending && keep(&n.payload)) + .filter(|n| { + n.id != root + && graph.root_of(n.id) == Some(root) + && n.state == State::Pending + && keep(&n.payload) + }) .map(|n| n.payload.as_str()) .collect() } @@ -263,19 +245,20 @@ fn declared_resources(q: &JobQueue, node_id: hive_jobq::NodeId) -> Vec .collect() } -/// The resources declared by **every** node of `kind`, one row per +/// The resources declared by **every** node of `kind` under `dag`, one row per /// node, sorted so the rows read as a set rather than an insertion order. /// /// The per-agent templates emit several nodes of one kind — one per agent — and /// what makes them concurrent is that each holds only its *own* agent's lease. /// That is a statement about the whole family, so it needs all the rows, not /// [`declared_resources`]'s single node. -fn declared_resources_of_kind(q: &JobQueue, kind: &str) -> Vec> { +fn declared_resources_of_kind(q: &JobQueue, dag: u64, kind: &str) -> Vec> { let sched = q.sched().lock().expect("job_queue mutex poisoned"); let graph = sched.graph(); + let root = graph.resolve_id(dag).expect("dag id is a real node id"); let mut rows: Vec> = graph .nodes() - .filter(|n| n.payload.as_str() == kind) + .filter(|n| n.id != root && graph.root_of(n.id) == Some(root) && n.payload.as_str() == kind) .map(|n| { n.deps .iter() @@ -296,8 +279,8 @@ fn declared_resources_of_kind(q: &JobQueue, kind: &str) -> Vec> { /// cannot tell them apart — two `set_wanted` rows look identical. Slicing by /// agent is what makes "the fresh agent goes straight to reconcile while the /// stale one rebuilds first" expressible as a declared shape. -fn declared_shape_for(q: &JobQueue, agent: &str) -> Vec { - declared_shape_filtered(q, &|kind: &NodeKind| kind.agent() == agent) +fn declared_shape_for(q: &JobQueue, dag: u64, agent: &str) -> Vec { + declared_shape_filtered(q, dag, &|kind: &NodeKind| kind.agent() == agent) } fn state_of(q: &JobQueue, dag_id: u64) -> State { @@ -324,68 +307,63 @@ fn dag_count(q: &JobQueue) -> usize { .count() } -// ---- insert (dedup removed — every insert is a fresh job) ---- -// -// `submit_assigns_distinct_ids` lived here and is gone. Its whole body was -// "two inserts get different container ids" — an assertion about the id -// allocation of a node type this issue deleted, and in any case `hive_jobq`'s -// property rather than c0re's. What the tests below keep is the part that was -// about c0re: **no dedup**, now read off the group count instead of off an id. +// ---- submit (dedup removed — every submit is a fresh DAG) ---- -/// Insert-time dedup was removed with the agent-per-node refactor (a -/// multi-agent job has no single agent to key a dedup on), so an identical -/// re-insert — same template + agent, still queued — now enqueues a distinct -/// group instead of collapsing into the pending one. Whether any dedup needs +#[test] +fn submit_assigns_distinct_ids() { + let q = JobQueue::new(1); + let first = submit(&q, "first", |builder| rebuild(builder, "agent-a")); + let second = submit(&q, "second", |builder| rebuild(builder, "agent-b")); + assert_ne!(first, second); + assert_eq!(dag_count(&q), 2); +} + +/// Submit-time dedup was removed with the agent-per-node refactor (a +/// multi-agent DAG has no single agent to key a dedup on), so an identical +/// resubmit — same template + agent, still queued — now enqueues a distinct +/// DAG instead of collapsing into the pending one. Whether any dedup needs /// reintroducing is tracked as a follow-up. #[test] fn identical_resubmit_is_a_distinct_dag() { let q = JobQueue::new(1); - let first = insert_named(&q, |builder| rebuild(builder, "agent-a")); - let after_first = dag_count(&q); - let resubmit = insert_named(&q, |builder| rebuild(builder, "agent-a")); - assert_ne!( - first, resubmit, - "no dedup: an identical re-insert declares its own nodes" - ); - // Counted as a doubling rather than against a literal: one rebuild is - // several group roots now, and pinning the number here would make this - // test fail on any shape change while saying nothing about dedup. - assert_eq!( - dag_count(&q), - after_first * 2, - "the re-insert added its own roots instead of collapsing into the pending ones" - ); + let first = submit(&q, "first", |builder| rebuild(builder, "agent-a")); + let resubmit = submit(&q, "again", |builder| rebuild(builder, "agent-a")); + assert_ne!(first, resubmit, "no dedup: identical resubmit is a new DAG"); + assert_eq!(dag_count(&q), 2); } #[test] fn distinct_submits_never_collapse() { let q = JobQueue::new(1); - let rebuild_a = insert_named(&q, |builder| rebuild(builder, "agent-a")); - let one_rebuild = dag_count(&q); - let rebuild_b = insert_named(&q, |builder| rebuild(builder, "agent-b")); - let two_rebuilds = dag_count(&q); - let restart_a = insert_named(&q, |builder| restart_online(builder, &["agent-a"], false)); + let rebuild_a = submit(&q, "r", |builder| rebuild(builder, "agent-a")); + let rebuild_b = submit(&q, "r", |builder| rebuild(builder, "agent-b")); + let restart_a = submit(&q, "r", |builder| { + restart_online(builder, &["agent-a"], false); + }); assert_ne!(rebuild_a, rebuild_b); assert_ne!(rebuild_a, restart_a); - assert_eq!( - two_rebuilds, - one_rebuild * 2, - "two rebuilds, nothing merged" - ); - assert_eq!( - dag_count(&q), - two_rebuilds + restart_a.len(), - "a restart of an agent that already has a queued rebuild is still its own group" - ); + assert_eq!(dag_count(&q), 3); } -// `resubmit_while_running_is_new_dag` lived here: the same two inserts as -// above, kept under a second name so a reader looking for "a config bump -// mid-build must not be swallowed" would find it. It asserted nothing the -// test above doesn't — `insert_job` never consults the state of an existing -// node, so "while running" could not change the outcome and was never staged. -// The scenario is named in that test's doc instead; a duplicate test is a -// second place for the same fact to rot. +#[test] +fn resubmit_while_running_is_new_dag() { + // The "while running" is not load-bearing and used to be staged by claiming + // a node first. `submit` appends a container and inserts the declared + // group; it never consults the state of any existing node, so whether an + // earlier DAG is running cannot change the outcome. What is actually being + // asserted — no dedup, ever — is `identical_resubmit_is_a_distinct_dag`. + // + // Kept as the *named* case because "a config bump mid-build must not be + // swallowed" is the scenario people worry about, and a reader looking for + // it should find it. + let q = JobQueue::new(1); + let a = submit(&q, "first", |builder| rebuild(builder, "agent-a")); + let again = submit(&q, "config bumped during build", |builder| { + rebuild(builder, "agent-a"); + }); + assert_ne!(a, again); + assert_eq!(dag_count(&q), 2); +} // ---- malformed specs: no longer expressible ---- // @@ -409,7 +387,7 @@ fn distinct_submits_never_collapse() { #[test] fn rebuild_chain_is_declared_serial() { // Was `rebuild_chain_claims_in_dep_order`, which drove the whole DAG to - // observe an order that is fully declared the moment `insert` returns. + // observe an order that is fully declared the moment `submit` returns. // // ⚠️ The old name was also wrong about the mechanism, and reading it rather // than the graph is how you'd stay wrong: **only part of this chain is dep @@ -423,11 +401,9 @@ fn rebuild_chain_is_declared_serial() { // The non-graceful shape asserted here has no quiesce chain, so nothing here // is concurrent — but the name would mislead about the graceful one. let q = JobQueue::new(1); - insert(&q, |builder| { - rebuild(builder, "agent-a"); - }); + let id = submit(&q, "r", |builder| rebuild(builder, "agent-a")); assert_eq!( - declared_shape(&q), + declared_shape(&q, id), vec![ row("meta_sync", None, &[]), // The brace: holds the lease + slot for everything nested below it. @@ -492,14 +468,20 @@ fn rebuild_chain_is_declared_serial() { #[test] fn graceful_rebuild_chain_drains_before_stopping() { let q = JobQueue::new(1); - insert(&q, |builder| { - templates::graceful_rebuild_nodes(builder, "agent-a", true, None); - }); + let id = q + .submit( + Source::AutoUpdate, + "sweep".to_owned(), + |builder: &JobBuilder| { + templates::graceful_rebuild_nodes(builder, "agent-a", true, None); + }, + ) + .expect("valid shape"); // Asserted as full rows, not just kinds: the kind list is identical whether // the quiesce chain runs beside the build or nested under it, so a // kind-only assertion cannot see the bug this shape exists to fix. assert_eq!( - declared_shape(&q), + declared_shape(&q, id), vec![ row("meta_sync", None, &[]), row("agent_window", None, &[("meta_sync", "done")]), @@ -548,11 +530,11 @@ fn non_graceful_rebuild_has_no_signal_or_drain() { // job keeps its nodes to itself and inserts them, so what it built is // observable where it matters — in what the scheduler runs. let q = JobQueue::new(1); - insert(&q, |builder| { + let id = submit(&q, "manual", |builder| { templates::rebuild_nodes(builder, "agent-a", true, None); }); assert_eq!( - declared_shape(&q) + declared_shape(&q, id) .iter() .map(|d| d.kind) .collect::>(), @@ -585,10 +567,8 @@ fn rebuild_chain_declares_its_resources_on_the_brace() { // `a_contended_resource_goes_to_the_oldest_waiter`. What is c0re's is // *which* nodes contend in the first place — a declaration, asserted here.) let q = JobQueue::new(1); - insert(&q, |builder| { - rebuild(builder, "agent-a"); - }); - let res = |kind: &str| declared_resources(&q, node_of(&q, kind)); + let id = submit(&q, "r", |builder| rebuild(builder, "agent-a")); + let res = |kind: &str| declared_resources(&q, node_of(&q, id, kind)); let agent = || Resource::Agent("agent-a".to_owned()); assert_eq!( @@ -630,25 +610,20 @@ fn rebuild_chain_declares_its_resources_on_the_brace() { // ---- per-agent lease ---- #[test] -fn multi_agent_restart_declares_concurrent_per_agent_subgraphs() { +fn multi_agent_restart_is_one_dag_with_concurrent_per_agent_subgraphs() { let q = JobQueue::new(4); - let roots = insert_named(&q, |builder| { - restart_online(builder, &["agent-a", "agent-b"], false) + let id = submit(&q, "hive-wide", |builder| { + restart_online(builder, &["agent-a", "agent-b"], false); }); - // Was "a hive-wide restart is ONE DAG": one container over both agents. - // With the container gone it is one *insert* over N independent groups — - // which is the same claim about the operator's action and a better one - // about the graph, since independence is what lets them run at once. - // Asserted against the named count rather than a literal: the point is - // that every root the job named is top-level, with nothing above it. - assert_eq!(dag_count(&q), roots.len(), "every named root is top-level"); + // A hive-wide restart is ONE DAG, not one-per-agent. + assert_eq!(dag_count(&q), 1); // Each agent's subgraph head (StopForUpdate, since both are running) is a // group root with no deps, so nothing orders them against each other; and // each declares only its OWN agent's lease, so nothing makes them contend. // Those two declared facts are what "they run concurrently" *means* here — // that a scheduler then does run independent, resource-disjoint roots at // once is hive_jobq's property, tested there. - let heads: Vec<_> = declared_shape(&q) + let heads: Vec<_> = declared_shape(&q, id) .into_iter() .filter(|d| d.kind == "stop_for_update") .collect(); @@ -661,7 +636,7 @@ fn multi_agent_restart_declares_concurrent_per_agent_subgraphs() { "both per-agent heads are independent group roots" ); assert_eq!( - declared_resources_of_kind(&q, "stop_for_update"), + declared_resources_of_kind(&q, id, "stop_for_update"), vec![ vec![Resource::Agent("agent-a".to_owned())], vec![Resource::Agent("agent-b".to_owned())], @@ -671,18 +646,18 @@ fn multi_agent_restart_declares_concurrent_per_agent_subgraphs() { } #[test] -fn multi_agent_stop_declares_concurrent_per_agent_subgraphs() { +fn multi_agent_stop_is_one_dag_with_concurrent_per_agent_subgraphs() { let q = JobQueue::new(4); - let roots = insert_named(&q, |builder| { - stop_online(builder, &["agent-a", "agent-b"], false) + let id = submit(&q, "hive-wide stop", |builder| { + stop_online(builder, &["agent-a", "agent-b"], false); }); - // See the restart case above for why this is a named-root count now. - assert_eq!(dag_count(&q), roots.len(), "every named root is top-level"); + // A hive-wide stop is ONE DAG, not one-per-agent. + assert_eq!(dag_count(&q), 1); // Same declared story as the restart case above: each agent's subgraph head // is a group root with no node-deps, holding only its own agent's lease. // Independent roots on disjoint resources is what "concurrently" means at // this layer — the running of them is hive_jobq's. - let heads: Vec<_> = declared_shape(&q) + let heads: Vec<_> = declared_shape(&q, id) .into_iter() .filter(|d| d.kind == "set_wanted") .collect(); @@ -692,7 +667,7 @@ fn multi_agent_stop_declares_concurrent_per_agent_subgraphs() { "both per-agent stop subgraph heads are independent group roots" ); assert_eq!( - declared_resources_of_kind(&q, "set_wanted"), + declared_resources_of_kind(&q, id, "set_wanted"), vec![ vec![Resource::Agent("agent-a".to_owned())], vec![Resource::Agent("agent-b".to_owned())], @@ -702,30 +677,27 @@ fn multi_agent_stop_declares_concurrent_per_agent_subgraphs() { } #[test] -fn multi_agent_start_folds_per_agent_stale_rebuild() { +fn multi_agent_start_one_dag_folds_per_agent_stale_rebuild() { let q = JobQueue::new(4); // fresh: offline + not stale → SetWanted → Reconcile. // stale: offline + stale → SetWanted → «rebuild subgraph». - let roots = insert_named(&q, |builder| { - power::start_nodes( + let id = submit(&q, "hive-wide start", |builder| { + submit::start_nodes( builder, &[ ("fresh".to_owned(), false, false), ("stale".to_owned(), false, true), ], - ) + ); }); - // One insert spanning both agents — and an *uneven* number of roots, which - // is the shape this test is about: the fresh agent names one, the stale one - // names four (its rebuild chains behind `SetWanted` rather than nesting - // under it, so the head alone would report the start done mid-rebuild). - assert_eq!(dag_count(&q), roots.len(), "every named root is top-level"); + // One DAG spanning both agents. + assert_eq!(dag_count(&q), 1); // The fold is a *declared* difference, readable the moment submit returns: // both agents get a `SetWanted(Up)` group root, but the fresh agent's // subgraph ends at the Reconcile behind it while the stale agent's carries // the whole rebuild chain in between. assert_eq!( - declared_shape_for(&q, "fresh"), + declared_shape_for(&q, id, "fresh"), vec![ row("set_wanted", None, &[]), row("reconcile", Some("set_wanted"), &[]), @@ -733,7 +705,7 @@ fn multi_agent_start_folds_per_agent_stale_rebuild() { "a fresh agent is intent + convergence, nothing in between" ); assert_eq!( - declared_shape_for(&q, "stale"), + declared_shape_for(&q, id, "stale"), vec![ row("set_wanted", None, &[]), row("meta_sync", None, &[("set_wanted", "done")]), @@ -769,32 +741,31 @@ fn offline_agents_skip_mechanical_nodes_but_keep_reconcile() { // read and node exec. let q = JobQueue::new(4); // Offline graceful stop → SetWanted(Off) → Reconcile (no Signal/Drain). - let stop = insert_named(&q, |builder| { - power::stop_nodes(builder, &[("down".to_owned(), false)], true) + let stop = submit(&q, "stop down", |builder| { + submit::stop_nodes(builder, &[("down".to_owned(), false)], true); }); // Offline restart → a lone Reconcile (no SetWanted, no StopForUpdate): // nothing to bounce, and restart never rewrites intent, so the tail // Reconcile converges the down agent to its existing `wanted`. - let restart = insert_named(&q, |builder| { - power::restart_nodes(builder, &[("down2".to_owned(), false)], true) + let restart = submit(&q, "restart down", |builder| { + submit::restart_nodes(builder, &[("down2".to_owned(), false)], true); }); - // The whole group, **root included** — the root is a work node now - // (`SetWanted` for the stop, the lone `Reconcile` for the restart), not a - // container to be filtered out. One id per agent, which is what the chain - // named. - let shape = |roots: &[u64]| -> Vec { - q.node_subtrees(roots) + let shape = |id: u64| -> Vec { + // The group's work nodes: its subtree minus the container itself, + // which the generic view carries as an ordinary node. + q.node_subtrees(&[id]) .iter() + .filter(|n| n.id != id) .map(|n| n.payload.label.clone()) .collect() }; assert_eq!( - shape(&stop), + shape(stop), vec!["set_wanted".to_owned(), "reconcile".to_owned()], "offline graceful stop skips the signal/drain quiesce, keeps Reconcile" ); assert_eq!( - shape(&restart), + shape(restart), vec!["reconcile".to_owned()], "offline restart is a lone Reconcile (no SetWanted head, nothing to stop)" ); @@ -811,16 +782,22 @@ fn boot_sweep_nodes_declare_their_own_resources() { // meta commit inside another node's staged deploy window. Nothing failed to // compile; only an exhaustive caller list would have caught it. let q = JobQueue::new(4); - insert(&q, |builder| { - crate::workers::auto_update::boot_nodes( - builder, - true, - vec!["stale-agent".to_owned()], - vec!["drifted-agent".to_owned()], - ); - }); + let id = q + .submit( + Source::AutoUpdate, + "boot".to_owned(), + |builder: &JobBuilder| { + crate::workers::auto_update::boot_nodes( + builder, + true, + vec!["stale-agent".to_owned()], + vec!["drifted-agent".to_owned()], + ); + }, + ) + .expect("valid shape"); - let mut lock = declared_resources(&q, node_of(&q, "meta_lock")); + let mut lock = declared_resources(&q, node_of(&q, id, "meta_lock")); lock.sort_by_key(|r| format!("{r:?}")); assert_eq!( lock, @@ -829,7 +806,7 @@ fn boot_sweep_nodes_declare_their_own_resources() { ); assert_eq!( - declared_resources(&q, node_of(&q, "reconcile")), + declared_resources(&q, node_of(&q, id, "reconcile")), vec![Resource::Agent("drifted-agent".to_owned())], "a boot Reconcile touches the container, so it holds that agent's lease" ); @@ -924,10 +901,8 @@ fn rebuild_reconcile_waits_for_the_whole_build_subtree() { // **did** flatten this chain, and the guarantee survives only because the // roll-up point moved with it. let q = JobQueue::new(1); - insert(&q, |builder| { - rebuild(builder, "agent-a"); - }); - let shape = declared_shape(&q); + let id = submit(&q, "r", |builder| rebuild(builder, "agent-a")); + let shape = declared_shape(&q, id); let parent_of = |kind: &str| { shape .iter() @@ -998,7 +973,7 @@ fn rebuild_reconcile_waits_for_the_whole_build_subtree() { #[test] fn a_fanned_out_mechanical_node_declares_its_agent_lease() { let q = JobQueue::new(4); - insert(&q, |builder| { + let id = submit(&q, "fan-out", |builder| { templates::fanned_out_mechanical( builder, NodeKind::Start { @@ -1006,9 +981,9 @@ fn a_fanned_out_mechanical_node_declares_its_agent_lease() { }, ); }); - assert_eq!(declared_shape(&q), vec![row("start", None, &[])]); + assert_eq!(declared_shape(&q, id), vec![row("start", None, &[])]); assert_eq!( - declared_resources(&q, node_of(&q, "start")), + declared_resources(&q, node_of(&q, id, "start")), vec![Resource::Agent("agent-a".to_owned())], "the fanned-out node carries the lease itself, rather than relying on \ whoever happened to fan it out" @@ -1033,13 +1008,19 @@ fn a_fanned_out_mechanical_node_declares_its_agent_lease() { fn a_meta_lock_grows_one_rebuild_subgraph_per_agent() { let q = JobQueue::new(4); let agents = vec!["alice".to_owned(), "bob".to_owned()]; - insert(&q, |builder| { - templates::grown_graceful_rebuilds(builder, &agents, true); - }); + let id = q + .submit( + Source::AutoUpdate, + "sweep".to_owned(), + |builder: &JobBuilder| { + templates::grown_graceful_rebuilds(builder, &agents, true); + }, + ) + .expect("valid shape"); // One chain per agent, each an independent group root — so the two rebuild // concurrently, each on its own lease. - let shape = declared_shape(&q); + let shape = declared_shape(&q, id); let heads: Vec<_> = shape .iter() .filter(|d| d.kind == "meta_sync") @@ -1065,40 +1046,21 @@ fn a_meta_lock_grows_one_rebuild_subgraph_per_agent() { #[test] fn cancel_clears_queued_dag() { let q = JobQueue::new(1); - let roots = insert_named(&q, |builder| rebuild(builder, "agent-a")); - let [head, brace, tail] = roots.as_slice() else { - panic!("a rebuild names three roots, got {roots:?}") - }; - assert!(q.cancel(*head), "fully-queued dag cancels"); - // The operator sees `Cancelled` the moment the cancel returns — a group must - // not read `Queued` back to the operator who just cancelled it (the - // dashboard renders this roll-up from the snapshot - // `post_rebuild_queue_cancel` emits synchronously). - assert_eq!(state_of(&q, *head), State::Cancelled, "no stale Queued gap"); - // Both `EmitRebuilt` tails go with it — the ok one is `AFTER_OK`, the - // failure one keys on elimination — so a rebuild that never ran emits - // nothing. - // - // 🎯 **But `Reconcile` survives, and that is the finding.** Its edge onto - // the brace accepts `done|failed|skipped`, and a cancel-cascade *skips* the - // brace rather than cancelling it — so the edge is satisfied and the tail - // stays claimable. With a container above them all, one cancel took the - // whole group; a job is its roots now, and dropping it means dropping every - // id the insert returned. That is what the ids are for. - assert_eq!( - pending_kinds(&q), - vec!["reconcile"], - "the convergence tail outlives its head's cancel" - ); + let id = submit(&q, "r", |builder| rebuild(builder, "agent-a")); + assert!(q.cancel(id), "fully-queued dag cancels"); + // The operator sees `Cancelled` the moment the cancel returns — the spared + // tail is still `Pending`, and a DAG must not read `Queued` back to the + // operator who just cancelled it (the dashboard renders this roll-up from + // the snapshot `post_rebuild_queue_cancel` emits synchronously). + assert_eq!(state_of(&q, id), State::Cancelled, "no stale Queued gap"); + // Neither `EmitRebuilt` tail accepts a *dropped* dependency — the ok one is + // `AFTER_OK`, the failure one keys on elimination — so both are cancelled + // with the work and **nothing is left that could still run**: no node is + // spared, so a rebuild that never ran emits nothing. assert!( - !q.cancel(*brace), - "the brace was eliminated with the head — there is nothing left to cancel" - ); - assert!(q.cancel(*tail), "the surviving tail cancels on its own id"); - assert!( - pending_kinds(&q).is_empty(), - "cancelling every named root leaves nothing alive, got {:?}", - pending_kinds(&q) + pending_kinds(&q, id).is_empty(), + "a dropped rebuild leaves nothing alive, got {:?}", + pending_kinds(&q, id) ); } @@ -1113,26 +1075,36 @@ fn cancel_clears_queued_dag() { #[test] fn cancel_drops_one_agents_branch_leaving_the_rest() { let q = JobQueue::new(2); - let roots = insert_named(&q, |builder| { - restart_online(builder, &["agent-a", "agent-b"], false) + let id = submit(&q, "r", |builder| { + restart_online(builder, &["agent-a", "agent-b"], false); }); - // One root per agent, **in the order the chain named them** — that ordering - // is `insert_job`'s contract, and it is what replaced digging the right - // subgraph out of a snapshot by matching on its payload's agent field. - let [a_root, _b_root] = roots.as_slice() else { - panic!("a two-agent restart names one root per agent, got {roots:?}") - }; + // Per-agent subgraphs hang directly off the container, one per agent. + // ⚠️ `parent` is the *graph* parent here, not a DAG-relative one: what + // the typed view called a parentless group root is a direct child of the + // container node in the generic view. + let snap = q.node_subtrees(&[id]); + let a_root = snap + .iter() + .find(|n| { + n.parent == Some(id) + && n.payload + .data + .get("agent") + .and_then(serde_json::Value::as_str) + == Some("agent-a") + }) + .expect("agent-a has a group root"); - assert!(q.cancel(*a_root), "an interior/group root cancels alone"); + assert!(q.cancel(a_root.id), "an interior/group root cancels alone"); // agent-a's subgraph is gone; agent-b's is untouched and still alive. assert!( - pending_kinds_for(&q, "agent-a").is_empty(), + pending_kinds_for(&q, id, "agent-a").is_empty(), "agent-a's branch was dropped whole, got {:?}", - pending_kinds_for(&q, "agent-a") + pending_kinds_for(&q, id, "agent-a") ); assert!( - !pending_kinds_for(&q, "agent-b").is_empty(), + !pending_kinds_for(&q, id, "agent-b").is_empty(), "agent-b's branch survives its sibling's cancel" ); } @@ -1152,26 +1124,23 @@ fn cancel_drops_one_agents_branch_leaving_the_rest() { /// deliberately-stopped as far as reconcile and crash-watch are concerned. #[test] fn cancelled_power_op_runs_no_compensating_node() { - /// Insert-cancel-assert for one power op. Taking the roots the job already - /// named is what removes the need to put three differently-typed recipes in - /// one array: each caller inserts its own, so no closure type has to be - /// erased to a boxed one. - fn assert_cancels_clean(q: &JobQueue, roots: &[u64], writes_intent: bool, case: &str) { - // Read the intent head off the inserted nodes rather than out of a - // spec: a declared job holds its own nodes and inserts them. The queue - // is fresh per case, so the whole graph is this one op. - let has_intent = declared_shape(q).iter().any(|d| d.kind == "set_wanted"); + /// Submit-cancel-assert for one power op. Taking the already-submitted DAG + /// id is what removes the need to put three differently-typed recipes in + /// one array: each caller submits its own spec, so no closure type has to + /// be erased to a boxed one. + fn assert_cancels_clean(q: &JobQueue, id: u64, writes_intent: bool, case: &str) { + // Read the intent head off the submitted DAG rather than out of the + // spec: a declared job holds its own nodes and inserts them. + let has_intent = declared_shape(q, id).iter().any(|d| d.kind == "set_wanted"); assert_eq!(has_intent, writes_intent, "{case}: intent head"); - for root in roots { - assert!(q.cancel(*root), "{case}: cancelled while queued"); - assert_eq!(state_of(q, *root), State::Cancelled); - } + assert!(q.cancel(id), "{case}: cancelled while queued"); + assert_eq!(state_of(q, id), State::Cancelled); // Nothing is left that *could* run. Asserting on the pending set rather // than on "what is ready this instant" also covers a node that is alive // but blocked — which is exactly what a leftover compensating node // would look like. assert_eq!( - pending_kinds(q), + pending_kinds(q, id), Vec::<&str>::new(), "{case}: a power op emits no tail node, so a cancelled one leaves nothing" ); @@ -1183,20 +1152,22 @@ fn cancelled_power_op_runs_no_compensating_node() { let case = format!("graceful={graceful} running={running}"); let q = JobQueue::new(1); - let roots = insert_named(&q, |builder| { - power::restart_nodes(builder, &targets, graceful) + let id = submit(&q, "bounce", |builder| { + submit::restart_nodes(builder, &targets, graceful); }); - assert_cancels_clean(&q, &roots, false, &format!("restart {case}")); + assert_cancels_clean(&q, id, false, &format!("restart {case}")); let q = JobQueue::new(1); - let roots = insert_named(&q, |builder| power::stop_nodes(builder, &targets, graceful)); - assert_cancels_clean(&q, &roots, true, &format!("stop {case}")); + let id = submit(&q, "stop", |builder| { + submit::stop_nodes(builder, &targets, graceful); + }); + assert_cancels_clean(&q, id, true, &format!("stop {case}")); let q = JobQueue::new(1); - let roots = insert_named(&q, |builder| { - power::start_nodes(builder, &[("agent-a".to_owned(), running, false)]) + let id = submit(&q, "start", |builder| { + submit::start_nodes(builder, &[("agent-a".to_owned(), running, false)]); }); - assert_cancels_clean(&q, &roots, true, &format!("start {case}")); + assert_cancels_clean(&q, id, true, &format!("start {case}")); } } } @@ -1214,21 +1185,16 @@ fn cancelled_power_op_runs_no_compensating_node() { #[test] fn cancelled_dag_still_runs_its_approval_tail() { let q = JobQueue::new(1); - insert(&q, |builder| { + let id = submit(&q, "approval #7", |builder| { templates::approval_deploy(builder, "agent-a", 7); }); - // Found by kind, not returned: `approval_deploy` deliberately names - // nothing, because nothing polls it — the approval row is how an operator - // follows a deploy, so the template is fire-and-forget in production and a - // test must not make it return an id it would otherwise have no use for. - let window = node_of(&q, "deploy_window").get(); - assert!(q.cancel(window), "fully-queued dag cancels"); + assert!(q.cancel(id), "fully-queued dag cancels"); // The `Cancelled` tail is the only node whose edge accepts a dropped // dependency, so it is the only one `cancel` spares — and *which* tail // survives is the whole assertion: the template emits one per outcome and // the spared one names how the approval row is about to be resolved. // Nothing computes it, so reading the survivor is reading the answer. - let spared = pending_payloads(&q); + let spared = pending_payloads(&q, id); assert!( matches!( spared.as_slice(), @@ -1239,20 +1205,18 @@ fn cancelled_dag_still_runs_its_approval_tail() { ), "only the cancelled-outcome tail is spared, got {spared:?}" ); - // ⚠️ `Cancelled`, and it reads terminal **while the spared tail is still - // pending** — the one place this differs from the container era, where the - // cancel landed on a node *above* the window and the window rolled up - // `Finishing`. Here the operator cancels the window itself, so its own - // state is `Cancelled` however its subtree is doing. Deliberately asserted - // rather than routed around: the group's card goes terminal while a - // bookkeeping node runs on. That is acceptable for the tail this test - // protects (it resolves the approval row and nothing waits on it), and it - // would not be for work an operator expects to still be watching. - assert_eq!(state_of(&q, window), State::Cancelled); + // ⚠️ `Finishing`, not `Cancelled` — and the change is a **fix**, not a + // regression. This used to read the host-side `DagView::rollup_state`, + // which flattened the spared tail away and reported the group settled + // while a node of it was still pending. The root's own state is the + // scheduler's answer: `Finishing` means "own logic done, children still + // running", and the tail this test exists to protect *is* such a child. + // A group that still has work to do does not read terminal. + assert_eq!(state_of(&q, id), State::Finishing); // An unrelated group landing in the same graph doesn't disturb this one's // state — a root rolls up its own subtree, not the graph. - let _other = insert_named(&q, |builder| rebuild(builder, "agent-b")); - assert_eq!(state_of(&q, window), State::Cancelled); + let _other = submit(&q, "r", |builder| rebuild(builder, "agent-b")); + assert_eq!(state_of(&q, id), State::Finishing); } // ---- approval deploy subtree ---- @@ -1265,12 +1229,12 @@ fn cancelled_dag_still_runs_its_approval_tail() { #[test] fn deploy_dag_runs_phases_in_order_and_tails_a_failed_apply() { let q = JobQueue::new(1); - insert(&q, |builder| { + let id = submit(&q, "approval #7", |builder| { templates::approval_deploy(builder, "agent-a", 7); }); assert_eq!( - declared_shape(&q), + declared_shape(&q, id), vec![ // The window is the group root and holds the meta window for the // whole subtree; the three phases are its sub-nodes. @@ -1323,12 +1287,12 @@ fn deploy_apply_grows_rebuild_subgraph_and_finalizes_after_it() { // children run (`a_completing_node_grows_the_work_it_declared`, // `parent_parks_in_finishing_until_children_roll_up`). let q = JobQueue::new(1); - insert(&q, |builder| { + let id = submit(&q, "deploy graft", |builder| { templates::deploy_rebuild_nodes(builder, "agent-a", 11); }); assert_eq!( - declared_shape(&q), + declared_shape(&q, id), vec![ row("meta_sync", None, &[]), row("agent_window", None, &[("meta_sync", "done")]), @@ -1503,11 +1467,11 @@ fn error_truncation_cuts_on_a_char_boundary() { #[test] fn graceful_stop_shape_signal_drain_reconcile() { let q = JobQueue::new(1); - insert(&q, |builder| { + let id = submit(&q, "graceful", |builder| { stop_online(builder, &["agent-a"], true); }); assert_eq!( - declared_shape(&q), + declared_shape(&q, id), vec![ // The whole stop hangs under `set_wanted`: the durable intent is // written first, and the mechanical steps are its sub-nodes. @@ -1530,8 +1494,8 @@ fn graceful_stop_shape_signal_drain_reconcile() { // `exec.rs`. assert_eq!( [ - declared_resources(&q, node_of(&q, "signal")), - declared_resources(&q, node_of(&q, "drain")), + declared_resources(&q, node_of(&q, id, "signal")), + declared_resources(&q, node_of(&q, id, "drain")), ], [vec![], vec![]], "the quiesce pair borrows the brace's lease and declares nothing" @@ -1541,11 +1505,11 @@ fn graceful_stop_shape_signal_drain_reconcile() { #[test] fn spawn_shape_provision_create_dropin_reconcile() { let q = JobQueue::new(1); - insert(&q, |builder| { + let id = submit(&q, "approval #7 spawn", |builder| { templates::spawn(builder, "newbie", 7); }); assert_eq!( - declared_shape(&q), + declared_shape(&q, id), vec![ row("provision", None, &[]), row("create", Some("provision"), &[]), @@ -1566,7 +1530,7 @@ fn spawn_shape_provision_create_dropin_reconcile() { #[test] fn perm_change_shape_prefixes_rebuild_chain() { let q = JobQueue::new(1); - insert(&q, |builder| { + let id = submit(&q, "perm", |builder| { templates::perm_change( builder, "agent-a", @@ -1577,7 +1541,7 @@ fn perm_change_shape_prefixes_rebuild_chain() { ); }); assert_eq!( - declared_shape(&q) + declared_shape(&q, id) .iter() .map(|d| d.kind) .collect::>(), @@ -1605,15 +1569,15 @@ fn reparent_shape_is_a_lone_agentless_meta_window_node() { // `MetaLock`, and it must declare the meta window — a topology commit // must not land inside another node's staged deploy window. let q = JobQueue::new(1); - insert(&q, |builder| { + let id = submit(&q, "set-parent", |builder| { templates::reparent(builder, vec![(ident("alice"), Some(ident("bob")))]); }); assert_eq!( - declared_shape(&q), + declared_shape(&q, id), vec![row("reparent", None, &[])], "one node, no rebuild subgraph" ); - let node = node_of(&q, "reparent"); + let node = node_of(&q, id, "reparent"); assert_eq!( declared_resources(&q, node), vec![Resource::MetaWindow], @@ -1629,15 +1593,15 @@ fn reparent_bulk_shape_carries_every_move_on_one_node() { // request is the reason a single node was chosen in the first place. let moves = vec![(ident("alice"), Some(ident("bob"))), (ident("carol"), None)]; let q = JobQueue::new(1); - insert(&q, |builder| { + let id = submit(&q, "set-parent-bulk", |builder| { templates::reparent(builder, moves.clone()); }); assert_eq!( - declared_shape(&q), + declared_shape(&q, id), vec![row("reparent", None, &[])], "one node for the whole request, not one per move" ); - let NodeKind::Reparent { moves: got } = payload_of(&q, "reparent") else { + let NodeKind::Reparent { moves: got } = payload_of(&q, id, "reparent") else { panic!("expected a Reparent node"); }; assert_eq!(got, moves, "every move rides the single node"); diff --git a/hive-c0re/src/server.rs b/hive-c0re/src/server.rs index 8e534ba1..26799a9a 100644 --- a/hive-c0re/src/server.rs +++ b/hive-c0re/src/server.rs @@ -180,19 +180,13 @@ async fn dispatch(req: &HostRequest, coord: Arc) -> HostResponse { // submit returns a DAG id immediately, the caller polls // `QueueDag` (`hivectl`'s wait/progress loop) for the // outcome instead of blocking here on the commit. - let inserted = coord.job_queue.insert_job(|b| { - vec![crate::job_queue::templates::reparent( - b, - vec![(child.clone(), new_parent.clone())], - )] - }); - match inserted { - Ok(ids) => { - coord.emit_rebuild_queue_snapshot(); - HostResponse::queued(ids.into_iter().map(hive_jobq::NodeId::get).collect()) - } - Err(e) => HostResponse::error(format!("queue reparent: {e}")), - } + let id = crate::job_queue::submit::reparent( + &coord, + vec![(child.clone(), new_parent.clone())], + crate::job_queue::Source::Manual, + "manual set-parent via hivectl".to_owned(), + ); + HostResponse::queued(vec![id]) } HostRequest::SetResourceLimits { name, @@ -783,36 +777,39 @@ enum Verb { } async fn submit_single(coord: &Arc, name: &str, verb: Verb) -> HostResponse { - use crate::job_queue::power; - // A single-target op is the N-target one with N = 1 — the shapes are - // identical, so there is no separate builder to keep in sync. - let targets = [name.to_owned()]; - let ids = match verb { + use crate::job_queue::{Source, submit}; + let id = match verb { Verb::Kill => { tracing::info!(%name, "kill"); - power::stop_many(coord, &targets, false).await + submit::stop( + coord, + name, + Source::Manual, + "manual kill via hivectl".to_owned(), + ) + .await } Verb::Restart => { tracing::info!(%name, "restart"); - power::restart_many(coord, &targets, false).await + submit::restart( + coord, + name, + Source::Manual, + "manual restart via hivectl".to_owned(), + ) + .await } Verb::Rebuild => { tracing::info!(%name, "rebuild"); - // Not a power op: a rebuild's shape doesn't depend on live state, - // so it is a plain template insert rather than a `*_many` gather. - let inserted = coord - .job_queue - .insert_job(|b| crate::job_queue::templates::rebuild(b, name, true)); - if inserted.is_ok() { - coord.emit_rebuild_queue_snapshot(); - } - inserted + submit::rebuild( + coord, + name, + Source::Manual, + "manual rebuild via hivectl".to_owned(), + ) } }; - match ids { - Ok(ids) => HostResponse::queued(ids.into_iter().map(hive_jobq::NodeId::get).collect()), - Err(e) => HostResponse::error(format!("queue insert failed: {e}")), - } + HostResponse::queued(vec![id]) } /// Stop the given `agents` (resolved logical names) then `infra` containers @@ -840,17 +837,27 @@ async fn handle_stop( let mut errors: Vec = Vec::new(); let mut queued: Vec = Vec::new(); - // One insert for all targeted agents — a per-agent stop subgraph each + // One DAG for all targeted agents — a per-agent stop subgraph each // (`SetWanted(Offline) → [Signal → Drain →] Reconcile`), independent - // roots that run concurrently on their own leases. + // roots that run concurrently on their own leases. A hive-wide + // `hivectl stop` is now a single DAG, not N. if !agents.is_empty() { - match crate::job_queue::power::stop_many(coord, agents, graceful).await { - Ok(ids) => { - queued.extend(ids.into_iter().map(hive_jobq::NodeId::get)); - ok_items.extend(agents.iter().cloned()); - } - Err(e) => errors.push(format!("queue stop: {e}")), - } + let reason = if graceful { + "manual via hivectl graceful stop" + } else { + "manual via hivectl stop" + }; + queued.push( + crate::job_queue::submit::stop_many( + coord, + agents, + graceful, + crate::job_queue::Source::Manual, + reason.to_owned(), + ) + .await, + ); + ok_items.extend(agents.iter().cloned()); } // Agents go down before infra so they're not mid-request against a @@ -937,13 +944,16 @@ async fn handle_start( // hivectl's wait loop. let mut queued: Vec = Vec::new(); if !agents.is_empty() { - match crate::job_queue::power::start_many(coord, agents).await { - Ok(ids) => { - queued.extend(ids.into_iter().map(hive_jobq::NodeId::get)); - ok_items.extend(agents.iter().cloned()); - } - Err(e) => errors.push(format!("queue start: {e}")), - } + queued.push( + crate::job_queue::submit::start_many( + coord, + agents, + crate::job_queue::Source::Manual, + "manual via hivectl start".to_owned(), + ) + .await, + ); + ok_items.extend(agents.iter().cloned()); } let mut resp = finish_lifecycle(ok_items, &errors); @@ -972,20 +982,28 @@ async fn handle_restart_scoped( let mut errors: Vec = Vec::new(); let mut queued: Vec = Vec::new(); - // One insert for all targeted agents — a per-agent restart subgraph each - // (`[Signal → Drain →] StopForUpdate → Reconcile`; no `SetWanted`, since a - // restart converges to the agent's *existing* intent rather than rewriting - // it), independent roots running concurrently on their own leases. No - // client-side stop-then-start composition — the whole restart survives a - // dropped connection because the graph owns it. + // One DAG for all targeted agents — a per-agent restart subgraph each + // (`SetWanted → [Signal → Drain →] StopForUpdate → Reconcile`), + // independent roots that run concurrently on their own leases. A + // hive-wide `hivectl restart` is now a single DAG, not N. No + // client-side stop-then-start composition — the whole restart survives + // a dropped connection because the DAG owns it. if !agents.is_empty() { - match crate::job_queue::power::restart_many(coord, &agents, graceful).await { - Ok(ids) => { - queued.extend(ids.into_iter().map(hive_jobq::NodeId::get)); - ok_items.extend(agents.iter().cloned()); - } - Err(e) => errors.push(format!("queue restart: {e}")), - } + queued.push( + crate::job_queue::submit::restart_many( + coord, + &agents, + graceful, + crate::job_queue::Source::Manual, + if graceful { + "manual via hivectl restart --graceful".to_owned() + } else { + "manual restart via hivectl restart".to_owned() + }, + ) + .await, + ); + ok_items.extend(agents.iter().cloned()); } for &container in &infra { diff --git a/hive-c0re/src/socket_server/lifecycle_handlers.rs b/hive-c0re/src/socket_server/lifecycle_handlers.rs index c0261405..46446c12 100644 --- a/hive-c0re/src/socket_server/lifecycle_handlers.rs +++ b/hive-c0re/src/socket_server/lifecycle_handlers.rs @@ -20,9 +20,13 @@ pub(super) async fn handle_start(coord: &Arc, agent: &str, name: &s // Persist `wanted = Up` and submit the Start DAG; the submit layer // upgrades a stale-rev start to a full rebuild so the container // runs current nix derivations before it starts. - if let Err(e) = crate::job_queue::power::start_many(coord, &[name.to_owned()]).await { - tracing::error!(%agent, %name, error = ?e, "start: insert failed"); - } + crate::job_queue::submit::start( + coord, + name, + crate::job_queue::Source::Manual, + format!("agent `{agent}` start tool"), + ) + .await; Response::Ok } @@ -44,9 +48,13 @@ pub(super) async fn handle_restart(coord: &Arc, agent: &str, name: return err; } tracing::info!(%agent, %name, "submit restart"); - if let Err(e) = crate::job_queue::power::restart_many(coord, &[name.to_owned()], false).await { - tracing::error!(%agent, %name, error = ?e, "restart: insert failed"); - } + crate::job_queue::submit::restart( + coord, + name, + crate::job_queue::Source::Manual, + format!("agent `{agent}` restart tool"), + ) + .await; Response::Ok } @@ -145,13 +153,12 @@ pub(super) fn handle_update(coord: &Arc, agent: &str, name: &str) - return err; } tracing::info!(%agent, %name, "submit rebuild"); - if let Err(e) = coord.job_queue.insert_job(|b| { - crate::job_queue::templates::rebuild(b, name, true); - Vec::new() - }) { - tracing::error!(%agent, %name, error = ?e, "update: insert failed"); - } - coord.emit_rebuild_queue_snapshot(); + crate::job_queue::submit::rebuild( + coord, + name, + crate::job_queue::Source::Manual, + format!("agent `{agent}` update tool"), + ); Response::Ok } diff --git a/hive-c0re/src/workers/auto_update.rs b/hive-c0re/src/workers/auto_update.rs index 0084b102..8132f0b7 100644 --- a/hive-c0re/src/workers/auto_update.rs +++ b/hive-c0re/src/workers/auto_update.rs @@ -109,11 +109,12 @@ pub async fn ensure_root_agent(coord: &Arc) -> Result<()> { tracing::warn!( "manager container exists but no applied flake — forcing rebuild to migrate" ); - if let Err(e) = coord.job_queue.insert_job(|b| { - crate::job_queue::templates::rebuild(b, MANAGER_NAME, true); - Vec::new() - }) { - tracing::warn!(error = ?e, "manager migration rebuild insert failed"); + if let Err(e) = coord.job_queue.submit( + crate::job_queue::Source::AutoUpdate, + "manager migration: no applied flake".to_owned(), + |b| crate::job_queue::templates::rebuild(b, MANAGER_NAME, true), + ) { + tracing::warn!(error = ?e, "manager migration rebuild submit failed"); } } else { tracing::debug!("manager container already present"); @@ -377,19 +378,18 @@ fn submit_boot_tree( n_deferred: usize, n_skipped: usize, ) { - // Fully-quiet boot (nothing stale, nothing drifted) inserts nothing. + use crate::job_queue::Source; + + // Fully-quiet boot (nothing stale, nothing drifted) submits nothing. if !any_stale && drifted.is_empty() { return; } - // The summary the sweep used to hand the container as its `reason` is a log - // line now: it was only ever stored on a node nobody read, and the counts - // are worth having where they can actually be seen. - tracing::info!( - rebuilds = fanout.len(), - reconciles = drifted.len(), - deferred = n_deferred, - up_to_date = n_skipped, - "boot: sweep" + let reason = format!( + "boot: {} rebuild(s), {} reconcile(s), {} deferred (offline), {} up-to-date", + fanout.len(), + drifted.len(), + n_deferred, + n_skipped, ); // The sweep's own rebuild subgraphs emit their `Rebuilt` events as they @@ -397,11 +397,10 @@ fn submit_boot_tree( // The subgraphs also carry their own per-agent crash-watch suppression // during their `Swap` (applied at claim time); a reconcile-only boot needs // no transient. - if let Err(e) = coord.job_queue.insert_job(|b| { + if let Err(e) = coord.job_queue.submit(Source::AutoUpdate, reason, |b| { boot_nodes(b, any_stale, fanout, drifted); - Vec::new() }) { - tracing::warn!(error = ?e, "boot: sweep DAG insert failed"); + tracing::warn!(error = ?e, "boot: sweep DAG submit failed"); } coord.emit_rebuild_queue_snapshot(); } diff --git a/hive-host-sock/src/jobs.rs b/hive-host-sock/src/jobs.rs index 7930807e..7502277b 100644 --- a/hive-host-sock/src/jobs.rs +++ b/hive-host-sock/src/jobs.rs @@ -1,11 +1,6 @@ -//! Vocabulary hive-c0re's job queue shares with its clients: what a -//! permission change carries ([`PermPayload`]) and the scheduler's -//! lifecycle [`State`]. -//! -//! A `Source` enum lived here too — where a job came from, rendered as the -//! "why" chip. It was a field on the DAG container, and it went with it: a -//! job is its nodes now, and a node says what it does rather than who asked -//! for it. +//! Vocabulary hive-c0re's job queue shares with its clients: where a job +//! came from ([`Source`]), what a permission change carries +//! ([`PermPayload`]), and the scheduler's lifecycle [`State`]. //! //! **The typed `DagView`/`NodeView` projection that used to live here is //! gone.** One graph is served one way now — `hive_jobq_wire`'s generic @@ -17,6 +12,37 @@ use serde::{Deserialize, Serialize}; +/// Where the submit request originated — drives the "why" chip on the +/// dashboard. +#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)] +#[serde(rename_all = "snake_case")] +pub enum Source { + /// Operator action (dashboard button, CLI, manager tool). + Manual, + /// Meta-update cascade rebuild (grown into the meta-update DAG). + MetaUpdate, + /// Boot-time submission (the boot sweep DAG + boot reconciles). + AutoUpdate, + /// Crash recovery path (future use). + CrashRecover, + /// Operator approved a pending `Approval` row; `approval_id` on + /// the DAG points back at the source row. + Approval, +} + +impl Source { + #[must_use] + pub fn as_str(self) -> &'static str { + match self { + Source::Manual => "manual", + Source::MetaUpdate => "meta_update", + Source::AutoUpdate => "auto_update", + Source::CrashRecover => "crash_recover", + Source::Approval => "approval", + } + } +} + pub use hive_jobq::State; /// Kind-specific payload for `Template::PermChange` DAGs.