Compare commits

..
5 changed files with 108 additions and 180 deletions

View file

@ -76,15 +76,12 @@ container build:
The power ops write the durable `wanted` intent via a head `SetWanted` 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 node (not a pre-submit side effect) — it holds the agent lease, so
intent-write + reconcile is atomic per-agent. `restart` takes an agent intent-write + reconcile is atomic per-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.
```text ```text
rebuild(a): Prebuild(a) → StopForUpdate(a) → Swap(a) →(after-any) Reconcile(a) rebuild(a): Prebuild(a) → StopForUpdate(a) → Swap(a) →(after-any) Reconcile(a)
graceful-stop(a): SetWanted(a,Off) → Signal(a) → Drain(a) → 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) 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) 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) stale-start(a): SetWanted(a,Up) → «rebuild subgraph» (prebuild noops — agent is down)

View file

@ -27,27 +27,12 @@ pub fn rebuild(coord: &Arc<Coordinator>, agent: &str, source: Source, reason: St
submit_and_emit(coord, templates::rebuild(agent, source, reason, None, true)) submit_and_emit(coord, templates::rebuild(agent, source, reason, None, true))
} }
/// Restart a single agent: mechanical stop + converge to `wanted = Up`. /// Restart: mechanical stop + converge to `wanted = Up`. The intent
/// The intent write is the template's head `SetWanted(Up)` node — it /// write is the template's head `SetWanted(Up)` node — it matters when
/// matters when `wanted` drifted `Offline` under a running agent (an /// `wanted` drifted `Offline` under a running agent (an operator asking
/// operator asking for a restart plainly wants it running, not a stop). /// 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<Coordinator>, agent: &str, source: Source, reason: String) -> u64 { pub fn restart(coord: &Arc<Coordinator>, agent: &str, source: Source, reason: String) -> u64 {
restart_many(coord, &[agent.to_owned()], false, source, reason) submit_and_emit(coord, templates::restart(agent, 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<Coordinator>,
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 /// Start: `SetWanted(Up)` (a DAG node now) then reconcile. A stale rev
@ -82,18 +67,17 @@ pub fn graceful_stop(coord: &Arc<Coordinator>, agent: &str, source: Source, reas
submit_and_emit(coord, templates::graceful_stop(agent, source, reason)) submit_and_emit(coord, templates::graceful_stop(agent, source, reason))
} }
/// Graceful restart of a single agent: signal → drain → mechanical stop → /// Graceful restart: signal → drain → mechanical stop → reconcile (starts
/// reconcile (starts it back up) — one atomic DAG, no client-side "await /// it back up) — one atomic DAG, no client-side "await the stop DAG then
/// the stop DAG then submit a start DAG" split. The head `SetWanted(Up)` /// submit a start DAG" split. The head `SetWanted(Up)` node writes the
/// node writes the intent as part of the DAG. Thin wrapper over /// intent as part of the DAG.
/// [`restart_many`] with `graceful = true`.
pub fn graceful_restart( pub fn graceful_restart(
coord: &Arc<Coordinator>, coord: &Arc<Coordinator>,
agent: &str, agent: &str,
source: Source, source: Source,
reason: String, reason: String,
) -> u64 { ) -> 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. /// Perm change: commit the JSON file(s) then rebuild.

View file

@ -4,11 +4,9 @@
//! `Vec<Node>` + `deps`). //! `Vec<Node>` + `deps`).
//! //!
//! Every node carries its own `agent` (there is no DAG-level agent) — the //! 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 //! `node` helper stamps the template's agent onto each. Today's templates
//! agent *list* and stamps each agent onto its own subgraph, so a hive-wide //! are single-agent (every node shares one agent); a future multi-agent
//! `hivectl restart` is ONE DAG with N independent per-agent subgraphs //! template would stamp different agents per subgraph.
//! (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 //! The power ops write the durable `wanted` intent via a head
//! `SetWanted(w)` node (not a pre-submit side effect); it holds the agent //! `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 /// Restart: write `wanted = Up` (head `SetWanted` node), mechanical stop,
/// independent subgraph — a head `SetWanted(Up)` (a root: no cross-agent /// then converge — a stop + start like the old `lifecycle::restart`
/// dep) then its restart chain — so all agents' restarts run concurrently, /// regardless of prior intent drift.
/// each taking its own agent lease. `graceful` inserts `Signal → Drain` pub fn restart(agent: &str, source: Source, reason: String) -> DagSpec {
/// 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)));
}
}
DagSpec { DagSpec {
template: if graceful { template: Template::Restart,
Template::GracefulRestart
} else {
Template::Restart
},
source, source,
reason, reason,
parent_id: None, parent_id: None,
@ -185,7 +159,38 @@ pub fn restart(agents: &[String], graceful: bool, source: Source, reason: String
inputs: Vec::new(), inputs: Vec::new(),
perm_payload: None, perm_payload: None,
transient: Some(TransientKind::Restarting), 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)),
],
} }
} }

View file

@ -67,12 +67,7 @@ fn distinct_submits_never_collapse() {
let b = submit(&q, rebuild("agent-b", "r")); let b = submit(&q, rebuild("agent-b", "r"));
let c = submit( let c = submit(
&q, &q,
templates::restart( templates::restart("agent-a", Source::Manual, "r".to_owned()),
&["agent-a".to_owned()],
false,
Source::Manual,
"r".to_owned(),
),
); );
assert_ne!(a, b); assert_ne!(a, b);
assert_ne!(a, c); assert_ne!(a, c);
@ -205,12 +200,7 @@ fn lease_serializes_two_lifecycle_dags_for_same_agent() {
let q = JobQueue::new(4); let q = JobQueue::new(4);
let restart = submit( let restart = submit(
&q, &q,
templates::restart( templates::restart("agent-a", Source::Manual, "restart".to_owned()),
&["agent-a".to_owned()],
false,
Source::Manual,
"restart".to_owned(),
),
); );
let stop = submit( let stop = submit(
&q, &q,
@ -294,60 +284,16 @@ fn agents_do_not_contend_on_each_others_leases() {
let q = JobQueue::new(4); let q = JobQueue::new(4);
submit( submit(
&q, &q,
templates::restart( templates::restart("agent-a", Source::Manual, "r".to_owned()),
&["agent-a".to_owned()],
false,
Source::Manual,
"r".to_owned(),
),
); );
submit( submit(
&q, &q,
templates::restart( templates::restart("agent-b", Source::Manual, "r".to_owned()),
&["agent-b".to_owned()],
false,
Source::Manual,
"r".to_owned(),
),
); );
let claims = q.claim_ready(); let claims = q.claim_ready();
assert_eq!(claims.len(), 2, "different agents run concurrently"); 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 ---- // ---- failure: cancel-downstream + AfterAny ----
#[test] #[test]
@ -568,12 +514,7 @@ fn terminal_dag_reported_exactly_once_and_lease_released() {
let q = JobQueue::new(1); let q = JobQueue::new(1);
let id = submit( let id = submit(
&q, &q,
templates::restart( templates::restart("agent-a", Source::Manual, "r".to_owned()),
&["agent-a".to_owned()],
false,
Source::Manual,
"r".to_owned(),
),
); );
// restart = SetWanted → StopForUpdate → Reconcile; not terminal until // restart = SetWanted → StopForUpdate → Reconcile; not terminal until
// the last node completes. // the last node completes.
@ -681,12 +622,7 @@ fn trim_keeps_terminal_parent_with_live_children() {
); );
let lock = claim_one(&q); let lock = claim_one(&q);
q.complete_node(meta, lock.node_id, Ok(())); q.complete_node(meta, lock.node_id, Ok(()));
let mut child_spec = templates::restart( let mut child_spec = templates::restart("agent-x", Source::MetaUpdate, "cascade".to_owned());
&["agent-x".to_owned()],
false,
Source::MetaUpdate,
"cascade".to_owned(),
);
child_spec.parent_id = Some(meta); child_spec.parent_id = Some(meta);
let child = submit(&q, child_spec); let child = submit(&q, child_spec);
// Churn > MAX_HISTORY_PER_TEMPLATE terminal meta_update DAGs. // Churn > MAX_HISTORY_PER_TEMPLATE terminal meta_update DAGs.

View file

@ -552,29 +552,28 @@ fn submit_single(coord: &Arc<Coordinator>, name: &str, verb: Verb) -> HostRespon
HostResponse::queued(vec![id]) HostResponse::queued(vec![id])
} }
/// Restart every container in **one** DAG — a per-agent restart subgraph /// Restart every container by submitting one restart DAG per agent —
/// each, running concurrently on their own leases (so unrelated agents' /// each serializes on its own lease, so unrelated agents' restarts
/// restarts overlap while nothing races an in-flight rebuild). Returns /// overlap while nothing races an in-flight rebuild. Returns once all
/// once queued; per-node progress surfaces on the single DAG. /// are queued; per-agent results surface on the queue.
async fn handle_restart_all(coord: &Arc<Coordinator>) -> Result<HostResponse> { async fn handle_restart_all(coord: &Arc<Coordinator>) -> Result<HostResponse> {
tracing::info!("restart-all"); tracing::info!("restart-all");
let containers = lifecycle::list().await?; let agents = lifecycle::list().await?;
let agents: Vec<String> = containers let mut ok_agents: Vec<String> = Vec::new();
.iter() let mut queued: Vec<u64> = Vec::new();
.filter_map(|a| a.strip_prefix(lifecycle::AGENT_PREFIX).map(str::to_owned)) for agent in &agents {
.collect(); let Some(logical) = agent.strip_prefix(lifecycle::AGENT_PREFIX) else {
let queued = if agents.is_empty() { continue;
Vec::new() };
} else { queued.push(crate::job_queue::submit::restart(
vec![crate::job_queue::submit::restart_many(
coord, coord,
&agents, logical,
false,
crate::job_queue::Source::Manual, crate::job_queue::Source::Manual,
"manual restart via hivectl restart-all".to_owned(), "manual restart via hivectl restart-all".to_owned(),
)] ));
}; ok_agents.push(logical.to_owned());
let mut resp = HostResponse::list(agents); }
let mut resp = HostResponse::list(ok_agents);
resp.queued_dags = Some(queued); resp.queued_dags = Some(queued);
Ok(resp) Ok(resp)
} }
@ -719,14 +718,16 @@ async fn handle_start(
} }
/// Restart containers hive-wide (`hivectl restart`) — the DAG-based /// Restart containers hive-wide (`hivectl restart`) — the DAG-based
/// sibling of [`handle_stop`]/[`handle_start`]. All targeted agents ride /// sibling of [`handle_stop`]/[`handle_start`], replacing the old
/// **one** DAG (a per-agent restart subgraph each: `SetWanted → [Signal → /// client-side stop-then-start composition (issue tracker "dagify hivectl
/// Drain →] StopForUpdate → Reconcile`, independent roots that run /// commands"). Each targeted agent gets exactly one atomic DAG submitted
/// concurrently on their own leases), not N separate DAGs — a hive-wide /// up front: the `Restart` template (mechanical stop + reconcile) in the
/// restart is one job. `graceful` prepends signal→drain per agent. No /// common case, or `GracefulRestart` (signal → drain → mechanical stop →
/// client-side stop-then-start composition, so a dropped `hivectl` /// reconcile) with `graceful` set — no "submit a DAG, wait for it, submit
/// connection never strands an agent. Infra containers have no lease/DAG /// another" composition on either path, so a dropped `hivectl` connection
/// and restart synchronously (stop then start). /// 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( async fn handle_restart_scoped(
coord: &Arc<Coordinator>, coord: &Arc<Coordinator>,
scope: &LifecycleScope, scope: &LifecycleScope,
@ -739,25 +740,30 @@ async fn handle_restart_scoped(
let mut errors: Vec<String> = Vec::new(); let mut errors: Vec<String> = Vec::new();
let mut queued: Vec<u64> = Vec::new(); let mut queued: Vec<u64> = Vec::new();
// One DAG for all targeted agents — a per-agent restart subgraph each // One atomic DAG per agent, submitted up front — `Restart`
// (`SetWanted → [Signal → Drain →] StopForUpdate → Reconcile`), // (mechanical stop + reconcile) or, with `--graceful`,
// independent roots that run concurrently on their own leases. A // `GracefulRestart` (signal → drain → mechanical stop → reconcile).
// hive-wide `hivectl restart` is now a single DAG, not N. No // No client- or server-side "submit a stop DAG, await it, then
// client-side stop-then-start composition — the whole restart survives // submit a start DAG" composition: that window is exactly the
// a dropped connection because the DAG owns it. // dropped-connection gap this DAG-based path exists to close.
if !agents.is_empty() { for agent in &agents {
queued.push(crate::job_queue::submit::restart_many( let id = if graceful {
crate::job_queue::submit::graceful_restart(
coord, coord,
&agents, agent,
graceful,
crate::job_queue::Source::Manual, crate::job_queue::Source::Manual,
if graceful { "manual via hivectl restart --graceful".to_owned(),
"manual via hivectl restart --graceful".to_owned() )
} else { } else {
"manual restart via hivectl restart".to_owned() crate::job_queue::submit::restart(
}, coord,
)); agent,
ok_items.extend(agents.iter().cloned()); crate::job_queue::Source::Manual,
"manual restart via hivectl restart".to_owned(),
)
};
queued.push(id);
ok_items.push(agent.clone());
} }
for &container in &infra { for &container in &infra {