From 7ac6819652e0e9cd417d9df3afcbbfc652584f39 Mon Sep 17 00:00:00 2001 From: atlas Date: Sat, 26 Sep 2026 15:18:41 +0200 Subject: [PATCH] hive-c0re: fail on a malformed agent name and on an unreadable container list An agent name that is not a valid Ident made `Coordinator::agent_paths` panic. Job payloads carry names as plain strings (the swarm's published wanted state is one source), and a panic inside a job-queue node never reaches `complete_growing`, so the node's resources (the deploy window included) were held until hive-c0re restarted. `agent_paths` now returns an error; the job-queue nodes, the admin-socket spawn and set-limits paths, the root-agent spawn and the dashboard set-limits handler propagate it. `lifecycle::list().await.unwrap_or_default()` turned a failed container list into "no agents": - meta-update cascade: the lock bump committed and zero rebuilds fanned out, reported as success. The cascade is now resolved before the lock bump and a list failure fails the node. - dashboard update-all: queued nothing and returned 200 "ok". Now 500 with the error. - container rescan: every row was emitted as removed and the cache emptied. Now the last snapshot stands; `hivectl status` gets an error. - dashboard journal: answered 404 "no managed container". Now 500. - spawn/rebuild port-collision check: silently skipped. Now fails. - startup migration: the per-agent phases ran over nothing, and phase 3 handed an empty agent list to `meta::sync_agents`, which renders the meta flake with exactly the agents it is given. Both now log the list failure and skip. The hive-jobq scheduler still leaks a node's resources on any executor panic; that root is not addressed here. Refs #4723 --- hive-c0re/src/container_view.rs | 10 +++-- hive-c0re/src/coordinator.rs | 50 +++++++++++++++++++----- hive-c0re/src/dashboard/journal.rs | 10 ++++- hive-c0re/src/dashboard/lifecycle_ops.rs | 23 ++++++----- hive-c0re/src/job_queue/exec.rs | 39 +++++++++--------- hive-c0re/src/lifecycle/mod.rs | 24 +++++++++--- hive-c0re/src/lifecycle/tests.rs | 15 +++++++ hive-c0re/src/migrate.rs | 28 +++++++++---- hive-c0re/src/server.rs | 19 +++++---- hive-c0re/src/workers/auto_update.rs | 2 +- 10 files changed, 153 insertions(+), 67 deletions(-) diff --git a/hive-c0re/src/container_view.rs b/hive-c0re/src/container_view.rs index 1a2866d6..655693a2 100644 --- a/hive-c0re/src/container_view.rs +++ b/hive-c0re/src/container_view.rs @@ -152,8 +152,12 @@ fn agent_url_for_domain(domain: &str, name: &str) -> String { /// fallback onto the hive-wide `agent_cpu_quota` / `agent_memory_max`, /// and those live on [`crate::coordinator::HiveEnv`], not on disk. Both callers already /// hold a `Coordinator`, so this is a parameter rather than a global. -pub async fn build_all(hive: &crate::coordinator::HiveEnv) -> Vec { - let raw = lifecycle::list().await.unwrap_or_default(); +/// +/// # Errors +/// +/// The container list could not be read. +pub async fn build_all(hive: &crate::coordinator::HiveEnv) -> anyhow::Result> { + let raw = lifecycle::list().await?; let locked = read_meta_locked_revs(); // Read once per scan rather than re-read for every container on every // SSE scan. An unreadable file shows every row at the hive defaults; @@ -254,7 +258,7 @@ pub async fn build_all(hive: &crate::coordinator::HiveEnv) -> Vec memory_max, }); } - out + Ok(out) } /// Host-side mirror of `hive_agent::login::has_session`. Returns true diff --git a/hive-c0re/src/coordinator.rs b/hive-c0re/src/coordinator.rs index 06c48abe..3a443d2c 100644 --- a/hive-c0re/src/coordinator.rs +++ b/hive-c0re/src/coordinator.rs @@ -558,21 +558,22 @@ impl Coordinator { /// from `crate::paths::agent_runtime_dir(name)` (pure path) or from /// `lifecycle::ensure_agent_runtime_dir(name)` when the dir must be /// created. All other paths are derived statically from `name`. - #[must_use] - pub fn agent_paths(name: &str, agent_dir: PathBuf) -> AgentPaths { - // `name` is validated upstream (spawn-approval / enqueue gate / the - // MANAGER_NAME const), so an invalid ident here is a construction - // bug. This is the step-3 boundary between the Ident-threaded path - // builders and the job_queue layer (threaded post hive-jobq cutover). + /// + /// # Errors + /// + /// `name` is not a valid [`hive_types::Ident`]. Job payloads carry agent + /// names as plain strings, one source being the swarm's published wanted + /// state, so this is an error for the caller to report, not an invariant. + pub fn agent_paths(name: &str, agent_dir: PathBuf) -> Result { let name = hive_types::Ident::parse(name) - .expect("agent_paths: name must be a valid ident (validated at spawn/enqueue)"); - AgentPaths { + .map_err(|e| anyhow::anyhow!("invalid agent name {name:?}: {e}"))?; + Ok(AgentPaths { agent: agent_dir, proposed: Self::agent_proposed_dir(&name), applied: crate::paths::applied_dir(name.as_str()), claude: Self::agent_claude_dir(&name), notes: Self::agent_notes_dir(&name), - } + }) } /// Emit a `RebuildQueueChanged` tick. Called from the queue mutation @@ -845,7 +846,15 @@ impl Coordinator { /// Cheap when nothing changed (one `nixos-container list` + a /// `HashMap` diff + zero emits). pub async fn rescan_containers_and_emit(self: &Arc) { - let fresh = container_view::build_all(&self.hive_env()).await; + // An unreadable list says nothing about which containers exist, so + // the last snapshot stands rather than every row being reported gone. + let fresh = match container_view::build_all(&self.hive_env()).await { + Ok(fresh) => fresh, + Err(e) => { + tracing::warn!(error = ?e, "container rescan: list failed; keeping last snapshot"); + return; + } + }; let mut last = self.last_containers.lock().await; let mut changed_or_new = Vec::new(); let mut removed = Vec::new(); @@ -1580,3 +1589,24 @@ mod rebuilt_todo_summary_tests { ); } } + +#[cfg(test)] +mod agent_paths_tests { + use super::Coordinator; + + /// A name from a job payload that is not an `Ident` fails the caller's + /// node instead of panicking the executor task. + #[test] + fn a_malformed_name_is_an_error_not_a_panic() { + let err = Coordinator::agent_paths("../etc", std::path::PathBuf::from("/nonexistent")) + .expect_err("a malformed agent name must be rejected"); + assert!(format!("{err:#}").contains("invalid agent name"), "{err:#}"); + } + + #[test] + fn a_valid_name_builds_its_paths() { + let paths = Coordinator::agent_paths("iris", std::path::PathBuf::from("/run/x")) + .expect("a valid agent name"); + assert_eq!(paths.agent, std::path::PathBuf::from("/run/x")); + } +} diff --git a/hive-c0re/src/dashboard/journal.rs b/hive-c0re/src/dashboard/journal.rs index f09908c5..6071ff98 100644 --- a/hive-c0re/src/dashboard/journal.rs +++ b/hive-c0re/src/dashboard/journal.rs @@ -59,6 +59,7 @@ pub(super) struct JournalQuery { (status = 200, description = "journal text", body = String, content_type = "text/plain"), (status = 400, description = "bad agent name or unknown unit"), (status = 404, description = "no such managed container"), + (status = 500, description = "container list unreadable"), ), tag = "journal" )] @@ -97,7 +98,14 @@ pub(super) async fn get_journal( // containers so we don't shell out with arbitrary input. let container = strip_container_prefix(name.as_str()); let prefixed = format!("{}{container}", lifecycle::AGENT_PREFIX); - let live = lifecycle::list().await.unwrap_or_default(); + let live = match lifecycle::list().await { + Ok(live) => live, + Err(e) => { + return Err(error_problem(&format!( + "journal: listing containers: {e:#}" + ))); + } + }; if !live.iter().any(|c| c == &prefixed) { return Err(ProblemDetails::from_status_code(StatusCode::NOT_FOUND) .with_detail(format!("journal: no managed container {prefixed:?}"))); diff --git a/hive-c0re/src/dashboard/lifecycle_ops.rs b/hive-c0re/src/dashboard/lifecycle_ops.rs index 3ba31e22..bef88cd1 100644 --- a/hive-c0re/src/dashboard/lifecycle_ops.rs +++ b/hive-c0re/src/dashboard/lifecycle_ops.rs @@ -392,7 +392,10 @@ pub(super) async fn post_resource_limits( } let agent_dir = crate::paths::agent_runtime_dir(ident.as_str()); let hive = state.coord.hive_env(); - let paths = crate::coordinator::Coordinator::agent_paths(ident.as_str(), agent_dir); + let paths = match crate::coordinator::Coordinator::agent_paths(ident.as_str(), agent_dir) { + Ok(p) => p, + Err(e) => return error_response(&format!("write_dropins {logical}: {e:#}")), + }; if let Err(e) = crate::lifecycle::write_dropins(ident.as_str(), &hive, &paths).await { return error_response(&format!("write_dropins {logical}: {e:#}")); } @@ -405,18 +408,18 @@ pub(super) async fn post_resource_limits( #[utoipa::path( post, path = "/api/update-all", - responses((status = 200, description = "rebuilds queued", body = String)), + responses( + (status = 200, description = "rebuilds queued", body = String), + (status = 500, description = "container list unreadable, nothing queued"), + ), tag = "lifecycle_ops" )] pub(super) async fn post_update_all(State(state): State) -> Response { - let containers = lifecycle::list().await.unwrap_or_default(); - for container in containers { - let Some(logical) = container - .strip_prefix(lifecycle::AGENT_PREFIX) - .map(str::to_owned) - else { - continue; - }; + let agents = match lifecycle::agent_names(lifecycle::list().await) { + Ok(agents) => agents, + Err(e) => return error_response(&format!("update-all: listing containers: {e:#}")), + }; + for logical in agents { if let Err(e) = state.coord.job_queue.insert_job(|b| { crate::job_queue::templates::rebuild(b, &logical, true); Vec::new() diff --git a/hive-c0re/src/job_queue/exec.rs b/hive-c0re/src/job_queue/exec.rs index 0b709e6a..6bf0bdde 100644 --- a/hive-c0re/src/job_queue/exec.rs +++ b/hive-c0re/src/job_queue/exec.rs @@ -399,7 +399,7 @@ async fn run_meta_sync(coord: &Arc, name: &str, relock: bool) -> Re // listener (event-driven: registered at start/create). let agent_dir = crate::paths::agent_runtime_dir(name); let hive = coord.hive_env(); - let paths = Coordinator::agent_paths(name, agent_dir); + let paths = Coordinator::agent_paths(name, agent_dir)?; crate::lifecycle::prepare_rebuild_dirs(name, &paths).await?; // Idempotent meta sync so a manual rebuild can also recover from a // divergent meta repo; then bump just this agent's input. `relock = @@ -442,7 +442,7 @@ async fn run_swap(coord: &Arc, name: &str, id: NodeId) -> Result<() // and listener were created earlier. Pure path accessor suffices. let agent_dir = crate::paths::agent_runtime_dir(name); let hive = coord.hive_env(); - let paths = Coordinator::agent_paths(name, agent_dir); + let paths = Coordinator::agent_paths(name, agent_dir)?; let result = crate::lifecycle::swap_update(name, &hive, &paths, Some(id.get())).await; // On success the Ok-only bookkeeping tail (rev marker, forge/matrix // sync, kick, rescan, snapshot) runs in the sibling `RebuildBookkeeping` node, @@ -492,7 +492,7 @@ async fn run_rebuild_bookkeeping(coord: &Arc, name: &str) -> Result async fn run_provision(coord: &Arc, name: &str) -> Result<()> { let agent_dir = crate::paths::agent_runtime_dir(name); let hive = coord.hive_env(); - let paths = Coordinator::agent_paths(name, agent_dir); + let paths = Coordinator::agent_paths(name, agent_dir)?; crate::lifecycle::provision_container(name, &hive, &paths).await?; Ok(()) } @@ -536,13 +536,15 @@ async fn run_meta_lock( return Ok(fanout.unwrap_or_default()); } let _progress = coord.meta_update_guard(); + // Resolved before the lock bump so a failed container list aborts + // the update before anything is committed. + let cascade = match fanout { + Some(list) => validate_agent_names(list), + None => meta_update_cascade_agents(inputs).await?, + }; crate::meta::lock_update(inputs).await?; // Lock file changed — meta-inputs panel re-renders. crate::dashboard::emit_meta_inputs_snapshot(coord); - let cascade = match fanout { - Some(list) => validate_agent_names(list), - None => meta_update_cascade_agents(inputs).await, - }; // Pull each cascade agent's own input too — an agent's config-repo main // is trusted, so a meta-input bump is a reasonable place to also catch // it up. Without this, a meta-input bump rebuilds every affected agent @@ -627,7 +629,7 @@ async fn run_start(coord: &Arc, name: &str) -> Result<()> { // omitting this becomes a compile error. let agent_dir = crate::paths::agent_runtime_dir(name); let hive = coord.hive_env(); - let paths = Coordinator::agent_paths(name, agent_dir); + let paths = Coordinator::agent_paths(name, agent_dir)?; let token = crate::lifecycle::converge_start_preamble(name, &hive, &paths).await?; crate::lifecycle::start_with_fallback(token).await?; // Bind the MCP listener immediately after starting the container. @@ -764,7 +766,7 @@ async fn run_write_dropin(coord: &Arc, name: &str) -> Result<()> { // on the upstream Prebuild/Start node). let agent_dir = crate::paths::agent_runtime_dir(name); let hive = coord.hive_env(); - let paths = Coordinator::agent_paths(name, agent_dir); + let paths = Coordinator::agent_paths(name, agent_dir)?; crate::lifecycle::write_dropins(name, &hive, &paths).await?; Ok(()) } @@ -870,26 +872,23 @@ async fn run_deploy_tail( /// cascade agent's name. The `touched_hyperhive` branch's names need no /// such filter: they come from `lifecycle::list()`, already real /// container names by construction, not parsed out of caller input. -pub async fn meta_update_cascade_agents(inputs: &[String]) -> Vec { +/// +/// # Errors +/// +/// The every-container branch could not list the containers. +pub async fn meta_update_cascade_agents(inputs: &[String]) -> Result> { let touched_hyperhive = inputs .iter() .any(|i| i == "hyperhive" || i.starts_with("hyperhive/")); let touched_agents = parse_agent_input_names(inputs); let mut names = if touched_hyperhive || inputs.is_empty() { - crate::lifecycle::list() - .await - .unwrap_or_default() - .into_iter() - .filter_map(|c| { - c.strip_prefix(crate::lifecycle::AGENT_PREFIX) - .map(str::to_owned) - }) - .collect() + crate::lifecycle::agent_names(crate::lifecycle::list().await) + .context("listing containers for the meta-update cascade")? } else { touched_agents }; names.sort(); - names + Ok(names) } /// Parse `agent-` input strings into validated agent names. Pure and diff --git a/hive-c0re/src/lifecycle/mod.rs b/hive-c0re/src/lifecycle/mod.rs index 998397a6..6843e732 100644 --- a/hive-c0re/src/lifecycle/mod.rs +++ b/hive-c0re/src/lifecycle/mod.rs @@ -269,9 +269,11 @@ fn validate(name: &str) -> Result<()> { /// would otherwise loop on `AddrInUse` forever; we surface the /// conflict here so spawn / rebuild fails loudly with an actionable /// message instead. -async fn port_collision(self_name: &str) -> Option { +async fn port_collision(self_name: &str) -> Result> { let port = agent_web_port(self_name); - let raw = list().await.unwrap_or_default(); + let raw = list() + .await + .context("listing containers for the port-collision check")?; for c in raw { let Some(other) = c.strip_prefix(AGENT_PREFIX) else { continue; @@ -280,10 +282,10 @@ async fn port_collision(self_name: &str) -> Option { continue; } if agent_web_port(other) == port && is_running(other).await { - return Some(other.to_owned()); + return Ok(Some(other.to_owned())); } } - None + Ok(None) } pub async fn spawn(name: &str, hive: &HiveEnv, paths: &AgentPaths) -> Result<()> { @@ -303,7 +305,7 @@ pub async fn spawn(name: &str, hive: &HiveEnv, paths: &AgentPaths) -> Result<()> /// create the container — that's `create_only` / the `Create` node. pub async fn provision_container(name: &str, hive: &HiveEnv, paths: &AgentPaths) -> Result<()> { validate(name)?; - if let Some(other) = port_collision(name).await { + if let Some(other) = port_collision(name).await? { bail!( "port {} is already taken by '{other}' — rename one of them and retry", agent_web_port(name) @@ -356,7 +358,7 @@ pub async fn create_container(name: &str, hive: &HiveEnv, paths: &AgentPaths) -> /// than the clear bail this replaces. pub async fn prepare_rebuild_dirs(name: &str, paths: &AgentPaths) -> Result<()> { validate(name)?; - if let Some(other) = port_collision(name).await { + if let Some(other) = port_collision(name).await? { bail!( "port {} is already taken by '{other}' — rename one of them and retry", agent_web_port(name) @@ -888,6 +890,16 @@ pub async fn list() -> Result> { .collect()) } +/// Logical agent names from a [`list`] result. An unreadable list stays an +/// error rather than becoming "no agents": a caller fanning work out over +/// every agent would otherwise report success after doing nothing. +pub fn agent_names(listed: Result>) -> Result> { + Ok(listed? + .into_iter() + .filter_map(|c| c.strip_prefix(AGENT_PREFIX).map(str::to_owned)) + .collect()) +} + /// Sync `/etc/tmpfiles.d/hyperhive-agents.conf` with the currently-known /// agent set (from `nixos-container list`). Strips the `h-` prefix to get /// logical names. Best-effort: errors are logged but never propagated — a diff --git a/hive-c0re/src/lifecycle/tests.rs b/hive-c0re/src/lifecycle/tests.rs index ce50093c..64fdfe48 100644 --- a/hive-c0re/src/lifecycle/tests.rs +++ b/hive-c0re/src/lifecycle/tests.rs @@ -377,3 +377,18 @@ fn failed_destroy_with_an_unreadable_list_fails() { .expect_err("an unreadable list must not pass for an absent container"); assert!(format!("{err:#}").contains("cannot confirm"), "{err:#}"); } + +#[test] +fn agent_names_strips_the_container_prefix() { + let names = agent_names(Ok(vec![format!("{AGENT_PREFIX}iris"), "other".to_owned()])) + .expect("a readable list"); + assert_eq!(names, vec!["iris"]); +} + +/// A failed list is not an empty hive: the error reaches the caller. +#[test] +fn agent_names_propagates_an_unreadable_list() { + let err = agent_names(Err(anyhow::anyhow!("connect to hive-priv socket"))) + .expect_err("an unreadable list must not read as zero agents"); + assert!(format!("{err:#}").contains("hive-priv"), "{err:#}"); +} diff --git a/hive-c0re/src/migrate.rs b/hive-c0re/src/migrate.rs index 92cb46a4..39b3d50c 100644 --- a/hive-c0re/src/migrate.rs +++ b/hive-c0re/src/migrate.rs @@ -63,7 +63,13 @@ pub async fn run(coord: &Arc) -> Result<()> { Err(e) => tracing::warn!(error = ?e, "clear stale meta lock failed"), } } - let names = enumerate_agents().await; + let names = match enumerate_agents().await { + Ok(names) => names, + Err(e) => { + tracing::warn!(error = ?e, "migration: container list failed; skipping per-agent phases"); + Vec::new() + } + }; tracing::info!(count = names.len(), "migration: scanning"); // Phase 0: move harness-owned files out of state/ into harness/. @@ -94,9 +100,15 @@ pub async fn run(coord: &Arc) -> Result<()> { // Phase 3: meta repo. tracing::debug!("migration: phase 3 (meta sync_agents)"); - let agents = lifecycle::agents_for_meta_listing() - .await - .unwrap_or_default(); + // `sync_agents` renders exactly the agents it is given, so an empty + // stand-in for an unreadable list would drop every agent from the flake. + let agents = match lifecycle::agents_for_meta_listing().await { + Ok(agents) => agents, + Err(e) => { + tracing::warn!(error = ?e, "migration: container list failed; skipping meta sync_agents"); + return Ok(()); + } + }; match tokio::time::timeout(GIT_TIMEOUT, meta::sync_agents(&coord.hive_env(), &agents)).await { Ok(Err(e)) => tracing::warn!(error = ?e, "migration: meta sync_agents failed"), Err(_) => { @@ -139,9 +151,9 @@ fn migrate_harness_files(name: &hive_types::Ident) { } } -async fn enumerate_agents() -> Vec { - let containers = lifecycle::list().await.unwrap_or_default(); - containers +async fn enumerate_agents() -> Result> { + let containers = lifecycle::list().await?; + Ok(containers .into_iter() .filter_map(|c| { let name = if c == MANAGER_CONTAINER { @@ -151,7 +163,7 @@ async fn enumerate_agents() -> Vec { }; hive_types::Ident::parse(name).ok() }) - .collect() + .collect()) } async fn migrate_applied_repo(name: &str) -> Result<()> { diff --git a/hive-c0re/src/server.rs b/hive-c0re/src/server.rs index 53f1bf10..f5faad0c 100644 --- a/hive-c0re/src/server.rs +++ b/hive-c0re/src/server.rs @@ -293,7 +293,7 @@ async fn handle_spawn(coord: &Arc, name: &str) -> Result) -> HostResponse { - let rows = crate::container_view::build_all(&coord.hive_env()) - .await - .into_iter() - .map(hive_sh4re::container::AgentStatusRow::from) - .collect(); - HostResponse::agent_statuses(rows) + match crate::container_view::build_all(&coord.hive_env()).await { + Ok(views) => HostResponse::agent_statuses( + views + .into_iter() + .map(hive_sh4re::container::AgentStatusRow::from) + .collect(), + ), + Err(e) => HostResponse::error(format!("listing containers: {e:#}")), + } } /// `SubscribeAgentStatus` — ack once, then push one single-row @@ -502,7 +505,7 @@ async fn handle_set_resource_limits( // in the JSON until the agent's next spawn or rebuild. let agent_dir = crate::paths::agent_runtime_dir(name.as_str()); let hive = coord.hive_env(); - let paths = Coordinator::agent_paths(name.as_str(), agent_dir); + let paths = Coordinator::agent_paths(name.as_str(), agent_dir)?; crate::lifecycle::write_dropins(name.as_str(), &hive, &paths).await?; let (cpu, mem) = crate::resource_limits::effective( diff --git a/hive-c0re/src/workers/auto_update.rs b/hive-c0re/src/workers/auto_update.rs index d2899ebc..f5422401 100644 --- a/hive-c0re/src/workers/auto_update.rs +++ b/hive-c0re/src/workers/auto_update.rs @@ -190,7 +190,7 @@ pub async fn ensure_root_agent(coord: &Arc) -> Result<()> { // ensure_agent_runtime_dir needed here. let runtime = crate::paths::agent_runtime_dir(MANAGER_NAME); let hive = coord.hive_env(); - let paths = Coordinator::agent_paths(MANAGER_NAME, runtime); + let paths = Coordinator::agent_paths(MANAGER_NAME, runtime)?; lifecycle::spawn(MANAGER_NAME, &hive, &paths).await?; seed_manager_tool_groups(); if let Err(e) = coord.power.set(MANAGER_NAME, crate::power::Wanted::Up) {