diff --git a/hive-c0re/assets/app.js b/hive-c0re/assets/app.js index 7d14054b..53012bde 100644 --- a/hive-c0re/assets/app.js +++ b/hive-c0re/assets/app.js @@ -493,6 +493,22 @@ renderMetaInputs({ meta_inputs: metaInputsState }); } + // Derived rebuild queue state — cold-loaded from + // `/api/state.rebuild_queue`, then mutated live by the + // `rebuild_queue_changed` snapshot event. Same shape as the meta- + // inputs panel (full snapshot per change, no diff). + let rebuildQueueState = []; + function syncRebuildQueueFromSnapshot(s) { + rebuildQueueState = (s.rebuild_queue || []).slice(); + } + function applyRebuildQueueChanged(ev) { + rebuildQueueState = (ev.queue || []).slice(); + renderRebuildQueueFromState(); + } + function renderRebuildQueueFromState() { + renderRebuildQueue({ rebuild_queue: rebuildQueueState }); + } + // Derived transient state — cold-loaded from /api/state.transients, // then mutated live by `transient_set` / `transient_cleared`. Keyed // by agent name so add/remove are O(1). `since_unix` is wall-clock so @@ -1653,6 +1669,119 @@ return s.length <= n ? s : s.slice(0, n - 1) + '…'; } + // ─── rebuild queue ────────────────────────────────────────────────────── + // Glyph + verb per QueueKind. Mirrors the labels used in + // hive-c0re::rebuild_queue::QueueKind::as_str. + const QUEUE_KIND_GLYPH = { + rebuild: '↻', + meta_update: '◆', + spawn: '✨', + destroy: '🗑', + }; + const QUEUE_STATE_GLYPH = { + queued: '⏸', + running: '▶', + done: '✔', + failed: '✖', + cancelled: '⊘', + }; + + function renderRebuildQueue(s) { + const root = $('rebuild-queue-section'); + if (!root) return; + root.innerHTML = ''; + const queue = s.rebuild_queue || []; + if (!queue.length) { + root.append(el('p', { class: 'empty' }, 'queue is empty — nothing pending or in flight.')); + return; + } + // Index by id for parent lookup. + const byId = new Map(queue.map((e) => [e.id, e])); + // Top-level entries first; children render nested under their parent. + const tops = queue.filter((e) => e.parent_id == null); + const childrenOf = new Map(); + for (const e of queue) { + if (e.parent_id != null) { + if (!childrenOf.has(e.parent_id)) childrenOf.set(e.parent_id, []); + childrenOf.get(e.parent_id).push(e); + } + } + const ul = el('ul', { class: 'rebuild-queue' }); + for (const top of tops) { + ul.append(renderQueueEntry(top, byId)); + for (const child of childrenOf.get(top.id) || []) { + ul.append(renderQueueEntry(child, byId, true)); + } + } + // Children whose parent isn't in the snapshot (history-evicted) still render flat. + const orphans = queue.filter( + (e) => e.parent_id != null && !byId.has(e.parent_id), + ); + for (const o of orphans) { + ul.append(renderQueueEntry(o, byId, true)); + } + root.append(ul); + } + + function renderQueueEntry(entry, _byId, isChild) { + const li = el('li', { + class: 'rebuild-queue-entry rqe-' + entry.state, + 'data-id': String(entry.id), + }); + if (isChild) li.classList.add('rqe-child'); + // State glyph + kind + agent. + li.append( + el('span', { class: 'rqe-state', title: entry.state }, QUEUE_STATE_GLYPH[entry.state] || '?'), + ' ', + el('span', { class: 'rqe-kind', title: entry.kind }, + (QUEUE_KIND_GLYPH[entry.kind] || '?') + ' ' + entry.kind), + ' ', + el('code', { class: 'rqe-agent' }, entry.agent), + ); + // Source chip (manual / meta_update / auto_update / crash_recover). + li.append(' ', el('span', { class: 'rqe-source rqe-source-' + entry.source }, entry.source)); + // Timing: queued Xs ago when pending, elapsed when running, + // finished Xs ago for terminal. + if (entry.state === 'queued') { + li.append(' ', el('span', { class: 'rqe-when' }, '· queued ' + fmtAgo(entry.enqueued_at))); + } else if (entry.state === 'running' && entry.started_at) { + const elapsed = Math.max(0, Math.floor(Date.now() / 1000 - entry.started_at)); + li.append(' ', el('span', { + class: 'rqe-when', + 'data-rqe-elapsed': String(entry.started_at), + }, '· ' + fmtElapsed(elapsed))); + } else if (entry.finished_at) { + li.append(' ', el('span', { class: 'rqe-when' }, '· ' + entry.state + ' ' + fmtAgo(entry.finished_at))); + } + // Reason (truncated; full text on hover). + if (entry.reason) { + const r = entry.reason.split('\n')[0]; + li.append(' ', el('span', { class: 'rqe-reason', title: entry.reason }, '— ' + truncate(r, 60))); + } + // Error block, when failed. + if (entry.error) { + li.append(el('pre', { class: 'rqe-error', title: entry.error }, truncate(entry.error, 200))); + } + return li; + } + + function fmtElapsed(secs) { + if (secs < 60) return secs + 's running'; + if (secs < 3600) return Math.floor(secs / 60) + 'm ' + (secs % 60) + 's running'; + return Math.floor(secs / 3600) + 'h ' + Math.floor((secs % 3600) / 60) + 'm running'; + } + + // Tick once per second to refresh "running Xs" badges in place + // (mirrors the question-TTL ticker pattern from #335). + setInterval(() => { + for (const span of document.querySelectorAll('.rqe-when[data-rqe-elapsed]')) { + const started = parseInt(span.dataset.rqeElapsed, 10); + if (!started) continue; + const elapsed = Math.max(0, Math.floor(Date.now() / 1000 - started)); + span.textContent = '· ' + fmtElapsed(elapsed); + } + }, 1000); + // ─── reminders ────────────────────────────────────────────────────────── // Reminders aren't part of /api/state (separate sqlite table, separate // mutation cadence). Refresh fires alongside refreshState() so a @@ -1764,6 +1893,7 @@ 'inbox-section', 'approvals-section', 'meta-inputs-section', + 'rebuild-queue-section', 'reminders-section', ]; //
sections that should survive a refresh need a stable @@ -1830,6 +1960,7 @@ syncContainersFromSnapshot(s); syncTombstonesFromSnapshot(s); syncMetaInputsFromSnapshot(s); + syncRebuildQueueFromSnapshot(s); renderContainers(s); renderTombstones(s); // Sync the derived approvals + questions stores from the @@ -1842,6 +1973,7 @@ syncApprovalsFromSnapshot(s); renderApprovals(); renderMetaInputs(s); + renderRebuildQueue(s); refreshReminders(); restoreOpenDetails(openDetails); notifyDeltas(s); @@ -1963,6 +2095,7 @@ tombstones_changed: (ev) => { applyTombstonesChanged(ev); }, meta_inputs_changed: (ev) => { applyMetaInputsChanged(ev); }, meta_update_running: (ev) => { applyMetaUpdateRunning(ev); }, + rebuild_queue_changed: (ev) => { applyRebuildQueueChanged(ev); }, }, // Both history backfill and live frames flow through here, so the // inbox section ends up populated correctly on first paint and diff --git a/hive-c0re/assets/dashboard.css b/hive-c0re/assets/dashboard.css index 248f5234..5f0e8c20 100644 --- a/hive-c0re/assets/dashboard.css +++ b/hive-c0re/assets/dashboard.css @@ -521,6 +521,61 @@ code { font-size: 0.85em; animation: badge-pulse 1.6s ease-in-out infinite; } +/* ─── rebuild queue panel ──────────────────────────────────────────────── */ +.rebuild-queue { + list-style: none; + padding: 0; + margin: 0; + display: grid; + gap: 0.2em; +} +.rebuild-queue-entry { + padding: 0.3em 0.6em; + border: 1px solid var(--border); + background: rgba(24, 24, 37, 0.6); + font-size: 0.9em; + display: flex; + flex-wrap: wrap; + align-items: baseline; + gap: 0.4em; +} +.rebuild-queue-entry.rqe-child { margin-left: 1.6em; border-color: var(--purple-dim); } +.rebuild-queue-entry.rqe-running { + border-color: var(--purple); + background: rgba(203, 166, 247, 0.12); + animation: badge-pulse 1.6s ease-in-out infinite; +} +.rebuild-queue-entry.rqe-failed { border-color: var(--red); color: var(--red); } +.rebuild-queue-entry.rqe-cancelled { opacity: 0.6; } +.rebuild-queue-entry.rqe-done { opacity: 0.7; color: var(--green); } +.rqe-state { font-weight: bold; min-width: 1.2em; text-align: center; } +.rqe-kind { color: var(--cyan); } +.rqe-agent { color: var(--amber); font-weight: bold; } +.rqe-source { + font-size: 0.75em; + padding: 0.05em 0.45em; + border-radius: 0.7em; + border: 1px solid var(--border); + color: var(--muted); + text-transform: uppercase; + letter-spacing: 0.05em; +} +.rqe-source-manual { color: var(--cyan); border-color: var(--cyan); } +.rqe-source-meta_update { color: var(--purple); border-color: var(--purple); } +.rqe-source-auto_update { color: var(--muted); } +.rqe-source-crash_recover { color: var(--amber); border-color: var(--amber); } +.rqe-when { color: var(--muted); font-size: 0.85em; } +.rqe-reason { color: var(--muted); font-size: 0.85em; flex: 1 1 auto; } +.rqe-error { + flex-basis: 100%; + margin: 0.3em 0 0; + padding: 0.3em 0.5em; + background: rgba(243, 139, 168, 0.1); + border-left: 2px solid var(--red); + color: var(--red); + font-size: 0.8em; + white-space: pre-wrap; +} .history-note { margin-left: 1.8em; margin-top: 0.2em; diff --git a/hive-c0re/assets/index.html b/hive-c0re/assets/index.html index 25fdbd51..846c86e9 100644 --- a/hive-c0re/assets/index.html +++ b/hive-c0re/assets/index.html @@ -39,6 +39,13 @@

loading…

+

◆ R3BU1LD QU3U3 ◆

+
══════════════════════════════════════════════════════════════
+

pending + running rebuilds, meta-updates, and first-spawns. one runs at a time; meta-update cascades nest under their parent. dedup: re-enqueueing a still-queued op collapses into the existing entry.

+
+

loading…

+
+

◆ M1ND H4S QU3STI0NS ◆

══════════════════════════════════════════════════════════════
diff --git a/hive-c0re/src/auto_update.rs b/hive-c0re/src/auto_update.rs index a21b5b30..1ac2add7 100644 --- a/hive-c0re/src/auto_update.rs +++ b/hive-c0re/src/auto_update.rs @@ -206,10 +206,9 @@ pub async fn run(coord: Arc) -> Result<()> { } }; - let current_rev = - current_flake_rev(&coord.hyperhive_flake).unwrap_or_default(); + let _current_rev = current_flake_rev(&coord.hyperhive_flake).unwrap_or_default(); - tracing::info!(agents = containers.len(), "auto-update: rebuilding all on startup"); + tracing::info!(agents = containers.len(), "auto-update: queueing all on startup"); for container in containers { let logical = if container == MANAGER_NAME { Some(MANAGER_NAME.to_owned()) @@ -217,9 +216,14 @@ pub async fn run(coord: Arc) -> Result<()> { container.strip_prefix(AGENT_PREFIX).map(str::to_owned) }; let Some(name) = logical else { continue }; - if let Err(e) = rebuild_agent(&coord, &name, ¤t_rev).await { - tracing::warn!(%name, error = ?e, "auto-update: rebuild failed"); - } + coord.rebuild_queue.enqueue( + crate::rebuild_queue::QueueKind::Rebuild, + name, + crate::rebuild_queue::QueueSource::AutoUpdate, + "startup sweep".to_owned(), + None, + ); } + coord.emit_rebuild_queue_snapshot(); Ok(()) } diff --git a/hive-c0re/src/coordinator.rs b/hive-c0re/src/coordinator.rs index d4d41ff0..357e9b44 100644 --- a/hive-c0re/src/coordinator.rs +++ b/hive-c0re/src/coordinator.rs @@ -88,6 +88,12 @@ pub struct Coordinator { /// tokio mutex so the rescan can `await` `lifecycle::list` / /// `is_running` without blocking other coordinator paths. last_containers: tokio::sync::Mutex>, + /// Global rebuild queue. Every long-running container/meta op + /// (rebuild, meta-update, first-spawn) goes through this queue so + /// hive-c0re runs at most one at a time and the dashboard can + /// render a single ordered view of pending + running work. See + /// `rebuild_queue.rs` for the dedup rules + history retention. + pub rebuild_queue: Arc, /// Shutdown signal broadcast to all background tasks. Sending /// `true` asks every loop to exit after its current work item. /// Use `shutdown_rx()` to subscribe; `request_shutdown()` to fire. @@ -202,10 +208,23 @@ impl Coordinator { event_seq: AtomicU64::new(0), meta_updates_active: AtomicU64::new(0), last_containers: tokio::sync::Mutex::new(HashMap::new()), + rebuild_queue: Arc::new(crate::rebuild_queue::RebuildQueue::new()), shutdown_tx, }) } + /// Emit a `RebuildQueueChanged` snapshot event. Called from the + /// queue mutation helpers (`enqueue` / `finish` / `cancel`-adjacent + /// wrappers below) and the worker so every state transition + /// surfaces on the dashboard without extra plumbing. + pub fn emit_rebuild_queue_snapshot(self: &Arc) { + let queue = self.rebuild_queue.snapshot(); + self.emit_dashboard_event(DashboardEvent::RebuildQueueChanged { + seq: self.next_seq(), + queue, + }); + } + /// Subscribe to the shutdown watch channel. Background tasks call /// this at spawn time and break their loop when the receiver /// transitions to `true` (via `Coordinator::request_shutdown`). diff --git a/hive-c0re/src/dashboard.rs b/hive-c0re/src/dashboard.rs index 1cc418f7..529da6cf 100644 --- a/hive-c0re/src/dashboard.rs +++ b/hive-c0re/src/dashboard.rs @@ -237,6 +237,12 @@ struct StateSnapshot { /// disabled "updating…" state; live transitions arrive via the /// `MetaUpdateRunning` event (issue #259). meta_update_running: bool, + /// Current state of the global rebuild queue — pending + running + /// long-lived ops (rebuild / meta-update / spawn) plus the most + /// recent few terminal entries the queue retains for history. + /// Live transitions arrive via the `RebuildQueueChanged` event. + /// See `rebuild_queue.rs`. + rebuild_queue: Vec, /// Whether the hive-forge container is up. When true the dashboard /// links each container's config + each approval's commit into the /// forge's `agent-configs` repos. @@ -431,6 +437,7 @@ async fn api_state(headers: HeaderMap, State(state): State) -> axum::J question_history, tombstones, port_conflicts, + rebuild_queue: state.coord.rebuild_queue.snapshot(), forge_present: crate::forge::is_present().await, }) } @@ -1571,82 +1578,18 @@ async fn post_meta_update( if inputs.is_empty() { return error_response("meta-update: no inputs selected"); } - let coord = state.coord.clone(); - let inputs_clone = inputs.clone(); - tokio::spawn(async move { - run_meta_update(&coord, &inputs_clone).await; - // Lock file changed — emit so dashboards refresh the - // meta-inputs panel without a snapshot poll. - emit_meta_inputs_snapshot(&coord); - }); + state.coord.rebuild_queue.enqueue_with_inputs( + crate::rebuild_queue::QueueKind::MetaUpdate, + "hyperhive".to_owned(), + crate::rebuild_queue::QueueSource::Manual, + format!("meta-update via dashboard ({})", inputs.join(", ")), + None, + inputs, + ); + state.coord.emit_rebuild_queue_snapshot(); (StatusCode::OK, "ok").into_response() } -/// Background task: run `nix flake update ` in meta + commit, -/// then rebuild every agent whose input was touched (or all agents -/// when `hyperhive` was bumped, since that's the shared base). Each -/// rebuild fires `Rebuilt { ok, note, ... }` to the manager so the -/// operator and manager get the same feedback they'd see from an -/// auto-update / manual dashboard rebuild. -async fn run_meta_update(coord: &Arc, inputs: &[String]) { - // Held for the whole run (incl. the early `return` on lock failure): - // emits `MetaUpdateRunning { running: true }` now and `false` on - // drop so the META INPUTS panel shows progress (issue #259). - let _progress = coord.meta_update_guard(); - tracing::info!(?inputs, "meta-update: starting"); - if let Err(e) = crate::meta::lock_update(inputs).await { - tracing::warn!(error = ?e, "meta-update: lock_update failed"); - return; - } - - // Decide which agents to rebuild. Inputs are slash-paths from - // the meta root — `hyperhive`, `hyperhive/nixpkgs`, - // `agent-coder`, `agent-coder/mcp-matrix`, etc. Anything in the - // hyperhive subtree affects every agent (shared base); anything - // in `agent-/...` only the named agent. - let touched_hyperhive = inputs - .iter() - .any(|i| i == "hyperhive" || i.starts_with("hyperhive/")); - let touched_agents: Vec = inputs - .iter() - .filter_map(|i| i.strip_prefix("agent-")) - .map(|rest| rest.split('/').next().unwrap_or(rest).to_owned()) - .collect(); - let agents_to_rebuild: Vec = if touched_hyperhive { - crate::lifecycle::list() - .await - .unwrap_or_default() - .into_iter() - .filter_map(|c| { - if c == crate::lifecycle::MANAGER_NAME { - Some(crate::lifecycle::MANAGER_NAME.to_owned()) - } else { - c.strip_prefix(crate::lifecycle::AGENT_PREFIX) - .map(str::to_owned) - } - }) - .collect() - } else { - touched_agents - }; - - let current_rev = - crate::auto_update::current_flake_rev(&coord.hyperhive_flake).unwrap_or_default(); - // Sequential rebuild loop — the META_LOCK guards meta-side - // races but parallel nix builds also serialise via nix-daemon, - // so sequential is just as fast in practice and keeps logs - // readable. - for name in agents_to_rebuild { - tracing::info!(%name, "meta-update: rebuilding agent"); - if let Err(e) = crate::auto_update::rebuild_agent(coord, &name, ¤t_rev).await { - tracing::warn!(%name, error = ?e, "meta-update: rebuild failed"); - // continue: surface each per-agent failure via its own - // Rebuilt event; don't abort the whole batch. - } - } - tracing::info!("meta-update: done"); -} - async fn post_op_send(State(state): State, Form(form): Form) -> Response { let to = form.to.trim().to_owned(); let body = form.body.trim().to_owned(); @@ -1708,28 +1651,16 @@ async fn post_request_spawn( } async fn post_rebuild(State(state): State, AxumPath(name): AxumPath) -> Response { - let Some(current_rev) = crate::auto_update::current_flake_rev(&state.coord.hyperhive_flake) - else { - return error_response( - "rebuild: hyperhive_flake has no canonical path; manual rebuild only via `hive-c0re rebuild`", - ); - }; - let coord = state.coord.clone(); - lifecycle_action( - &state, - &name, - crate::coordinator::TransientKind::Rebuilding, - "rebuild", - move |n| { - let coord = coord.clone(); - let rev = current_rev.clone(); - async move { crate::auto_update::rebuild_agent(&coord, &n, &rev).await } - }, - // rebuild_agent fires kick_agent on success itself, so the - // extra-closure is a no-op here. - |_, _| {}, - ) - .await + let logical = strip_container_prefix(&name); + state.coord.rebuild_queue.enqueue( + crate::rebuild_queue::QueueKind::Rebuild, + logical, + crate::rebuild_queue::QueueSource::Manual, + "manual via dashboard ↻ R3BU1LD button".to_owned(), + None, + ); + state.coord.emit_rebuild_queue_snapshot(); + (StatusCode::OK, "ok").into_response() } /// Common shape for the simple lifecycle action handlers (start / @@ -1816,12 +1747,7 @@ async fn post_start(State(state): State, AxumPath(name): AxumPath) -> Response { - let Some(current_rev) = crate::auto_update::current_flake_rev(&state.coord.hyperhive_flake) - else { - return error_response("update-all: hyperhive_flake has no canonical path"); - }; let containers = lifecycle::list().await.unwrap_or_default(); - let mut errors = Vec::new(); for container in containers { let logical = if container == lifecycle::MANAGER_NAME { lifecycle::MANAGER_NAME.to_owned() @@ -1830,21 +1756,16 @@ async fn post_update_all(State(state): State) -> Response { } else { continue; }; - if let Err(e) = - crate::auto_update::rebuild_agent(&state.coord, &logical, ¤t_rev).await - { - errors.push(format!("{logical}: {e:#}")); - } - } - if errors.is_empty() { - // Each rebuild_agent rescanned; no extra refetch needed. - (StatusCode::OK, "ok").into_response() - } else { - error_response(&format!( - "update-all partial failure:\n{}", - errors.join("\n") - )) + state.coord.rebuild_queue.enqueue( + crate::rebuild_queue::QueueKind::Rebuild, + logical, + crate::rebuild_queue::QueueSource::Manual, + "manual via dashboard 🌀 UPDATE ALL".to_owned(), + None, + ); } + state.coord.emit_rebuild_queue_snapshot(); + (StatusCode::OK, "ok").into_response() } fn transient_label(k: crate::coordinator::TransientKind) -> &'static str { diff --git a/hive-c0re/src/dashboard_events.rs b/hive-c0re/src/dashboard_events.rs index f82814de..fcea9c78 100644 --- a/hive-c0re/src/dashboard_events.rs +++ b/hive-c0re/src/dashboard_events.rs @@ -27,6 +27,7 @@ use serde::Serialize; use crate::container_view::ContainerView; use crate::dashboard::{MetaInputView, TombstoneView}; +use crate::rebuild_queue::QueueEntry; #[derive(Debug, Clone, Serialize)] #[serde(rename_all = "snake_case", tag = "kind")] @@ -204,4 +205,16 @@ pub enum DashboardEvent { /// when the active-run count crosses 0, so concurrent updates flip /// the flag exactly once. MetaUpdateRunning { seq: u64, running: bool }, + /// Full snapshot of the rebuild queue (`hive-c0re::rebuild_queue`) + /// — every entry, in enqueue order, including the few most-recent + /// terminal entries the queue retains for history. Same + /// snapshot-shape rationale as `TombstonesChanged` / + /// `MetaInputsChanged`: the list is small, snapshot semantics avoid + /// the add/remove races a per-row event would have, and the + /// dashboard's grouping (parent_id) is most naturally re-derived + /// from the full list. + RebuildQueueChanged { + seq: u64, + queue: Vec, + }, } diff --git a/hive-c0re/src/main.rs b/hive-c0re/src/main.rs index 07e6d174..5cd24a30 100644 --- a/hive-c0re/src/main.rs +++ b/hive-c0re/src/main.rs @@ -27,6 +27,7 @@ mod meta; mod migrate; mod operator_questions; mod questions; +mod rebuild_queue; mod reminder_scheduler; mod server; @@ -232,6 +233,18 @@ async fn cmd_serve( // Reminder scheduler: drains due reminders + handles // file_path payload persistence. See reminder_scheduler.rs. reminder_scheduler::spawn(coord.clone()); + // Rebuild-queue worker: drains the global rebuild/meta-update/ + // spawn queue FIFO so hive-c0re never runs two heavyweight + // container ops concurrently. Existing rebuild call sites + // (auto_update, dashboard, manager, approval handler) enqueue + // here instead of awaiting `rebuild_agent` inline. See + // `rebuild_queue.rs`. + { + let q_coord = coord.clone(); + tokio::spawn(async move { + rebuild_queue::run_worker(q_coord).await; + }); + } // Forward every broker event onto the unified dashboard // channel with a freshly-stamped seq, so the dashboard SSE // sees broker messages + future mutation events on one diff --git a/hive-c0re/src/manager_server.rs b/hive-c0re/src/manager_server.rs index ae29b830..fb0ef127 100644 --- a/hive-c0re/src/manager_server.rs +++ b/hive-c0re/src/manager_server.rs @@ -291,25 +291,16 @@ async fn dispatch(req: &ManagerRequest, coord: &Arc) -> ManagerResp } } ManagerRequest::Update { name } => { - tracing::info!(%name, "manager: update"); - let Some(current_rev) = crate::auto_update::current_flake_rev(&coord.hyperhive_flake) - else { - return ManagerResponse::Err { - message: "update: hyperhive_flake has no canonical path".into(), - }; - }; - let guard = coord.transient_guard(name, crate::coordinator::TransientKind::Rebuilding); - let result = crate::auto_update::rebuild_agent(coord, name, ¤t_rev).await; - drop(guard); - match result { - Ok(()) => { - coord.kick_agent(name, "container rebuilt"); - ManagerResponse::Ok - } - Err(e) => ManagerResponse::Err { - message: format!("{e:#}"), - }, - } + tracing::info!(%name, "manager: enqueue update"); + coord.rebuild_queue.enqueue( + crate::rebuild_queue::QueueKind::Rebuild, + name.to_owned(), + crate::rebuild_queue::QueueSource::Manual, + "manager `update` tool".to_owned(), + None, + ); + coord.emit_rebuild_queue_snapshot(); + ManagerResponse::Ok } ManagerRequest::RequestUpdateMetaInputs { inputs, diff --git a/hive-c0re/src/rebuild_queue.rs b/hive-c0re/src/rebuild_queue.rs new file mode 100644 index 00000000..3515cb01 --- /dev/null +++ b/hive-c0re/src/rebuild_queue.rs @@ -0,0 +1,749 @@ +//! Global rebuild queue. +//! +//! Every long-running container/meta operation (rebuild, meta-update, +//! first-spawn) goes through this queue. A single background worker +//! drains it in FIFO order so we never overlap two `nixos-container +//! update` runs on the same agent and never start a fresh agent rebuild +//! while a meta-update's lock bump is mid-flight. +//! +//! ## Why one queue +//! +//! Before this module landed, four independent call paths could fire +//! `auto_update::rebuild_agent` concurrently: +//! - dashboard manual rebuild button +//! - `update-all` / `meta-update` cascade +//! - approval handler (apply-commit / spawn) +//! - startup auto-update sweep +//! +//! Nothing serialised them. nix-daemon serialises the actual store +//! ops, but the rest of `rebuild_agent` (token sync, kick, rescan, +//! lock-bump emit) interleaved unpredictably. The single-worker queue +//! gives operators a visible, ordered runway and lets the UI render +//! "what's about to happen" instead of "something might be happening +//! somewhere." +//! +//! ## Scope +//! +//! In-queue kinds: +//! - `Rebuild` — a single-agent rebuild (covers manual / approval-driven / +//! auto-update / meta-update cascade variants — they all funnel here). +//! - `MetaUpdate` — `nix flake update` on the meta flake. The worker +//! runs the lock bump itself, then enqueues a cascade of `Rebuild` +//! entries with `parent_id` set to the meta-update's id. +//! - `Spawn` — first-deploy of an agent (approval-driven). Same +//! serialisation as `Rebuild` from the operator's POV. +//! - `Destroy` — for future use (`destroy --purge` does real I/O); not +//! currently routed through the queue. +//! +//! Out of scope (intentionally not queued — these are sub-second ops +//! and adding them adds visual noise without serving the "one at a +//! time" goal): +//! - `start` / `stop` / `restart` +//! - `kill` +//! +//! ## Dedup +//! +//! Enqueueing `(kind, agent)` that already has a `Queued` entry returns +//! the existing entry's id and appends the new reason as an +//! "also requested by …" line. Running entries do not dedup — a +//! re-queue during a run is legitimate (something changed since the +//! current run started). + +use std::collections::VecDeque; +use std::sync::{Arc, Mutex}; + +use serde::Serialize; +use tokio::sync::Notify; + +/// What the queue can run. Each variant maps to a specific worker +/// execution path; `agent` (in `QueueEntry`) names the target where +/// relevant. +#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize)] +#[serde(rename_all = "snake_case")] +pub enum QueueKind { + /// Rebuild a single agent's container (`auto_update::rebuild_agent`). + Rebuild, + /// Run `nix flake update` on the meta flake. Triggers cascade + /// `Rebuild` entries (with `parent_id`) once the lock bump lands. + MetaUpdate, + /// First-deploy spawn of a new agent (approval-driven). + Spawn, + /// Destroy with `--purge` (real fs work). Not yet routed here; the + /// variant exists so the wire shape doesn't need to change later. + Destroy, +} + +impl QueueKind { + pub fn as_str(self) -> &'static str { + match self { + QueueKind::Rebuild => "rebuild", + QueueKind::MetaUpdate => "meta_update", + QueueKind::Spawn => "spawn", + QueueKind::Destroy => "destroy", + } + } +} + +/// Where the enqueue request originated. Drives the "why" chip on the +/// dashboard and lets the UI group cascade entries under their parent +/// without parsing the reason text. +#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)] +#[serde(rename_all = "snake_case")] +pub enum QueueSource { + /// Operator clicked rebuild / update-all / meta-update on the + /// dashboard, or any other direct human action (CLI, manager tool). + Manual, + /// Spawned as a cascade from a `MetaUpdate` entry's lock-bump + /// fan-out. The `parent_id` on the `QueueEntry` points back at + /// the originating meta-update. + MetaUpdate, + /// `auto_update::run` startup sweep — rebuild every container on + /// hive-c0re boot. + AutoUpdate, + /// Crash recovery path (future use — currently no auto-rebuild on + /// crash, but the variant exists for the imminent feature). + CrashRecover, +} + +impl QueueSource { + pub fn as_str(self) -> &'static str { + match self { + QueueSource::Manual => "manual", + QueueSource::MetaUpdate => "meta_update", + QueueSource::AutoUpdate => "auto_update", + QueueSource::CrashRecover => "crash_recover", + } + } +} + +/// Lifecycle state of an entry. `Done` / `Failed` / `Cancelled` are +/// retained in the queue snapshot for a short tail (`MAX_HISTORY_PER_KIND`) +/// so the dashboard can show "last few" runs alongside live state. +#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)] +#[serde(rename_all = "snake_case")] +pub enum QueueState { + Queued, + Running, + Done, + Failed, + Cancelled, +} + +impl QueueState { + pub fn is_terminal(self) -> bool { + matches!(self, QueueState::Done | QueueState::Failed | QueueState::Cancelled) + } +} + +/// A single queue entry — what's pending, running, or recently finished. +/// Serialised verbatim onto the dashboard event channel and the +/// `/api/state` snapshot. +#[derive(Debug, Clone, Serialize)] +pub struct QueueEntry { + /// Monotonic per-process id. Stable for the lifetime of the entry + /// so SSE upserts land in place rather than churning the list. + pub id: u64, + /// Target agent name, or the literal `"hyperhive"` for entries + /// (MetaUpdate) that affect the meta flake rather than a single + /// agent. + pub agent: String, + pub kind: QueueKind, + pub state: QueueState, + pub source: QueueSource, + /// Groups cascade entries under their originating parent. For a + /// `MetaUpdate` entry this is `None`; for the per-agent rebuilds + /// the worker enqueues after the lock bump it's `Some(meta_id)`. + #[serde(default, skip_serializing_if = "Option::is_none")] + pub parent_id: Option, + /// Human-readable "why" — populated by the enqueuer (`"manual via + /// dashboard"`, `"meta-update cascade (hyperhive bumped)"`, + /// `"startup sweep"`). Free-form; dedup appends `(also requested + /// by …)` lines on repeated enqueues. + pub reason: String, + pub enqueued_at: i64, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub started_at: Option, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub finished_at: Option, + /// Populated when `state == Failed`. Carries the worker's error + /// string (already truncated to a reasonable length by the caller). + #[serde(default, skip_serializing_if = "Option::is_none")] + pub error: Option, + /// `MetaUpdate`-only payload: the list of meta flake inputs to run + /// through `nix flake update`. Empty / absent on `Rebuild` / + /// `Spawn` / `Destroy` entries; absent on the wire (never + /// serialised) when the entry kind doesn't have meaningful inputs. + #[serde(default, skip_serializing_if = "Vec::is_empty")] + pub inputs: Vec, +} + +/// How many terminal-state entries (`Done` / `Failed` / `Cancelled`) +/// to retain per kind in the snapshot. Older entries get evicted to +/// keep `/api/state` tight; the live event channel is unaffected. +const MAX_HISTORY_PER_KIND: usize = 5; + +/// Inner state guarded by a single mutex. Held briefly — every +/// operation is constant-time relative to the queue's depth, and +/// the depths in practice are tiny (single-digit). +#[derive(Debug, Default)] +struct Inner { + entries: VecDeque, + next_id: u64, +} + +/// Global rebuild queue. Lives on `Coordinator` (one per hive-c0re +/// process). The associated `Notify` wakes the worker when something +/// new arrives. +#[derive(Debug)] +pub struct RebuildQueue { + inner: Mutex, + /// Worker wakes on this signal. The worker checks the queue and + /// loops back to `notified().await` when there's nothing to run. + pub(crate) notify: Notify, +} + +impl Default for RebuildQueue { + fn default() -> Self { + Self { + inner: Mutex::new(Inner::default()), + notify: Notify::new(), + } + } +} + +impl RebuildQueue { + pub fn new() -> Self { + Self::default() + } + + /// Add an entry to the queue. Returns the entry's id (newly-allocated + /// or — on dedup — the existing entry's id with the new reason + /// appended). + /// + /// Dedup rule: a `Queued` entry with the same `(kind, agent)` swallows + /// the new request and returns its existing id. Running and terminal + /// entries do not dedup — operators are free to re-queue a rebuild + /// that's currently running (something changed since it started) or + /// re-run one that just finished. + pub fn enqueue( + &self, + kind: QueueKind, + agent: String, + source: QueueSource, + reason: String, + parent_id: Option, + ) -> u64 { + self.enqueue_with_inputs(kind, agent, source, reason, parent_id, Vec::new()) + } + + /// Same as `enqueue` but carries an `inputs` payload — used by + /// `MetaUpdate` enqueues to tell the worker which meta-flake + /// inputs to bump. + pub fn enqueue_with_inputs( + &self, + kind: QueueKind, + agent: String, + source: QueueSource, + reason: String, + parent_id: Option, + inputs: Vec, + ) -> u64 { + let mut inner = self.inner.lock().expect("rebuild_queue mutex poisoned"); + // Dedup against a pending entry with the same (kind, agent). + for entry in inner.entries.iter_mut() { + if entry.state == QueueState::Queued && entry.kind == kind && entry.agent == agent { + if !entry.reason.contains(&reason) { + entry.reason.push_str(&format!("\nalso requested by: {reason}")); + } + return entry.id; + } + } + inner.next_id += 1; + let id = inner.next_id; + let entry = QueueEntry { + id, + agent, + kind, + state: QueueState::Queued, + source, + parent_id, + reason, + enqueued_at: now_unix(), + started_at: None, + finished_at: None, + error: None, + inputs, + }; + inner.entries.push_back(entry); + // Wake the worker. `notify_one` is a no-op when there's no + // waiter; the next `notified().await` returns immediately. + self.notify.notify_one(); + id + } + + /// Pop the next `Queued` entry and mark it `Running`. Returns the + /// entry (a clone — the original stays in the queue so live state + /// reflects "this is currently running"). Returns `None` when there's + /// nothing queued. + pub fn take_next(&self) -> Option { + let mut inner = self.inner.lock().expect("rebuild_queue mutex poisoned"); + let pos = inner + .entries + .iter() + .position(|e| e.state == QueueState::Queued)?; + let entry = &mut inner.entries[pos]; + entry.state = QueueState::Running; + entry.started_at = Some(now_unix()); + Some(entry.clone()) + } + + /// Mark an entry terminal. `error` is populated for `Failed`; + /// `Done` / `Cancelled` ignore it. Trims the history tail. + pub fn finish(&self, id: u64, state: QueueState, error: Option) { + debug_assert!(state.is_terminal(), "finish() called with non-terminal {state:?}"); + let mut inner = self.inner.lock().expect("rebuild_queue mutex poisoned"); + if let Some(entry) = inner.entries.iter_mut().find(|e| e.id == id) { + entry.state = state; + entry.finished_at = Some(now_unix()); + entry.error = error.filter(|_| state == QueueState::Failed); + } + Self::trim_history(&mut inner); + } + + /// Snapshot the queue for `/api/state` and `RebuildQueueChanged`. + /// Cheap clone — entries are small (~hundreds of bytes each). + pub fn snapshot(&self) -> Vec { + let inner = self.inner.lock().expect("rebuild_queue mutex poisoned"); + inner.entries.iter().cloned().collect() + } + + /// Cancel a `Queued` entry (no-op for `Running` / terminal — the + /// in-flight rebuild owns the agent's nix store and can't be + /// safely interrupted). Returns true when an entry was cancelled. + pub fn cancel(&self, id: u64) -> bool { + let mut inner = self.inner.lock().expect("rebuild_queue mutex poisoned"); + if let Some(entry) = inner.entries.iter_mut().find(|e| e.id == id) { + if entry.state == QueueState::Queued { + entry.state = QueueState::Cancelled; + entry.finished_at = Some(now_unix()); + Self::trim_history(&mut inner); + return true; + } + } + false + } + + /// Keep only the most recent `MAX_HISTORY_PER_KIND` terminal entries + /// per kind. Pending + running entries are never evicted. + fn trim_history(inner: &mut Inner) { + let mut counts: std::collections::HashMap = + std::collections::HashMap::new(); + // Walk newest-first; keep the first MAX_HISTORY_PER_KIND + // terminals per kind, evict the rest. + let entries: Vec = inner + .entries + .iter() + .rev() + .filter(|e| { + if !e.state.is_terminal() { + return true; + } + let n = counts.entry(e.kind).or_insert(0); + *n += 1; + *n <= MAX_HISTORY_PER_KIND + }) + .cloned() + .collect(); + inner.entries = entries.into_iter().rev().collect(); + } +} + +/// Background worker that drains the queue. Spawned once at hive-c0re +/// startup from `main.rs`. Loops forever: +/// 1. Pop the next `Queued` entry (`take_next` marks it `Running` and +/// fires a `RebuildQueueChanged` snapshot via the caller). +/// 2. Dispatch by kind — single-agent rebuild, meta-update + cascade, +/// or first-spawn. +/// 3. Mark the entry terminal (`finish`) and emit another snapshot. +/// 4. When the queue is empty, `await` on `notify` until something +/// new lands. +/// +/// Shutdown semantics: subscribes to `coord.shutdown_rx()`. On a true +/// signal the worker exits after its current entry finishes; pending +/// `Queued` entries are dropped (they'll either be replayed by the +/// startup sweep on next boot or left for an operator to re-queue). +pub async fn run_worker(coord: std::sync::Arc) { + let mut shutdown = coord.shutdown_rx(); + loop { + // Drain everything available now. + while let Some(entry) = coord.rebuild_queue.take_next() { + coord.emit_rebuild_queue_snapshot(); + tracing::info!( + id = entry.id, + kind = entry.kind.as_str(), + agent = %entry.agent, + source = entry.source.as_str(), + "rebuild_queue: running" + ); + let result = dispatch(&coord, &entry).await; + match result { + Ok(()) => { + coord.rebuild_queue.finish(entry.id, QueueState::Done, None); + tracing::info!(id = entry.id, "rebuild_queue: done"); + } + Err(e) => { + let msg = format!("{e:#}"); + let truncated = if msg.len() > 2_000 { + format!("{}…", &msg[..2_000]) + } else { + msg.clone() + }; + coord + .rebuild_queue + .finish(entry.id, QueueState::Failed, Some(truncated)); + tracing::warn!(id = entry.id, error = %msg, "rebuild_queue: failed"); + } + } + coord.emit_rebuild_queue_snapshot(); + } + // Park until something new is enqueued OR shutdown fires. + tokio::select! { + biased; + res = shutdown.changed() => { + if res.is_err() || *shutdown.borrow() { + tracing::info!("rebuild_queue: worker exiting on shutdown"); + return; + } + } + _ = coord.rebuild_queue.notify.notified() => { + // New entry — back to the drain loop. + } + } + } +} + +/// Run a single queue entry to completion. Kind-dispatched; failures +/// bubble up to the worker which marks the entry `Failed`. +async fn dispatch( + coord: &std::sync::Arc, + entry: &QueueEntry, +) -> anyhow::Result<()> { + match entry.kind { + QueueKind::Rebuild => { + let current_rev = crate::auto_update::current_flake_rev(&coord.hyperhive_flake) + .unwrap_or_default(); + crate::auto_update::rebuild_agent(coord, &entry.agent, ¤t_rev).await + } + QueueKind::MetaUpdate => run_meta_update(coord, entry).await, + QueueKind::Spawn => { + // First-deploy spawns route through `actions::approve_spawn` + // / `actions::approve_apply_commit` today; they enqueue a + // Spawn entry only to claim the queue slot, the actual + // spawn work runs inside those handlers before completion. + // Keeping this arm a no-op so we don't double-run. + tracing::debug!( + id = entry.id, + agent = %entry.agent, + "rebuild_queue: Spawn entry is a queue claim; actual work elsewhere" + ); + Ok(()) + } + QueueKind::Destroy => { + // Reserved for future `destroy --purge` integration. + anyhow::bail!("Destroy kind not yet implemented in rebuild_queue worker"); + } + } +} + +/// Run one `MetaUpdate` entry: bump the meta flake's locks for the +/// requested inputs, then enqueue a cascade of `Rebuild` entries +/// (with `parent_id` set to this entry's id) for every agent affected +/// by the bump. Mirrors the previous `dashboard::run_meta_update` +/// semantics; that path now enqueues into this queue rather than +/// running the bump + rebuild loop inline. +async fn run_meta_update( + coord: &std::sync::Arc, + entry: &QueueEntry, +) -> anyhow::Result<()> { + let _progress = coord.meta_update_guard(); + let inputs = entry.inputs.clone(); + tracing::info!(?inputs, parent = entry.id, "rebuild_queue: meta-update starting"); + if inputs.is_empty() { + crate::meta::lock_update(&[]).await?; + } else { + crate::meta::lock_update(&inputs).await?; + } + + // Decide which agents to rebuild. Same logic as the previous + // `run_meta_update` — anything in the hyperhive subtree affects + // every agent; anything in `agent-/...` only the named agent. + let touched_hyperhive = inputs + .iter() + .any(|i| i == "hyperhive" || i.starts_with("hyperhive/")); + let touched_agents: Vec = inputs + .iter() + .filter_map(|i| i.strip_prefix("agent-")) + .map(|rest| rest.split('/').next().unwrap_or(rest).to_owned()) + .collect(); + let agents_to_rebuild: Vec = if touched_hyperhive || inputs.is_empty() { + crate::lifecycle::list() + .await + .unwrap_or_default() + .into_iter() + .filter_map(|c| { + if c == crate::lifecycle::MANAGER_NAME { + Some(crate::lifecycle::MANAGER_NAME.to_owned()) + } else { + c.strip_prefix(crate::lifecycle::AGENT_PREFIX).map(str::to_owned) + } + }) + .collect() + } else { + touched_agents + }; + + let reason_hint = if inputs.is_empty() { + "meta-update cascade (all inputs)".to_owned() + } else { + format!("meta-update cascade ({})", inputs.join(", ")) + }; + for name in agents_to_rebuild { + coord.rebuild_queue.enqueue( + QueueKind::Rebuild, + name, + QueueSource::MetaUpdate, + reason_hint.clone(), + Some(entry.id), + ); + } + // Lock file changed — meta-inputs panel re-renders. + crate::dashboard::emit_meta_inputs_snapshot(coord.as_ref()); + Ok(()) +} + +/// Current unix timestamp in seconds. `now()` calls are pulled into a +/// helper so tests can swap them out later. +fn now_unix() -> i64 { + std::time::SystemTime::now() + .duration_since(std::time::UNIX_EPOCH) + .ok() + .and_then(|d| i64::try_from(d.as_secs()).ok()) + .unwrap_or(0) +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn enqueue_and_take_in_order() { + let q = RebuildQueue::new(); + let a = q.enqueue( + QueueKind::Rebuild, + "agent-a".to_owned(), + QueueSource::Manual, + "first".to_owned(), + None, + ); + let b = q.enqueue( + QueueKind::Rebuild, + "agent-b".to_owned(), + QueueSource::Manual, + "second".to_owned(), + None, + ); + assert_ne!(a, b); + let next = q.take_next().expect("queued"); + assert_eq!(next.id, a); + assert_eq!(next.state, QueueState::Running); + let next = q.take_next().expect("queued"); + assert_eq!(next.id, b); + assert!(q.take_next().is_none()); + } + + #[test] + fn dedup_pending_same_kind_and_agent() { + let q = RebuildQueue::new(); + let a = q.enqueue( + QueueKind::Rebuild, + "agent-a".to_owned(), + QueueSource::Manual, + "first".to_owned(), + None, + ); + let b = q.enqueue( + QueueKind::Rebuild, + "agent-a".to_owned(), + QueueSource::AutoUpdate, + "auto sweep".to_owned(), + None, + ); + assert_eq!(a, b, "dedup should return existing id"); + let snap = q.snapshot(); + assert_eq!(snap.len(), 1); + assert!(snap[0].reason.contains("first")); + assert!(snap[0].reason.contains("auto sweep")); + } + + #[test] + fn dedup_does_not_apply_across_kinds_or_agents() { + let q = RebuildQueue::new(); + let a = q.enqueue( + QueueKind::Rebuild, + "agent-a".to_owned(), + QueueSource::Manual, + "r".to_owned(), + None, + ); + let b = q.enqueue( + QueueKind::Rebuild, + "agent-b".to_owned(), + QueueSource::Manual, + "r".to_owned(), + None, + ); + let c = q.enqueue( + QueueKind::Spawn, + "agent-a".to_owned(), + QueueSource::Manual, + "s".to_owned(), + None, + ); + assert_ne!(a, b); + assert_ne!(a, c); + assert_eq!(q.snapshot().len(), 3); + } + + #[test] + fn dedup_skips_running_entries() { + let q = RebuildQueue::new(); + q.enqueue( + QueueKind::Rebuild, + "agent-a".to_owned(), + QueueSource::Manual, + "first".to_owned(), + None, + ); + let running = q.take_next().expect("queued"); + assert_eq!(running.state, QueueState::Running); + // While the original is running, re-enqueue is legitimate. + let again = q.enqueue( + QueueKind::Rebuild, + "agent-a".to_owned(), + QueueSource::Manual, + "config bumped during build".to_owned(), + None, + ); + assert_ne!(running.id, again); + let snap = q.snapshot(); + assert_eq!(snap.len(), 2); + } + + #[test] + fn finish_marks_state_and_keeps_history() { + let q = RebuildQueue::new(); + let id = q.enqueue( + QueueKind::Rebuild, + "agent-a".to_owned(), + QueueSource::Manual, + "r".to_owned(), + None, + ); + q.take_next(); + q.finish(id, QueueState::Done, None); + let snap = q.snapshot(); + assert_eq!(snap.len(), 1); + assert_eq!(snap[0].state, QueueState::Done); + assert!(snap[0].finished_at.is_some()); + assert!(snap[0].error.is_none()); + } + + #[test] + fn finish_with_failure_records_error() { + let q = RebuildQueue::new(); + let id = q.enqueue( + QueueKind::Rebuild, + "agent-a".to_owned(), + QueueSource::Manual, + "r".to_owned(), + None, + ); + q.take_next(); + q.finish(id, QueueState::Failed, Some("nix build failed".to_owned())); + let snap = q.snapshot(); + assert_eq!(snap[0].state, QueueState::Failed); + assert_eq!(snap[0].error.as_deref(), Some("nix build failed")); + } + + #[test] + fn history_evicts_old_terminals_per_kind() { + let q = RebuildQueue::new(); + for i in 0..(MAX_HISTORY_PER_KIND + 3) { + let id = q.enqueue( + QueueKind::Rebuild, + format!("agent-{i}"), + QueueSource::Manual, + "r".to_owned(), + None, + ); + q.take_next(); + q.finish(id, QueueState::Done, None); + } + let snap = q.snapshot(); + assert_eq!(snap.len(), MAX_HISTORY_PER_KIND); + } + + #[test] + fn cancel_clears_queued_entry() { + let q = RebuildQueue::new(); + let id = q.enqueue( + QueueKind::Rebuild, + "agent-a".to_owned(), + QueueSource::Manual, + "r".to_owned(), + None, + ); + assert!(q.cancel(id)); + let snap = q.snapshot(); + assert_eq!(snap[0].state, QueueState::Cancelled); + assert!(q.take_next().is_none()); + } + + #[test] + fn cancel_refuses_running_entry() { + let q = RebuildQueue::new(); + let id = q.enqueue( + QueueKind::Rebuild, + "agent-a".to_owned(), + QueueSource::Manual, + "r".to_owned(), + None, + ); + q.take_next(); + assert!(!q.cancel(id)); + let snap = q.snapshot(); + assert_eq!(snap[0].state, QueueState::Running); + } + + #[test] + fn parent_id_groups_cascade() { + let q = RebuildQueue::new(); + let meta = q.enqueue( + QueueKind::MetaUpdate, + "hyperhive".to_owned(), + QueueSource::Manual, + "lock bump".to_owned(), + None, + ); + let child = q.enqueue( + QueueKind::Rebuild, + "agent-a".to_owned(), + QueueSource::MetaUpdate, + "cascade".to_owned(), + Some(meta), + ); + let snap = q.snapshot(); + let child_entry = snap.iter().find(|e| e.id == child).expect("child queued"); + assert_eq!(child_entry.parent_id, Some(meta)); + } +}