workers: skip a cycle instead of spawning/alerting on an unreadable container list

crash_watch's 10s poll and auto_update's ensure_root_agent both read
lifecycle::list().await.unwrap_or_default(), which turned a failed read
into 'zero containers'. In crash_watch that made every previously-running
agent look like it crashed simultaneously (prev.difference(current) over
an empty current), and left prev empty for the next cycle too, so a
second wave of false 'agent logged in' / 'agent needs login' events fired
against the next successful read. In ensure_root_agent it read as
'manager container missing' and called lifecycle::spawn on a manager
that might already exist.

Both sites now treat a list error as its own outcome: log it at warn and
skip the cycle's decision entirely. crash_watch leaves prev exactly as
the last good read produced it. ensure_root_agent attempts no spawn.

Factors each site's decision into a pure helper (plan_cycle /
plan_root_agent) matching the check_not_live / confirm_gone_after_failed_destroy
pattern, with unit tests for the error case, a control for the readable
case, and (for crash_watch) an invert-proof run locally against the old
unwrap_or_default logic before reverting.
This commit is contained in:
atlas 2026-09-24 12:50:07 +02:00 • committed by mara
commit f62ec349d5
2 changed files with 172 additions and 43 deletions

View file

@ -81,6 +81,36 @@ fn ruthless() -> bool {
}
}
/// What `ensure_root_agent` does with a `lifecycle::list()` result, before
/// it even gets to check whether the manager is present.
#[derive(Debug, PartialEq, Eq)]
enum RootAgentPlan {
/// The manager is in this readable list.
Present,
/// The manager is absent from this readable list: spawn it.
Absent,
/// The list could not be read. That is not "manager absent" — it
/// proves nothing either way, so the auto-spawn decision is skipped
/// for this attempt rather than risking a spawn over a manager that's
/// actually there.
Unreadable,
}
/// Turn a `lifecycle::list()` result into `ensure_root_agent`'s plan.
fn plan_root_agent(list_result: anyhow::Result<Vec<String>>) -> RootAgentPlan {
match list_result {
Ok(names)
if names
.iter()
.any(|c| c.strip_prefix(AGENT_PREFIX) == Some(MANAGER_NAME)) =>
{
RootAgentPlan::Present
}
Ok(_) => RootAgentPlan::Absent,
Err(_) => RootAgentPlan::Unreadable,
}
}
/// Auto-create the manager container on startup if it isn't already there.
/// hive-c0re manages the manager end-to-end: operators no longer declare
/// `containers.h-ruth` in their host NixOS config. Bypasses the approval
@ -98,12 +128,25 @@ pub async fn ensure_root_agent(coord: &Arc<Coordinator>) -> Result<()> {
// container already exists this is the only run that can still seed the
// grant, and it has to land before her next rebuild bakes the binds.
seed_manager_capabilities();
let existing = lifecycle::list().await.unwrap_or_default();
let list_result = lifecycle::list().await;
if let Err(e) = &list_result {
tracing::warn!(
error = ?e,
"manager container list unreadable — skipping auto-spawn check for this attempt"
);
}
let plan = plan_root_agent(list_result);
let current_rev = current_flake_rev(&coord.hyperhive_flake);
if existing
.iter()
.any(|c| c.strip_prefix(AGENT_PREFIX) == Some(MANAGER_NAME))
{
if plan == RootAgentPlan::Unreadable {
// An unreadable list is not "manager absent" — spawning on it would
// both waste `provision_container`'s work and hit `nixos-container
// create`'s "already exists" failure if the manager is actually
// there. Skip the decision this attempt; there is no periodic
// retry for this boot-time call, so the next chance is the next
// `hive-c0re` restart.
return Ok(());
}
if plan == RootAgentPlan::Present {
// Container exists already. If it predates the unified lifecycle
// (no applied flake on disk) we must rebuild — otherwise it's
// running whatever the host-declarative config was at create
@ -503,9 +546,38 @@ fn submit_boot_tree(
#[cfg(test)]
mod tests {
use super::{BootAction, boot_action, manager_seed_caps, should_seed_manager_caps};
use super::{
BootAction, MANAGER_NAME, RootAgentPlan, boot_action, manager_seed_caps, plan_root_agent,
should_seed_manager_caps,
};
use crate::lifecycle::AGENT_PREFIX;
use crate::power::Wanted;
// -----------------------------------------------------------------------
// `ensure_root_agent`'s spawn-or-not decision. An unreadable list must
// never be read as "manager absent" — see `plan_root_agent`'s doc.
/// The regression this fix exists for: a read failure must not spawn.
#[test]
fn unreadable_list_skips_without_spawning() {
let result: anyhow::Result<Vec<String>> =
Err(anyhow::anyhow!("connect to hive-priv socket"));
assert_eq!(plan_root_agent(result), RootAgentPlan::Unreadable);
}
/// Control: a genuinely empty, readable list still plans a spawn.
#[test]
fn readable_list_missing_manager_spawns() {
let result: anyhow::Result<Vec<String>> = Ok(Vec::new());
assert_eq!(plan_root_agent(result), RootAgentPlan::Absent);
}
#[test]
fn readable_list_with_manager_present_is_a_noop() {
let result: anyhow::Result<Vec<String>> = Ok(vec![format!("{AGENT_PREFIX}{MANAGER_NAME}")]);
assert_eq!(plan_root_agent(result), RootAgentPlan::Present);
}
// -----------------------------------------------------------------------
// Root-agent capability seed. `capabilities_path()` resolves under
// `paths::meta_root()`, which is hardcoded to `/var/lib/hyperhive` with no

View file

@ -28,46 +28,57 @@ pub fn spawn(coord: Arc<Coordinator>) {
let mut prev_sub_agents: HashSet<String> = HashSet::new();
let mut seeded = false;
loop {
let raw = lifecycle::list().await.unwrap_or_default();
let mut current_running = HashSet::new();
let mut current_logged_in = HashSet::new();
let mut sub_agents: Vec<String> = Vec::new();
for c in &raw {
let Some(logical) = c.strip_prefix(AGENT_PREFIX) else {
continue;
};
let logical = logical.to_owned();
sub_agents.push(logical.clone());
if lifecycle::is_running(&logical).await {
current_running.insert(logical.clone());
}
if hive_types::Ident::parse(&logical)
.is_ok_and(|id| claude_has_session(&Coordinator::agent_claude_dir(&id)))
{
current_logged_in.insert(logical.clone());
}
let list_result = lifecycle::list().await;
if let Err(e) = &list_result {
tracing::warn!(
error = ?e,
"crash watcher: container list unreadable — skipping this cycle"
);
}
if let CyclePlan::Process(raw) = plan_cycle(list_result) {
let mut current_running = HashSet::new();
let mut current_logged_in = HashSet::new();
let mut sub_agents: Vec<String> = Vec::new();
for c in &raw {
let Some(logical) = c.strip_prefix(AGENT_PREFIX) else {
continue;
};
let logical = logical.to_owned();
sub_agents.push(logical.clone());
if lifecycle::is_running(&logical).await {
current_running.insert(logical.clone());
}
if hive_types::Ident::parse(&logical)
.is_ok_and(|id| claude_has_session(&Coordinator::agent_claude_dir(&id)))
{
current_logged_in.insert(logical.clone());
}
}
if seeded {
emit_crash_transitions(&coord, &prev_running, &current_running);
emit_login_transitions(
&prev_logged_in,
&current_logged_in,
&sub_agents,
&prev_sub_agents,
)
.await;
if seeded {
emit_crash_transitions(&coord, &prev_running, &current_running);
emit_login_transitions(
&prev_logged_in,
&current_logged_in,
&sub_agents,
&prev_sub_agents,
)
.await;
}
// Periodic container rescan — catches state flips that
// happen outside our mutation surface (operator runs
// `nixos-container stop` over ssh, agent logs in via its
// own web UI, etc.) so the dashboard converges within one
// POLL_INTERVAL. Idempotent + cheap when nothing changed.
coord.rescan_containers_and_emit().await;
prev_running = current_running;
prev_logged_in = current_logged_in;
prev_sub_agents = sub_agents.into_iter().collect();
seeded = true;
}
// Periodic container rescan — catches state flips that
// happen outside our mutation surface (operator runs
// `nixos-container stop` over ssh, agent logs in via its
// own web UI, etc.) so the dashboard converges within one
// POLL_INTERVAL. Idempotent + cheap when nothing changed.
coord.rescan_containers_and_emit().await;
prev_running = current_running;
prev_logged_in = current_logged_in;
prev_sub_agents = sub_agents.into_iter().collect();
seeded = true;
// On `CyclePlan::Skip`, `prev_*` and `seeded` are left exactly as
// the last good cycle produced them, so the next successful read
// diffs against real prior state instead of an empty one.
tokio::select! {
() = tokio::time::sleep(POLL_INTERVAL) => {}
@ -80,6 +91,29 @@ pub fn spawn(coord: Arc<Coordinator>) {
});
}
/// What this poll cycle does with `lifecycle::list()`'s result.
enum CyclePlan {
/// The list read; diff `Vec` against `prev` as usual.
Process(Vec<String>),
/// The read failed: skip the cycle outright.
Skip,
}
/// Turn a `lifecycle::list()` result into this poll's plan. An unreadable
/// list is not "zero containers" — treating it as such made every
/// previously-running agent land in `prev.difference(current)`, i.e. look
/// like it crashed simultaneously with every other agent. So a read failure
/// plans `Skip`: no diff against `prev`, no crash or login events this
/// cycle, and (at the call site) `prev` stays exactly what the last good
/// read left it, so the next successful read compares against real state
/// rather than an empty one.
fn plan_cycle(list_result: anyhow::Result<Vec<String>>) -> CyclePlan {
match list_result {
Ok(raw) => CyclePlan::Process(raw),
Err(_) => CyclePlan::Skip,
}
}
fn emit_crash_transitions(coord: &Coordinator, prev: &HashSet<String>, current: &HashSet<String>) {
let transients = coord.transient_snapshot();
// Operator actions whose RAII guard already cleared but only just;
@ -177,6 +211,29 @@ async fn emit_login_transitions(
mod tests {
use super::*;
/// The regression this whole fix exists for: a read failure must not
/// reach `emit_crash_transitions` at all, so `prev` is left untouched
/// for the next cycle to diff against.
#[test]
fn list_error_skips_the_cycle() {
let result: anyhow::Result<Vec<String>> =
Err(anyhow::anyhow!("connect to hive-priv socket"));
assert!(matches!(plan_cycle(result), CyclePlan::Skip));
}
/// Control: a real drop from running to absent is a *readable*, empty
/// list — it must still be processed (and so still reach
/// `emit_crash_transitions`, which fires the crash event), unlike a
/// read failure.
#[test]
fn readable_empty_list_still_processes() {
let result: anyhow::Result<Vec<String>> = Ok(Vec::new());
match plan_cycle(result) {
CyclePlan::Process(raw) => assert!(raw.is_empty()),
CyclePlan::Skip => panic!("a readable empty list must not be treated as unreadable"),
}
}
#[test]
fn deliberate_when_either_source_says_so() {
// Active guard, and the race repro: a lifecycle action completes and