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
This commit is contained in:
parent
ee25b7de20
commit
7ac6819652
10 changed files with 153 additions and 67 deletions
|
|
@ -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`,
|
/// fallback onto the hive-wide `agent_cpu_quota` / `agent_memory_max`,
|
||||||
/// and those live on [`crate::coordinator::HiveEnv`], not on disk. Both callers already
|
/// 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.
|
/// hold a `Coordinator`, so this is a parameter rather than a global.
|
||||||
pub async fn build_all(hive: &crate::coordinator::HiveEnv) -> Vec<ContainerView> {
|
///
|
||||||
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<Vec<ContainerView>> {
|
||||||
|
let raw = lifecycle::list().await?;
|
||||||
let locked = read_meta_locked_revs();
|
let locked = read_meta_locked_revs();
|
||||||
// Read once per scan rather than re-read for every container on every
|
// 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;
|
// 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<ContainerView>
|
||||||
memory_max,
|
memory_max,
|
||||||
});
|
});
|
||||||
}
|
}
|
||||||
out
|
Ok(out)
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Host-side mirror of `hive_agent::login::has_session`. Returns true
|
/// Host-side mirror of `hive_agent::login::has_session`. Returns true
|
||||||
|
|
|
||||||
|
|
@ -558,21 +558,22 @@ impl Coordinator {
|
||||||
/// from `crate::paths::agent_runtime_dir(name)` (pure path) or from
|
/// from `crate::paths::agent_runtime_dir(name)` (pure path) or from
|
||||||
/// `lifecycle::ensure_agent_runtime_dir(name)` when the dir must be
|
/// `lifecycle::ensure_agent_runtime_dir(name)` when the dir must be
|
||||||
/// created. All other paths are derived statically from `name`.
|
/// created. All other paths are derived statically from `name`.
|
||||||
#[must_use]
|
///
|
||||||
pub fn agent_paths(name: &str, agent_dir: PathBuf) -> AgentPaths {
|
/// # Errors
|
||||||
// `name` is validated upstream (spawn-approval / enqueue gate / the
|
///
|
||||||
// MANAGER_NAME const), so an invalid ident here is a construction
|
/// `name` is not a valid [`hive_types::Ident`]. Job payloads carry agent
|
||||||
// bug. This is the step-3 boundary between the Ident-threaded path
|
/// names as plain strings, one source being the swarm's published wanted
|
||||||
// builders and the job_queue layer (threaded post hive-jobq cutover).
|
/// state, so this is an error for the caller to report, not an invariant.
|
||||||
|
pub fn agent_paths(name: &str, agent_dir: PathBuf) -> Result<AgentPaths> {
|
||||||
let name = hive_types::Ident::parse(name)
|
let name = hive_types::Ident::parse(name)
|
||||||
.expect("agent_paths: name must be a valid ident (validated at spawn/enqueue)");
|
.map_err(|e| anyhow::anyhow!("invalid agent name {name:?}: {e}"))?;
|
||||||
AgentPaths {
|
Ok(AgentPaths {
|
||||||
agent: agent_dir,
|
agent: agent_dir,
|
||||||
proposed: Self::agent_proposed_dir(&name),
|
proposed: Self::agent_proposed_dir(&name),
|
||||||
applied: crate::paths::applied_dir(name.as_str()),
|
applied: crate::paths::applied_dir(name.as_str()),
|
||||||
claude: Self::agent_claude_dir(&name),
|
claude: Self::agent_claude_dir(&name),
|
||||||
notes: Self::agent_notes_dir(&name),
|
notes: Self::agent_notes_dir(&name),
|
||||||
}
|
})
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Emit a `RebuildQueueChanged` tick. Called from the queue mutation
|
/// Emit a `RebuildQueueChanged` tick. Called from the queue mutation
|
||||||
|
|
@ -845,7 +846,15 @@ impl Coordinator {
|
||||||
/// Cheap when nothing changed (one `nixos-container list` + a
|
/// Cheap when nothing changed (one `nixos-container list` + a
|
||||||
/// `HashMap` diff + zero emits).
|
/// `HashMap` diff + zero emits).
|
||||||
pub async fn rescan_containers_and_emit(self: &Arc<Self>) {
|
pub async fn rescan_containers_and_emit(self: &Arc<Self>) {
|
||||||
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 last = self.last_containers.lock().await;
|
||||||
let mut changed_or_new = Vec::new();
|
let mut changed_or_new = Vec::new();
|
||||||
let mut removed = 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"));
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
|
||||||
|
|
@ -59,6 +59,7 @@ pub(super) struct JournalQuery {
|
||||||
(status = 200, description = "journal text", body = String, content_type = "text/plain"),
|
(status = 200, description = "journal text", body = String, content_type = "text/plain"),
|
||||||
(status = 400, description = "bad agent name or unknown unit"),
|
(status = 400, description = "bad agent name or unknown unit"),
|
||||||
(status = 404, description = "no such managed container"),
|
(status = 404, description = "no such managed container"),
|
||||||
|
(status = 500, description = "container list unreadable"),
|
||||||
),
|
),
|
||||||
tag = "journal"
|
tag = "journal"
|
||||||
)]
|
)]
|
||||||
|
|
@ -97,7 +98,14 @@ pub(super) async fn get_journal(
|
||||||
// containers so we don't shell out with arbitrary input.
|
// containers so we don't shell out with arbitrary input.
|
||||||
let container = strip_container_prefix(name.as_str());
|
let container = strip_container_prefix(name.as_str());
|
||||||
let prefixed = format!("{}{container}", lifecycle::AGENT_PREFIX);
|
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) {
|
if !live.iter().any(|c| c == &prefixed) {
|
||||||
return Err(ProblemDetails::from_status_code(StatusCode::NOT_FOUND)
|
return Err(ProblemDetails::from_status_code(StatusCode::NOT_FOUND)
|
||||||
.with_detail(format!("journal: no managed container {prefixed:?}")));
|
.with_detail(format!("journal: no managed container {prefixed:?}")));
|
||||||
|
|
|
||||||
|
|
@ -392,7 +392,10 @@ pub(super) async fn post_resource_limits(
|
||||||
}
|
}
|
||||||
let agent_dir = crate::paths::agent_runtime_dir(ident.as_str());
|
let agent_dir = crate::paths::agent_runtime_dir(ident.as_str());
|
||||||
let hive = state.coord.hive_env();
|
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 {
|
if let Err(e) = crate::lifecycle::write_dropins(ident.as_str(), &hive, &paths).await {
|
||||||
return error_response(&format!("write_dropins {logical}: {e:#}"));
|
return error_response(&format!("write_dropins {logical}: {e:#}"));
|
||||||
}
|
}
|
||||||
|
|
@ -405,18 +408,18 @@ pub(super) async fn post_resource_limits(
|
||||||
#[utoipa::path(
|
#[utoipa::path(
|
||||||
post,
|
post,
|
||||||
path = "/api/update-all",
|
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"
|
tag = "lifecycle_ops"
|
||||||
)]
|
)]
|
||||||
pub(super) async fn post_update_all(State(state): State<AppState>) -> Response {
|
pub(super) async fn post_update_all(State(state): State<AppState>) -> Response {
|
||||||
let containers = lifecycle::list().await.unwrap_or_default();
|
let agents = match lifecycle::agent_names(lifecycle::list().await) {
|
||||||
for container in containers {
|
Ok(agents) => agents,
|
||||||
let Some(logical) = container
|
Err(e) => return error_response(&format!("update-all: listing containers: {e:#}")),
|
||||||
.strip_prefix(lifecycle::AGENT_PREFIX)
|
|
||||||
.map(str::to_owned)
|
|
||||||
else {
|
|
||||||
continue;
|
|
||||||
};
|
};
|
||||||
|
for logical in agents {
|
||||||
if let Err(e) = state.coord.job_queue.insert_job(|b| {
|
if let Err(e) = state.coord.job_queue.insert_job(|b| {
|
||||||
crate::job_queue::templates::rebuild(b, &logical, true);
|
crate::job_queue::templates::rebuild(b, &logical, true);
|
||||||
Vec::new()
|
Vec::new()
|
||||||
|
|
|
||||||
|
|
@ -399,7 +399,7 @@ async fn run_meta_sync(coord: &Arc<Coordinator>, name: &str, relock: bool) -> Re
|
||||||
// listener (event-driven: registered at start/create).
|
// listener (event-driven: registered at start/create).
|
||||||
let agent_dir = crate::paths::agent_runtime_dir(name);
|
let agent_dir = crate::paths::agent_runtime_dir(name);
|
||||||
let hive = coord.hive_env();
|
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?;
|
crate::lifecycle::prepare_rebuild_dirs(name, &paths).await?;
|
||||||
// Idempotent meta sync so a manual rebuild can also recover from a
|
// Idempotent meta sync so a manual rebuild can also recover from a
|
||||||
// divergent meta repo; then bump just this agent's input. `relock =
|
// divergent meta repo; then bump just this agent's input. `relock =
|
||||||
|
|
@ -442,7 +442,7 @@ async fn run_swap(coord: &Arc<Coordinator>, name: &str, id: NodeId) -> Result<()
|
||||||
// and listener were created earlier. Pure path accessor suffices.
|
// and listener were created earlier. Pure path accessor suffices.
|
||||||
let agent_dir = crate::paths::agent_runtime_dir(name);
|
let agent_dir = crate::paths::agent_runtime_dir(name);
|
||||||
let hive = coord.hive_env();
|
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;
|
let result = crate::lifecycle::swap_update(name, &hive, &paths, Some(id.get())).await;
|
||||||
// On success the Ok-only bookkeeping tail (rev marker, forge/matrix
|
// On success the Ok-only bookkeeping tail (rev marker, forge/matrix
|
||||||
// sync, kick, rescan, snapshot) runs in the sibling `RebuildBookkeeping` node,
|
// sync, kick, rescan, snapshot) runs in the sibling `RebuildBookkeeping` node,
|
||||||
|
|
@ -492,7 +492,7 @@ async fn run_rebuild_bookkeeping(coord: &Arc<Coordinator>, name: &str) -> Result
|
||||||
async fn run_provision(coord: &Arc<Coordinator>, name: &str) -> Result<()> {
|
async fn run_provision(coord: &Arc<Coordinator>, name: &str) -> Result<()> {
|
||||||
let agent_dir = crate::paths::agent_runtime_dir(name);
|
let agent_dir = crate::paths::agent_runtime_dir(name);
|
||||||
let hive = coord.hive_env();
|
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?;
|
crate::lifecycle::provision_container(name, &hive, &paths).await?;
|
||||||
Ok(())
|
Ok(())
|
||||||
}
|
}
|
||||||
|
|
@ -536,13 +536,15 @@ async fn run_meta_lock(
|
||||||
return Ok(fanout.unwrap_or_default());
|
return Ok(fanout.unwrap_or_default());
|
||||||
}
|
}
|
||||||
let _progress = coord.meta_update_guard();
|
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?;
|
crate::meta::lock_update(inputs).await?;
|
||||||
// Lock file changed — meta-inputs panel re-renders.
|
// Lock file changed — meta-inputs panel re-renders.
|
||||||
crate::dashboard::emit_meta_inputs_snapshot(coord);
|
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
|
// 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
|
// 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
|
// it up. Without this, a meta-input bump rebuilds every affected agent
|
||||||
|
|
@ -627,7 +629,7 @@ async fn run_start(coord: &Arc<Coordinator>, name: &str) -> Result<()> {
|
||||||
// omitting this becomes a compile error.
|
// omitting this becomes a compile error.
|
||||||
let agent_dir = crate::paths::agent_runtime_dir(name);
|
let agent_dir = crate::paths::agent_runtime_dir(name);
|
||||||
let hive = coord.hive_env();
|
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?;
|
let token = crate::lifecycle::converge_start_preamble(name, &hive, &paths).await?;
|
||||||
crate::lifecycle::start_with_fallback(token).await?;
|
crate::lifecycle::start_with_fallback(token).await?;
|
||||||
// Bind the MCP listener immediately after starting the container.
|
// Bind the MCP listener immediately after starting the container.
|
||||||
|
|
@ -764,7 +766,7 @@ async fn run_write_dropin(coord: &Arc<Coordinator>, name: &str) -> Result<()> {
|
||||||
// on the upstream Prebuild/Start node).
|
// on the upstream Prebuild/Start node).
|
||||||
let agent_dir = crate::paths::agent_runtime_dir(name);
|
let agent_dir = crate::paths::agent_runtime_dir(name);
|
||||||
let hive = coord.hive_env();
|
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?;
|
crate::lifecycle::write_dropins(name, &hive, &paths).await?;
|
||||||
Ok(())
|
Ok(())
|
||||||
}
|
}
|
||||||
|
|
@ -870,26 +872,23 @@ async fn run_deploy_tail(
|
||||||
/// cascade agent's name. The `touched_hyperhive` branch's names need no
|
/// cascade agent's name. The `touched_hyperhive` branch's names need no
|
||||||
/// such filter: they come from `lifecycle::list()`, already real
|
/// such filter: they come from `lifecycle::list()`, already real
|
||||||
/// container names by construction, not parsed out of caller input.
|
/// container names by construction, not parsed out of caller input.
|
||||||
pub async fn meta_update_cascade_agents(inputs: &[String]) -> Vec<String> {
|
///
|
||||||
|
/// # Errors
|
||||||
|
///
|
||||||
|
/// The every-container branch could not list the containers.
|
||||||
|
pub async fn meta_update_cascade_agents(inputs: &[String]) -> Result<Vec<String>> {
|
||||||
let touched_hyperhive = inputs
|
let touched_hyperhive = inputs
|
||||||
.iter()
|
.iter()
|
||||||
.any(|i| i == "hyperhive" || i.starts_with("hyperhive/"));
|
.any(|i| i == "hyperhive" || i.starts_with("hyperhive/"));
|
||||||
let touched_agents = parse_agent_input_names(inputs);
|
let touched_agents = parse_agent_input_names(inputs);
|
||||||
let mut names = if touched_hyperhive || inputs.is_empty() {
|
let mut names = if touched_hyperhive || inputs.is_empty() {
|
||||||
crate::lifecycle::list()
|
crate::lifecycle::agent_names(crate::lifecycle::list().await)
|
||||||
.await
|
.context("listing containers for the meta-update cascade")?
|
||||||
.unwrap_or_default()
|
|
||||||
.into_iter()
|
|
||||||
.filter_map(|c| {
|
|
||||||
c.strip_prefix(crate::lifecycle::AGENT_PREFIX)
|
|
||||||
.map(str::to_owned)
|
|
||||||
})
|
|
||||||
.collect()
|
|
||||||
} else {
|
} else {
|
||||||
touched_agents
|
touched_agents
|
||||||
};
|
};
|
||||||
names.sort();
|
names.sort();
|
||||||
names
|
Ok(names)
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Parse `agent-<name>` input strings into validated agent names. Pure and
|
/// Parse `agent-<name>` input strings into validated agent names. Pure and
|
||||||
|
|
|
||||||
|
|
@ -269,9 +269,11 @@ fn validate(name: &str) -> Result<()> {
|
||||||
/// would otherwise loop on `AddrInUse` forever; we surface the
|
/// would otherwise loop on `AddrInUse` forever; we surface the
|
||||||
/// conflict here so spawn / rebuild fails loudly with an actionable
|
/// conflict here so spawn / rebuild fails loudly with an actionable
|
||||||
/// message instead.
|
/// message instead.
|
||||||
async fn port_collision(self_name: &str) -> Option<String> {
|
async fn port_collision(self_name: &str) -> Result<Option<String>> {
|
||||||
let port = agent_web_port(self_name);
|
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 {
|
for c in raw {
|
||||||
let Some(other) = c.strip_prefix(AGENT_PREFIX) else {
|
let Some(other) = c.strip_prefix(AGENT_PREFIX) else {
|
||||||
continue;
|
continue;
|
||||||
|
|
@ -280,10 +282,10 @@ async fn port_collision(self_name: &str) -> Option<String> {
|
||||||
continue;
|
continue;
|
||||||
}
|
}
|
||||||
if agent_web_port(other) == port && is_running(other).await {
|
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<()> {
|
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.
|
/// create the container — that's `create_only` / the `Create` node.
|
||||||
pub async fn provision_container(name: &str, hive: &HiveEnv, paths: &AgentPaths) -> Result<()> {
|
pub async fn provision_container(name: &str, hive: &HiveEnv, paths: &AgentPaths) -> Result<()> {
|
||||||
validate(name)?;
|
validate(name)?;
|
||||||
if let Some(other) = port_collision(name).await {
|
if let Some(other) = port_collision(name).await? {
|
||||||
bail!(
|
bail!(
|
||||||
"port {} is already taken by '{other}' — rename one of them and retry",
|
"port {} is already taken by '{other}' — rename one of them and retry",
|
||||||
agent_web_port(name)
|
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.
|
/// than the clear bail this replaces.
|
||||||
pub async fn prepare_rebuild_dirs(name: &str, paths: &AgentPaths) -> Result<()> {
|
pub async fn prepare_rebuild_dirs(name: &str, paths: &AgentPaths) -> Result<()> {
|
||||||
validate(name)?;
|
validate(name)?;
|
||||||
if let Some(other) = port_collision(name).await {
|
if let Some(other) = port_collision(name).await? {
|
||||||
bail!(
|
bail!(
|
||||||
"port {} is already taken by '{other}' — rename one of them and retry",
|
"port {} is already taken by '{other}' — rename one of them and retry",
|
||||||
agent_web_port(name)
|
agent_web_port(name)
|
||||||
|
|
@ -888,6 +890,16 @@ pub async fn list() -> Result<Vec<String>> {
|
||||||
.collect())
|
.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<Vec<String>>) -> Result<Vec<String>> {
|
||||||
|
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
|
/// Sync `/etc/tmpfiles.d/hyperhive-agents.conf` with the currently-known
|
||||||
/// agent set (from `nixos-container list`). Strips the `h-` prefix to get
|
/// agent set (from `nixos-container list`). Strips the `h-` prefix to get
|
||||||
/// logical names. Best-effort: errors are logged but never propagated — a
|
/// logical names. Best-effort: errors are logged but never propagated — a
|
||||||
|
|
|
||||||
|
|
@ -377,3 +377,18 @@ fn failed_destroy_with_an_unreadable_list_fails() {
|
||||||
.expect_err("an unreadable list must not pass for an absent container");
|
.expect_err("an unreadable list must not pass for an absent container");
|
||||||
assert!(format!("{err:#}").contains("cannot confirm"), "{err:#}");
|
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:#}");
|
||||||
|
}
|
||||||
|
|
|
||||||
|
|
@ -63,7 +63,13 @@ pub async fn run(coord: &Arc<Coordinator>) -> Result<()> {
|
||||||
Err(e) => tracing::warn!(error = ?e, "clear stale meta lock failed"),
|
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");
|
tracing::info!(count = names.len(), "migration: scanning");
|
||||||
|
|
||||||
// Phase 0: move harness-owned files out of state/ into harness/.
|
// Phase 0: move harness-owned files out of state/ into harness/.
|
||||||
|
|
@ -94,9 +100,15 @@ pub async fn run(coord: &Arc<Coordinator>) -> Result<()> {
|
||||||
|
|
||||||
// Phase 3: meta repo.
|
// Phase 3: meta repo.
|
||||||
tracing::debug!("migration: phase 3 (meta sync_agents)");
|
tracing::debug!("migration: phase 3 (meta sync_agents)");
|
||||||
let agents = lifecycle::agents_for_meta_listing()
|
// `sync_agents` renders exactly the agents it is given, so an empty
|
||||||
.await
|
// stand-in for an unreadable list would drop every agent from the flake.
|
||||||
.unwrap_or_default();
|
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 {
|
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"),
|
Ok(Err(e)) => tracing::warn!(error = ?e, "migration: meta sync_agents failed"),
|
||||||
Err(_) => {
|
Err(_) => {
|
||||||
|
|
@ -139,9 +151,9 @@ fn migrate_harness_files(name: &hive_types::Ident) {
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
async fn enumerate_agents() -> Vec<hive_types::Ident> {
|
async fn enumerate_agents() -> Result<Vec<hive_types::Ident>> {
|
||||||
let containers = lifecycle::list().await.unwrap_or_default();
|
let containers = lifecycle::list().await?;
|
||||||
containers
|
Ok(containers
|
||||||
.into_iter()
|
.into_iter()
|
||||||
.filter_map(|c| {
|
.filter_map(|c| {
|
||||||
let name = if c == MANAGER_CONTAINER {
|
let name = if c == MANAGER_CONTAINER {
|
||||||
|
|
@ -151,7 +163,7 @@ async fn enumerate_agents() -> Vec<hive_types::Ident> {
|
||||||
};
|
};
|
||||||
hive_types::Ident::parse(name).ok()
|
hive_types::Ident::parse(name).ok()
|
||||||
})
|
})
|
||||||
.collect()
|
.collect())
|
||||||
}
|
}
|
||||||
|
|
||||||
async fn migrate_applied_repo(name: &str) -> Result<()> {
|
async fn migrate_applied_repo(name: &str) -> Result<()> {
|
||||||
|
|
|
||||||
|
|
@ -293,7 +293,7 @@ async fn handle_spawn(coord: &Arc<Coordinator>, name: &str) -> Result<HostRespon
|
||||||
tracing::info!(%name, "spawn");
|
tracing::info!(%name, "spawn");
|
||||||
let agent_dir = crate::paths::agent_runtime_dir(name);
|
let agent_dir = crate::paths::agent_runtime_dir(name);
|
||||||
let hive = coord.hive_env();
|
let hive = coord.hive_env();
|
||||||
let paths = Coordinator::agent_paths(name, agent_dir);
|
let paths = Coordinator::agent_paths(name, agent_dir)?;
|
||||||
// lifecycle::spawn creates the runtime dir internally before start.
|
// lifecycle::spawn creates the runtime dir internally before start.
|
||||||
// MCP listener registration is event-driven: bind immediately on
|
// MCP listener registration is event-driven: bind immediately on
|
||||||
// success so the harness can connect on its first turn without
|
// success so the harness can connect on its first turn without
|
||||||
|
|
@ -375,12 +375,15 @@ async fn handle_set_paused(
|
||||||
|
|
||||||
/// Collect per-agent status rows for `hivectl status` and the dashboard.
|
/// Collect per-agent status rows for `hivectl status` and the dashboard.
|
||||||
async fn handle_agent_status(coord: &Arc<Coordinator>) -> HostResponse {
|
async fn handle_agent_status(coord: &Arc<Coordinator>) -> HostResponse {
|
||||||
let rows = crate::container_view::build_all(&coord.hive_env())
|
match crate::container_view::build_all(&coord.hive_env()).await {
|
||||||
.await
|
Ok(views) => HostResponse::agent_statuses(
|
||||||
|
views
|
||||||
.into_iter()
|
.into_iter()
|
||||||
.map(hive_sh4re::container::AgentStatusRow::from)
|
.map(hive_sh4re::container::AgentStatusRow::from)
|
||||||
.collect();
|
.collect(),
|
||||||
HostResponse::agent_statuses(rows)
|
),
|
||||||
|
Err(e) => HostResponse::error(format!("listing containers: {e:#}")),
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
/// `SubscribeAgentStatus` — ack once, then push one single-row
|
/// `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.
|
// in the JSON until the agent's next spawn or rebuild.
|
||||||
let agent_dir = crate::paths::agent_runtime_dir(name.as_str());
|
let agent_dir = crate::paths::agent_runtime_dir(name.as_str());
|
||||||
let hive = coord.hive_env();
|
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?;
|
crate::lifecycle::write_dropins(name.as_str(), &hive, &paths).await?;
|
||||||
|
|
||||||
let (cpu, mem) = crate::resource_limits::effective(
|
let (cpu, mem) = crate::resource_limits::effective(
|
||||||
|
|
|
||||||
|
|
@ -190,7 +190,7 @@ pub async fn ensure_root_agent(coord: &Arc<Coordinator>) -> Result<()> {
|
||||||
// ensure_agent_runtime_dir needed here.
|
// ensure_agent_runtime_dir needed here.
|
||||||
let runtime = crate::paths::agent_runtime_dir(MANAGER_NAME);
|
let runtime = crate::paths::agent_runtime_dir(MANAGER_NAME);
|
||||||
let hive = coord.hive_env();
|
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?;
|
lifecycle::spawn(MANAGER_NAME, &hive, &paths).await?;
|
||||||
seed_manager_tool_groups();
|
seed_manager_tool_groups();
|
||||||
if let Err(e) = coord.power.set(MANAGER_NAME, crate::power::Wanted::Up) {
|
if let Err(e) = coord.power.set(MANAGER_NAME, crate::power::Wanted::Up) {
|
||||||
|
|
|
||||||
Loading…
Reference in a new issue