303 lines
12 KiB
Rust
303 lines
12 KiB
Rust
//! Startup auto-update: on `hive-c0re serve` boot, rebuild every known
|
|
//! container unconditionally. `nixos-container update` is a no-op at the
|
|
//! nix level when nothing changed (same store path), so the cost is low
|
|
//! and avoids rev-marker staleness (all agents always need an update pass
|
|
//! when any meta commit lands). See `docs/coordinator.md::Auto-update sweep`.
|
|
|
|
use std::path::{Path, PathBuf};
|
|
use std::sync::Arc;
|
|
|
|
use anyhow::{Context, Result};
|
|
|
|
use crate::coordinator::Coordinator;
|
|
use crate::lifecycle::{self, AGENT_PREFIX, MANAGER_NAME};
|
|
|
|
/// Marker file recording the hyperhive rev a sub-agent's container was last
|
|
/// built against. Sibling of `applied/<name>/` (rather than inside it) to
|
|
/// keep it out of the applied repo's git history. Uses a leading dot so a
|
|
/// glob over `applied/*` doesn't include it.
|
|
pub fn rev_marker_path(name: &str) -> PathBuf {
|
|
PathBuf::from(format!("/var/lib/hyperhive/applied/.{name}.hyperhive-rev"))
|
|
}
|
|
|
|
/// Resolve the current rev of `hyperhive_flake`. For a path on disk we
|
|
/// canonicalize (following symlinks) so a /etc/hyperhive → /nix/store/...
|
|
/// update yields a different string. For anything else we return None.
|
|
#[must_use]
|
|
pub fn current_flake_rev(hyperhive_flake: &str) -> Option<String> {
|
|
let path = Path::new(hyperhive_flake);
|
|
if !path.exists() {
|
|
return None;
|
|
}
|
|
std::fs::canonicalize(path)
|
|
.ok()
|
|
.map(|p| p.display().to_string())
|
|
}
|
|
|
|
/// Returns true when the applied repo has commits that have not yet been
|
|
/// deployed (i.e. the applied HEAD differs from the sha currently locked in
|
|
/// meta's flake.lock). This is the semantic the dashboard `needs_update` chip
|
|
/// conveys: "there is a config change ready to apply via rebuild."
|
|
#[must_use]
|
|
pub fn agent_config_pending(name: &str, deployed_sha: Option<&str>) -> bool {
|
|
let applied_head = std::process::Command::new("git")
|
|
.args([
|
|
"-C",
|
|
&format!("/var/lib/hyperhive/applied/{name}"),
|
|
"rev-parse",
|
|
"HEAD",
|
|
])
|
|
.output()
|
|
.ok()
|
|
.filter(|o| o.status.success())
|
|
.and_then(|o| String::from_utf8(o.stdout).ok())
|
|
.map(|s| s.trim().to_owned());
|
|
|
|
match (applied_head.as_deref(), deployed_sha) {
|
|
(Some(head), Some(sha)) => !head.starts_with(sha) && !sha.starts_with(head),
|
|
_ => false,
|
|
}
|
|
}
|
|
|
|
/// Rebuild one sub-agent and refresh its marker. Used by both the startup
|
|
/// scanner and the dashboard's manual "update" button so the two paths
|
|
/// can't diverge.
|
|
///
|
|
/// `queue_entry_id` is `Some(id)` when the rebuild was dispatched from
|
|
/// the `rebuild_queue` worker (lets the function annotate its phase via
|
|
/// `coord.set_queue_step`) and `None` when called directly (e.g. the
|
|
/// manager-migration nudge in `ensure_manager`).
|
|
pub async fn rebuild_agent(
|
|
coord: &Arc<Coordinator>,
|
|
name: &str,
|
|
current_rev: &str,
|
|
queue_entry_id: Option<u64>,
|
|
) -> Result<()> {
|
|
tracing::info!(%name, rev = %current_rev, "rebuild agent");
|
|
let agent_dir = coord
|
|
.ensure_runtime(name)
|
|
.with_context(|| format!("ensure_runtime {name}"))?;
|
|
let applied_dir = Coordinator::agent_applied_dir(name);
|
|
let claude_dir = Coordinator::agent_claude_dir(name);
|
|
let notes_dir = Coordinator::agent_notes_dir(name);
|
|
// Suppress crash_watch during the stop+start window inside
|
|
// lifecycle::rebuild. Dashboard rebuilds already do this via
|
|
// lifecycle_action; this catches the auto-update scan + any
|
|
// other direct caller.
|
|
let guard = coord.transient_guard(name, crate::coordinator::TransientKind::Rebuilding);
|
|
let result = lifecycle::rebuild(
|
|
name,
|
|
&coord.hyperhive_flake,
|
|
&agent_dir,
|
|
&applied_dir,
|
|
&claude_dir,
|
|
¬es_dir,
|
|
coord.dashboard_port,
|
|
&coord.operator_pronouns,
|
|
&coord.context_window_tokens,
|
|
&|step| coord.set_queue_step(queue_entry_id, step),
|
|
)
|
|
.await;
|
|
drop(guard);
|
|
match &result {
|
|
Ok(()) => {
|
|
if let Err(e) = std::fs::write(rev_marker_path(name), current_rev) {
|
|
tracing::warn!(%name, error = ?e, "write rev marker failed");
|
|
}
|
|
coord.notify_manager(&hive_sh4re::HelperEvent::Rebuilt {
|
|
agent: name.to_owned(),
|
|
ok: true,
|
|
note: None,
|
|
sha: None,
|
|
tag: None,
|
|
});
|
|
coord.set_queue_step(queue_entry_id, "forge sync");
|
|
// Run the full forge sync on every successful rebuild so
|
|
// the rebuild path is equivalent to the hive-c0re startup
|
|
// sweep: token, config-repo mirror, meta read access, and
|
|
// meta remote are all kept in sync. Recovers missing tokens
|
|
// (e.g. first-spawn seeding failed transiently) without
|
|
// requiring a full hive-c0re restart.
|
|
crate::forge::sync_agent(name, crate::forge::core_token().as_deref()).await;
|
|
// Mirror the matrix side of the startup sweep: if hive-matrix
|
|
// is present, ensure this agent has a registered user + token.
|
|
// Idempotent (skips if token file already exists). Keeps the
|
|
// rebuild path equivalent to the startup sweep for newly-spawned
|
|
// agents that missed ensure_all().
|
|
crate::matrix::sync_agent_standalone(name).await;
|
|
// Wake the agent on its next turn so claude sees a
|
|
// "you were rebuilt — check /state/ for notes, --continue
|
|
// session intact" hint. Covers dashboard rebuild, admin
|
|
// CLI rebuild, auto-update startup scan, and the
|
|
// dashboard's meta-input update path — all of which
|
|
// route through rebuild_agent.
|
|
coord.kick_agent(name, "container rebuilt");
|
|
// Container state (needs_update, deployed_sha) may have
|
|
// shifted — rescan so dashboards drop the "needs update"
|
|
// chip without waiting for the next /api/state poll.
|
|
coord.rescan_containers_and_emit().await;
|
|
// Lock bump → meta-inputs panel needs to re-render.
|
|
crate::dashboard::emit_meta_inputs_snapshot(coord);
|
|
}
|
|
Err(e) => {
|
|
coord.notify_manager(&hive_sh4re::HelperEvent::Rebuilt {
|
|
agent: name.to_owned(),
|
|
ok: false,
|
|
note: Some(format!("{e:#}")),
|
|
sha: None,
|
|
tag: None,
|
|
});
|
|
coord.rescan_containers_and_emit().await;
|
|
}
|
|
}
|
|
result
|
|
}
|
|
|
|
/// Auto-create the manager container on startup if it isn't already there.
|
|
/// hive-c0re manages `root` end-to-end: operators no
|
|
/// longer declare `containers.root` in their host NixOS config. Bypasses
|
|
/// the approval queue — manager is required infrastructure. Idempotent.
|
|
pub async fn ensure_manager(coord: &Arc<Coordinator>) -> Result<()> {
|
|
let existing = lifecycle::list().await.unwrap_or_default();
|
|
let current_rev = current_flake_rev(&coord.hyperhive_flake);
|
|
if existing.iter().any(|c| c == MANAGER_NAME) {
|
|
// 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
|
|
// time, with a wrong systemd unit and port.
|
|
let applied_flake = Coordinator::agent_applied_dir(MANAGER_NAME).join("flake.nix");
|
|
if !applied_flake.exists()
|
|
&& let Some(rev) = current_rev.as_ref()
|
|
{
|
|
tracing::warn!(
|
|
"manager container exists but no applied flake — forcing rebuild to migrate"
|
|
);
|
|
let coord_clone = coord.clone();
|
|
if let Err(e) = rebuild_agent(&coord_clone, MANAGER_NAME, rev.as_str(), None).await {
|
|
tracing::warn!(error = ?e, "manager migration rebuild failed");
|
|
}
|
|
} else {
|
|
tracing::debug!("manager container already present");
|
|
}
|
|
return Ok(());
|
|
}
|
|
tracing::info!("manager container missing — spawning");
|
|
let runtime = coord.ensure_runtime(MANAGER_NAME)?;
|
|
let proposed = Coordinator::agent_proposed_dir(MANAGER_NAME);
|
|
let applied = Coordinator::agent_applied_dir(MANAGER_NAME);
|
|
let claude_dir = Coordinator::agent_claude_dir(MANAGER_NAME);
|
|
let notes_dir = Coordinator::agent_notes_dir(MANAGER_NAME);
|
|
lifecycle::spawn(
|
|
MANAGER_NAME,
|
|
&coord.hyperhive_flake,
|
|
&runtime,
|
|
&proposed,
|
|
&applied,
|
|
&claude_dir,
|
|
¬es_dir,
|
|
coord.dashboard_port,
|
|
&coord.operator_pronouns,
|
|
&coord.context_window_tokens,
|
|
)
|
|
.await?;
|
|
if let Some(rev) = current_rev {
|
|
let _ = std::fs::write(rev_marker_path(MANAGER_NAME), &rev);
|
|
}
|
|
Ok(())
|
|
}
|
|
|
|
/// Sort `names` in-place so parents precede their children in the topology.
|
|
/// Uses BFS from root agents (depth 0). Agents absent from `topo` sort last,
|
|
/// alphabetically within their tier. Stable within each depth tier.
|
|
pub fn topology_sort(names: &mut Vec<String>, topo: &std::collections::BTreeMap<String, Option<String>>) {
|
|
use std::collections::{HashMap, VecDeque};
|
|
// Build depth map using owned clones so the borrow on `names` is released
|
|
// before the sort_by mutable borrow.
|
|
let name_set: Vec<String> = names.clone();
|
|
let mut depth: HashMap<String, usize> = HashMap::new();
|
|
let mut queue: VecDeque<String> = VecDeque::new();
|
|
// Seed roots: entries with no parent, or names not present in topo at all.
|
|
for name in &name_set {
|
|
if topo.get(name).map_or(true, |p| p.is_none()) {
|
|
depth.insert(name.clone(), 0);
|
|
queue.push_back(name.clone());
|
|
}
|
|
}
|
|
// BFS to assign depths to children.
|
|
while let Some(parent) = queue.pop_front() {
|
|
let d = depth[&parent] + 1;
|
|
for name in &name_set {
|
|
let is_child = topo.get(name).and_then(|p| p.as_deref()) == Some(parent.as_str());
|
|
if is_child && !depth.contains_key(name) {
|
|
depth.insert(name.clone(), d);
|
|
queue.push_back(name.clone());
|
|
}
|
|
}
|
|
}
|
|
names.sort_by(|a, b| {
|
|
let da = depth.get(a).copied().unwrap_or(usize::MAX);
|
|
let db = depth.get(b).copied().unwrap_or(usize::MAX);
|
|
da.cmp(&db).then(a.cmp(b))
|
|
});
|
|
}
|
|
|
|
/// Rebuild every container on startup. Enqueues a `StartupSweep` parent
|
|
/// entry (agent = `"hyperhive"`) followed by per-agent `Rebuild` children
|
|
/// linked via `parent_id`. The dashboard renders them nested so the operator
|
|
/// can see at a glance "boot N agents, here is each rebuild's status".
|
|
/// Returns Ok even if some rebuilds failed.
|
|
pub async fn run(coord: Arc<Coordinator>) -> Result<()> {
|
|
let containers = match lifecycle::list().await {
|
|
Ok(c) => c,
|
|
Err(e) => {
|
|
tracing::warn!(error = ?e, "auto-update: nixos-container list failed");
|
|
return Ok(());
|
|
}
|
|
};
|
|
|
|
// Enqueue the parent sweep entry. The worker processes it trivially
|
|
// (no-op dispatch) so it completes quickly; its purpose is to give the
|
|
// dashboard a "why" header for the per-agent child rebuilds below.
|
|
let sweep_id = coord.rebuild_queue.enqueue(
|
|
crate::rebuild_queue::QueueKind::StartupSweep,
|
|
"hyperhive".to_owned(),
|
|
crate::rebuild_queue::QueueSource::AutoUpdate,
|
|
format!("startup sweep ({} containers)", containers.len()),
|
|
None,
|
|
);
|
|
|
|
tracing::info!(
|
|
agents = containers.len(),
|
|
sweep_id,
|
|
"auto-update: queueing all on startup"
|
|
);
|
|
|
|
// Resolve container names to logical agent names, then sort by
|
|
// topology depth so parents are always rebuilt before their
|
|
// children. Root agents (depth 0) go first; agents absent from
|
|
// the topology file sort last (stable, alphabetical within tier).
|
|
let mut logical_names: Vec<String> = containers
|
|
.iter()
|
|
.filter_map(|c| {
|
|
if c == MANAGER_NAME {
|
|
Some(MANAGER_NAME.to_owned())
|
|
} else {
|
|
c.strip_prefix(AGENT_PREFIX).map(str::to_owned)
|
|
}
|
|
})
|
|
.collect();
|
|
let topo = crate::topology::read();
|
|
topology_sort(&mut logical_names, &topo);
|
|
for name in logical_names {
|
|
coord.rebuild_queue.enqueue(
|
|
crate::rebuild_queue::QueueKind::Rebuild,
|
|
name,
|
|
crate::rebuild_queue::QueueSource::StartupSweep,
|
|
"startup sweep".to_owned(),
|
|
Some(sweep_id),
|
|
);
|
|
}
|
|
coord.emit_rebuild_queue_snapshot();
|
|
Ok(())
|
|
}
|
|
|