diff --git a/hive-c0re/assets/app.js b/hive-c0re/assets/app.js index 53012bde..7d14054b 100644 --- a/hive-c0re/assets/app.js +++ b/hive-c0re/assets/app.js @@ -493,22 +493,6 @@ 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 @@ -1669,119 +1653,6 @@ 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 @@ -1893,7 +1764,6 @@ 'inbox-section', 'approvals-section', 'meta-inputs-section', - 'rebuild-queue-section', 'reminders-section', ]; //
sections that should survive a refresh need a stable @@ -1960,7 +1830,6 @@ syncContainersFromSnapshot(s); syncTombstonesFromSnapshot(s); syncMetaInputsFromSnapshot(s); - syncRebuildQueueFromSnapshot(s); renderContainers(s); renderTombstones(s); // Sync the derived approvals + questions stores from the @@ -1973,7 +1842,6 @@ syncApprovalsFromSnapshot(s); renderApprovals(); renderMetaInputs(s); - renderRebuildQueue(s); refreshReminders(); restoreOpenDetails(openDetails); notifyDeltas(s); @@ -2095,7 +1963,6 @@ 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 5f0e8c20..248f5234 100644 --- a/hive-c0re/assets/dashboard.css +++ b/hive-c0re/assets/dashboard.css @@ -521,61 +521,6 @@ 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 846c86e9..25fdbd51 100644 --- a/hive-c0re/assets/index.html +++ b/hive-c0re/assets/index.html @@ -39,13 +39,6 @@

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 1ac2add7..a21b5b30 100644 --- a/hive-c0re/src/auto_update.rs +++ b/hive-c0re/src/auto_update.rs @@ -206,9 +206,10 @@ 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: queueing all on startup"); + tracing::info!(agents = containers.len(), "auto-update: rebuilding all on startup"); for container in containers { let logical = if container == MANAGER_NAME { Some(MANAGER_NAME.to_owned()) @@ -216,14 +217,9 @@ pub async fn run(coord: Arc) -> Result<()> { container.strip_prefix(AGENT_PREFIX).map(str::to_owned) }; let Some(name) = logical else { continue }; - coord.rebuild_queue.enqueue( - crate::rebuild_queue::QueueKind::Rebuild, - name, - crate::rebuild_queue::QueueSource::AutoUpdate, - "startup sweep".to_owned(), - None, - ); + if let Err(e) = rebuild_agent(&coord, &name, ¤t_rev).await { + tracing::warn!(%name, error = ?e, "auto-update: rebuild failed"); + } } - coord.emit_rebuild_queue_snapshot(); Ok(()) } diff --git a/hive-c0re/src/coordinator.rs b/hive-c0re/src/coordinator.rs index 357e9b44..d4d41ff0 100644 --- a/hive-c0re/src/coordinator.rs +++ b/hive-c0re/src/coordinator.rs @@ -88,12 +88,6 @@ 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. @@ -208,23 +202,10 @@ 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 529da6cf..1cc418f7 100644 --- a/hive-c0re/src/dashboard.rs +++ b/hive-c0re/src/dashboard.rs @@ -237,12 +237,6 @@ 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. @@ -437,7 +431,6 @@ 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, }) } @@ -1578,18 +1571,82 @@ async fn post_meta_update( if inputs.is_empty() { return error_response("meta-update: no inputs selected"); } - 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(); + 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); + }); (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(); @@ -1651,16 +1708,28 @@ async fn post_request_spawn( } async fn post_rebuild(State(state): State, AxumPath(name): AxumPath) -> Response { - 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() + 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 } /// Common shape for the simple lifecycle action handlers (start / @@ -1747,7 +1816,12 @@ 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() @@ -1756,16 +1830,21 @@ async fn post_update_all(State(state): State) -> Response { } else { continue; }; - state.coord.rebuild_queue.enqueue( - crate::rebuild_queue::QueueKind::Rebuild, - logical, - crate::rebuild_queue::QueueSource::Manual, - "manual via dashboard 🌀 UPDATE ALL".to_owned(), - None, - ); + 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.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 fcea9c78..f82814de 100644 --- a/hive-c0re/src/dashboard_events.rs +++ b/hive-c0re/src/dashboard_events.rs @@ -27,7 +27,6 @@ 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")] @@ -205,16 +204,4 @@ 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 5cd24a30..07e6d174 100644 --- a/hive-c0re/src/main.rs +++ b/hive-c0re/src/main.rs @@ -27,7 +27,6 @@ mod meta; mod migrate; mod operator_questions; mod questions; -mod rebuild_queue; mod reminder_scheduler; mod server; @@ -233,18 +232,6 @@ 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 fb0ef127..ae29b830 100644 --- a/hive-c0re/src/manager_server.rs +++ b/hive-c0re/src/manager_server.rs @@ -291,16 +291,25 @@ async fn dispatch(req: &ManagerRequest, coord: &Arc) -> ManagerResp } } ManagerRequest::Update { name } => { - 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 + 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:#}"), + }, + } } ManagerRequest::RequestUpdateMetaInputs { inputs, diff --git a/hive-c0re/src/rebuild_queue.rs b/hive-c0re/src/rebuild_queue.rs deleted file mode 100644 index 3515cb01..00000000 --- a/hive-c0re/src/rebuild_queue.rs +++ /dev/null @@ -1,749 +0,0 @@ -//! 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)); - } -}