From 1739716fa20ea34d8922083651d8dc28be9f3921 Mon Sep 17 00:00:00 2001 From: atlas Date: Tue, 14 Jul 2026 23:00:10 +0200 Subject: [PATCH] feat(#2448): emit one multi-agent DAG for hivectl restart / restart-all MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit A hive-wide restart was N separate single-agent DAGs (one submit::restart per agent). Now that agent is per-node (#2445), make it ONE DAG with a per-agent restart subgraph each. - templates::restart takes an agent list: each agent gets an independent subgraph (a head SetWanted(Up) root, then its restart chain), so the N subgraphs run concurrently on their own leases. One agent = the ordinary single-agent restart; unifies the old restart + graceful_restart fns. - submit::restart / graceful_restart stay as single-agent wrappers over the new submit::restart_many(agents, graceful). - server.rs handle_restart_all + handle_restart_scoped submit one restart_many call instead of looping per agent. Infra containers unchanged (no lease/DAG, synchronous). Scope: restart + restart-all only. Broad stop+start is the same pattern (stop/start templates take agent lists) — a follow-up increment. --- hive-c0re/src/job_queue/submit.rs | 36 ++++++++--- hive-c0re/src/job_queue/templates.rs | 73 ++++++++++----------- hive-c0re/src/server.rs | 96 +++++++++++++--------------- 3 files changed, 105 insertions(+), 100 deletions(-) diff --git a/hive-c0re/src/job_queue/submit.rs b/hive-c0re/src/job_queue/submit.rs index 17c83ba6..93b718ef 100644 --- a/hive-c0re/src/job_queue/submit.rs +++ b/hive-c0re/src/job_queue/submit.rs @@ -27,12 +27,27 @@ pub fn rebuild(coord: &Arc, agent: &str, source: Source, reason: St submit_and_emit(coord, templates::rebuild(agent, source, reason, None, true)) } -/// Restart: mechanical stop + converge to `wanted = Up`. The intent -/// write is the template's head `SetWanted(Up)` node — it matters when -/// `wanted` drifted `Offline` under a running agent (an operator asking -/// for a restart plainly wants it running, not a stop). +/// Restart a single agent: mechanical stop + converge to `wanted = Up`. +/// The intent write is the template's head `SetWanted(Up)` node — it +/// matters when `wanted` drifted `Offline` under a running agent (an +/// operator asking for a restart plainly wants it running, not a stop). +/// Thin wrapper over [`restart_many`] with a one-agent slice. pub fn restart(coord: &Arc, agent: &str, source: Source, reason: String) -> u64 { - submit_and_emit(coord, templates::restart(agent, source, reason)) + restart_many(coord, &[agent.to_owned()], false, source, reason) +} + +/// Restart `agents` (one or many) in a **single** DAG — one per-agent +/// subgraph each, running concurrently on their own leases. `graceful` +/// prepends the signal→drain quiesce per agent. The whole hive-wide +/// `hivectl restart` / `restart-all` is now one DAG instead of N. +pub fn restart_many( + coord: &Arc, + agents: &[String], + graceful: bool, + source: Source, + reason: String, +) -> u64 { + submit_and_emit(coord, templates::restart(agents, graceful, source, reason)) } /// Start: `SetWanted(Up)` (a DAG node now) then reconcile. A stale rev @@ -67,17 +82,18 @@ pub fn graceful_stop(coord: &Arc, agent: &str, source: Source, reas submit_and_emit(coord, templates::graceful_stop(agent, source, reason)) } -/// Graceful restart: signal → drain → mechanical stop → reconcile (starts -/// it back up) — one atomic DAG, no client-side "await the stop DAG then -/// submit a start DAG" split. The head `SetWanted(Up)` node writes the -/// intent as part of the DAG. +/// Graceful restart of a single agent: signal → drain → mechanical stop → +/// reconcile (starts it back up) — one atomic DAG, no client-side "await +/// the stop DAG then submit a start DAG" split. The head `SetWanted(Up)` +/// node writes the intent as part of the DAG. Thin wrapper over +/// [`restart_many`] with `graceful = true`. pub fn graceful_restart( coord: &Arc, agent: &str, source: Source, reason: String, ) -> u64 { - submit_and_emit(coord, templates::graceful_restart(agent, source, reason)) + restart_many(coord, &[agent.to_owned()], true, source, reason) } /// Perm change: commit the JSON file(s) then rebuild. diff --git a/hive-c0re/src/job_queue/templates.rs b/hive-c0re/src/job_queue/templates.rs index 37d029c5..be258043 100644 --- a/hive-c0re/src/job_queue/templates.rs +++ b/hive-c0re/src/job_queue/templates.rs @@ -4,9 +4,11 @@ //! `Vec` + `deps`). //! //! Every node carries its own `agent` (there is no DAG-level agent) — the -//! `node` helper stamps the template's agent onto each. Today's templates -//! are single-agent (every node shares one agent); a future multi-agent -//! template would stamp different agents per subgraph. +//! `node` helper stamps the template's agent onto each. `restart` takes an +//! agent *list* and stamps each agent onto its own subgraph, so a hive-wide +//! `hivectl restart` is ONE DAG with N independent per-agent subgraphs +//! (each a root chain, run concurrently on its own lease) rather than N +//! separate DAGs. The other templates are still single-agent. //! //! The power ops write the durable `wanted` intent via a head //! `SetWanted(w)` node (not a pre-submit side effect); it holds the agent @@ -146,37 +148,36 @@ pub fn graceful_stop(agent: &str, source: Source, reason: String) -> DagSpec { } } -/// Restart: write `wanted = Up` (head `SetWanted` node), mechanical stop, -/// then converge — a stop + start like the old `lifecycle::restart` -/// regardless of prior intent drift. -pub fn restart(agent: &str, source: Source, reason: String) -> DagSpec { - DagSpec { - template: Template::Restart, - source, - reason, - parent_id: None, - approval_id: None, - inputs: Vec::new(), - perm_payload: None, - transient: Some(TransientKind::Restarting), - nodes: vec![ - node(agent, NodeKind::SetWanted { up: true }, Vec::new()), - node(agent, NodeKind::StopForUpdate, after_ok(0)), - node(agent, NodeKind::Reconcile, after_ok(1)), - ], +/// Restart one or more agents in a **single** DAG. Each agent gets an +/// independent subgraph — a head `SetWanted(Up)` (a root: no cross-agent +/// dep) then its restart chain — so all agents' restarts run concurrently, +/// each taking its own agent lease. `graceful` inserts `Signal → Drain` +/// before the mechanical `StopForUpdate` (each agent's subgraph is still one +/// atomic restart). One agent = the ordinary single-agent restart; many = +/// a hive-wide `hivectl restart` as one DAG instead of N separate ones. +pub fn restart(agents: &[String], graceful: bool, source: Source, reason: String) -> DagSpec { + let mut nodes = Vec::new(); + for agent in agents { + let base = u32::try_from(nodes.len()).unwrap_or(u32::MAX); + // Head of this agent's subgraph — a root (empty deps), so the N + // per-agent subgraphs are independent and run concurrently. + nodes.push(node(agent, NodeKind::SetWanted { up: true }, Vec::new())); + if graceful { + nodes.push(node(agent, NodeKind::Signal, after_ok(base))); + nodes.push(node(agent, NodeKind::Drain, after_ok(base + 1))); + nodes.push(node(agent, NodeKind::StopForUpdate, after_ok(base + 2))); + nodes.push(node(agent, NodeKind::Reconcile, after_ok(base + 3))); + } else { + nodes.push(node(agent, NodeKind::StopForUpdate, after_ok(base))); + nodes.push(node(agent, NodeKind::Reconcile, after_ok(base + 1))); + } } -} - -/// Graceful restart: write `wanted = Up` (head `SetWanted`), signal → -/// drain → mechanical stop → converge. One atomic DAG start to finish -/// (no client- or server-side "submit one DAG, await it, submit the next" -/// composition): the `Drain` node is the same bounded harness-checkpoint -/// wait `graceful_stop` uses, then `StopForUpdate` (mechanical, ignores -/// `wanted`) and the tail `Reconcile` (converges to `wanted = Up`, i.e. -/// starts it back up) chain exactly like `restart`'s tail. -pub fn graceful_restart(agent: &str, source: Source, reason: String) -> DagSpec { DagSpec { - template: Template::GracefulRestart, + template: if graceful { + Template::GracefulRestart + } else { + Template::Restart + }, source, reason, parent_id: None, @@ -184,13 +185,7 @@ pub fn graceful_restart(agent: &str, source: Source, reason: String) -> DagSpec inputs: Vec::new(), perm_payload: None, transient: Some(TransientKind::Restarting), - nodes: vec![ - node(agent, NodeKind::SetWanted { up: true }, Vec::new()), - node(agent, NodeKind::Signal, after_ok(0)), - node(agent, NodeKind::Drain, after_ok(1)), - node(agent, NodeKind::StopForUpdate, after_ok(2)), - node(agent, NodeKind::Reconcile, after_ok(3)), - ], + nodes, } } diff --git a/hive-c0re/src/server.rs b/hive-c0re/src/server.rs index 58da91a7..a50f9481 100644 --- a/hive-c0re/src/server.rs +++ b/hive-c0re/src/server.rs @@ -552,28 +552,29 @@ fn submit_single(coord: &Arc, name: &str, verb: Verb) -> HostRespon HostResponse::queued(vec![id]) } -/// Restart every container by submitting one restart DAG per agent — -/// each serializes on its own lease, so unrelated agents' restarts -/// overlap while nothing races an in-flight rebuild. Returns once all -/// are queued; per-agent results surface on the queue. +/// Restart every container in **one** DAG — a per-agent restart subgraph +/// each, running concurrently on their own leases (so unrelated agents' +/// restarts overlap while nothing races an in-flight rebuild). Returns +/// once queued; per-node progress surfaces on the single DAG. async fn handle_restart_all(coord: &Arc) -> Result { tracing::info!("restart-all"); - let agents = lifecycle::list().await?; - let mut ok_agents: Vec = Vec::new(); - let mut queued: Vec = Vec::new(); - for agent in &agents { - let Some(logical) = agent.strip_prefix(lifecycle::AGENT_PREFIX) else { - continue; - }; - queued.push(crate::job_queue::submit::restart( + let containers = lifecycle::list().await?; + let agents: Vec = containers + .iter() + .filter_map(|a| a.strip_prefix(lifecycle::AGENT_PREFIX).map(str::to_owned)) + .collect(); + let queued = if agents.is_empty() { + Vec::new() + } else { + vec![crate::job_queue::submit::restart_many( coord, - logical, + &agents, + false, crate::job_queue::Source::Manual, "manual restart via hivectl restart-all".to_owned(), - )); - ok_agents.push(logical.to_owned()); - } - let mut resp = HostResponse::list(ok_agents); + )] + }; + let mut resp = HostResponse::list(agents); resp.queued_dags = Some(queued); Ok(resp) } @@ -718,16 +719,14 @@ async fn handle_start( } /// Restart containers hive-wide (`hivectl restart`) — the DAG-based -/// sibling of [`handle_stop`]/[`handle_start`], replacing the old -/// client-side stop-then-start composition (issue tracker "dagify hivectl -/// commands"). Each targeted agent gets exactly one atomic DAG submitted -/// up front: the `Restart` template (mechanical stop + reconcile) in the -/// common case, or `GracefulRestart` (signal → drain → mechanical stop → -/// reconcile) with `graceful` set — no "submit a DAG, wait for it, submit -/// another" composition on either path, so a dropped `hivectl` connection -/// never strands an agent, the same way `handle_restart_all` already -/// avoids it. Infra containers have no lease/DAG and restart -/// synchronously (stop then start). +/// sibling of [`handle_stop`]/[`handle_start`]. All targeted agents ride +/// **one** DAG (a per-agent restart subgraph each: `SetWanted → [Signal → +/// Drain →] StopForUpdate → Reconcile`, independent roots that run +/// concurrently on their own leases), not N separate DAGs — a hive-wide +/// restart is one job. `graceful` prepends signal→drain per agent. No +/// client-side stop-then-start composition, so a dropped `hivectl` +/// connection never strands an agent. Infra containers have no lease/DAG +/// and restart synchronously (stop then start). async fn handle_restart_scoped( coord: &Arc, scope: &LifecycleScope, @@ -740,30 +739,25 @@ async fn handle_restart_scoped( let mut errors: Vec = Vec::new(); let mut queued: Vec = Vec::new(); - // One atomic DAG per agent, submitted up front — `Restart` - // (mechanical stop + reconcile) or, with `--graceful`, - // `GracefulRestart` (signal → drain → mechanical stop → reconcile). - // No client- or server-side "submit a stop DAG, await it, then - // submit a start DAG" composition: that window is exactly the - // dropped-connection gap this DAG-based path exists to close. - for agent in &agents { - let id = if graceful { - crate::job_queue::submit::graceful_restart( - coord, - agent, - crate::job_queue::Source::Manual, - "manual via hivectl restart --graceful".to_owned(), - ) - } else { - crate::job_queue::submit::restart( - coord, - agent, - crate::job_queue::Source::Manual, - "manual restart via hivectl restart".to_owned(), - ) - }; - queued.push(id); - ok_items.push(agent.clone()); + // 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() { + 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() + }, + )); + ok_items.extend(agents.iter().cloned()); } for &container in &infra {