diff --git a/docs/coordinator.md b/docs/coordinator.md index 7b0d6cff..8670e92a 100644 --- a/docs/coordinator.md +++ b/docs/coordinator.md @@ -76,15 +76,12 @@ container build: The power ops write the durable `wanted` intent via a head `SetWanted` node (not a pre-submit side effect) — it holds the agent lease, so -intent-write + reconcile is atomic per-agent. `restart` takes an agent -*list*: a hive-wide `hivectl restart` / `restart-all` is ONE DAG with a -per-agent restart subgraph each (independent roots, run concurrently on -their own leases), not N separate DAGs. +intent-write + reconcile is atomic per-agent. ```text rebuild(a): Prebuild(a) → StopForUpdate(a) → Swap(a) →(after-any) Reconcile(a) graceful-stop(a): SetWanted(a,Off) → Signal(a) → Drain(a) → Reconcile(a) -restart(a..): per agent: SetWanted(a,Up) → StopForUpdate(a) → Reconcile(a) (N subgraphs, 1 DAG) +restart(a): SetWanted(a,Up) → StopForUpdate(a) → Reconcile(a) graceful-restart(a): SetWanted(a,Up) → Signal(a) → Drain(a) → StopForUpdate(a) → Reconcile(a) start(a): SetWanted(a,Up) → Reconcile(a) (stale rev ⇒ stale-start below) stale-start(a): SetWanted(a,Up) → «rebuild subgraph» (prebuild noops — agent is down) diff --git a/hive-c0re/src/job_queue/submit.rs b/hive-c0re/src/job_queue/submit.rs index 93b718ef..17c83ba6 100644 --- a/hive-c0re/src/job_queue/submit.rs +++ b/hive-c0re/src/job_queue/submit.rs @@ -27,27 +27,12 @@ pub fn rebuild(coord: &Arc, agent: &str, source: Source, reason: St submit_and_emit(coord, templates::rebuild(agent, source, reason, None, true)) } -/// 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. +/// 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). pub fn restart(coord: &Arc, agent: &str, source: Source, reason: String) -> u64 { - 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)) + submit_and_emit(coord, templates::restart(agent, source, reason)) } /// Start: `SetWanted(Up)` (a DAG node now) then reconcile. A stale rev @@ -82,18 +67,17 @@ pub fn graceful_stop(coord: &Arc, agent: &str, source: Source, reas submit_and_emit(coord, templates::graceful_stop(agent, source, reason)) } -/// 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`. +/// 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. pub fn graceful_restart( coord: &Arc, agent: &str, source: Source, reason: String, ) -> u64 { - restart_many(coord, &[agent.to_owned()], true, source, reason) + submit_and_emit(coord, templates::graceful_restart(agent, 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 be258043..37d029c5 100644 --- a/hive-c0re/src/job_queue/templates.rs +++ b/hive-c0re/src/job_queue/templates.rs @@ -4,11 +4,9 @@ //! `Vec` + `deps`). //! //! Every node carries its own `agent` (there is no DAG-level agent) — the -//! `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. +//! `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. //! //! The power ops write the durable `wanted` intent via a head //! `SetWanted(w)` node (not a pre-submit side effect); it holds the agent @@ -148,36 +146,12 @@ pub fn graceful_stop(agent: &str, source: Source, reason: String) -> DagSpec { } } -/// 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))); - } - } +/// 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: if graceful { - Template::GracefulRestart - } else { - Template::Restart - }, + template: Template::Restart, source, reason, parent_id: None, @@ -185,7 +159,38 @@ pub fn restart(agents: &[String], graceful: bool, source: Source, reason: String inputs: Vec::new(), perm_payload: None, transient: Some(TransientKind::Restarting), - nodes, + nodes: vec![ + node(agent, NodeKind::SetWanted { up: true }, Vec::new()), + node(agent, NodeKind::StopForUpdate, after_ok(0)), + node(agent, NodeKind::Reconcile, after_ok(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, + 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::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)), + ], } } diff --git a/hive-c0re/src/job_queue/tests.rs b/hive-c0re/src/job_queue/tests.rs index 9a691ce8..df26da4e 100644 --- a/hive-c0re/src/job_queue/tests.rs +++ b/hive-c0re/src/job_queue/tests.rs @@ -67,12 +67,7 @@ fn distinct_submits_never_collapse() { let b = submit(&q, rebuild("agent-b", "r")); let c = submit( &q, - templates::restart( - &["agent-a".to_owned()], - false, - Source::Manual, - "r".to_owned(), - ), + templates::restart("agent-a", Source::Manual, "r".to_owned()), ); assert_ne!(a, b); assert_ne!(a, c); @@ -205,12 +200,7 @@ fn lease_serializes_two_lifecycle_dags_for_same_agent() { let q = JobQueue::new(4); let restart = submit( &q, - templates::restart( - &["agent-a".to_owned()], - false, - Source::Manual, - "restart".to_owned(), - ), + templates::restart("agent-a", Source::Manual, "restart".to_owned()), ); let stop = submit( &q, @@ -294,60 +284,16 @@ fn agents_do_not_contend_on_each_others_leases() { let q = JobQueue::new(4); submit( &q, - templates::restart( - &["agent-a".to_owned()], - false, - Source::Manual, - "r".to_owned(), - ), + templates::restart("agent-a", Source::Manual, "r".to_owned()), ); submit( &q, - templates::restart( - &["agent-b".to_owned()], - false, - Source::Manual, - "r".to_owned(), - ), + templates::restart("agent-b", Source::Manual, "r".to_owned()), ); let claims = q.claim_ready(); assert_eq!(claims.len(), 2, "different agents run concurrently"); } -#[test] -fn multi_agent_restart_is_one_dag_with_concurrent_per_agent_subgraphs() { - let q = JobQueue::new(4); - let id = submit( - &q, - templates::restart( - &["agent-a".to_owned(), "agent-b".to_owned()], - false, - Source::Manual, - "hive-wide".to_owned(), - ), - ); - // A hive-wide restart is ONE DAG, not one-per-agent. - assert_eq!(q.snapshot().len(), 1); - // Each agent's subgraph head (SetWanted) is a root, so both are - // claimable at once — each takes its OWN agent's lease (no contention - // across distinct agents), all inside the single DAG. - let claims = q.claim_ready(); - assert!(claims.iter().all(|c| c.dag_id == id)); - let mut heads: Vec<(&str, &str, bool)> = claims - .iter() - .map(|c| (c.agent.as_str(), c.kind.as_str(), c.lease_acquired)) - .collect(); - heads.sort_unstable(); - assert_eq!( - heads, - vec![ - ("agent-a", "set_wanted", true), - ("agent-b", "set_wanted", true), - ], - "both per-agent subgraphs start concurrently, each acquiring its own lease" - ); -} - // ---- failure: cancel-downstream + AfterAny ---- #[test] @@ -568,12 +514,7 @@ fn terminal_dag_reported_exactly_once_and_lease_released() { let q = JobQueue::new(1); let id = submit( &q, - templates::restart( - &["agent-a".to_owned()], - false, - Source::Manual, - "r".to_owned(), - ), + templates::restart("agent-a", Source::Manual, "r".to_owned()), ); // restart = SetWanted → StopForUpdate → Reconcile; not terminal until // the last node completes. @@ -681,12 +622,7 @@ fn trim_keeps_terminal_parent_with_live_children() { ); let lock = claim_one(&q); q.complete_node(meta, lock.node_id, Ok(())); - let mut child_spec = templates::restart( - &["agent-x".to_owned()], - false, - Source::MetaUpdate, - "cascade".to_owned(), - ); + let mut child_spec = templates::restart("agent-x", Source::MetaUpdate, "cascade".to_owned()); child_spec.parent_id = Some(meta); let child = submit(&q, child_spec); // Churn > MAX_HISTORY_PER_TEMPLATE terminal meta_update DAGs. diff --git a/hive-c0re/src/server.rs b/hive-c0re/src/server.rs index a50f9481..58da91a7 100644 --- a/hive-c0re/src/server.rs +++ b/hive-c0re/src/server.rs @@ -552,29 +552,28 @@ fn submit_single(coord: &Arc, name: &str, verb: Verb) -> HostRespon HostResponse::queued(vec![id]) } -/// 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. +/// 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. async fn handle_restart_all(coord: &Arc) -> Result { tracing::info!("restart-all"); - 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( + 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( coord, - &agents, - false, + logical, crate::job_queue::Source::Manual, "manual restart via hivectl restart-all".to_owned(), - )] - }; - let mut resp = HostResponse::list(agents); + )); + ok_agents.push(logical.to_owned()); + } + let mut resp = HostResponse::list(ok_agents); resp.queued_dags = Some(queued); Ok(resp) } @@ -719,14 +718,16 @@ async fn handle_start( } /// Restart containers hive-wide (`hivectl restart`) — the DAG-based -/// 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). +/// 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). async fn handle_restart_scoped( coord: &Arc, scope: &LifecycleScope, @@ -739,25 +740,30 @@ async fn handle_restart_scoped( let mut errors: Vec = Vec::new(); let mut queued: Vec = Vec::new(); - // 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()); + // 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()); } for &container in &infra {