diff --git a/docs/approvals.md b/docs/approvals.md index e5305769..527f402b 100644 --- a/docs/approvals.md +++ b/docs/approvals.md @@ -478,16 +478,12 @@ 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, 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. + 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. 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 871a45f8..7c4ff619 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/`. -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. +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). ## Destroy vs purge diff --git a/docs/web-ui/dashboard.md b/docs/web-ui/dashboard.md index a8cfaff2..7556155b 100644 --- a/docs/web-ui/dashboard.md +++ b/docs/web-ui/dashboard.md @@ -852,29 +852,12 @@ 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 **transient** - (`transientsState`) — covers the create-and-start window where the + 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 container literally isn't up yet, before any backend state event - 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 + has fired. (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 @@ -1368,9 +1351,7 @@ 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. `transient_kind` is an **open** - display string (the running node's own tag), not a fixed - enum — render it, don't branch on it. + client-side, no polling. - `container_state_changed` (container: ContainerView) / `container_removed` (name) — per-row container mutations, emitted by `Coordinator::rescan_containers_and_emit` from diff --git a/frontend/packages/agent/src/agent.css b/frontend/packages/agent/src/agent.css index 79f32625..2b9d7007 100644 --- a/frontend/packages/agent/src/agent.css +++ b/frontend/packages/agent/src/agent.css @@ -530,12 +530,6 @@ 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); @@ -661,11 +655,10 @@ pre.diff { text-decoration-color: var(--muted); } -/* 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. */ +/* "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. */ .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 b53ecaeb..729e4b06 100644 --- a/frontend/packages/agent/src/app.js +++ b/frontend/packages/agent/src/app.js @@ -861,67 +861,8 @@ 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 - * (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. */ + * (subsystem, summary, source, age_seconds). */ function buildTodosList(todos) { const wrap = el('div', { class: 'agent-inbox' }); if (!todos.length) { @@ -929,7 +870,6 @@ 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'; @@ -940,13 +880,8 @@ 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( - cb, ' ', - el('label', { for: cbId, class: 'inbox-from' }, label), ' ', + el('span', { 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 52cb4e75..de060f92 100644 --- a/hive-agent/src/web_ui/actions.rs +++ b/hive-agent/src/web_ui/actions.rs @@ -1,5 +1,4 @@ -//! Operator action POST handlers (send, cancel, compact, model, effort, -//! reset, todos mark-done). +//! Operator action POST handlers (send, cancel, compact, model, effort, reset). use axum::{ Form, @@ -163,41 +162,3 @@ 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 4693c361..5b06e15d 100644 --- a/hive-agent/src/web_ui/mod.rs +++ b/hive-agent/src/web_ui/mod.rs @@ -116,7 +116,6 @@ 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)); diff --git a/hive-c0re/src/actions.rs b/hive-c0re/src/actions.rs index d3eb90c4..4dadeace 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; +use crate::coordinator::{Coordinator, TransientKind}; use crate::lifecycle; /// Approve a pending request. Marks the approval row durably, then @@ -882,10 +882,7 @@ 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 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); + let guard = coord.transient_guard(name, TransientKind::Destroying); 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 4b9e1366..7029cbf3 100644 --- a/hive-c0re/src/coordinator.rs +++ b/hive-c0re/src/coordinator.rs @@ -8,7 +8,6 @@ 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; @@ -119,10 +118,7 @@ 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). - /// 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>, + 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 @@ -137,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 @@ -361,75 +357,27 @@ 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 { - /// 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, - /// 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, + pub kind: TransientKind, + pub since: std::time::Instant, } -/// 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 { +/// 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 { coord: Arc, name: String, } -impl Drop for CrashWatchSuppression { +impl Drop for TransientGuard { fn drop(&mut self) { - 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())); + self.coord.clear_transient(&self.name); } } @@ -460,6 +408,38 @@ 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 @@ -568,7 +548,7 @@ impl Coordinator { agent_io_weight, model_prices, agents: Mutex::new(HashMap::new()), - crash_suppressed: Mutex::new(HashSet::new()), + transient: Mutex::new(HashMap::new()), recent_transient: Mutex::new(HashMap::new()), recent_crashes: Mutex::new(HashMap::new()), graceful_stop_pending: Mutex::new(HashSet::new()), @@ -1106,9 +1086,22 @@ impl Coordinator { self.agents.lock().unwrap().keys().cloned().collect() } - /// 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) { + /// 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, kind: TransientKind) { + self.transient.lock().unwrap().insert( + name.to_owned(), + TransientState { + kind, + since: std::time::Instant::now(), + }, + ); // 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. @@ -1120,31 +1113,32 @@ impl Coordinator { self.emit_dashboard_event(DashboardEvent::TransientSet { seq: self.next_seq(), name: name.to_owned(), - transient_kind: label, + transient_kind: kind.as_str(), since_unix, }); } - /// Emit the "a pill went away" edge **and stamp the tombstone the crash - /// 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: - /// 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(), - }); + /// 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) { + let removed = self.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.kind, 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, @@ -1211,19 +1205,20 @@ impl Coordinator { result } - /// 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 { + /// 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 { 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, (deliberate, _))| (k.clone(), *deliberate)) + .map(|(k, (kind, _))| (k.clone(), *kind)) .collect() } @@ -1254,62 +1249,21 @@ impl Coordinator { map.iter().map(|(k, v)| (k.clone(), v.len())).collect() } - /// 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. - /// - /// 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() - .unwrap() - .insert(name.to_owned()); - CrashWatchSuppression { + /// 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. + pub fn transient_guard(self: &Arc, name: &str, kind: TransientKind) -> TransientGuard { + self.set_transient(name, kind); + TransientGuard { 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, 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 { - 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() + self.transient.lock().unwrap().clone() } /// 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 04c43935..f86000cb 100644 --- a/hive-c0re/src/dashboard/state_snapshot.rs +++ b/hive-c0re/src/dashboard/state_snapshot.rs @@ -214,8 +214,7 @@ struct PortConflict { #[derive(Serialize)] struct TransientView { name: String, - /// Owned: the label is the running node's wire tag, not one of a fixed set. - kind: String, + kind: &'static str, secs: u64, } @@ -534,18 +533,26 @@ fn build_transient_views( .filter(|(name, _)| !containers.iter().any(|c| &c.name == *name)) .map(|(name, st)| TransientView { name: name.clone(), - kind: st.label.clone(), - // 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(), + kind: transient_label(st.kind), + 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 8c6c8db0..b5451b05 100644 --- a/hive-c0re/src/dashboard_events.rs +++ b/hive-c0re/src/dashboard_events.rs @@ -144,15 +144,9 @@ pub enum DashboardEvent { TransientSet { seq: u64, name: String, - /// 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, + /// Lifecycle kind: `"spawning"` / `"starting"` / `"stopping"` / + /// `"restarting"` / `"rebuilding"` / `"destroying"`. + transient_kind: &'static str, since_unix: i64, }, /// The matching lifecycle action resolved (success or failure). @@ -386,7 +380,7 @@ mod tests { DashboardEvent::TransientSet { seq: 1, name: "x".into(), - transient_kind: "rebuilding".into(), + transient_kind: "rebuilding", since_unix: 0, }, DashboardEvent::TransientCleared { diff --git a/hive-c0re/src/job_queue/exec.rs b/hive-c0re/src/job_queue/exec.rs index 7dcdf4dd..0b73ac76 100644 --- a/hive-c0re/src/job_queue/exec.rs +++ b/hive-c0re/src/job_queue/exec.rs @@ -418,11 +418,13 @@ async fn run_reconcile(coord: &Arc, claim: &Claim) -> Result, claim: &Claim) -> Result { let name = &claim.agent; - // 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. + // 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)); // 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 — @@ -447,8 +449,10 @@ async fn run_start(coord: &Arc, claim: &Claim) -> Result, claim: &Claim) -> Result { let name = &claim.agent; - // See `run_start`: no node-local guard — `Stop` reports `Stopping` from its - // own kind now. + let _guard = claim + .transient + .is_none() + .then(|| coord.transient_guard(name, crate::coordinator::TransientKind::Stopping)); 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 ada7faad..9300b526 100644 --- a/hive-c0re/src/job_queue/mod.rs +++ b/hive-c0re/src/job_queue/mod.rs @@ -47,6 +47,7 @@ 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; @@ -59,27 +60,6 @@ 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)] @@ -90,6 +70,10 @@ 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 @@ -107,6 +91,7 @@ struct NodeRuntime { struct DagMeta { source: Source, reason: String, + transient: Option, created_at: i64, } @@ -230,6 +215,7 @@ impl JobQueue { NodeKind::Dag { source: spec.source, reason: spec.reason, + transient: spec.transient, created_at: now_unix(), }, Vec::new(), @@ -307,11 +293,15 @@ 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. @@ -420,59 +410,24 @@ impl JobQueue { .map(ToOwned::to_owned) } - /// `(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` / - /// `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. - /// - /// `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. - /// - /// 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. + /// 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. #[must_use] - pub fn running_transients(&self) -> Vec { + pub fn held_transients(&self) -> Vec<(u64, String, TransientKind)> { let inner = self.lock(); inner .sched - .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(RunningTransient { - agent, - 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())), - }) + .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)) }) .collect() } @@ -526,6 +481,7 @@ impl QueueInner { let NodeKind::Dag { source, reason, + transient, created_at, } = &self.sched.graph().node(container)?.payload else { @@ -534,6 +490,7 @@ 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 96fba3c3..fa208985 100644 --- a/hive-c0re/src/job_queue/model.rs +++ b/hive-c0re/src/job_queue/model.rs @@ -15,6 +15,8 @@ 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 @@ -283,6 +285,7 @@ pub enum NodeKind { Dag { source: Source, reason: String, + transient: Option, created_at: i64, }, } @@ -390,44 +393,6 @@ 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 @@ -498,5 +463,8 @@ 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 105d5326..d235653b 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-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::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. +//! 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. //! //! 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; +use std::collections::{HashMap, HashSet}; use std::sync::Arc; use super::Claim; @@ -54,12 +54,8 @@ struct NodeDone { pub async fn run_worker(coord: Arc) { let mut shutdown = coord.shutdown_rx(); let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel::(); - // 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(); + // (DAG id, agent) → transient guard held for that agent's lease window. + let mut transients: HashMap<(u64, String), crate::coordinator::TransientGuard> = HashMap::new(); loop { // Checked every iteration, not just in the `select!` below — a // continuous stream of ready claims never reaches the `select!`, so @@ -150,53 +146,18 @@ fn handle_completion(coord: &Arc, done: NodeDone) { coord.emit_rebuild_queue_snapshot(); } -/// 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 -/// `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 -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. `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.emit_transient_cleared(agent, *deliberate); - } - still - }); - - for t in running { - if prev.get(&t.agent).map(|(l, _)| l) == Some(&t.label) { - continue; - } - coord.emit_transient_set(&t.agent, t.label.clone()); - prev.insert(t.agent, (t.label, t.takes_container_down)); +/// 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)`. +fn reconcile_transients( + coord: &Arc, + transients: &mut HashMap<(u64, String), crate::coordinator::TransientGuard>, +) { + 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)); } } diff --git a/hive-c0re/src/job_queue/submit.rs b/hive-c0re/src/job_queue/submit.rs index 43955ec9..9ced4137 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; +use crate::coordinator::{Coordinator, TransientKind}; use crate::lifecycle; fn submit_and_emit(coord: &Arc, spec: super::DagSpec) -> u64 { @@ -198,10 +198,16 @@ 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(source: Source, reason: String, nodes: Vec) -> DagSpec { +fn power_dag( + transient: TransientKind, + source: Source, + reason: String, + nodes: Vec, +) -> DagSpec { DagSpec { source, reason, + transient: Some(transient), nodes, } } @@ -222,15 +228,18 @@ pub(crate) fn stop_spec( .iter() .map(|(agent, running)| stop_chain(agent, graceful, *running)) .collect(); - power_dag(source, reason, concat_subgraphs(chains)) + power_dag( + TransientKind::Stopping, + source, + reason, + concat_subgraphs(chains), + ) } /// Assemble the start DAG from explicit `(agent, running, stale)` targets. -/// -/// 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. +/// 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. pub(crate) fn start_spec( targets: &[(String, bool, bool)], source: Source, @@ -240,7 +249,13 @@ pub(crate) fn start_spec( .iter() .map(|(agent, running, stale)| start_chain(agent, *running, *stale)) .collect(); - power_dag(source, reason, concat_subgraphs(chains)) + 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)) } /// Assemble the restart DAG from explicit `(agent, running)` targets. @@ -254,7 +269,12 @@ pub(crate) fn restart_spec( .iter() .map(|(agent, running)| restart_chain(agent, graceful, *running)) .collect(); - power_dag(source, reason, concat_subgraphs(chains)) + power_dag( + TransientKind::Restarting, + 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 70c6d6a9..19abd54b 100644 --- a/hive-c0re/src/job_queue/templates.rs +++ b/hive-c0re/src/job_queue/templates.rs @@ -34,6 +34,7 @@ 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 @@ -366,6 +367,7 @@ pub fn rebuild(agent: &str, source: Source, reason: String, relock: bool) -> Dag DagSpec { source, reason, + transient: Some(TransientKind::Rebuilding), nodes, } } @@ -400,6 +402,7 @@ 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 { @@ -448,10 +451,16 @@ 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) -> DagSpec { +pub fn reconcile_only( + agent: &str, + source: Source, + reason: String, + transient: Option, +) -> DagSpec { DagSpec { source, reason, + transient, nodes: vec![node( NodeKind::Reconcile { agent: agent.to_owned(), @@ -476,6 +485,7 @@ 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![ @@ -519,6 +529,7 @@ pub fn perm_change(agent: &str, source: Source, reason: String, payload: PermPay DagSpec { source, reason, + transient: Some(TransientKind::Rebuilding), nodes, } } @@ -557,6 +568,7 @@ pub fn meta_update( DagSpec { source, reason, + transient: Some(TransientKind::Rebuilding), nodes, } } @@ -578,6 +590,7 @@ 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 49d55bd4..5fd97c49 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()), + templates::reconcile_only("agent-a", Source::Manual, "stop".to_owned(), None), ); // 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()), + templates::reconcile_only("agent-a", Source::Manual, "stop".to_owned(), None), ); // 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,55 +789,26 @@ 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 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() { +fn meta_update_carries_rebuilding_transient_and_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); @@ -1010,7 +981,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()), + templates::reconcile_only("agent-a", Source::Manual, "start".to_owned(), None), ); let c = claim_one(&q); q.complete_node(c.node_id, Err("start failed".to_owned())); @@ -1161,7 +1132,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()), + templates::reconcile_only("agent-a", Source::Manual, "stop".to_owned(), None), ); let c = claim_one(&q); assert_eq!(c.dag_id, next); @@ -1430,7 +1401,12 @@ 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()), + templates::reconcile_only( + &format!("agent-{i}"), + Source::Manual, + "start".to_owned(), + None, + ), ); 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 f80a30fa..355bba5d 100644 --- a/hive-c0re/src/migrate.rs +++ b/hive-c0re/src/migrate.rs @@ -1,22 +1,8 @@ -//! 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`. +//! 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`. use std::path::Path; use std::sync::Arc; @@ -28,17 +14,21 @@ 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 timeout for the blocking startup convergence. `run` is +/// Per-shellout timeouts for the blocking startup migration. `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 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. +/// 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. 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 @@ -105,6 +95,45 @@ 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. + let guard = + coord.transient_guard(name.as_str(), crate::coordinator::TransientKind::Rebuilding); + 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(()) } @@ -139,6 +168,105 @@ 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"); + let _guard = coord.transient_guard(MANAGER_NAME, crate::coordinator::TransientKind::Rebuilding); + + // 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 @@ -198,6 +326,57 @@ 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 b98c2c64..e657ed39 100644 --- a/hive-c0re/src/paths.rs +++ b/hive-c0re/src/paths.rs @@ -248,6 +248,18 @@ 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 599d6796..36e42477 100644 --- a/hive-c0re/src/workers/auto_update.rs +++ b/hive-c0re/src/workers/auto_update.rs @@ -23,7 +23,6 @@ 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/... @@ -147,7 +146,6 @@ 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"); } @@ -157,35 +155,6 @@ 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. @@ -384,6 +353,7 @@ 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 73999f1c..49111296 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; +use crate::coordinator::{Coordinator, TransientKind}; use crate::lifecycle::{self, AGENT_PREFIX}; const POLL_INTERVAL: Duration = Duration::from_secs(10); @@ -87,13 +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) { - // 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 active = transients.get(stopped).map(|st| st.kind); let recently_cleared = recent.get(stopped).copied(); if is_deliberate_stop(active, recently_cleared) { continue; @@ -107,19 +101,26 @@ fn emit_crash_transitions(coord: &Coordinator, prev: &HashSet, current: } } -/// 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) +/// 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) } fn emit_login_transitions( @@ -167,12 +168,29 @@ mod tests { use super::*; #[test] - 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))); + 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:?}"); + } } #[test] @@ -182,16 +200,13 @@ mod tests { } #[test] - 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))); + 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"); + } } }