Compare commits

..
22 changed files with 553 additions and 661 deletions

View file

@ -478,16 +478,12 @@ each phase is a no-op once already applied. Behaviour:
dirs, or claude creds.
- Meta-flake phase: rewrites each `applied/<n>/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#<n>`, guarded by a marker file so it
ran once per hive. It is gone: containers have been rendered
onto `meta#<n>` 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#<n>`. 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

View file

@ -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#<n>` 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#<n>`. Removing it forces a re-run on
next hive-c0re start (idempotent — only the actual repoint
step would re-fire).
## Destroy vs purge

View file

@ -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

View file

@ -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;

View file

@ -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 || ''),
);

View file

@ -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<MarkTodosDoneForm>) -> Response {
let ids: Vec<i64> = form
.ids
.split(',')
.filter_map(|s| s.trim().parse::<i64>().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()
}

View file

@ -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));

View file

@ -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<Coordinator>, 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);

View file

@ -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<HashSet<String>>,
transient: Mutex<HashMap<String, TransientState>>,
/// 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<HashMap<String, (bool, std::time::Instant)>>,
recent_transient: Mutex<HashMap<String, (TransientKind, std::time::Instant)>>,
/// 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<Utc>,
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<Coordinator>` 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<Coordinator>,
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<String, bool> {
/// 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<String, TransientKind> {
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<Self>, 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<Self>, 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<String, TransientState> {
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

View file

@ -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

View file

@ -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 {

View file

@ -418,11 +418,13 @@ async fn run_reconcile(coord: &Arc<Coordinator>, claim: &Claim) -> Result<NodeOu
/// out when it observes `wanted = Up` and the container down.
async fn run_start(coord: &Arc<Coordinator>, claim: &Claim) -> Result<NodeOutput> {
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<Coordinator>, claim: &Claim) -> Result<NodeOutput
/// out when it observes `wanted = Offline` and the container up.
async fn run_stop(coord: &Arc<Coordinator>, claim: &Claim) -> Result<NodeOutput> {
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 {

View file

@ -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<Utc>,
}
/// 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<TransientKind>,
}
/// 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<TransientKind>,
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<RunningTransient> {
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,
})
}

View file

@ -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<TransientKind>,
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<crate::coordinator::TransientKind>,
pub nodes: Vec<NodeSpec>,
}

View file

@ -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<Coordinator>) {
let mut shutdown = coord.shutdown_rx();
let (tx, mut rx) = tokio::sync::mpsc::unbounded_channel::<NodeDone>();
// 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<String, (String, bool)> = 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<Coordinator>, 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<Coordinator>, prev: &mut HashMap<String, (String, bool)>) {
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<Coordinator>,
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));
}
}

View file

@ -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<Coordinator>, spec: super::DagSpec) -> u64 {
@ -198,10 +198,16 @@ fn concat_subgraphs(chains: Vec<Vec<NodeSpec>>) -> Vec<NodeSpec> {
/// 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<NodeSpec>) -> DagSpec {
fn power_dag(
transient: TransientKind,
source: Source,
reason: String,
nodes: Vec<NodeSpec>,
) -> 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`].

View file

@ -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<TransientKind>,
) -> 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())],
}
}

View file

@ -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

View file

@ -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<Coordinator>) -> 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<Coordinator>) {
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<hive_types::Ident> {
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<String> = 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

View file

@ -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 {

View file

@ -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<Coordinator>) -> 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<Coordinator>) -> 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<String> = 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) {

View file

@ -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<String>, 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<String>, 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<bool>, recently_cleared: Option<bool>) -> 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<TransientKind>,
recently_cleared: Option<TransientKind>,
) -> 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");
}
}
}