From d3d73b5ffb3dd6d57e0b0a0cb15a856b9587eefa Mon Sep 17 00:00:00 2001 From: atlas Date: Sat, 1 Aug 2026 16:06:06 +0200 Subject: [PATCH 1/9] refactor(#2815): derive the transient pill from the running node MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The dashboard pill was declared once per DAG at submit time, so a rebuild reported `rebuilding` for its entire life — through the prebuild, the stop, the swap, the tail and the reconcile. It named the intent of the request, not what was happening. It is now read off the nodes actually running. A node lights a pill when it is `Running` and declares the agent's own resource. Declaring is the test, not targeting: `Prebuild` and `MetaSync` name an agent but are lease-exempt on purpose (the container keeps serving), so they must not light one. It is also not the lease *owner* — `resource_state()` answers "who holds the slot", which is a different question from "what is running", and a descendant that borrows an ancestor's grant never appears in that map. `TransientKind` is gone entirely rather than being re-derived. The label is the node's own wire tag (`NodeKind::as_str`) — the same vocabulary `NodeView.kind` already ships, so a pill and a DAG node name an operation identically and there is no second taxonomy to keep in step. Work with no node behind it (destroy, migration) supplies its own literal. `DagSpec::transient`, `Claim::transient`, `DagMeta::transient` and `NodeKind::Dag`'s `transient` field all go with it. ## the safety half, which is deliberately not the display half `crash_watch::is_deliberate_stop` used to match a `TransientKind` to decide whether a vanished container was intentional or a crash. That made a pill's display vocabulary decide an alerting question, so renaming or adding a label would silently move the alerting boundary. `TransientState` now carries two independent fields: `label` (rendered, nothing branches on it) and `deliberate_stop` (read only by the crash watcher). The producer sets the second, because the producer is the only thing that knows — it is not recoverable from the first. For queue work that value is `NodeKind::takes_container_down()`, and it is emphatically not "holds a lease": `Create` and `Start` hold the agent's lease exactly like `Stop` does, and a container dying *while starting* is a real crash that must keep reporting as one. The default is `false` on purpose — a wrong `false` costs a spurious crash event, a wrong `true` swallows a real crash silently. ## known cost, accepted on the issue A restart no longer reads `restarting`. No `NodeKind` is unique to a restart — `restart_chain` reuses `Signal` / `StopForUpdate` / `Drain` / `Reconcile` — because "restart" is a property of the DAG's shape, not of any node. A restart now reads `signal` / `stop_for_update`, then the agent returns. `Start` / `Stop` / `PostSwap` run inside a lease-holding ancestor and re-declare nothing, so they light no pill and the agent reads idle for those windows. Closing that is the resources-where-constructed work (#2818), not this change. Checked with clippy (`--all-targets -D warnings`), `cargo test -p hive-c0re -p hive-jobq` (321 + 40 passed) and `nix fmt`. --- hive-c0re/src/actions.rs | 6 +- hive-c0re/src/coordinator.rs | 105 +++++++++++----------- hive-c0re/src/dashboard/state_snapshot.rs | 19 +--- hive-c0re/src/dashboard_events.rs | 14 ++- hive-c0re/src/job_queue/exec.rs | 18 ++-- hive-c0re/src/job_queue/mod.rs | 75 ++++++++++------ hive-c0re/src/job_queue/model.rs | 44 +++++++-- hive-c0re/src/job_queue/scheduler.rs | 55 ++++++++---- hive-c0re/src/job_queue/submit.rs | 40 +++------ hive-c0re/src/job_queue/templates.rs | 15 +--- hive-c0re/src/job_queue/tests.rs | 70 ++++++++++----- hive-c0re/src/migrate.rs | 9 +- hive-c0re/src/workers/auto_update.rs | 1 - hive-c0re/src/workers/crash_watch.rs | 85 +++++++----------- 14 files changed, 295 insertions(+), 261 deletions(-) diff --git a/hive-c0re/src/actions.rs b/hive-c0re/src/actions.rs index 4dadeace..b2a354ea 100644 --- a/hive-c0re/src/actions.rs +++ b/hive-c0re/src/actions.rs @@ -8,7 +8,7 @@ use std::sync::Arc; use anyhow::{Context as _, Result, bail}; use hive_sh4re::{ApprovalKind, ApprovalStatus, HelperEvent}; -use crate::coordinator::{Coordinator, TransientKind}; +use crate::coordinator::Coordinator; use crate::lifecycle; /// Approve a pending request. Marks the approval row durably, then @@ -882,7 +882,9 @@ pub async fn destroy(coord: &Arc, name: &str, purge: bool) -> Resul tracing::info!(%name, purge, "destroy"); // Guard auto-clears on the success path's final scope exit and on // every early-return / cancellation along the way. - let guard = coord.transient_guard(name, TransientKind::Destroying); + // Destroy is not a queue node, so it names its own label. `true`: the + // container is going away, so its disappearance must not read as a crash. + let guard = coord.transient_guard(name, "destroying", true); lifecycle::destroy(name).await?; coord.unregister_agent(name); let runtime = crate::paths::agent_runtime_dir(name); diff --git a/hive-c0re/src/coordinator.rs b/hive-c0re/src/coordinator.rs index 7029cbf3..196b95ac 100644 --- a/hive-c0re/src/coordinator.rs +++ b/hive-c0re/src/coordinator.rs @@ -133,7 +133,7 @@ pub struct Coordinator { /// agents whose tombstone is still inside the grace window. Crash /// watcher consults both this and the active map before declaring /// a stop deliberate. - recent_transient: Mutex>, + recent_transient: Mutex>, /// Timestamps of recent unexpected container crashes, keyed by agent. /// Fed by `crash_watch` each time it classifies a stop as a crash (so /// a crash-looping container — which `Restart=on-failure` flips back @@ -357,9 +357,30 @@ pub struct AgentPaths { /// Per-agent in-progress state that the dashboard surfaces between approve /// click and container ready. +/// +/// The two fields answer genuinely different questions and are set +/// independently on purpose. There used to be a single `TransientKind` enum +/// serving both, which meant a display concern and a safety decision shared one +/// vocabulary and moved together. #[derive(Debug, Clone)] pub struct TransientState { - pub kind: TransientKind, + /// What the dashboard pill renders. For queue-driven work this is the + /// running node's own wire tag ([`crate::job_queue::NodeKind::as_str`]) — + /// the same vocabulary the DAG view ships, so a pill and a node name an + /// operation identically. Work with no node behind it (destroy, migration) + /// supplies its own. + /// + /// Display only. Nothing branches on it — match on a string and this + /// becomes a taxonomy again, silently. + pub label: String, + /// Whether the container going down is **expected**, i.e. this operation + /// takes it down on purpose. Read by the crash watcher to tell a + /// deliberate stop from a crash, so a wrong value here either raises a + /// false alarm or swallows a real one. + /// + /// Set by whoever creates the transient, which is the only place that + /// actually knows — it is not recoverable from `label`. + pub deliberate_stop: bool, pub since: std::time::Instant, } @@ -408,38 +429,6 @@ impl Drop for MetaUpdateGuard { } } -#[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Serialize)] -#[serde(rename_all = "snake_case")] -pub enum TransientKind { - /// `lifecycle::spawn` is running (nixos-container create + update + start). - Spawning, - /// `lifecycle::start` is running. - Starting, - /// `lifecycle::kill` is running. - Stopping, - /// A restart (`lifecycle::kill` then `lifecycle::start`) is running. - Restarting, - /// `lifecycle::rebuild` is running (nixos-container update). - Rebuilding, - /// `actions::destroy` is running. - Destroying, -} - -impl TransientKind { - /// Wire/UI label. Matches the strings the dashboard already - /// renders in the transient spinner. - pub fn as_str(self) -> &'static str { - match self { - TransientKind::Spawning => "spawning", - TransientKind::Starting => "starting", - TransientKind::Stopping => "stopping", - TransientKind::Restarting => "restarting", - TransientKind::Rebuilding => "rebuilding", - TransientKind::Destroying => "destroying", - } - } -} - /// Field-named payload for [`Coordinator::emit_approval_resolved`]. /// Mirrors the `ApprovalResolved` dashboard-event fields. `agent` /// borrows from the caller; `approval_kind` / `status` are @@ -1094,11 +1083,12 @@ impl Coordinator { /// is cancelled (HTTP request aborted, runtime shutdown mid-rebuild, /// panic). A bare set with no guaranteed clear would leak the transient /// and leave the dashboard stuck in "rebuilding…" forever. - fn set_transient(&self, name: &str, kind: TransientKind) { + fn set_transient(&self, name: &str, label: String, deliberate_stop: bool) { self.transient.lock().unwrap().insert( name.to_owned(), TransientState { - kind, + label: label.clone(), + deliberate_stop, since: std::time::Instant::now(), }, ); @@ -1113,7 +1103,7 @@ impl Coordinator { self.emit_dashboard_event(DashboardEvent::TransientSet { seq: self.next_seq(), name: name.to_owned(), - transient_kind: kind.as_str(), + transient_kind: label, since_unix, }); } @@ -1130,10 +1120,10 @@ impl Coordinator { // spurious ContainerCrash on every operator stop/restart. // Old entries get reaped lazily on read so the map doesn't // grow unbounded. - self.recent_transient - .lock() - .unwrap() - .insert(name.to_owned(), (state.kind, std::time::Instant::now())); + self.recent_transient.lock().unwrap().insert( + name.to_owned(), + (state.deliberate_stop, std::time::Instant::now()), + ); self.emit_dashboard_event(DashboardEvent::TransientCleared { seq: self.next_seq(), name: name.to_owned(), @@ -1205,20 +1195,19 @@ impl Coordinator { result } - /// Set of agents whose transient was cleared within the last - /// `grace` seconds — i.e. agents the operator just acted on, - /// whose stop the crash watcher should NOT classify as a crash. - /// Lazily reaps entries older than `grace` so the map stays - /// bounded by the active agent count. - pub fn recent_transient_within( - &self, - grace: std::time::Duration, - ) -> HashMap { + /// Per-agent `deliberate_stop` for transients cleared within the last + /// `grace` seconds — i.e. agents the operator just acted on, whose stop the + /// crash watcher should NOT classify as a crash. Lazily reaps entries older + /// than `grace` so the map stays bounded by the active agent count. + /// + /// Carries only the safety bit, not the display label: nothing downstream + /// should be able to re-derive a stop/crash decision from a pill's wording. + pub fn recent_transient_within(&self, grace: std::time::Duration) -> HashMap { let now = std::time::Instant::now(); let mut map = self.recent_transient.lock().unwrap(); map.retain(|_, (_, ts)| now.duration_since(*ts) <= grace); map.iter() - .map(|(k, (kind, _))| (k.clone(), *kind)) + .map(|(k, (deliberate, _))| (k.clone(), *deliberate)) .collect() } @@ -1254,8 +1243,18 @@ impl Coordinator { /// cancelled or panic between set and clear (HTTP handlers, spawned /// tasks). The guard's `Drop` runs even on task cancellation, so /// the dashboard's spinner can't get pinned forever. - pub fn transient_guard(self: &Arc, name: &str, kind: TransientKind) -> TransientGuard { - self.set_transient(name, kind); + /// + /// `label` is what the pill renders; `deliberate_stop` says whether this + /// operation takes the container down on purpose, and is what the crash + /// watcher reads. Only the caller knows the second one — it is not + /// recoverable from the first. + pub fn transient_guard( + self: &Arc, + name: &str, + label: impl Into, + deliberate_stop: bool, + ) -> TransientGuard { + self.set_transient(name, label.into(), deliberate_stop); TransientGuard { coord: self.clone(), name: name.to_owned(), diff --git a/hive-c0re/src/dashboard/state_snapshot.rs b/hive-c0re/src/dashboard/state_snapshot.rs index f86000cb..0016487a 100644 --- a/hive-c0re/src/dashboard/state_snapshot.rs +++ b/hive-c0re/src/dashboard/state_snapshot.rs @@ -214,7 +214,8 @@ struct PortConflict { #[derive(Serialize)] struct TransientView { name: String, - kind: &'static str, + /// Owned: the label is the running node's wire tag, not one of a fixed set. + kind: String, secs: u64, } @@ -533,26 +534,12 @@ fn build_transient_views( .filter(|(name, _)| !containers.iter().any(|c| &c.name == *name)) .map(|(name, st)| TransientView { name: name.clone(), - kind: transient_label(st.kind), + kind: st.label.clone(), secs: st.since.elapsed().as_secs(), }) .collect() } -fn transient_label(k: crate::coordinator::TransientKind) -> &'static str { - use crate::coordinator::TransientKind::{ - Destroying, Rebuilding, Restarting, Spawning, Starting, Stopping, - }; - match k { - Spawning => "spawning", - Starting => "starting", - Stopping => "stopping", - Restarting => "restarting", - Rebuilding => "rebuilding", - Destroying => "destroying", - } -} - /// Render each pending approval into its dashboard view (short sha for /// `MergeConfigPr`, just the name for `Spawn`). /// Project a resolved sqlite row into the lean shape the dashboard diff --git a/hive-c0re/src/dashboard_events.rs b/hive-c0re/src/dashboard_events.rs index b5451b05..8c6c8db0 100644 --- a/hive-c0re/src/dashboard_events.rs +++ b/hive-c0re/src/dashboard_events.rs @@ -144,9 +144,15 @@ pub enum DashboardEvent { TransientSet { seq: u64, name: String, - /// Lifecycle kind: `"spawning"` / `"starting"` / `"stopping"` / - /// `"restarting"` / `"rebuilding"` / `"destroying"`. - transient_kind: &'static str, + /// What the pill renders. For queue-driven work this is the running + /// node's own wire tag (`"swap"`, `"create"`, `"stop_for_update"`, …) — + /// the same vocabulary the DAG view ships. Work with no node behind it + /// (destroy, migration) supplies its own (`"destroying"`, + /// `"rebuilding"`). + /// + /// Owned rather than `&'static str`: a label now comes from the node + /// that happens to be running, not from a fixed set. + transient_kind: String, since_unix: i64, }, /// The matching lifecycle action resolved (success or failure). @@ -380,7 +386,7 @@ mod tests { DashboardEvent::TransientSet { seq: 1, name: "x".into(), - transient_kind: "rebuilding", + transient_kind: "rebuilding".into(), since_unix: 0, }, DashboardEvent::TransientCleared { diff --git a/hive-c0re/src/job_queue/exec.rs b/hive-c0re/src/job_queue/exec.rs index 0b73ac76..7dcdf4dd 100644 --- a/hive-c0re/src/job_queue/exec.rs +++ b/hive-c0re/src/job_queue/exec.rs @@ -418,13 +418,11 @@ async fn run_reconcile(coord: &Arc, claim: &Claim) -> Result, claim: &Claim) -> Result { let name = &claim.agent; - // Node-local transient only when the DAG holds none (the - // boot-reconcile template); a rebuild/spawn/etc. DAG's lease-window - // transient already covers this node. - let _guard = claim - .transient - .is_none() - .then(|| coord.transient_guard(name, crate::coordinator::TransientKind::Starting)); + // No node-local transient guard: the pill is derived from the running node + // set, and `Start` reports `Starting` via `NodeKind::transient_kind`. This + // used to take one "only when the DAG holds none", which was a second + // derivation covering the gap left by a DAG-level declaration that couldn't + // describe a sub-step. // Run the typed start preamble: ensures the runtime dir exists and // writes the nspawn/resource-limits drop-ins. The returned // StartableAgent token is the only way to call start_with_fallback — @@ -449,10 +447,8 @@ async fn run_start(coord: &Arc, claim: &Claim) -> Result, claim: &Claim) -> Result { let name = &claim.agent; - let _guard = claim - .transient - .is_none() - .then(|| coord.transient_guard(name, crate::coordinator::TransientKind::Stopping)); + // See `run_start`: no node-local guard — `Stop` reports `Stopping` from its + // own kind now. crate::lifecycle::kill(name).await?; coord.unregister_agent(name); coord.notify_manager(&hive_sh4re::HelperEvent::Killed { diff --git a/hive-c0re/src/job_queue/mod.rs b/hive-c0re/src/job_queue/mod.rs index 9300b526..1bf2fbf9 100644 --- a/hive-c0re/src/job_queue/mod.rs +++ b/hive-c0re/src/job_queue/mod.rs @@ -47,7 +47,6 @@ use hive_jobq::{Dep, Graph, NodeId}; use hive_sh4re::wire_time::now_unix; use tokio::sync::Notify; -use crate::coordinator::TransientKind; pub use hive_jobq::TerminalState; pub use model::{DagSpec, DagView, NodeKind, NodeSpec, PermPayload, Source, State}; use resource::Resource; @@ -70,10 +69,6 @@ pub struct Claim { /// The agent this node targets (its own, not a DAG-level field). Empty for /// the agentless [`NodeKind::MetaLock`] + [`NodeKind::Dag`] container nodes. pub agent: String, - /// Transient pill kind for the lease window (from the spec). Whether the - /// pill is currently shown is derived from live lease ownership - /// ([`JobQueue::held_transients`]), not a per-claim edge. - pub transient: Option, } /// Per-node runtime metadata the crate graph doesn't carry. Lifecycle @@ -91,7 +86,6 @@ struct NodeRuntime { struct DagMeta { source: Source, reason: String, - transient: Option, created_at: i64, } @@ -215,7 +209,6 @@ impl JobQueue { NodeKind::Dag { source: spec.source, reason: spec.reason, - transient: spec.transient, created_at: now_unix(), }, Vec::new(), @@ -293,15 +286,11 @@ impl JobQueue { let Some(container) = inner.sched.graph().root_of(id) else { continue; }; - let Some(meta) = inner.dag_meta(container) else { - continue; - }; claims.push(Claim { dag_id: container.get(), node_id: id, kind, agent, - transient: meta.transient, }); // `started_at` is stamped on the graph `Node` by the scheduler's // transition to `Running` — no host-side copy needed. @@ -410,24 +399,58 @@ impl JobQueue { .map(ToOwned::to_owned) } - /// The `(dag_id, agent, kind)` triples for every per-agent lease currently - /// held by a DAG that carries a transient pill — the live transient-pill - /// set, a pull query over crate resource ownership (replaces the old - /// lease-release event stream). A DAG with no transient kind is omitted. + /// The `(agent, label)` pairs for the live transient-pill set — derived from + /// the nodes **actually running**, not from an intent a template declared at + /// submit time. A rebuild used to report `rebuilding` for its whole life: + /// through the prebuild, the stop, the swap, the tail and the reconcile. + /// + /// A node lights a pill when it is `Running` **and declares the agent's + /// resource itself**. Declaring is the test, not targeting: `Prebuild` and + /// `MetaSync` name an agent but are lease-exempt on purpose — the container + /// keeps serving right through them — so they must not light one. It is also + /// not the *lease owner*: `resource_state()` answers "who holds the slot", + /// a different question from "what is running". + /// + /// The label is the node's own wire tag ([`NodeKind::as_str`]) — the same + /// vocabulary [`NodeView::kind`] already ships, so a pill and a DAG node + /// name an operation identically and there is no second taxonomy to keep in + /// step. + /// + /// Consequence, by design: `Start` / `Stop` / `PostSwap` run *inside* a + /// lease-holding ancestor and re-declare nothing, so they light no pill and + /// the agent reads idle for those windows. Closing that is the point of the + /// resources-where-constructed work, not of this function. + /// + /// An agent's lease is cap-1, so at most one pair per agent. + /// The third element is [`NodeKind::takes_container_down`] — the crash + /// watcher's input, carried alongside the label rather than inferred from + /// it (a `Start` pill and a `Stop` pill are both pills; only one of them + /// means a vanished container is expected). #[must_use] - pub fn held_transients(&self) -> Vec<(u64, String, TransientKind)> { + pub fn held_transients(&self) -> Vec<(String, String, bool)> { let inner = self.lock(); inner .sched - .resource_state() - .into_iter() - .filter_map(|(res, holder)| { - let Resource::Agent(agent) = res else { - return None; - }; - let container = inner.sched.graph().root_of(holder)?; - let kind = inner.dag_meta(container)?.transient?; - Some((container.get(), agent, kind)) + .graph() + .nodes() + .filter(|n| matches!(n.state, State::Running)) + .filter_map(|n| { + let agent = n + .payload + .resource_deps() + .into_iter() + .find_map(|d| match d { + Dep::Resource { + name: Resource::Agent(a), + .. + } => Some(a), + _ => None, + })?; + Some(( + agent, + n.payload.as_str().to_owned(), + n.payload.takes_container_down(), + )) }) .collect() } @@ -481,7 +504,6 @@ impl QueueInner { let NodeKind::Dag { source, reason, - transient, created_at, } = &self.sched.graph().node(container)?.payload else { @@ -490,7 +512,6 @@ impl QueueInner { Some(DagMeta { source: *source, reason: reason.clone(), - transient: *transient, created_at: *created_at, }) } diff --git a/hive-c0re/src/job_queue/model.rs b/hive-c0re/src/job_queue/model.rs index fa208985..96fba3c3 100644 --- a/hive-c0re/src/job_queue/model.rs +++ b/hive-c0re/src/job_queue/model.rs @@ -15,8 +15,6 @@ pub use hive_host_sock::jobs::{DagView, NodeId, PermPayload, Source, State}; use serde::Serialize; -use crate::coordinator::TransientKind; - use hive_jobq::{DepWhen, TerminalState}; /// A dependency edge (intra-DAG only — cross-DAG ordering comes from @@ -285,7 +283,6 @@ pub enum NodeKind { Dag { source: Source, reason: String, - transient: Option, created_at: i64, }, } @@ -393,6 +390,44 @@ impl NodeKind { ) } + /// Whether running this node is *expected* to take the agent's container + /// down. Feeds `TransientState::deliberate_stop`, which the crash watcher + /// reads to tell an intentional stop from a crash. + /// + /// This is a **safety** question, not a display one — it decides whether a + /// vanished container raises an alert. It is deliberately not derived from + /// the pill label: a label is free to be renamed or added without moving + /// the alerting boundary, and only the operation itself knows its intent. + /// + /// Default is `false`, and that asymmetry is the point. A wrong `false` + /// costs a spurious crash event; a wrong `true` **swallows a real crash** + /// silently. So a kind earns `true` by being listed here, and anything new + /// is noisy-but-safe until someone decides otherwise. + #[must_use] + pub fn takes_container_down(&self) -> bool { + matches!( + self, + // Explicit stops, and the quiesce steps that precede one. + NodeKind::Stop { .. } + | NodeKind::StopForUpdate { .. } + | NodeKind::Signal { .. } + | NodeKind::Drain { .. } + | NodeKind::SetWanted { up: false, .. } + // The rebuild's own machinery: the container is down across the + // swap and the drop-in write that reconfigures it. + | NodeKind::Swap { .. } + | NodeKind::WriteDropin { .. } + ) + // Everything else is `false` on purpose, including the ones that would + // be easy to wave through: + // - `Create` / `Start` / `SetWanted{up}` bring a container UP. A + // container disappearing *while starting* is a genuine crash and has + // to keep reporting as one. + // - `Reconcile` is a planner; it fans out `Start` / `Stop`, which carry + // their own answer. + // - `DeployWindow` brackets a deploy without itself stopping anything. + } + /// Kinds that **mutate the meta repo** and so must hold the global /// [`Resource::MetaWindow`](super::resource::Resource::MetaWindow) for /// their duration: no two meta mutations may interleave, because a commit @@ -463,8 +498,5 @@ pub struct DagSpec { pub source: Source, /// Free-form "why". pub reason: String, - /// Dashboard transient pill (and crash-watch suppression) held for - /// the lease window — from lease acquisition to DAG terminal. - pub transient: Option, pub nodes: Vec, } diff --git a/hive-c0re/src/job_queue/scheduler.rs b/hive-c0re/src/job_queue/scheduler.rs index d235653b..24cb8000 100644 --- a/hive-c0re/src/job_queue/scheduler.rs +++ b/hive-c0re/src/job_queue/scheduler.rs @@ -3,12 +3,12 @@ //! and on any completion re-evaluate. Concurrency comes from the build-slot //! count, not multiple workers. //! -//! Owns the per-DAG transient guard (dashboard pill + crash-watch suppression) -//! that the sync queue core can't hold itself. The guard set is *reconciled* -//! from live lease ownership ([`super::JobQueue::held_transients`]) each loop: -//! a `(dag, agent)` pill exists for exactly as long as that agent's lease is -//! held, so it appears when the agent's owner node starts and disappears when -//! its subgraph settles — one pill per agent a DAG touches. +//! Owns the per-agent transient guard (dashboard pill + crash-watch +//! suppression) that the sync queue core can't hold itself. The guard set is +//! *reconciled* each loop from [`super::JobQueue::held_transients`], which +//! reports what is **running right now** under each held agent lease — so the +//! label tracks the DAG's progress (signal → swap → reconcile) instead of +//! repeating one intent the template declared before any of it started. //! //! Per-DAG terminal work (approval resolution, `Rebuilt`) is not drained here: //! it runs as the DAG's focused terminal node (`ResolveApproval` / @@ -19,7 +19,7 @@ //! its `Start`/`Stop`) flows through `NodeOutput.append_subgraph`, applied //! before the emitting node completes — see `handle_completion`. -use std::collections::{HashMap, HashSet}; +use std::collections::HashMap; use std::sync::Arc; use super::Claim; @@ -54,8 +54,10 @@ struct NodeDone { pub async fn run_worker(coord: Arc) { let mut shutdown = coord.shutdown_rx(); let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel::(); - // (DAG id, agent) → transient guard held for that agent's lease window. - let mut transients: HashMap<(u64, String), crate::coordinator::TransientGuard> = HashMap::new(); + // Keyed by agent (its lease is cap-1, so one pill each); the label rides + // along so a change of label can be detected and the guard swapped. + let mut transients: HashMap = + HashMap::new(); loop { // Checked every iteration, not just in the `select!` below — a // continuous stream of ready claims never reaches the `select!`, so @@ -146,18 +148,35 @@ fn handle_completion(coord: &Arc, done: NodeDone) { coord.emit_rebuild_queue_snapshot(); } -/// Reconcile the transient-guard set against live lease ownership: drop pills -/// whose lease is no longer held, create one for each newly-held `(dag, agent)`. +/// Reconcile the transient-guard set against the live pill set: drop guards for +/// pills that are no longer current, create one for each newly-current +/// `(agent, label)`. Keyed by agent — an agent's lease is cap-1, so it has at +/// most one pill. +/// +/// `deliberate_stop` rides along per node ([`NodeKind::takes_container_down`]) +/// rather than being blanket-`true` for anything holding a lease: `Create` and +/// `Start` hold the agent's lease too, and a container vanishing *while +/// starting* is a real crash that must keep reporting as one. +/// +/// [`NodeKind::takes_container_down`]: super::NodeKind::takes_container_down +/// +/// ⚠️ **`retain` must run to completion before anything is created.** +/// `TransientGuard::drop` calls `clear_transient(agent)` — keyed by agent alone, +/// with no notion of *which* label it was clearing. So when a pill's label +/// changes for the same agent (which is now routine: the label follows the +/// running node as a DAG advances), creating the new guard first and dropping +/// the old second would clear the pill that was just set. Dropping first is what +/// makes the swap safe. fn reconcile_transients( coord: &Arc, - transients: &mut HashMap<(u64, String), crate::coordinator::TransientGuard>, + transients: &mut HashMap, ) { let held = coord.job_queue.held_transients(); - let keys: HashSet<(u64, String)> = held.iter().map(|(d, a, _)| (*d, a.clone())).collect(); - transients.retain(|k, _| keys.contains(k)); - for (dag_id, agent, kind) in held { - transients - .entry((dag_id, agent.clone())) - .or_insert_with(|| coord.transient_guard(&agent, kind)); + transients.retain(|agent, (label, _)| held.iter().any(|(a, l, _)| a == agent && l == label)); + for (agent, label, takes_down) in held { + transients.entry(agent.clone()).or_insert_with(|| { + let guard = coord.transient_guard(&agent, label.clone(), takes_down); + (label, guard) + }); } } diff --git a/hive-c0re/src/job_queue/submit.rs b/hive-c0re/src/job_queue/submit.rs index 9ced4137..43955ec9 100644 --- a/hive-c0re/src/job_queue/submit.rs +++ b/hive-c0re/src/job_queue/submit.rs @@ -27,7 +27,7 @@ use std::sync::Arc; use super::model::{DagSpec, Dep, NodeKind, NodeSpec}; use super::templates::{RebuildOpts, after_ok, child, node, rebuild_nodes}; use super::{Source, templates}; -use crate::coordinator::{Coordinator, TransientKind}; +use crate::coordinator::Coordinator; use crate::lifecycle; fn submit_and_emit(coord: &Arc, spec: super::DagSpec) -> u64 { @@ -198,16 +198,10 @@ fn concat_subgraphs(chains: Vec>) -> Vec { /// Wrap assembled power-op `nodes` in a `DagSpec`. No tail node: a power op's /// effect is its nodes (`SetWanted` + `Reconcile`), with nothing left to do once /// they settle. -fn power_dag( - transient: TransientKind, - source: Source, - reason: String, - nodes: Vec, -) -> DagSpec { +fn power_dag(source: Source, reason: String, nodes: Vec) -> DagSpec { DagSpec { source, reason, - transient: Some(transient), nodes, } } @@ -228,18 +222,15 @@ pub(crate) fn stop_spec( .iter() .map(|(agent, running)| stop_chain(agent, graceful, *running)) .collect(); - power_dag( - TransientKind::Stopping, - source, - reason, - concat_subgraphs(chains), - ) + power_dag(source, reason, concat_subgraphs(chains)) } /// Assemble the start DAG from explicit `(agent, running, stale)` targets. -/// Transient is `Rebuilding` when any down+stale agent grew a rebuild -/// subgraph (crash-watch suppression during its Swap), else `Starting`; -/// applied per-agent at claim time, so each agent still shows its own pill. +/// +/// No DAG-level pill: each agent's dashboard label is derived from the node +/// 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_spec( targets: &[(String, bool, bool)], source: Source, @@ -249,13 +240,7 @@ pub(crate) fn start_spec( .iter() .map(|(agent, running, stale)| start_chain(agent, *running, *stale)) .collect(); - let any_rebuild = targets.iter().any(|(_, running, stale)| !running && *stale); - let transient = if any_rebuild { - TransientKind::Rebuilding - } else { - TransientKind::Starting - }; - power_dag(transient, source, reason, concat_subgraphs(chains)) + power_dag(source, reason, concat_subgraphs(chains)) } /// Assemble the restart DAG from explicit `(agent, running)` targets. @@ -269,12 +254,7 @@ pub(crate) fn restart_spec( .iter() .map(|(agent, running)| restart_chain(agent, graceful, *running)) .collect(); - power_dag( - TransientKind::Restarting, - source, - reason, - concat_subgraphs(chains), - ) + power_dag(source, reason, concat_subgraphs(chains)) } /// Restart a single agent. Thin wrapper over [`restart_many`]. diff --git a/hive-c0re/src/job_queue/templates.rs b/hive-c0re/src/job_queue/templates.rs index 19abd54b..70c6d6a9 100644 --- a/hive-c0re/src/job_queue/templates.rs +++ b/hive-c0re/src/job_queue/templates.rs @@ -34,7 +34,6 @@ use anyhow::{Result, bail}; use hive_jobq::{DepWhen, TerminalState}; use super::model::{DagSpec, Dep, NodeKind, NodeSpec, PermPayload, Source}; -use crate::coordinator::TransientKind; /// After-ok edge on the previous node — the common chain link. Shared with /// the async power-op builders in `submit.rs` (which assemble per-agent @@ -367,7 +366,6 @@ pub fn rebuild(agent: &str, source: Source, reason: String, relock: bool) -> Dag DagSpec { source, reason, - transient: Some(TransientKind::Rebuilding), nodes, } } @@ -402,7 +400,6 @@ pub fn approval_deploy(agent: &str, approval_id: i64, reason: String) -> DagSpec DagSpec { source: Source::Approval, reason, - transient: Some(TransientKind::Rebuilding), nodes: vec![ node( NodeKind::DeployWindow { @@ -451,16 +448,10 @@ pub fn approval_deploy(agent: &str, approval_id: i64, reason: String) -> DagSpec /// single-node lifecycle DAGs that exercise per-agent lease serialization /// in the queue tests); production paths no longer emit a bare reconcile. #[cfg(test)] -pub fn reconcile_only( - agent: &str, - source: Source, - reason: String, - transient: Option, -) -> DagSpec { +pub fn reconcile_only(agent: &str, source: Source, reason: String) -> DagSpec { DagSpec { source, reason, - transient, nodes: vec![node( NodeKind::Reconcile { agent: agent.to_owned(), @@ -485,7 +476,6 @@ pub fn spawn(agent: &str, approval_id: i64, reason: String) -> DagSpec { DagSpec { source: Source::Approval, reason, - transient: Some(TransientKind::Spawning), nodes: { let a = || agent.to_owned(); vec![ @@ -529,7 +519,6 @@ pub fn perm_change(agent: &str, source: Source, reason: String, payload: PermPay DagSpec { source, reason, - transient: Some(TransientKind::Rebuilding), nodes, } } @@ -568,7 +557,6 @@ pub fn meta_update( DagSpec { source, reason, - transient: Some(TransientKind::Rebuilding), nodes, } } @@ -590,7 +578,6 @@ pub fn reparent( DagSpec { source, reason, - transient: None, nodes: vec![node(NodeKind::Reparent { moves }, Vec::new())], } } diff --git a/hive-c0re/src/job_queue/tests.rs b/hive-c0re/src/job_queue/tests.rs index 5fd97c49..49d55bd4 100644 --- a/hive-c0re/src/job_queue/tests.rs +++ b/hive-c0re/src/job_queue/tests.rs @@ -246,7 +246,7 @@ fn graceful_rebuild_chain_drains_before_stopping() { let spec = DagSpec { source: Source::AutoUpdate, reason: "sweep".to_owned(), - transient: None, + nodes: templates::rebuild_nodes( "agent-a", templates::RebuildOpts { @@ -430,7 +430,7 @@ fn lease_serializes_two_lifecycle_dags_for_same_agent() { let restart = submit(&q, restart_online(&["agent-a"], false, "restart")); let stop = submit( &q, - templates::reconcile_only("agent-a", Source::Manual, "stop".to_owned(), None), + templates::reconcile_only("agent-a", Source::Manual, "stop".to_owned()), ); // Restart's first node (StopForUpdate) takes the lease; stop's // Reconcile must wait even though slots are free. @@ -461,7 +461,7 @@ fn lease_exempt_prebuild_overlaps_other_dag_on_same_agent() { submit(&q, rebuild("agent-a", "rebuild")); submit( &q, - templates::reconcile_only("agent-a", Source::Manual, "stop".to_owned(), None), + templates::reconcile_only("agent-a", Source::Manual, "stop".to_owned()), ); // Both DAGs' heads are lease-independent of each other: the rebuild's // MetaSync (meta window) and the stop's Reconcile (agent lease). @@ -723,7 +723,7 @@ fn append_subgraph_roots_on_emitter_and_rebases_local_deps() { let spec = DagSpec { source: Source::AutoUpdate, reason: "sweep".to_owned(), - transient: None, + nodes: vec![NodeSpec { kind: NodeKind::MetaLock { sweep: true, @@ -789,26 +789,55 @@ fn drain_meta_syncs(q: &JobQueue) -> Vec<(String, String)> { rest } +/// Crash-watch suppression for a cascade rebuild, which the deleted half of +/// `meta_update_grows_cascade_in_dag` used to assert via `DagSpec::transient`. +/// +/// The property is unchanged — a container going down under a rebuild must not +/// read as a crash — but it is no longer a DAG-level declaration: each node +/// answers for itself, so the assertion moves to the nodes a cascade actually +/// runs. Kept as its own test rather than dropped, because it is the *property* +/// that mattered, not the field that used to carry it. #[test] -fn meta_update_carries_rebuilding_transient_and_grows_cascade_in_dag() { +fn rebuild_chain_nodes_suppress_crash_watch() { + for kind in [ + NodeKind::StopForUpdate { + agent: "a".to_owned(), + }, + NodeKind::Swap { + agent: "a".to_owned(), + }, + NodeKind::Drain { + agent: "a".to_owned(), + }, + ] { + assert!( + kind.takes_container_down(), + "{} must suppress crash-watch — a rebuild takes the container down \ + on purpose", + kind.as_str() + ); + } + // The counter-case, and the reason this can't be "any node in a rebuild": + // the tail brings the container back up, so a container that dies there + // really did crash. + assert!( + !NodeKind::Start { + agent: "a".to_owned() + } + .takes_container_down() + ); +} + +#[test] +fn meta_update_grows_cascade_in_dag() { // The meta-update `MetaLock` grows one rebuild subgraph per affected // agent into its OWN DAG (via append_subgraph), not child DAGs. - // The DAG carries `Rebuilding` so the folded rebuilds keep crash-watch - // suppression (the property the old child Rebuild DAGs had via their own - // transient). let spec = templates::meta_update( vec!["nixpkgs".to_owned()], Source::Manual, "bump".to_owned(), None, ); - assert!( - matches!( - spec.transient, - Some(crate::coordinator::TransientKind::Rebuilding) - ), - "meta-update DAG must carry Rebuilding so cascade rebuilds get suppression" - ); let q = JobQueue::new(4); let id = submit(&q, spec); let meta_lock = claim_one(&q); @@ -981,7 +1010,7 @@ fn failed_reconcile_marks_dag_failed() { let q = JobQueue::new(1); let id = submit( &q, - templates::reconcile_only("agent-a", Source::Manual, "start".to_owned(), None), + templates::reconcile_only("agent-a", Source::Manual, "start".to_owned()), ); let c = claim_one(&q); q.complete_node(c.node_id, Err("start failed".to_owned())); @@ -1132,7 +1161,7 @@ fn dag_settles_terminal_and_releases_lease_after_work() { // immediately. let next = submit( &q, - templates::reconcile_only("agent-a", Source::Manual, "stop".to_owned(), None), + templates::reconcile_only("agent-a", Source::Manual, "stop".to_owned()), ); let c = claim_one(&q); assert_eq!(c.dag_id, next); @@ -1401,12 +1430,7 @@ fn history_evicts_oldest_terminals_past_flat_cap() { for i in 0..(MAX_HISTORY_DAGS + OVERFLOW) { let id = submit( &q, - templates::reconcile_only( - &format!("agent-{i}"), - Source::Manual, - "start".to_owned(), - None, - ), + templates::reconcile_only(&format!("agent-{i}"), Source::Manual, "start".to_owned()), ); let c = claim_one(&q); // Fail the single work node so the DAG *lingers*: a fully-`Done` DAG diff --git a/hive-c0re/src/migrate.rs b/hive-c0re/src/migrate.rs index 355bba5d..38a5d5e2 100644 --- a/hive-c0re/src/migrate.rs +++ b/hive-c0re/src/migrate.rs @@ -108,8 +108,10 @@ pub async fn run(coord: &Arc) -> Result<()> { // update activation triggers. Without this, crash_watch // would fire ContainerCrash for every agent here and the // manager would spuriously try to recover them. - let guard = - coord.transient_guard(name.as_str(), crate::coordinator::TransientKind::Rebuilding); + // No queue node behind this one — migration repoints containers + // directly — so the label is supplied here. `true`: the repoint takes + // the container down, which is the whole reason for the guard. + let guard = coord.transient_guard(name.as_str(), "rebuilding", true); let result = repoint_container(name.as_str()).await; drop(guard); if let Err(e) = result { @@ -201,7 +203,8 @@ async fn rename_manager_container(coord: &Arc) { return; } tracing::info!("migration phase 5: renaming root container to h-root"); - let _guard = coord.transient_guard(MANAGER_NAME, crate::coordinator::TransientKind::Rebuilding); + // `true`: the old container is stopped immediately below. + let _guard = coord.transient_guard(MANAGER_NAME, "rebuilding", true); // Stop the old container. Abort if stop fails — continuing with a // running `root` and then starting `h-root` risks two manager diff --git a/hive-c0re/src/workers/auto_update.rs b/hive-c0re/src/workers/auto_update.rs index 36e42477..3fc78e61 100644 --- a/hive-c0re/src/workers/auto_update.rs +++ b/hive-c0re/src/workers/auto_update.rs @@ -353,7 +353,6 @@ fn submit_boot_tree( // Rebuilding when the sweep will grow rebuild subgraphs (per-agent // crash-watch suppression during their Swap, applied at claim time); // a reconcile-only boot needs no transient. - transient: any_stale.then_some(crate::coordinator::TransientKind::Rebuilding), nodes, }; if let Err(e) = coord.job_queue.submit(spec) { diff --git a/hive-c0re/src/workers/crash_watch.rs b/hive-c0re/src/workers/crash_watch.rs index 49111296..9b74f826 100644 --- a/hive-c0re/src/workers/crash_watch.rs +++ b/hive-c0re/src/workers/crash_watch.rs @@ -8,7 +8,7 @@ use std::sync::Arc; use std::time::Duration; use crate::container_view::claude_has_session; -use crate::coordinator::{Coordinator, TransientKind}; +use crate::coordinator::Coordinator; use crate::lifecycle::{self, AGENT_PREFIX}; const POLL_INTERVAL: Duration = Duration::from_secs(10); @@ -87,7 +87,7 @@ fn emit_crash_transitions(coord: &Coordinator, prev: &HashSet, current: // guard between two crash-watch polls. let recent = coord.recent_transient_within(RECENT_TRANSIENT_GRACE); for stopped in prev.difference(current) { - let active = transients.get(stopped).map(|st| st.kind); + let active = transients.get(stopped).map(|st| st.deliberate_stop); let recently_cleared = recent.get(stopped).copied(); if is_deliberate_stop(active, recently_cleared) { continue; @@ -101,26 +101,19 @@ fn emit_crash_transitions(coord: &Coordinator, prev: &HashSet, current: } } -/// Pure classifier: did the operator stop / restart / destroy / -/// rebuild this container, or did it crash? Splits the matcher out -/// so it has a focused unit test without needing a Coordinator -/// fixture. `active` is the currently-set transient (if any), -/// `recently_cleared` is one whose RAII guard dropped within the -/// grace window. -fn is_deliberate_stop( - active: Option, - recently_cleared: Option, -) -> bool { - let is_op_kind = |kind: TransientKind| { - matches!( - kind, - TransientKind::Stopping - | TransientKind::Restarting - | TransientKind::Destroying - | TransientKind::Rebuilding - ) - }; - active.is_some_and(is_op_kind) || recently_cleared.is_some_and(is_op_kind) +/// Pure classifier: did an operation take this container down on purpose, or +/// did it crash? Splits the matcher out so it has a focused unit test without +/// needing a Coordinator fixture. `active` is the currently-set transient's +/// `deliberate_stop` (if any), `recently_cleared` is one whose RAII guard +/// dropped within the grace window. +/// +/// This reads the flag the transient's creator set and nothing else. It used to +/// match a `TransientKind`, which meant a pill's *display* vocabulary decided a +/// crash-alert question — so renaming or adding a label silently moved the +/// alerting boundary. Whoever starts the operation knows whether the container +/// is meant to go down; nothing downstream can re-derive it. +fn is_deliberate_stop(active: Option, recently_cleared: Option) -> bool { + active.unwrap_or(false) || recently_cleared.unwrap_or(false) } fn emit_login_transitions( @@ -168,29 +161,12 @@ mod tests { use super::*; #[test] - fn deliberate_when_active_transient_is_operator_kind() { - for kind in [ - TransientKind::Stopping, - TransientKind::Restarting, - TransientKind::Destroying, - TransientKind::Rebuilding, - ] { - assert!(is_deliberate_stop(Some(kind), None), "{kind:?}"); - } - } - - #[test] - fn deliberate_when_recent_transient_is_operator_kind() { - // Race repros: lifecycle action completes + drops the guard - // between two polls. recent_transient catches it. - for kind in [ - TransientKind::Stopping, - TransientKind::Restarting, - TransientKind::Destroying, - TransientKind::Rebuilding, - ] { - assert!(is_deliberate_stop(None, Some(kind)), "{kind:?}"); - } + fn deliberate_when_either_source_says_so() { + // Active guard, and the race repro: a lifecycle action completes and + // drops its guard between two polls, so only `recent` still carries it. + assert!(is_deliberate_stop(Some(true), None)); + assert!(is_deliberate_stop(None, Some(true))); + assert!(is_deliberate_stop(Some(true), Some(true))); } #[test] @@ -200,13 +176,16 @@ mod tests { } #[test] - fn not_deliberate_when_only_spawning_starting() { - // Spawning/Starting are never paired with a "stopped" transition - // — they're starts. If we see one alongside a stop, it's - // unrelated (e.g. just-started container died), still a crash. - for kind in [TransientKind::Spawning, TransientKind::Starting] { - assert!(!is_deliberate_stop(Some(kind), None), "{kind:?} active"); - assert!(!is_deliberate_stop(None, Some(kind)), "{kind:?} recent"); - } + fn not_deliberate_when_the_operation_was_bringing_the_container_up() { + // The case that used to be spelled `Spawning` / `Starting`: an + // operation IS in flight, but it is not one that takes the container + // down, so a container that vanishes under it really did crash. + // + // This is why `deliberate_stop` is carried rather than inferred from + // the pill — `Create` and `Start` hold the agent's lease exactly like + // `Stop` does, so "has a pill" cannot answer this. + assert!(!is_deliberate_stop(Some(false), None)); + assert!(!is_deliberate_stop(None, Some(false))); + assert!(!is_deliberate_stop(Some(false), Some(false))); } } From 6458c039a01cec14a962a0605c1ec34f9a52ecc2 Mon Sep 17 00:00:00 2001 From: atlas Date: Sat, 1 Aug 2026 16:23:10 +0200 Subject: [PATCH 2/9] docs(#2815): the transient pill's vocabulary is open, not a fixed set MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit argus flagged on !2910 that the dashboard doc still described the pill's old shape. It did, in two ways that now teach the wrong thing: - it listed a fixed vocabulary, where the label is now the running node's own wire tag — the `NodeKind::as_str` strings `NodeView.kind` already carries. A client that switches on specific values is now wrong, and `restarting` in particular no longer exists at all. - it described the transient as operator-initiated ("set the moment the operator clicks"), which was the distinction between it and the rebuild-queue fallback. That is no longer true: the transient is derived from the running node, so worker-driven work the operator never clicked lights the same pill. The queue-`kind` half of that section is untouched — that path did not change and its vocabulary is still fixed. Also notes on the wire-event list that `transient_kind` is a display string to render, not an enum to branch on, since that is the property a client would otherwise have to infer from a now-open set. Docs only, no code change. --- docs/web-ui/dashboard.md | 31 +++++++++++++++++++++++++------ 1 file changed, 25 insertions(+), 6 deletions(-) diff --git a/docs/web-ui/dashboard.md b/docs/web-ui/dashboard.md index 7556155b..a8cfaff2 100644 --- a/docs/web-ui/dashboard.md +++ b/docs/web-ui/dashboard.md @@ -852,12 +852,29 @@ fetch entirely. icon, so it's obvious at a glance which container is actually moving. **Pending-state derivation:** the pill is sourced from two - separate stores in priority order. (1) The operator-initiated - **transient** (`transientsState`) is set on the dashboard the - moment the operator clicks start / stop / restart / rebuild / - destroy / spawn — covers the create-and-start window where the + separate stores in priority order. (1) The **transient** + (`transientsState`) — covers the create-and-start window where the container literally isn't up yet, before any backend state event - has fired. (2) If no transient is set, the **rebuild-queue + has fired. + + A transient is **derived from the job-queue node currently + running** against that agent, not declared per request, so its + label follows the operation as it progresses (a rebuild reads + `stop_for_update`, then `swap`, then `reconcile` rather than one + constant `rebuilding` for its whole life). Two consequences for + anything rendering it: + + - The label vocabulary is **open** — it is the node's own wire tag + (`NodeKind::as_str`, the same strings `NodeView.kind` carries), + not a fixed set. Treat it as an opaque display string; do not + switch on specific values. `restarting` in particular no longer + exists, because no node kind is unique to a restart. + - It is **not** exclusively operator-initiated. Work the operator + never clicked (a meta-update cascade, a crash-recover rebuild) + lights the same pill, since it is the running node that sets it. + + Ops with no queue node behind them (destroy, migration) supply + their own label directly. (2) If no transient is set, the **rebuild-queue entry** for this agent is consulted (`rebuildQueueState`); this covers worker-driven ops — meta-update cascades, crash-recover rebuilds, approval-driven rebuilds — that the operator didn't @@ -1351,7 +1368,9 @@ payload): - `transient_set` (name, transient_kind, since_unix) / `transient_cleared` (name) — lifecycle action spinners. The client ticks the elapsed-seconds badge off `since_unix` - client-side, no polling. + client-side, no polling. `transient_kind` is an **open** + display string (the running node's own tag), not a fixed + enum — render it, don't branch on it. - `container_state_changed` (container: ContainerView) / `container_removed` (name) — per-row container mutations, emitted by `Coordinator::rescan_containers_and_emit` from From 17500a391d3bbd57fd02c0733fb7bb3aa5437878 Mon Sep 17 00:00:00 2001 From: atlas Date: Sat, 1 Aug 2026 16:39:38 +0200 Subject: [PATCH 3/9] refactor(#2815): held_transients -> running_transients MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit mara on !2910: "rename now, we will see if we can remove it later when some of the users have been removed or work differently." Nothing is held. The old name described a transient the DAG declared and kept for its whole lifetime — precisely the thing this PR replaces — so it outlived its own meaning the moment the derivation landed. The value is recomputed from the running set on every call. Kept as a function rather than inlined at its single call site, per the above: removing it is a later step that depends on its users changing, not something this PR should force. Rename plus its two references (the call in `reconcile_transients` and the module doc link). No behaviour change; the doc comment records what the old name meant so the rename doesn't erase the reason for it. Checked with clippy (`--all-targets -D warnings`), `cargo test -p hive-c0re -p hive-jobq` (322 + 41 passed) and `nix fmt`. --- hive-c0re/src/job_queue/mod.rs | 43 ++++++++++++---------------- hive-c0re/src/job_queue/scheduler.rs | 8 +++--- 2 files changed, 23 insertions(+), 28 deletions(-) diff --git a/hive-c0re/src/job_queue/mod.rs b/hive-c0re/src/job_queue/mod.rs index 1bf2fbf9..71f6263e 100644 --- a/hive-c0re/src/job_queue/mod.rs +++ b/hive-c0re/src/job_queue/mod.rs @@ -399,35 +399,30 @@ impl JobQueue { .map(ToOwned::to_owned) } - /// The `(agent, label)` pairs for the live transient-pill set — derived from - /// the nodes **actually running**, not from an intent a template declared at - /// submit time. A rebuild used to report `rebuilding` for its whole life: - /// through the prebuild, the stop, the swap, the tail and the reconcile. + /// `(agent, label, takes_container_down)` for the live transient-pill set, + /// recomputed from the nodes **actually running** — not from an intent a + /// template declared at submit time. (A rebuild used to report `rebuilding` + /// for its whole life: prebuild, stop, swap, tail and reconcile alike.) /// /// A node lights a pill when it is `Running` **and declares the agent's - /// resource itself**. Declaring is the test, not targeting: `Prebuild` and - /// `MetaSync` name an agent but are lease-exempt on purpose — the container - /// keeps serving right through them — so they must not light one. It is also - /// not the *lease owner*: `resource_state()` answers "who holds the slot", - /// a different question from "what is running". + /// resource itself**. Declaring is the test, not targeting — `Prebuild` / + /// `MetaSync` name an agent but are lease-exempt on purpose, since the + /// container keeps serving through them. Nor is it the lease *owner*: + /// `resource_state()` answers "who holds the slot", a different question. /// - /// The label is the node's own wire tag ([`NodeKind::as_str`]) — the same - /// vocabulary [`NodeView::kind`] already ships, so a pill and a DAG node - /// name an operation identically and there is no second taxonomy to keep in - /// step. + /// `label` is the node's own wire tag ([`NodeKind::as_str`]), the vocabulary + /// [`NodeView::kind`] already ships, so a pill and a DAG node name an + /// operation identically. `takes_container_down` is the crash watcher's + /// input, carried rather than inferred from the label — a `Start` pill and a + /// `Stop` pill are both pills; only one means a vanished container is + /// expected. /// - /// Consequence, by design: `Start` / `Stop` / `PostSwap` run *inside* a - /// lease-holding ancestor and re-declare nothing, so they light no pill and - /// the agent reads idle for those windows. Closing that is the point of the - /// resources-where-constructed work, not of this function. - /// - /// An agent's lease is cap-1, so at most one pair per agent. - /// The third element is [`NodeKind::takes_container_down`] — the crash - /// watcher's input, carried alongside the label rather than inferred from - /// it (a `Start` pill and a `Stop` pill are both pills; only one of them - /// means a vanished container is expected). + /// By design, `Start` / `Stop` / `PostSwap` run inside a lease-holding + /// ancestor and re-declare nothing, so they light no pill; closing that is + /// the resources-where-constructed work, not this function. An agent's lease + /// is cap-1, so at most one entry per agent. #[must_use] - pub fn held_transients(&self) -> Vec<(String, String, bool)> { + pub fn running_transients(&self) -> Vec<(String, String, bool)> { let inner = self.lock(); inner .sched diff --git a/hive-c0re/src/job_queue/scheduler.rs b/hive-c0re/src/job_queue/scheduler.rs index 24cb8000..2df092b6 100644 --- a/hive-c0re/src/job_queue/scheduler.rs +++ b/hive-c0re/src/job_queue/scheduler.rs @@ -5,7 +5,7 @@ //! //! Owns the per-agent transient guard (dashboard pill + crash-watch //! suppression) that the sync queue core can't hold itself. The guard set is -//! *reconciled* each loop from [`super::JobQueue::held_transients`], which +//! *reconciled* each loop from [`super::JobQueue::running_transients`], which //! reports what is **running right now** under each held agent lease — so the //! label tracks the DAG's progress (signal → swap → reconcile) instead of //! repeating one intent the template declared before any of it started. @@ -171,9 +171,9 @@ fn reconcile_transients( coord: &Arc, transients: &mut HashMap, ) { - let held = coord.job_queue.held_transients(); - transients.retain(|agent, (label, _)| held.iter().any(|(a, l, _)| a == agent && l == label)); - for (agent, label, takes_down) in held { + let running = coord.job_queue.running_transients(); + transients.retain(|agent, (label, _)| running.iter().any(|(a, l, _)| a == agent && l == label)); + for (agent, label, takes_down) in running { transients.entry(agent.clone()).or_insert_with(|| { let guard = coord.transient_guard(&agent, label.clone(), takes_down); (label, guard) From 56202065d59ec879f1ebac649435d1840fb5ba12 Mon Sep 17 00:00:00 2001 From: atlas Date: Sat, 1 Aug 2026 17:30:05 +0200 Subject: [PATCH 4/9] refactor(#2815): the scheduler publishes pill edges, it doesn't own pills MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit mara on !2910: "transient guard as well - should be removable now?" — for the queue path, yes. `set_transient`'s own doc explained why the RAII guard existed: a cancelled future must not leak an imperatively-set transient and pin the dashboard on "rebuilding…" forever. That cannot happen to a derived set. `running_transients()` is recomputed from the graph every loop, so a node that stops running stops appearing — there is nothing to own and nothing to leak. So the scheduler no longer holds a guard per pill. It keeps the previous derived value and publishes the transitions, which is the one thing a derived read cannot express: the dashboard wants `TransientSet` / `TransientCleared` edges, and the crash watcher wants the *moment* a pill cleared, since its grace window is what stops an operator stop from reading as a crash. That also retires a hazard rather than restating it. The old code carried a warning that stale guards had to be dropped before new ones were created, because `TransientGuard::drop` clears by agent with no notion of which label it was clearing — so a same-agent label change could clear the pill it had just set. With no guards there is no ordering to get wrong; clears are emitted before sets so a relabel reads as clear-then-set rather than two overlapping pills. `set_transient` / `clear_transient` become `pub(crate)`. The guard stays for destroy and migration, which have no node behind them and where the cancellation concern is real. Checked with clippy (`--all-targets -D warnings`), `cargo test -p hive-c0re -p hive-jobq` (322 + 41 passed) and `nix fmt`. --- hive-c0re/src/coordinator.rs | 27 +++++++----- hive-c0re/src/job_queue/scheduler.rs | 66 +++++++++++++++++----------- 2 files changed, 57 insertions(+), 36 deletions(-) diff --git a/hive-c0re/src/coordinator.rs b/hive-c0re/src/coordinator.rs index 196b95ac..7325d92e 100644 --- a/hive-c0re/src/coordinator.rs +++ b/hive-c0re/src/coordinator.rs @@ -1077,13 +1077,17 @@ impl Coordinator { /// Mark an agent as in-progress (only one state per agent for now). /// - /// Private on purpose: the RAII [`TransientGuard`] (via - /// [`Coordinator::transient_guard`]) is the only door, so the paired - /// `clear_transient` always runs on drop even if the surrounding future - /// is cancelled (HTTP request aborted, runtime shutdown mid-rebuild, - /// panic). A bare set with no guaranteed clear would leak the transient - /// and leave the dashboard stuck in "rebuilding…" forever. - fn set_transient(&self, name: &str, label: String, deliberate_stop: bool) { + /// Two callers, for two different reasons: + /// - **Work with a queue node behind it** — the job-queue scheduler, which + /// publishes the edges of a *derived* set + /// ([`crate::job_queue::JobQueue::running_transients`]). It needs no guard: + /// nothing is owned, and a node that stops running stops appearing. + /// - **Work with no node** (destroy, migration) — via the RAII + /// [`TransientGuard`], so the paired `clear_transient` still runs if the + /// surrounding future is cancelled (HTTP request aborted, shutdown + /// mid-rebuild, panic). There a bare set really would leak the transient + /// and pin the dashboard on "rebuilding…" forever. + pub(crate) fn set_transient(&self, name: &str, label: String, deliberate_stop: bool) { self.transient.lock().unwrap().insert( name.to_owned(), TransientState { @@ -1108,10 +1112,11 @@ impl Coordinator { }); } - /// Clear an agent's transient state. Private: only reachable through - /// [`TransientGuard`]'s `Drop`, which guarantees it runs (see - /// [`Coordinator::set_transient`]). - fn clear_transient(&self, name: &str) { + /// Clear an agent's transient state. Reached either from + /// [`TransientGuard`]'s `Drop` (the no-node callers) or from the scheduler + /// when a derived pill stops being current — see + /// [`Coordinator::set_transient`]. + pub(crate) fn clear_transient(&self, name: &str) { let removed = self.transient.lock().unwrap().remove(name); if let Some(state) = removed { // Stamp the tombstone so the crash watcher can still see diff --git a/hive-c0re/src/job_queue/scheduler.rs b/hive-c0re/src/job_queue/scheduler.rs index 2df092b6..492aa55d 100644 --- a/hive-c0re/src/job_queue/scheduler.rs +++ b/hive-c0re/src/job_queue/scheduler.rs @@ -54,10 +54,12 @@ struct NodeDone { pub async fn run_worker(coord: Arc) { let mut shutdown = coord.shutdown_rx(); let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel::(); - // Keyed by agent (its lease is cap-1, so one pill each); the label rides - // along so a change of label can be detected and the guard swapped. - let mut transients: HashMap = - HashMap::new(); + // Last derived pill set we published, keyed by agent (its lease is cap-1, + // so one pill each). Purely the previous value of a *derived* quantity — + // it exists to spot transitions, since the dashboard wants edges + // (`TransientSet` / `TransientCleared`) and the crash watcher wants the + // moment of the clear. Nothing owns a pill; nothing can leak one. + let mut transients: HashMap = HashMap::new(); loop { // Checked every iteration, not just in the `select!` below — a // continuous stream of ready claims never reaches the `select!`, so @@ -148,10 +150,23 @@ fn handle_completion(coord: &Arc, done: NodeDone) { coord.emit_rebuild_queue_snapshot(); } -/// Reconcile the transient-guard set against the live pill set: drop guards for -/// pills that are no longer current, create one for each newly-current -/// `(agent, label)`. Keyed by agent — an agent's lease is cap-1, so it has at -/// most one pill. +/// Publish the transitions between the previously-derived pill set and the +/// current one. `prev` is last loop's derived value, keyed by agent (an agent's +/// lease is cap-1, so at most one pill each). +/// +/// The pill set itself isn't owned or stored here — it is +/// [`super::JobQueue::running_transients`], recomputed from the graph. What this +/// publishes is the **edges**, which a derived read can't express on its own: +/// the dashboard wants `TransientSet` / `TransientCleared` events, and the crash +/// watcher wants the *moment* a pill cleared (its grace window is what keeps an +/// operator stop from reading as a crash). +/// +/// There used to be an RAII `TransientGuard` per pill here, and a hazard note +/// about dropping stale guards before creating new ones or a same-agent label +/// change would clear the pill it had just set. Both are gone: a guard exists so +/// a cancelled future can't *leak* an imperatively-set transient, and a derived +/// set has nothing to leak — a node that stops running simply stops appearing. +/// (Destroy and migration still take guards; they have no node behind them.) /// /// `deliberate_stop` rides along per node ([`NodeKind::takes_container_down`]) /// rather than being blanket-`true` for anything holding a lease: `Create` and @@ -159,24 +174,25 @@ fn handle_completion(coord: &Arc, done: NodeDone) { /// starting* is a real crash that must keep reporting as one. /// /// [`NodeKind::takes_container_down`]: super::NodeKind::takes_container_down -/// -/// ⚠️ **`retain` must run to completion before anything is created.** -/// `TransientGuard::drop` calls `clear_transient(agent)` — keyed by agent alone, -/// with no notion of *which* label it was clearing. So when a pill's label -/// changes for the same agent (which is now routine: the label follows the -/// running node as a DAG advances), creating the new guard first and dropping -/// the old second would clear the pill that was just set. Dropping first is what -/// makes the swap safe. -fn reconcile_transients( - coord: &Arc, - transients: &mut HashMap, -) { +fn reconcile_transients(coord: &Arc, prev: &mut HashMap) { let running = coord.job_queue.running_transients(); - transients.retain(|agent, (label, _)| running.iter().any(|(a, l, _)| a == agent && l == label)); + + // Cleared: in `prev`, gone (or relabelled) now. Emitted before the sets + // below so a same-agent label change reads as clear-then-set rather than + // two overlapping pills. + prev.retain(|agent, label| { + let still = running.iter().any(|(a, l, _)| a == agent && l == label); + if !still { + coord.clear_transient(agent); + } + still + }); + for (agent, label, takes_down) in running { - transients.entry(agent.clone()).or_insert_with(|| { - let guard = coord.transient_guard(&agent, label.clone(), takes_down); - (label, guard) - }); + if prev.get(&agent) == Some(&label) { + continue; + } + coord.set_transient(&agent, label.clone(), takes_down); + prev.insert(agent, label); } } From 884e39ba636c012616e3ae70ca81954704b96671 Mon Sep 17 00:00:00 2001 From: atlas Date: Sat, 1 Aug 2026 17:47:08 +0200 Subject: [PATCH 5/9] refactor(#2815): derive the transient snapshot, don't mirror it MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit mara on !2910: "why is set_transient still a thing if it completely derives from nodes?" It was still a thing because the scheduler mirrored the derived set into a stored map that every consumer read — derived state computed once and then cached, with the reconciliation loop existing only to keep the cache honest. `transient_snapshot()` now derives: `running_transients()` off the live graph, with the handful of entries that have no node behind them (destroy, migration) overlaid on top. There is no cached copy left to go stale or disagree with what is running. `set_transient` / `clear_transient` split by what they actually do: `set_manual_transient` / `clear_manual_transient` own the stored map for the no-node callers, and `emit_transient_set` / `emit_transient_cleared` publish the edges both paths need. Two things had to survive, and both are edges rather than state: - The dashboard's `TransientSet` / `TransientCleared` events. The scheduler carries the previous derived value and emits the diff. - The crash watcher's grace window. `recent_transient_within` answers "was a transient cleared just now?", which is what stops a deliberate stop from reading as a crash on the next 10s poll — a derived read of current state cannot answer it, so the clear still stamps. The scheduler keeps `deliberate_stop` alongside the label precisely so it is available at clear time: the node it came from is, by definition, no longer running to be asked. `TransientState::since` becomes wall-clock and, for derived entries, is the node's own `started_at` — the true start of the operation rather than the moment a watcher first noticed it, which is what the old guard-creation timestamp actually measured. `running_transients` returns a named `RunningTransient` rather than a 4-tuple; two of its fields are strings and one is a bool whose meaning is not guessable at a call site. Note for anyone reaching for a timestamp here: chrono is vendored with `default-features = false`, so there is no `Utc::now()`. The workspace convention is `wire_time::now_unix()` / `from_secs()`. Checked with clippy (`--all-targets -D warnings`), `cargo test -p hive-c0re -p hive-jobq` (322 + 41 passed) and `nix fmt`. --- hive-c0re/src/coordinator.rs | 128 +++++++++++++++------- hive-c0re/src/dashboard/state_snapshot.rs | 8 +- hive-c0re/src/job_queue/mod.rs | 37 ++++++- hive-c0re/src/job_queue/scheduler.rs | 24 ++-- 4 files changed, 142 insertions(+), 55 deletions(-) diff --git a/hive-c0re/src/coordinator.rs b/hive-c0re/src/coordinator.rs index 7325d92e..ade13914 100644 --- a/hive-c0re/src/coordinator.rs +++ b/hive-c0re/src/coordinator.rs @@ -8,6 +8,7 @@ use std::sync::atomic::{AtomicU64, Ordering}; use std::sync::{Arc, Mutex}; use anyhow::{Context, Result}; +use chrono::{DateTime, Utc}; use tokio::sync::{broadcast, watch}; use crate::approvals::Approvals; @@ -118,7 +119,9 @@ pub struct Coordinator { /// Agents whose lifecycle action (currently just spawn) is in flight. /// Read by the dashboard to render a spinner; cleared when the action /// resolves (success or failure). - transient: Mutex>, + /// Transients for work with NO queue node behind it (destroy, migration). + /// Queue-driven pills are derived, not stored — see `transient_snapshot`. + manual_transient: Mutex>, /// Tombstone for transients that have JUST been cleared. The /// crash watcher polls every 10s and would race the /// drop-clears-immediately path of `TransientGuard`: an operator @@ -381,7 +384,10 @@ pub struct TransientState { /// Set by whoever creates the transient, which is the only place that /// actually knows — it is not recoverable from `label`. pub deliberate_stop: bool, - pub since: std::time::Instant, + /// When the operation started. Wall-clock rather than `Instant` because a + /// derived entry takes it from the node's own `started_at` — the true start + /// of the work, not the moment a watcher first noticed it. + pub since: DateTime, } /// RAII handle returned by `Coordinator::transient_guard`. Cleared on @@ -398,7 +404,7 @@ pub struct TransientGuard { impl Drop for TransientGuard { fn drop(&mut self) { - self.coord.clear_transient(&self.name); + self.coord.clear_manual_transient(&self.name); } } @@ -537,7 +543,7 @@ impl Coordinator { agent_io_weight, model_prices, agents: Mutex::new(HashMap::new()), - transient: Mutex::new(HashMap::new()), + manual_transient: Mutex::new(HashMap::new()), recent_transient: Mutex::new(HashMap::new()), recent_crashes: Mutex::new(HashMap::new()), graceful_stop_pending: Mutex::new(HashSet::new()), @@ -1077,25 +1083,31 @@ impl Coordinator { /// Mark an agent as in-progress (only one state per agent for now). /// - /// Two callers, for two different reasons: - /// - **Work with a queue node behind it** — the job-queue scheduler, which - /// publishes the edges of a *derived* set - /// ([`crate::job_queue::JobQueue::running_transients`]). It needs no guard: - /// nothing is owned, and a node that stops running stops appearing. - /// - **Work with no node** (destroy, migration) — via the RAII - /// [`TransientGuard`], so the paired `clear_transient` still runs if the - /// surrounding future is cancelled (HTTP request aborted, shutdown - /// mid-rebuild, panic). There a bare set really would leak the transient - /// and pin the dashboard on "rebuilding…" forever. - pub(crate) fn set_transient(&self, name: &str, label: String, deliberate_stop: bool) { - self.transient.lock().unwrap().insert( + /// Record a transient for work that has **no queue node behind it** — + /// destroy and migration. Reached only through the RAII [`TransientGuard`], + /// so the paired clear still runs if the surrounding future is cancelled + /// (HTTP request aborted, shutdown mid-rebuild, panic); a bare set here + /// really would leak the transient and pin the dashboard on "rebuilding…". + /// + /// Queue-driven work does **not** come through here. Its pills are derived + /// from the running graph ([`crate::job_queue::JobQueue::running_transients`]) + /// and merged in by [`Coordinator::transient_snapshot`] — nothing is stored, + /// so nothing can go stale or leak. + fn set_manual_transient(&self, name: &str, label: String, deliberate_stop: bool) { + self.manual_transient.lock().unwrap().insert( name.to_owned(), TransientState { label: label.clone(), deliberate_stop, - since: std::time::Instant::now(), + since: hive_sh4re::wire_time::from_secs(hive_sh4re::wire_time::now_unix()), }, ); + self.emit_transient_set(name, label); + } + + /// Emit the "a pill appeared" edge. Shared by the manual path and the + /// job-queue scheduler, which publishes the transitions of its derived set. + pub(crate) fn emit_transient_set(&self, name: &str, label: String) { // Live-update dashboards. `since_unix` is wall-clock so the // browser can tick "Ns spawning…" without polling. The // intra-process map keeps using `Instant` for monotonicity. @@ -1112,30 +1124,37 @@ impl Coordinator { }); } - /// Clear an agent's transient state. Reached either from - /// [`TransientGuard`]'s `Drop` (the no-node callers) or from the scheduler - /// when a derived pill stops being current — see - /// [`Coordinator::set_transient`]. - pub(crate) fn clear_transient(&self, name: &str) { - let removed = self.transient.lock().unwrap().remove(name); + /// Drop a manual transient (see [`Coordinator::set_manual_transient`]). + /// Reached only from [`TransientGuard`]'s `Drop`, which guarantees it runs. + fn clear_manual_transient(&self, name: &str) { + let removed = self.manual_transient.lock().unwrap().remove(name); if let Some(state) = removed { - // Stamp the tombstone so the crash watcher can still see - // "operator kicked this off recently" on its next 10s poll - // — without this, the clear-then-poll race produced a - // spurious ContainerCrash on every operator stop/restart. - // Old entries get reaped lazily on read so the map doesn't - // grow unbounded. - self.recent_transient.lock().unwrap().insert( - name.to_owned(), - (state.deliberate_stop, std::time::Instant::now()), - ); - self.emit_dashboard_event(DashboardEvent::TransientCleared { - seq: self.next_seq(), - name: name.to_owned(), - }); + self.emit_transient_cleared(name, state.deliberate_stop); } } + /// Emit the "a pill went away" edge **and stamp the tombstone the crash + /// watcher reads**. Shared by the manual path and the job-queue scheduler. + /// + /// 🚨 The stamp is not bookkeeping. Without it the clear-then-poll race + /// produced a spurious `ContainerCrash` on **every** operator stop/restart: + /// the transient is gone by the time the 10s poll looks, so a deliberate + /// stop is indistinguishable from a crash. `recent_transient_within` is what + /// closes that window, which is why the clear has to be an *event* — a + /// derived read of current state cannot answer "was one here a moment ago?". + /// + /// Old entries are reaped lazily on read, so the map stays bounded. + pub(crate) fn emit_transient_cleared(&self, name: &str, deliberate_stop: bool) { + self.recent_transient.lock().unwrap().insert( + name.to_owned(), + (deliberate_stop, std::time::Instant::now()), + ); + self.emit_dashboard_event(DashboardEvent::TransientCleared { + seq: self.next_seq(), + name: name.to_owned(), + }); + } + /// Mark `name` as having a graceful stop in progress. While set, /// `socket_server::handle_recv` returns `Response::GracefulStop` for /// this agent instead of polling the broker (the inbound fence). @@ -1259,15 +1278,46 @@ impl Coordinator { label: impl Into, deliberate_stop: bool, ) -> TransientGuard { - self.set_transient(name, label.into(), deliberate_stop); + self.set_manual_transient(name, label.into(), deliberate_stop); TransientGuard { coord: self.clone(), name: name.to_owned(), } } + /// Every live transient, keyed by agent. + /// + /// **Derived on read**, not stored: the queue-driven pills come straight + /// from the running graph, so there is no cached copy to go stale, leak, or + /// disagree with what is actually running. The only stored entries are the + /// handful with no node behind them (destroy, migration), overlaid on top — + /// they win, since an agent being destroyed is the more urgent truth than + /// whatever node was mid-flight when it started. + #[must_use] pub fn transient_snapshot(&self) -> HashMap { - self.transient.lock().unwrap().clone() + let mut out: HashMap = self + .job_queue + .running_transients() + .into_iter() + .map(|t| { + ( + t.agent, + TransientState { + label: t.label, + deliberate_stop: t.takes_container_down, + since: t.since, + }, + ) + }) + .collect(); + out.extend( + self.manual_transient + .lock() + .unwrap() + .iter() + .map(|(k, v)| (k.clone(), v.clone())), + ); + out } /// Drop a system message into the given agent's inbox. Wakes the diff --git a/hive-c0re/src/dashboard/state_snapshot.rs b/hive-c0re/src/dashboard/state_snapshot.rs index 0016487a..04c43935 100644 --- a/hive-c0re/src/dashboard/state_snapshot.rs +++ b/hive-c0re/src/dashboard/state_snapshot.rs @@ -535,7 +535,13 @@ fn build_transient_views( .map(|(name, st)| TransientView { name: name.clone(), kind: st.label.clone(), - secs: st.since.elapsed().as_secs(), + // Clamped at 0: `since` is wall-clock now (the node's own + // `started_at`), so a backwards clock adjustment could otherwise + // render a negative age. + secs: (hive_sh4re::wire_time::from_secs(hive_sh4re::wire_time::now_unix()) - st.since) + .num_seconds() + .max(0) + .cast_unsigned(), }) .collect() } diff --git a/hive-c0re/src/job_queue/mod.rs b/hive-c0re/src/job_queue/mod.rs index 71f6263e..ada7faad 100644 --- a/hive-c0re/src/job_queue/mod.rs +++ b/hive-c0re/src/job_queue/mod.rs @@ -59,6 +59,27 @@ const MAX_HISTORY_DAGS: usize = 50; /// Cap on stored node error strings. const MAX_ERROR_LEN: usize = 2_000; +/// One live transient pill, derived from a running node. +/// +/// A named struct rather than a tuple because three of its four fields are +/// easy to confuse at a call site: two are strings and two answer questions +/// nobody should have to guess at ("is this the agent or the label?", "does +/// this bool mean deliberate or running?"). +#[derive(Debug, Clone)] +pub struct RunningTransient { + /// The agent whose lease the node declared. + pub agent: String, + /// The node's own wire tag, rendered as the pill. + pub label: String, + /// Whether this operation is expected to take the container down — the + /// crash watcher's input. See [`NodeKind::takes_container_down`]. + pub takes_container_down: bool, + /// When the node started running, so the dashboard can tick elapsed + /// seconds. Taken from the node itself, which is the true start of the + /// operation rather than the moment a watcher noticed it. + pub since: DateTime, +} + /// A node claimed for execution — everything the executor needs, snapshotted at /// claim time. #[derive(Debug, Clone)] @@ -422,7 +443,7 @@ impl JobQueue { /// the resources-where-constructed work, not this function. An agent's lease /// is cap-1, so at most one entry per agent. #[must_use] - pub fn running_transients(&self) -> Vec<(String, String, bool)> { + pub fn running_transients(&self) -> Vec { let inner = self.lock(); inner .sched @@ -441,11 +462,17 @@ impl JobQueue { } => Some(a), _ => None, })?; - Some(( + Some(RunningTransient { agent, - n.payload.as_str().to_owned(), - n.payload.takes_container_down(), - )) + label: n.payload.as_str().to_owned(), + takes_container_down: n.payload.takes_container_down(), + // `started_at` is set when a node enters `Running`, and this + // only sees `Running` nodes — the fallback is unreachable in + // practice, and "just now" is the honest answer if it isn't. + since: n + .started_at + .unwrap_or_else(|| hive_sh4re::wire_time::from_secs(now_unix())), + }) }) .collect() } diff --git a/hive-c0re/src/job_queue/scheduler.rs b/hive-c0re/src/job_queue/scheduler.rs index 492aa55d..105d5326 100644 --- a/hive-c0re/src/job_queue/scheduler.rs +++ b/hive-c0re/src/job_queue/scheduler.rs @@ -59,7 +59,7 @@ pub async fn run_worker(coord: Arc) { // it exists to spot transitions, since the dashboard wants edges // (`TransientSet` / `TransientCleared`) and the crash watcher wants the // moment of the clear. Nothing owns a pill; nothing can leak one. - let mut transients: HashMap = HashMap::new(); + let mut transients: HashMap = HashMap::new(); loop { // Checked every iteration, not just in the `select!` below — a // continuous stream of ready claims never reaches the `select!`, so @@ -174,25 +174,29 @@ fn handle_completion(coord: &Arc, done: NodeDone) { /// starting* is a real crash that must keep reporting as one. /// /// [`NodeKind::takes_container_down`]: super::NodeKind::takes_container_down -fn reconcile_transients(coord: &Arc, prev: &mut HashMap) { +fn reconcile_transients(coord: &Arc, prev: &mut HashMap) { let running = coord.job_queue.running_transients(); // Cleared: in `prev`, gone (or relabelled) now. Emitted before the sets // below so a same-agent label change reads as clear-then-set rather than - // two overlapping pills. - prev.retain(|agent, label| { - let still = running.iter().any(|(a, l, _)| a == agent && l == label); + // two overlapping pills. `deliberate_stop` is carried in `prev` precisely + // so it is still available *here* — the node it came from is, by + // definition, no longer running to be asked. + prev.retain(|agent, (label, deliberate)| { + let still = running + .iter() + .any(|t| &t.agent == agent && &t.label == label); if !still { - coord.clear_transient(agent); + coord.emit_transient_cleared(agent, *deliberate); } still }); - for (agent, label, takes_down) in running { - if prev.get(&agent) == Some(&label) { + for t in running { + if prev.get(&t.agent).map(|(l, _)| l) == Some(&t.label) { continue; } - coord.set_transient(&agent, label.clone(), takes_down); - prev.insert(agent, label); + coord.emit_transient_set(&t.agent, t.label.clone()); + prev.insert(t.agent, (t.label, t.takes_container_down)); } } From 7f920718f2adebd6584286e8d5219658f02500a7 Mon Sep 17 00:00:00 2001 From: atlas Date: Sat, 1 Aug 2026 18:15:27 +0200 Subject: [PATCH 6/9] refactor(#2815): drop the imperative transient path entirely MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit mara on !2910: "also remove imperative path for the things that are not nodes yet, file follow up issue to fix that". `TransientGuard`, the stored map and both manual set/clear are gone. `transient_snapshot()` is derived and nothing else — destroy and migration show no pill, because there is no node to derive one from. The pill returns for free when they become nodes. What those three guards were actually doing, though, was suppressing the crash watcher, not drawing a pill. `migrate.rs` said so in its own comment: without it, `crash_watch` fires `ContainerCrash` for every migrated agent and the manager tries to recover containers that were stopped on purpose. Destroy is the same — the container disappears deliberately and nothing in the graph says so. Deleting them outright would therefore have traded a dashboard pill for false crash alerts on every destroy and every migration. So the suppression survives as its own thing, `suppress_crash_watch`, with a name that says what it is. It is still RAII, and still held for the operation rather than stamped once, because the crash watcher's grace window is finite and a destroy is not — a single tombstone would expire mid-operation. The drop stamps the tombstone, covering the poll that lands just after. That leaves RAII in the codebase for exactly one purpose instead of two. Untangling the pill from the suppression is what made the transient layer deletable at all. Follow-up issue for making destroy + migration real queue nodes to follow; at that point this guard goes too. Checked with clippy (`--all-targets -D warnings`), `cargo test -p hive-c0re -p hive-jobq` (322 + 41 passed) and `nix fmt`. --- hive-c0re/src/actions.rs | 7 +- hive-c0re/src/coordinator.rs | 140 +++++++++++++-------------- hive-c0re/src/migrate.rs | 9 +- hive-c0re/src/workers/crash_watch.rs | 8 +- 4 files changed, 80 insertions(+), 84 deletions(-) diff --git a/hive-c0re/src/actions.rs b/hive-c0re/src/actions.rs index b2a354ea..d3eb90c4 100644 --- a/hive-c0re/src/actions.rs +++ b/hive-c0re/src/actions.rs @@ -882,9 +882,10 @@ pub async fn destroy(coord: &Arc, name: &str, purge: bool) -> Resul tracing::info!(%name, purge, "destroy"); // Guard auto-clears on the success path's final scope exit and on // every early-return / cancellation along the way. - // Destroy is not a queue node, so it names its own label. `true`: the - // container is going away, so its disappearance must not read as a crash. - let guard = coord.transient_guard(name, "destroying", true); + // Destroy has no queue node behind it, so nothing in the graph says this + // container is going away on purpose — without this the crash watcher + // reports every destroy as a crash and the manager tries to recover it. + let guard = coord.suppress_crash_watch(name); lifecycle::destroy(name).await?; coord.unregister_agent(name); let runtime = crate::paths::agent_runtime_dir(name); diff --git a/hive-c0re/src/coordinator.rs b/hive-c0re/src/coordinator.rs index ade13914..d03ec18a 100644 --- a/hive-c0re/src/coordinator.rs +++ b/hive-c0re/src/coordinator.rs @@ -119,9 +119,10 @@ pub struct Coordinator { /// Agents whose lifecycle action (currently just spawn) is in flight. /// Read by the dashboard to render a spinner; cleared when the action /// resolves (success or failure). - /// Transients for work with NO queue node behind it (destroy, migration). - /// Queue-driven pills are derived, not stored — see `transient_snapshot`. - manual_transient: Mutex>, + /// Agents whose container is being taken down by work with **no queue node + /// behind it** (destroy, migration), so the crash watcher must not report + /// the disappearance as a crash. Not a pill — see [`CrashWatchSuppression`]. + crash_suppressed: Mutex>, /// Tombstone for transients that have JUST been cleared. The /// crash watcher polls every 10s and would race the /// drop-clears-immediately path of `TransientGuard`: an operator @@ -390,21 +391,45 @@ pub struct TransientState { pub since: DateTime, } -/// RAII handle returned by `Coordinator::transient_guard`. Cleared on -/// drop — including drop-via-cancellation, the path that bare -/// `set_transient` / `clear_transient` pairs leaked through. Holds an -/// `Arc` so the guard is freely returnable / movable. -#[must_use = "the guard clears the transient when dropped; bind it for the operation's \ - duration (`let _guard = coord.transient_guard(...)`). An unbound call drops \ - it immediately and un-sets the transient at once — the exact footgun this guards against."] -pub struct TransientGuard { +/// RAII handle returned by [`Coordinator::suppress_crash_watch`]. While held, +/// the crash watcher treats this container disappearing as **expected**. +/// +/// This is *not* a dashboard pill. Transients are derived from running queue +/// nodes and nothing stores them. But destroy and migration take a container +/// down without a node behind them, so nothing in the graph says the +/// disappearance was intended — and without that, `crash_watch` fires a +/// `ContainerCrash` for every destroy and every migrated agent, and the manager +/// tries to "recover" containers that were removed on purpose. +/// +/// It is held rather than stamped once because +/// [`crate::workers::crash_watch`]'s grace window is finite and these +/// operations are not: a long destroy would outlive a single tombstone. The +/// tombstone is stamped on drop, covering the poll that lands just after. +/// +/// Goes away entirely once destroy + migration are real queue nodes. +#[must_use = "suppression lasts as long as the guard; bind it for the operation's duration \ + (`let _guard = coord.suppress_crash_watch(...)`). An unbound call drops it \ + immediately and the very next poll can report a deliberate stop as a crash."] +pub struct CrashWatchSuppression { coord: Arc, name: String, } -impl Drop for TransientGuard { +impl Drop for CrashWatchSuppression { fn drop(&mut self) { - self.coord.clear_manual_transient(&self.name); + self.coord + .crash_suppressed + .lock() + .unwrap() + .remove(&self.name); + // Tombstone the release so the next poll — which may land in the + // window between the container going away and this guard dropping — + // still reads the stop as deliberate. + self.coord + .recent_transient + .lock() + .unwrap() + .insert(self.name.clone(), (true, std::time::Instant::now())); } } @@ -543,7 +568,7 @@ impl Coordinator { agent_io_weight, model_prices, agents: Mutex::new(HashMap::new()), - manual_transient: Mutex::new(HashMap::new()), + crash_suppressed: Mutex::new(HashSet::new()), recent_transient: Mutex::new(HashMap::new()), recent_crashes: Mutex::new(HashMap::new()), graceful_stop_pending: Mutex::new(HashSet::new()), @@ -1081,32 +1106,8 @@ impl Coordinator { self.agents.lock().unwrap().keys().cloned().collect() } - /// Mark an agent as in-progress (only one state per agent for now). - /// - /// Record a transient for work that has **no queue node behind it** — - /// destroy and migration. Reached only through the RAII [`TransientGuard`], - /// so the paired clear still runs if the surrounding future is cancelled - /// (HTTP request aborted, shutdown mid-rebuild, panic); a bare set here - /// really would leak the transient and pin the dashboard on "rebuilding…". - /// - /// Queue-driven work does **not** come through here. Its pills are derived - /// from the running graph ([`crate::job_queue::JobQueue::running_transients`]) - /// and merged in by [`Coordinator::transient_snapshot`] — nothing is stored, - /// so nothing can go stale or leak. - fn set_manual_transient(&self, name: &str, label: String, deliberate_stop: bool) { - self.manual_transient.lock().unwrap().insert( - name.to_owned(), - TransientState { - label: label.clone(), - deliberate_stop, - since: hive_sh4re::wire_time::from_secs(hive_sh4re::wire_time::now_unix()), - }, - ); - self.emit_transient_set(name, label); - } - - /// Emit the "a pill appeared" edge. Shared by the manual path and the - /// job-queue scheduler, which publishes the transitions of its derived set. + /// Emit the "a pill appeared" edge, for the job-queue scheduler publishing + /// the transitions of its derived set. pub(crate) fn emit_transient_set(&self, name: &str, label: String) { // Live-update dashboards. `since_unix` is wall-clock so the // browser can tick "Ns spawning…" without polling. The @@ -1124,17 +1125,8 @@ impl Coordinator { }); } - /// Drop a manual transient (see [`Coordinator::set_manual_transient`]). - /// Reached only from [`TransientGuard`]'s `Drop`, which guarantees it runs. - fn clear_manual_transient(&self, name: &str) { - let removed = self.manual_transient.lock().unwrap().remove(name); - if let Some(state) = removed { - self.emit_transient_cleared(name, state.deliberate_stop); - } - } - /// Emit the "a pill went away" edge **and stamp the tombstone the crash - /// watcher reads**. Shared by the manual path and the job-queue scheduler. + /// watcher reads**, for the job-queue scheduler. /// /// 🚨 The stamp is not bookkeeping. Without it the clear-then-poll race /// produced a spurious `ContainerCrash` on **every** operator stop/restart: @@ -1272,31 +1264,37 @@ impl Coordinator { /// operation takes the container down on purpose, and is what the crash /// watcher reads. Only the caller knows the second one — it is not /// recoverable from the first. - pub fn transient_guard( - self: &Arc, - name: &str, - label: impl Into, - deliberate_stop: bool, - ) -> TransientGuard { - self.set_manual_transient(name, label.into(), deliberate_stop); - TransientGuard { + pub fn suppress_crash_watch(self: &Arc, name: &str) -> CrashWatchSuppression { + self.crash_suppressed + .lock() + .unwrap() + .insert(name.to_owned()); + CrashWatchSuppression { coord: self.clone(), name: name.to_owned(), } } + /// Whether a no-node operation is currently taking this container down. + #[must_use] + pub fn crash_watch_suppressed(&self, name: &str) -> bool { + self.crash_suppressed.lock().unwrap().contains(name) + } + /// Every live transient, keyed by agent. /// - /// **Derived on read**, not stored: the queue-driven pills come straight - /// from the running graph, so there is no cached copy to go stale, leak, or - /// disagree with what is actually running. The only stored entries are the - /// handful with no node behind them (destroy, migration), overlaid on top — - /// they win, since an agent being destroyed is the more urgent truth than - /// whatever node was mid-flight when it started. + /// **Derived on read, stored nowhere.** Straight off the running graph, so + /// there is no cached copy to go stale, leak, or disagree with what is + /// actually running. + /// + /// Work with no queue node behind it (destroy, migration) therefore shows + /// **no pill** — there is nothing in the graph to derive one from. Its + /// crash-watch suppression is a separate, narrower thing + /// ([`Coordinator::suppress_crash_watch`]); the pill comes back for free + /// once those become real nodes. #[must_use] pub fn transient_snapshot(&self) -> HashMap { - let mut out: HashMap = self - .job_queue + self.job_queue .running_transients() .into_iter() .map(|t| { @@ -1309,15 +1307,7 @@ impl Coordinator { }, ) }) - .collect(); - out.extend( - self.manual_transient - .lock() - .unwrap() - .iter() - .map(|(k, v)| (k.clone(), v.clone())), - ); - out + .collect() } /// Drop a system message into the given agent's inbox. Wakes the diff --git a/hive-c0re/src/migrate.rs b/hive-c0re/src/migrate.rs index 38a5d5e2..56f1f0bb 100644 --- a/hive-c0re/src/migrate.rs +++ b/hive-c0re/src/migrate.rs @@ -109,9 +109,8 @@ pub async fn run(coord: &Arc) -> Result<()> { // would fire ContainerCrash for every agent here and the // manager would spuriously try to recover them. // No queue node behind this one — migration repoints containers - // directly — so the label is supplied here. `true`: the repoint takes - // the container down, which is the whole reason for the guard. - let guard = coord.transient_guard(name.as_str(), "rebuilding", true); + // directly — so nothing in the graph marks the stop as intended. + let guard = coord.suppress_crash_watch(name.as_str()); let result = repoint_container(name.as_str()).await; drop(guard); if let Err(e) = result { @@ -203,8 +202,8 @@ async fn rename_manager_container(coord: &Arc) { return; } tracing::info!("migration phase 5: renaming root container to h-root"); - // `true`: the old container is stopped immediately below. - let _guard = coord.transient_guard(MANAGER_NAME, "rebuilding", true); + // The old container is stopped immediately below, on purpose. + let _guard = coord.suppress_crash_watch(MANAGER_NAME); // Stop the old container. Abort if stop fails — continuing with a // running `root` and then starting `h-root` risks two manager diff --git a/hive-c0re/src/workers/crash_watch.rs b/hive-c0re/src/workers/crash_watch.rs index 9b74f826..73999f1c 100644 --- a/hive-c0re/src/workers/crash_watch.rs +++ b/hive-c0re/src/workers/crash_watch.rs @@ -87,7 +87,13 @@ fn emit_crash_transitions(coord: &Coordinator, prev: &HashSet, current: // guard between two crash-watch polls. let recent = coord.recent_transient_within(RECENT_TRANSIENT_GRACE); for stopped in prev.difference(current) { - let active = transients.get(stopped).map(|st| st.deliberate_stop); + // Two sources, because a container can go down on purpose either way: + // a running queue node that declared it takes the container down, or a + // no-node operation (destroy, migration) holding a suppression guard. + let active = transients + .get(stopped) + .map(|st| st.deliberate_stop) + .or_else(|| coord.crash_watch_suppressed(stopped).then_some(true)); let recently_cleared = recent.get(stopped).copied(); if is_deliberate_stop(active, recently_cleared) { continue; From ff2721cc0f60fff6c19adb2eb90e78238cbcb0e7 Mon Sep 17 00:00:00 2001 From: atlas Date: Sat, 1 Aug 2026 19:13:50 +0200 Subject: [PATCH 7/9] docs(#2815): fix suppress_crash_watch's stale doc comment MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The doc still described a `label`/`deliberate_stop` parameter pair inherited from `transient_guard`, which this function replaced and whose signature it does not share — it only takes `name`. Rewritten to say what it does and, more usefully, what must not come through it: the queue answers the same question from the node itself via `NodeKind::takes_container_down`, so this path is only for the two operations that have no node behind them yet. --- hive-c0re/src/coordinator.rs | 20 +++++++++++--------- 1 file changed, 11 insertions(+), 9 deletions(-) diff --git a/hive-c0re/src/coordinator.rs b/hive-c0re/src/coordinator.rs index d03ec18a..4b9e1366 100644 --- a/hive-c0re/src/coordinator.rs +++ b/hive-c0re/src/coordinator.rs @@ -1254,16 +1254,18 @@ impl Coordinator { map.iter().map(|(k, v)| (k.clone(), v.len())).collect() } - /// Set a transient state and return a guard that clears it on drop. - /// Use this from any path where the surrounding future could be - /// cancelled or panic between set and clear (HTTP handlers, spawned - /// tasks). The guard's `Drop` runs even on task cancellation, so - /// the dashboard's spinner can't get pinned forever. + /// Tell the crash watcher that `name`'s container is going down **on + /// purpose**, for the lifetime of the returned guard. See + /// [`CrashWatchSuppression`] for why this exists at all. /// - /// `label` is what the pill renders; `deliberate_stop` says whether this - /// operation takes the container down on purpose, and is what the crash - /// watcher reads. Only the caller knows the second one — it is not - /// recoverable from the first. + /// Only for the operations with no queue node behind them. Anything the + /// job queue runs answers this from the node itself + /// ([`crate::job_queue::NodeKind::takes_container_down`]) and must not come + /// through here. + /// + /// The guard's `Drop` runs even on task cancellation, so an aborted HTTP + /// request or a panic mid-destroy can't leave a container permanently + /// exempt from crash reporting. pub fn suppress_crash_watch(self: &Arc, name: &str) -> CrashWatchSuppression { self.crash_suppressed .lock() From 41c1b1a3fbf6b23380de7ab09d64ad94ef23069e Mon Sep 17 00:00:00 2001 From: iris Date: Sat, 1 Aug 2026 19:11:46 +0200 Subject: [PATCH 8/9] agent web UI: bulk mark-done for the todos flyout MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The per-agent web UI todos flyout (loose-ends v2) had no mark-done affordance at all — dismissing a todo was only possible via the cancel_loose_end MCP tool, one id at a time. Add a checkbox per row, a select-all/select-none/mark-done bulk row, and a new POST /api/todos/mark-done handler that loops the existing single-id MarkTodoDone request over the in-agent socket (no new wire request type needed — the todos list is small, so N same-host round-trips is cheap). Fixes #2917 --- frontend/packages/agent/src/agent.css | 15 ++++-- frontend/packages/agent/src/app.js | 69 ++++++++++++++++++++++++++- hive-agent/src/web_ui/actions.rs | 41 +++++++++++++++- hive-agent/src/web_ui/mod.rs | 1 + 4 files changed, 119 insertions(+), 7 deletions(-) diff --git a/frontend/packages/agent/src/agent.css b/frontend/packages/agent/src/agent.css index 2b9d7007..79f32625 100644 --- a/frontend/packages/agent/src/agent.css +++ b/frontend/packages/agent/src/agent.css @@ -530,6 +530,12 @@ pre.diff { .agent-inbox .inbox-ts { color: var(--muted); font-size: 0.9em; margin-left: 0.5em; } .agent-inbox .inbox-from { color: var(--amber); } .agent-inbox .inbox-sep { color: var(--muted); margin-left: 0.4em; } +/* Todos flyout per-row checkbox (bulk mark-done) — sits inline before the + existing `.inbox-from` label, same row. */ +.agent-inbox .todo-cb { + vertical-align: middle; + margin-right: 0.2em; +} .agent-inbox .inbox-body { display: block; color: var(--fg); @@ -655,10 +661,11 @@ pre.diff { text-decoration-color: var(--muted); } -/* "mark all read" header row sits above the recent-messages list - in the inbox side-panel flyout. Same look as the answer-form - button (mauve hover, bg-elev background) so they read as part - of the same affordance family. */ +/* Bulk-action header row: "mark all read" above the recent-messages + list in the inbox flyout, and "select all / select none / mark + done" above the list in the todos flyout — same classes, shared + look (mauve hover, bg-elev background) so both read as part of + the same affordance family as the answer-form button. */ .agent-inbox .inbox-mark-all-row { display: flex; gap: 0.6em; diff --git a/frontend/packages/agent/src/app.js b/frontend/packages/agent/src/app.js index 729e4b06..b53ecaeb 100644 --- a/frontend/packages/agent/src/app.js +++ b/frontend/packages/agent/src/app.js @@ -861,8 +861,67 @@ window.marked = marked; return wrap; } + /** Bulk "mark done" row for the todos flyout: select all / select none + * + a mark-done button, disabled until at least one row is checked. + * POSTs the checked ids (comma-joined into one field, same shape as + * `hive-c0re`'s meta-inputs bulk form — axum's `Form` extractor doesn't + * natively decode repeated same-name keys) to this agent's own + * `/api/todos/mark-done`, then calls `refreshTodos()` on success so the + * flyout reloads without the now-dismissed rows. `wrap` is the panel + * root — bulk buttons read/toggle the checkboxes it contains. */ + function buildTodosMarkDoneRow(wrap) { + const status = el('span', { class: 'inbox-mark-status' }); + const selAll = el('button', { type: 'button', class: 'inbox-mark-all-btn' }, 'select all'); + const selNone = el('button', { type: 'button', class: 'inbox-mark-all-btn' }, 'select none'); + const markBtn = el('button', { + type: 'button', class: 'inbox-mark-all-btn', disabled: '', + }, '✓ mark done'); + const checkboxes = () => Array.from(wrap.querySelectorAll('input[data-todo-id]')); + const refreshDisabled = () => { + const any = checkboxes().some((cb) => cb.checked); + if (any) markBtn.removeAttribute('disabled'); + else markBtn.setAttribute('disabled', ''); + }; + selAll.addEventListener('click', () => { + checkboxes().forEach((cb) => { cb.checked = true; }); + refreshDisabled(); + }); + selNone.addEventListener('click', () => { + checkboxes().forEach((cb) => { cb.checked = false; }); + refreshDisabled(); + }); + wrap.addEventListener('change', (e) => { + if (e.target.matches('input[data-todo-id]')) refreshDisabled(); + }); + markBtn.addEventListener('click', () => { + const ids = checkboxes().filter((cb) => cb.checked).map((cb) => cb.dataset.todoId); + if (!ids.length) return; + status.textContent = 'marking…'; + asyncBtn(markBtn, async () => { + try { + const resp = await fetch('api/todos/mark-done', { + method: 'POST', + headers: { 'Content-Type': 'application/x-www-form-urlencoded' }, + body: 'ids=' + encodeURIComponent(ids.join(',')), + }); + if (resp.ok) { + status.textContent = '✓ marked done'; + refreshTodos(); + } else { + status.textContent = 'failed: ' + (await resp.text()); + } + } catch (err) { + status.textContent = 'failed: ' + err; + } + }); + }); + return el('div', { class: 'inbox-mark-all-row' }, selAll, selNone, markBtn, status); + } + /** Build the todos side-panel list. Each entry is a LooseEnd::Todo - * (subsystem, summary, source, age_seconds). */ + * (id, subsystem, summary, source, age_seconds). A checkbox per row + * plus the bulk row above lets the operator dismiss several at once + * instead of one `cancel_loose_end` call at a time. */ function buildTodosList(todos) { const wrap = el('div', { class: 'agent-inbox' }); if (!todos.length) { @@ -870,6 +929,7 @@ window.marked = marked; 'no todos — all subsystem queues are clear.')); return wrap; } + wrap.append(buildTodosMarkDoneRow(wrap)); const list = el('ul'); const fmtAge = (s) => { if (s < 60) return s + 's'; @@ -880,8 +940,13 @@ window.marked = marked; for (const t of todos) { const li = el('li'); const label = t.source ? t.subsystem + ' · ' + t.source : t.subsystem; + const cbId = 'todo-cb-' + t.id; + const cb = el('input', { + type: 'checkbox', id: cbId, class: 'todo-cb', 'data-todo-id': String(t.id), + }); li.append( - el('span', { class: 'inbox-from' }, label), ' ', + cb, ' ', + el('label', { for: cbId, class: 'inbox-from' }, label), ' ', el('span', { class: 'inbox-ts' }, fmtAge(t.age_seconds || 0) + ' ago'), el('div', { class: 'inbox-body' }, t.summary || ''), ); diff --git a/hive-agent/src/web_ui/actions.rs b/hive-agent/src/web_ui/actions.rs index de060f92..52cb4e75 100644 --- a/hive-agent/src/web_ui/actions.rs +++ b/hive-agent/src/web_ui/actions.rs @@ -1,4 +1,5 @@ -//! Operator action POST handlers (send, cancel, compact, model, effort, reset). +//! Operator action POST handlers (send, cancel, compact, model, effort, +//! reset, todos mark-done). use axum::{ Form, @@ -162,3 +163,41 @@ pub(super) async fn post_set_effort( tracing::info!(%level, "operator set effort"); (axum::http::StatusCode::OK, "ok").into_response() } + +#[derive(Deserialize)] +pub(super) struct MarkTodosDoneForm { + /// Comma-separated todo ids. Same "one field, JS joins the checked + /// boxes" shape as `hive-c0re`'s `meta_inputs::MetaUpdateForm` — axum's + /// `Form` extractor doesn't natively decode repeated same-name keys. + ids: String, +} + +/// `POST /api/todos/mark-done` — dismiss one or more of this agent's own +/// todos (loose-ends v2) from the todos flyout. Loops a `MarkTodoDone` call +/// per id over the in-agent socket rather than adding a new bulk request to +/// `hive-agent-sock`: the todos list is small (single-digit rows most of the +/// time), so N same-host socket round-trips isn't a real cost, and it keeps +/// the wire protocol's `Request` enum — already used by the `cancel_loose_end` +/// MCP tool — unchanged. Unknown/already-acked ids just don't add to the +/// `acked` count (same "acking twice is not a new action" semantics as the +/// single-id path); a request with no ids or where every id fails to parse +/// is rejected as a client error rather than silently acking nothing. +pub(super) async fn post_mark_todos_done(Form(form): Form) -> Response { + let ids: Vec = form + .ids + .split(',') + .filter_map(|s| s.trim().parse::().ok()) + .collect(); + if ids.is_empty() { + return error_response(StatusCode::BAD_REQUEST, "mark-done: no todo ids selected"); + } + let mut acked = 0u64; + for id in ids { + if let Some(hive_agent_sock::Response::Acked { count }) = + crate::todo_server::dial(&hive_agent_sock::Request::MarkTodoDone { id }).await + { + acked += count; + } + } + axum::Json(serde_json::json!({ "acked": acked })).into_response() +} diff --git a/hive-agent/src/web_ui/mod.rs b/hive-agent/src/web_ui/mod.rs index 5b06e15d..4693c361 100644 --- a/hive-agent/src/web_ui/mod.rs +++ b/hive-agent/src/web_ui/mod.rs @@ -116,6 +116,7 @@ pub async fn serve( .route("/api/new-session", post(actions::post_new_session)) .route("/api/logout", post(auth::post_logout)) .route("/api/todos", get(stats::api_todos)) + .route("/api/todos/mark-done", post(actions::post_mark_todos_done)) .route("/api/stats", get(stats::api_stats)) .route("/screen/ws", get(screen::screen_ws)) .route("/icon", get(screen::serve_icon)); From e02ac1e86eef3a4d46871f35e4df5dbf404ed037 Mon Sep 17 00:00:00 2001 From: atlas Date: Sat, 1 Aug 2026 21:01:13 +0200 Subject: [PATCH 9/9] refactor(#2916): drop the two obsolete startup migrations MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Phase 4 (repoint every container onto `meta#`) and phase 5 (rename the `root` container to `h-root`) were marker-guarded one-shots for layouts no live hive still has: containers are rendered onto `meta#` at creation, and the `h-` prefix has been the naming for far longer than any deployment predates. A one-shot nobody can still trigger is dead weight, so both are gone along with `repoint_container`, `rename_manager_container`, `CONTAINER_TIMEOUT` and the two marker paths. Phase 6 was not obsolete, only misplaced. Ruth's tool groups are now seeded by `ensure_root_agent` on the one path that creates her, rather than re-asserted on every hive-c0re boot. The skip-if-already-set guard survives the move: a destroy+recreate under the same name must not reset an operator's chosen group set back to MANAGER_DEFAULT. That also settles a latent bug. Phase 4's marker check was a `return`, not a skip, so on any hive carrying the marker phases 5 and 6 never ran at all — the tool-group backfill, whose whole job was preventing a silent privilege downgrade, has not executed here in a long time. Moving it to create-time removes the question rather than answering it. What stays is convergence: three unguarded, idempotent phases that re-run each boot and no-op once their state is right. The module doc now names the three categories so the next person can tell which kind they're adding. --- docs/approvals.md | 16 +- docs/persistence.md | 10 +- hive-c0re/src/migrate.rs | 229 +++------------------------ hive-c0re/src/paths.rs | 12 -- hive-c0re/src/workers/auto_update.rs | 31 ++++ 5 files changed, 70 insertions(+), 228 deletions(-) diff --git a/docs/approvals.md b/docs/approvals.md index 527f402b..e5305769 100644 --- a/docs/approvals.md +++ b/docs/approvals.md @@ -478,12 +478,16 @@ each phase is a no-op once already applied. Behaviour: dirs, or claude creds. - Meta-flake phase: rewrites each `applied//flake.nix` to the module-only boilerplate, wires the `applied` remote in - each proposed repo, bootstraps the meta repo from the - current agent list, and `nixos-container update`s every - container at `meta#`. The expensive last step is - guarded by `/var/lib/hyperhive/.meta-migration-done` so - it only runs once across hive-c0re restarts. Set - `HIVE_SKIP_META_MIGRATION=1` on the service to defer. + each proposed repo, and bootstraps the meta repo from the + current agent list. Set `HIVE_SKIP_META_MIGRATION=1` on the + service to defer. + + A further step used to `nixos-container update` every + container onto `meta#`, guarded by a marker file so it + ran once per hive. It is gone: containers have been rendered + onto `meta#` at creation for long enough that no live hive + needs the repoint, and a one-shot nobody can still trigger is + dead weight. Same for the `root` → `h-root` container rename. No state loss in either migration. claude creds, /state/ notes, the events DB, proposed history, and applied history diff --git a/docs/persistence.md b/docs/persistence.md index 7c4ff619..871a45f8 100644 --- a/docs/persistence.md +++ b/docs/persistence.md @@ -324,11 +324,11 @@ Contents: The root agent has the meta dir RO-mounted at `/meta/`. -Marker file `/var/lib/hyperhive/.meta-migration-done` is -written by the startup migration after every container has -been repointed at `meta#`. Removing it forces a re-run on -next hive-c0re start (idempotent — only the actual repoint -step would re-fire). +There is no longer a `.meta-migration-done` marker: the +one-shot container repoint it guarded has been removed, since +containers are rendered onto `meta#` at creation. A stale +marker file left over from an older hive is inert and can be +deleted. ## Destroy vs purge diff --git a/hive-c0re/src/migrate.rs b/hive-c0re/src/migrate.rs index 56f1f0bb..f80a30fa 100644 --- a/hive-c0re/src/migrate.rs +++ b/hive-c0re/src/migrate.rs @@ -1,8 +1,22 @@ -//! Startup auto-migration. Six idempotent phases: applied repo, -//! proposed repo, meta repo, container repoint, root→h-root rename, -//! and manager tool-groups backfill. -//! Kill-switch: `HIVE_SKIP_META_MIGRATION=1`. Full migration sequence -//! and phase details: `docs/approvals.md::Migration from the pre-tag`. +//! Startup convergence. Three phases, all idempotent and unguarded: +//! harness files, applied + proposed repos, meta repo. They re-run every +//! boot on purpose — each one is a no-op once its state is already +//! correct. +//! +//! Deliberately *not* here, and the distinction is the point: +//! +//! - **One-shot, marker-guarded migrations.** Two used to live here +//! (repointing containers onto the meta flake, renaming `root` to +//! `h-root`); both targeted layouts no live hive still has. Add one +//! only if it cannot be expressed as convergence, and expect to delete +//! it once every hive has passed it. +//! - **Create-time setup.** Ruth's tool groups were backfilled here on +//! every boot; they are now seeded where she is created +//! (`workers::auto_update::ensure_root_agent`). A thing that is true +//! from birth does not need re-asserting each morning. +//! +//! Kill-switch: `HIVE_SKIP_META_MIGRATION=1`. Full sequence and phase +//! details: `docs/approvals.md::Migration from the pre-tag`. use std::path::Path; use std::sync::Arc; @@ -14,21 +28,17 @@ use tokio::process::Command; use crate::coordinator::Coordinator; use crate::lifecycle::{self, AGENT_PREFIX, MANAGER_CONTAINER, MANAGER_NAME}; use crate::meta; -use crate::tool_groups; const KILL_SWITCH: &str = "HIVE_SKIP_META_MIGRATION"; -/// Per-shellout timeouts for the blocking startup migration. `run` is +/// Per-shellout timeout for the blocking startup convergence. `run` is /// awaited *before* the daemon starts serving (main.rs), so any child /// process that wedges here freezes the whole daemon — admin socket + -/// dashboard included — with no diagnostics: a git/container shellout was -/// observed blocked for 86min under a concurrent `nixos-rebuild`. Every -/// shellout now runs under a timeout that kills the child on elapse, so a -/// stuck migration degrades to a logged warning instead of a hung boot. -/// Git ops are quick; `nixos-container update` can legitimately trigger a -/// nix build, so it gets a much longer budget. +/// dashboard included — with no diagnostics: a git shellout was observed +/// blocked for 86min under a concurrent `nixos-rebuild`. Every shellout +/// runs under a timeout that kills the child on elapse, so a stuck phase +/// degrades to a logged warning instead of a hung boot. const GIT_TIMEOUT: Duration = Duration::from_mins(2); -const CONTAINER_TIMEOUT: Duration = Duration::from_mins(10); /// Substring that identifies the *current* agent flake boilerplate. /// Bumped whenever the template changes so the startup migration @@ -95,46 +105,6 @@ pub async fn run(coord: &Arc) -> Result<()> { Ok(Ok(())) => {} } - // Phase 4: container repoint, guarded by marker. - if crate::paths::meta_migration_marker().exists() { - tracing::debug!("migration: phase 4 marker present, skipping repoint"); - return Ok(()); - } - tracing::debug!("migration: phase 4 (container repoint)"); - let mut all_ok = true; - for name in &names { - // Mark Rebuilding so the crash watcher skips this container - // during the brief stop+start window the nixos-container - // update activation triggers. Without this, crash_watch - // would fire ContainerCrash for every agent here and the - // manager would spuriously try to recover them. - // No queue node behind this one — migration repoints containers - // directly — so nothing in the graph marks the stop as intended. - let guard = coord.suppress_crash_watch(name.as_str()); - let result = repoint_container(name.as_str()).await; - drop(guard); - if let Err(e) = result { - tracing::warn!(%name, error = ?e, "migration: container repoint failed"); - all_ok = false; - } - } - if all_ok - && !names.is_empty() - && let Err(e) = std::fs::write(crate::paths::meta_migration_marker(), b"done\n") - { - tracing::warn!(error = ?e, "migration: write repoint marker failed"); - } - - // Phase 5: rename `root` nixos-container to `h-root` for naming - // consistency with sub-agents. Guarded by marker; skipped on - // fresh installs (conf file absent) and after first successful run. - rename_manager_container(coord).await; - - // Phase 6: ensure ruth has explicit tool groups so removing the - // role-based fallback (Role::Manager → MANAGER_DEFAULT) doesn't - // silently strip her privileged tools on next rebuild. - backfill_manager_tool_groups(&names); - Ok(()) } @@ -169,106 +139,6 @@ fn migrate_harness_files(name: &hive_types::Ident) { } } -/// Phase 5: rename the `root` nixos-container to `h-root` so the -/// manager container name is consistent with the `h-` prefix used by -/// all sub-agents. Idempotent and marker-guarded. Steps: -/// -/// 1. Check `/etc/nixos-containers/root.conf` exists (old name present). -/// 2. Stop the `root` container. -/// 3. Copy `root.conf` → `h-root.conf`. -/// 4. Move `/var/lib/nixos-containers/root/` → `h-root/` (if present). -/// 5. `systemctl daemon-reload` so systemd sees the new unit name. -/// 6. `nixos-container start h-root`. -/// 7. Write the done marker. -/// -/// Best-effort: logs warnings on failure. A failed rename leaves both -/// conf files present; on the next hive-c0re start the marker is -/// absent so the phase retries. -async fn rename_manager_container(coord: &Arc) { - if crate::paths::hroot_rename_marker().exists() { - return; - } - let old_conf = std::path::PathBuf::from("/etc/nixos-containers/root.conf"); - let new_conf = std::path::PathBuf::from("/etc/nixos-containers/h-root.conf"); - if !old_conf.exists() { - // Fresh install — root container was never created under the old name. - let _ = std::fs::write(crate::paths::hroot_rename_marker(), b"done\n"); - return; - } - if new_conf.exists() { - // Already renamed (but marker was lost — write it and return). - tracing::info!("migration phase 5: h-root.conf already present, marking done"); - let _ = std::fs::write(crate::paths::hroot_rename_marker(), b"done\n"); - return; - } - tracing::info!("migration phase 5: renaming root container to h-root"); - // The old container is stopped immediately below, on purpose. - let _guard = coord.suppress_crash_watch(MANAGER_NAME); - - // Stop the old container. Abort if stop fails — continuing with a - // running `root` and then starting `h-root` risks two manager - // instances racing for the same broker / state files. - match Command::new("nixos-container") - .args(["stop", "root"]) - .status() - .await - { - Ok(s) if s.success() => {} - Ok(s) => { - tracing::warn!(status = %s, "migration phase 5: nixos-container stop root failed — aborting"); - return; - } - Err(e) => { - tracing::warn!(error = ?e, "migration phase 5: nixos-container stop root failed — aborting"); - return; - } - } - - // Copy conf file. - if let Err(e) = std::fs::copy(&old_conf, &new_conf) { - tracing::warn!(error = ?e, "migration phase 5: copy root.conf failed — aborting"); - return; - } - - // Move rootfs if it exists (may be absent for ephemeral containers). - let old_rootfs = std::path::PathBuf::from("/var/lib/nixos-containers/root"); - let new_rootfs = std::path::PathBuf::from("/var/lib/nixos-containers/h-root"); - if old_rootfs.exists() - && !new_rootfs.exists() - && let Err(e) = std::fs::rename(&old_rootfs, &new_rootfs) - { - tracing::warn!(error = ?e, "migration phase 5: rename rootfs failed (non-fatal)"); - } - - // Daemon reload so systemd picks up the new container@h-root unit. - if let Err(e) = Command::new("systemctl") - .args(["daemon-reload"]) - .status() - .await - { - tracing::warn!(error = ?e, "migration phase 5: systemctl daemon-reload failed"); - } - - // Start the renamed container. - if let Err(e) = Command::new("nixos-container") - .args(["start", "h-root"]) - .status() - .await - { - tracing::warn!(error = ?e, "migration phase 5: nixos-container start h-root failed"); - return; - } - - tracing::info!("migration phase 5: root container renamed to h-root"); - let _ = std::fs::write(crate::paths::hroot_rename_marker(), b"done\n"); - // Clean up the old conf file so `nixos-container list` doesn't show - // a stale stopped `root` entry. Best-effort; a failure here is - // harmless — h-root is already running and the marker is written. - if let Err(e) = std::fs::remove_file(&old_conf) { - tracing::warn!(error = ?e, "migration phase 5: remove old root.conf failed (non-fatal)"); - } -} - async fn enumerate_agents() -> Vec { let containers = lifecycle::list().await.unwrap_or_default(); containers @@ -328,57 +198,6 @@ async fn migrate_applied_repo(name: &str) -> Result<()> { Ok(()) } -async fn repoint_container(name: &str) -> Result<()> { - let container = lifecycle::container_name(name); - let flake_ref = format!("{}#{name}", crate::paths::meta_root().display()); - let mut cmd = Command::new("nixos-container"); - cmd.args(["update", &container, "--flake", &flake_ref]); - let out = output_with_timeout( - cmd, - CONTAINER_TIMEOUT, - &format!("nixos-container update {container}"), - ) - .await?; - if !out.status.success() { - anyhow::bail!( - "nixos-container update {container} exited {}: {}", - out.status, - String::from_utf8_lossy(&out.stderr).trim() - ); - } - tracing::info!(%name, %container, "migration: container repointed at meta"); - Ok(()) -} - -/// Phase 6: if ruth is a deployed agent and has no explicit entry in -/// `tool-groups.json`, set her groups to `MANAGER_DEFAULT` (all groups). -/// Idempotent — skips when entry already present. Prevents a silent tool -/// downgrade when upgrading from a build that relied on the manager-flavor -/// fallback in `effective_tool_groups()`. -fn backfill_manager_tool_groups(names: &[hive_types::Ident]) { - if !names.iter().any(|n| n.as_str() == MANAGER_NAME) { - return; // ruth not deployed — nothing to backfill - } - let existing = tool_groups::groups_for(MANAGER_NAME); - if !existing.is_empty() { - tracing::debug!("migration: ruth already has explicit tool groups — skipping backfill"); - return; - } - let all_groups: Vec = hive_sh4re::ToolGroup::MANAGER_DEFAULT - .iter() - .map(|g| g.as_str().to_owned()) - .collect(); - match tool_groups::set_groups(MANAGER_NAME, &all_groups) { - Ok(()) => tracing::info!( - "migration: backfilled ruth's tool groups to MANAGER_DEFAULT (all groups)" - ), - Err(e) => tracing::warn!( - error = ?e, - "migration: failed to backfill ruth's tool groups — she may lose privileged tools on next rebuild" - ), - } -} - /// Run a command to completion under a timeout, capturing its output. On /// timeout the child is killed (`kill_on_drop`) and an error is returned, /// so a wedged shellout can never freeze startup migration. `what` is a diff --git a/hive-c0re/src/paths.rs b/hive-c0re/src/paths.rs index e657ed39..b98c2c64 100644 --- a/hive-c0re/src/paths.rs +++ b/hive-c0re/src/paths.rs @@ -248,18 +248,6 @@ pub fn matrix_register_token() -> PathBuf { state_root().join("matrix-register-token") } -/// `.meta-migration-done` — one-shot marker: legacy meta layout migrated. -#[must_use] -pub fn meta_migration_marker() -> PathBuf { - state_root().join(".meta-migration-done") -} - -/// `.hroot-rename-done` — one-shot marker: legacy hive-root rename applied. -#[must_use] -pub fn hroot_rename_marker() -> PathBuf { - state_root().join(".hroot-rename-done") -} - /// `/run/hyperhive` — the runtime root (host admin socket + per-agent dirs). #[must_use] pub fn runtime_root() -> PathBuf { diff --git a/hive-c0re/src/workers/auto_update.rs b/hive-c0re/src/workers/auto_update.rs index 3fc78e61..599d6796 100644 --- a/hive-c0re/src/workers/auto_update.rs +++ b/hive-c0re/src/workers/auto_update.rs @@ -23,6 +23,7 @@ use anyhow::Result; use crate::coordinator::Coordinator; use crate::lifecycle::{self, AGENT_PREFIX, MANAGER_NAME}; +use crate::tool_groups; /// Resolve the current rev of `hyperhive_flake`. For a path on disk we /// canonicalize (following symlinks) so a /etc/hyperhive → /nix/store/... @@ -146,6 +147,7 @@ pub async fn ensure_root_agent(coord: &Arc) -> Result<()> { let hive = coord.hive_env(); let paths = Coordinator::agent_paths(MANAGER_NAME, runtime); lifecycle::spawn(MANAGER_NAME, &hive, &paths).await?; + seed_manager_tool_groups(); if let Err(e) = coord.power.set(MANAGER_NAME, crate::power::Wanted::Up) { tracing::warn!(error = ?e, "agent_power: set manager wanted=up failed"); } @@ -155,6 +157,35 @@ pub async fn ensure_root_agent(coord: &Arc) -> Result<()> { Ok(()) } +/// Give ruth her privileged tool groups on the one path that creates her. +/// +/// `effective_tool_groups()` has no manager-flavour fallback, so an agent +/// with no entry in `tool-groups.json` is an agent with no privileged +/// tools. Ruth needs hers from her first turn, and this is the only place +/// she is brought into existence — so it is written once, here, rather +/// than re-checked on every hive-c0re boot. +/// +/// Skips a name that already has an entry: a destroy+recreate under the +/// same name must not silently reset an operator's chosen group set back +/// to the default. +fn seed_manager_tool_groups() { + if !tool_groups::groups_for(MANAGER_NAME).is_empty() { + tracing::debug!("manager tool groups already set — leaving as-is"); + return; + } + let all_groups: Vec = hive_sh4re::ToolGroup::MANAGER_DEFAULT + .iter() + .map(|g| g.as_str().to_owned()) + .collect(); + match tool_groups::set_groups(MANAGER_NAME, &all_groups) { + Ok(()) => tracing::info!("seeded ruth's tool groups to MANAGER_DEFAULT (all groups)"), + Err(e) => tracing::warn!( + error = ?e, + "failed to seed ruth's tool groups — she will start without privileged tools" + ), + } +} + /// Sort `names` in-place so parents precede their children in the topology. /// Uses BFS from root agents (depth 0). Agents absent from `topo` sort last, /// alphabetically within their tier. Stable within each depth tier.