hive-c0re: fail on an unparseable resource-limits or topology file, write both atomically
resource-limits.json and topology.json were read with parse errors folded into an empty map, and written in place with std::fs::write. One truncated resource-limits.json followed by a single set_limits call rewrote the file with only that agent's entry, erasing every other agent's CPU and memory overrides without a log line. topology.json had the same shape: reconcile rebuilt it from the live set, losing pending (provisioned, never spawned) names. - agent_config::read_map / write_map are generic over the stored type. tool-groups and capabilities behave as before. - resource_limits::read / effective return an error for an existing but unreadable file; a missing file is still the empty map. set_limits fails without writing on such a file, and writes atomically. - topology: reconcile fails without writing on an unreadable file and writes atomically. all_agents logs the error and returns no agents, so a ManageRootAgent holder starts without cross-agent mounts. Read-path behaviour on an unreadable resource-limits.json, per caller: - write_dropins (every spawn / swap / WriteDropin): logs the error and keeps the limits drop-in already under /run; the agent still starts. With no drop-in yet (first start since boot) it writes the hive defaults, because no drop-in means an uncapped container. - render_flake: propagates, so sync_agents (and spawn/rebuild/destroy jobs) fail. An empty map would give tighter-capped agents the hive memoryMaxBytes. - container_view::build_all: logs the error each scan and renders the rows at the hive defaults (no ContainerView wire change). - set_resource_limits reply: propagates. Closes #4731
This commit is contained in:
parent
97cf8a1b2b
commit
6fac00dcc5
8 changed files with 332 additions and 119 deletions
|
|
@ -14,18 +14,20 @@ pub mod resource_limits;
|
||||||
pub mod tool_groups;
|
pub mod tool_groups;
|
||||||
pub mod topology;
|
pub mod topology;
|
||||||
|
|
||||||
use std::collections::BTreeMap;
|
|
||||||
use std::io::{self, Write as _};
|
use std::io::{self, Write as _};
|
||||||
use std::path::Path;
|
use std::path::Path;
|
||||||
|
|
||||||
/// Read a per-agent `name → [string]` JSON map. A missing file is the
|
use serde::Serialize;
|
||||||
/// empty map; any other read failure, or content that doesn't parse, is
|
use serde::de::DeserializeOwned;
|
||||||
/// an error, so a read-modify-write can't write back a map that has lost
|
|
||||||
/// every other agent's entry.
|
/// Read one of these JSON registries. A missing file is `T::default()`
|
||||||
fn read_map(path: &Path) -> io::Result<BTreeMap<String, Vec<String>>> {
|
/// (the empty map or set); any other read failure, or content that
|
||||||
|
/// doesn't parse, is an error, so a read-modify-write can't write back a
|
||||||
|
/// registry that has lost every other agent's entry.
|
||||||
|
fn read_map<T: DeserializeOwned + Default>(path: &Path) -> io::Result<T> {
|
||||||
let raw = match std::fs::read_to_string(path) {
|
let raw = match std::fs::read_to_string(path) {
|
||||||
Ok(raw) => raw,
|
Ok(raw) => raw,
|
||||||
Err(e) if e.kind() == io::ErrorKind::NotFound => return Ok(BTreeMap::new()),
|
Err(e) if e.kind() == io::ErrorKind::NotFound => return Ok(T::default()),
|
||||||
Err(e) => {
|
Err(e) => {
|
||||||
return Err(io::Error::new(
|
return Err(io::Error::new(
|
||||||
e.kind(),
|
e.kind(),
|
||||||
|
|
@ -41,13 +43,13 @@ fn read_map(path: &Path) -> io::Result<BTreeMap<String, Vec<String>>> {
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Serialise `map` as pretty JSON and replace `path` with it atomically:
|
/// Serialise `value` as pretty JSON and replace `path` with it atomically:
|
||||||
/// a temp file in the same directory (`rename(2)` is only atomic within a
|
/// a temp file in the same directory (`rename(2)` is only atomic within a
|
||||||
/// filesystem), fsynced, renamed over `path`, then the directory fsynced
|
/// filesystem), fsynced, renamed over `path`, then the directory fsynced
|
||||||
/// so the rename itself survives a crash. A reader sees the old file or
|
/// so the rename itself survives a crash. A reader sees the old file or
|
||||||
/// the new one, never a truncated one.
|
/// the new one, never a truncated one.
|
||||||
fn write_map(path: &Path, map: &BTreeMap<String, Vec<String>>) -> io::Result<()> {
|
fn write_map<T: Serialize + ?Sized>(path: &Path, value: &T) -> io::Result<()> {
|
||||||
let text = serde_json::to_string_pretty(map)
|
let text = serde_json::to_string_pretty(value)
|
||||||
.map_err(|e| io::Error::new(io::ErrorKind::InvalidData, e))?;
|
.map_err(|e| io::Error::new(io::ErrorKind::InvalidData, e))?;
|
||||||
let (Some(dir), Some(name)) = (path.parent(), path.file_name()) else {
|
let (Some(dir), Some(name)) = (path.parent(), path.file_name()) else {
|
||||||
return Err(io::Error::new(
|
return Err(io::Error::new(
|
||||||
|
|
@ -67,8 +69,12 @@ fn write_map(path: &Path, map: &BTreeMap<String, Vec<String>>) -> io::Result<()>
|
||||||
|
|
||||||
#[cfg(test)]
|
#[cfg(test)]
|
||||||
mod tests {
|
mod tests {
|
||||||
|
use std::collections::BTreeMap;
|
||||||
|
|
||||||
use super::*;
|
use super::*;
|
||||||
|
|
||||||
|
type Map = BTreeMap<String, Vec<String>>;
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
fn write_map_round_trips_and_leaves_no_temp_file() {
|
fn write_map_round_trips_and_leaves_no_temp_file() {
|
||||||
let dir = tempfile::tempdir().expect("tempdir");
|
let dir = tempfile::tempdir().expect("tempdir");
|
||||||
|
|
@ -77,7 +83,7 @@ mod tests {
|
||||||
write_map(&path, &map).expect("first write");
|
write_map(&path, &map).expect("first write");
|
||||||
let replaced = BTreeMap::from([("bob".to_owned(), vec!["inbox".to_owned()])]);
|
let replaced = BTreeMap::from([("bob".to_owned(), vec!["inbox".to_owned()])]);
|
||||||
write_map(&path, &replaced).expect("replacing write");
|
write_map(&path, &replaced).expect("replacing write");
|
||||||
assert_eq!(read_map(&path).expect("read"), replaced);
|
assert_eq!(read_map::<Map>(&path).expect("read"), replaced);
|
||||||
let entries: Vec<_> = std::fs::read_dir(dir.path())
|
let entries: Vec<_> = std::fs::read_dir(dir.path())
|
||||||
.expect("read_dir")
|
.expect("read_dir")
|
||||||
.map(|e| e.expect("entry").file_name())
|
.map(|e| e.expect("entry").file_name())
|
||||||
|
|
@ -88,7 +94,8 @@ mod tests {
|
||||||
#[test]
|
#[test]
|
||||||
fn read_map_of_a_missing_file_is_empty() {
|
fn read_map_of_a_missing_file_is_empty() {
|
||||||
let dir = tempfile::tempdir().expect("tempdir");
|
let dir = tempfile::tempdir().expect("tempdir");
|
||||||
let map = read_map(&dir.path().join("absent.json")).expect("missing file is not an error");
|
let map =
|
||||||
|
read_map::<Map>(&dir.path().join("absent.json")).expect("missing file is not an error");
|
||||||
assert!(map.is_empty());
|
assert!(map.is_empty());
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -97,7 +104,7 @@ mod tests {
|
||||||
let dir = tempfile::tempdir().expect("tempdir");
|
let dir = tempfile::tempdir().expect("tempdir");
|
||||||
let path = dir.path().join("map.json");
|
let path = dir.path().join("map.json");
|
||||||
std::fs::write(&path, "{\n \"alice\": [\"messag").expect("seed");
|
std::fs::write(&path, "{\n \"alice\": [\"messag").expect("seed");
|
||||||
let err = read_map(&path).expect_err("truncated JSON must not read as a map");
|
let err = read_map::<Map>(&path).expect_err("truncated JSON must not read as a map");
|
||||||
assert_eq!(err.kind(), io::ErrorKind::InvalidData);
|
assert_eq!(err.kind(), io::ErrorKind::InvalidData);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -28,7 +28,7 @@
|
||||||
//! a resource *cap* should not be sourced from the capped party.
|
//! a resource *cap* should not be sourced from the capped party.
|
||||||
|
|
||||||
use std::collections::BTreeMap;
|
use std::collections::BTreeMap;
|
||||||
use std::path::PathBuf;
|
use std::path::{Path, PathBuf};
|
||||||
|
|
||||||
const RESOURCE_LIMITS_FILE: &str = "resource-limits.json";
|
const RESOURCE_LIMITS_FILE: &str = "resource-limits.json";
|
||||||
|
|
||||||
|
|
@ -58,28 +58,35 @@ pub fn resource_limits_path() -> PathBuf {
|
||||||
crate::paths::meta_root().join(RESOURCE_LIMITS_FILE)
|
crate::paths::meta_root().join(RESOURCE_LIMITS_FILE)
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Read the per-agent limit map. Returns an empty map when the file is
|
/// Read the per-agent limit map. An absent file is the empty map —
|
||||||
/// absent or unparsable — callers treat a missing entry as "hive-wide
|
/// callers treat a missing entry as "hive-wide defaults". A file that
|
||||||
/// defaults", which is also the safe failure mode for a malformed file.
|
/// exists but can't be read or parsed is an error: an override can be
|
||||||
#[must_use]
|
/// tighter than the hive default, so reading it as empty is not a safe
|
||||||
pub fn read() -> BTreeMap<String, AgentLimits> {
|
/// fallback.
|
||||||
let path = resource_limits_path();
|
pub fn read() -> std::io::Result<BTreeMap<String, AgentLimits>> {
|
||||||
let Ok(raw) = std::fs::read_to_string(&path) else {
|
super::read_map(&resource_limits_path())
|
||||||
return BTreeMap::new();
|
|
||||||
};
|
|
||||||
serde_json::from_str(&raw).unwrap_or_default()
|
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Resolve the effective values for an agent, filling each unset field
|
/// Resolve the effective values for an agent from the override file at
|
||||||
|
/// `path` (normally [`resource_limits_path`]), filling each unset field
|
||||||
/// from the hive-wide default. This is the single place the fallback
|
/// from the hive-wide default. This is the single place the fallback
|
||||||
/// rule lives; `write_dropins` calls it and passes the result straight
|
/// rule lives; `write_dropins` calls it and passes the result straight
|
||||||
/// to systemd.
|
/// to systemd. Errors as [`read`] does.
|
||||||
///
|
///
|
||||||
/// Reads the override file. Use [`effective_from`] when resolving more
|
/// Use [`effective_from`] when resolving more than one agent in a row.
|
||||||
/// than one agent in a row.
|
pub fn effective(
|
||||||
#[must_use]
|
path: &Path,
|
||||||
pub fn effective(name: &str, hive_cpu_quota: &str, hive_memory_max: &str) -> (String, String) {
|
name: &str,
|
||||||
effective_from(&read(), name, hive_cpu_quota, hive_memory_max)
|
hive_cpu_quota: &str,
|
||||||
|
hive_memory_max: &str,
|
||||||
|
) -> std::io::Result<(String, String)> {
|
||||||
|
let limits = super::read_map(path)?;
|
||||||
|
Ok(effective_from(
|
||||||
|
&limits,
|
||||||
|
name,
|
||||||
|
hive_cpu_quota,
|
||||||
|
hive_memory_max,
|
||||||
|
))
|
||||||
}
|
}
|
||||||
|
|
||||||
/// [`effective`] against an already-loaded map — the multi-agent form.
|
/// [`effective`] against an already-loaded map — the multi-agent form.
|
||||||
|
|
@ -200,42 +207,27 @@ pub fn parse_bytes(value: &str) -> Option<u64> {
|
||||||
u64::try_from(bytes).ok()
|
u64::try_from(bytes).ok()
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Persist the full map. Sorted JSON output keeps meta-repo diffs
|
/// Set one agent's overrides and persist the map atomically. An entry
|
||||||
/// minimal.
|
/// with both fields unset is removed rather than stored, so "reset to
|
||||||
|
/// hive defaults" and "never configured" are the same state on disk.
|
||||||
///
|
///
|
||||||
/// # Errors
|
/// # Errors
|
||||||
///
|
///
|
||||||
/// Returns an `io::Error` when the meta dir can't be created or the
|
/// Fails without writing when the existing file can't be read or parsed
|
||||||
/// file can't be written (permissions, disk full). Serialization
|
/// (`InvalidData` for a parse failure), and otherwise when the meta dir
|
||||||
/// failure is surfaced as `InvalidData`, though it can't happen for
|
/// or the file can't be written.
|
||||||
/// this type — it's a plain map of strings.
|
pub fn set_limits(name: &str, limits: &AgentLimits) -> std::io::Result<()> {
|
||||||
pub fn write(map: &BTreeMap<String, AgentLimits>) -> std::io::Result<()> {
|
set_limits_at(&resource_limits_path(), name, limits)
|
||||||
let path = resource_limits_path();
|
|
||||||
if let Some(parent) = path.parent() {
|
|
||||||
std::fs::create_dir_all(parent)?;
|
|
||||||
}
|
|
||||||
let text = serde_json::to_string_pretty(map)
|
|
||||||
.map_err(|e| std::io::Error::new(std::io::ErrorKind::InvalidData, e))?;
|
|
||||||
std::fs::write(&path, format!("{text}\n"))
|
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Set one agent's overrides and persist. An entry with both fields
|
fn set_limits_at(path: &Path, name: &str, limits: &AgentLimits) -> std::io::Result<()> {
|
||||||
/// unset is removed rather than stored, so "reset to hive defaults" and
|
let mut current: BTreeMap<String, AgentLimits> = super::read_map(path)?;
|
||||||
/// "never configured" are the same state on disk.
|
|
||||||
///
|
|
||||||
/// # Errors
|
|
||||||
///
|
|
||||||
/// Propagates whatever the `write` function fails with. An unreadable or malformed
|
|
||||||
/// existing file is *not* an error — the `read` function degrades to an empty map,
|
|
||||||
/// so this call rewrites the file from scratch.
|
|
||||||
pub fn set_limits(name: &str, limits: &AgentLimits) -> std::io::Result<()> {
|
|
||||||
let mut current = read();
|
|
||||||
if limits.is_empty() {
|
if limits.is_empty() {
|
||||||
current.remove(name);
|
current.remove(name);
|
||||||
} else {
|
} else {
|
||||||
current.insert(name.to_owned(), limits.clone());
|
current.insert(name.to_owned(), limits.clone());
|
||||||
}
|
}
|
||||||
write(¤t)
|
super::write_map(path, ¤t)
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Validate a systemd `CPUQuota=` value. Percentages only, and values
|
/// Validate a systemd `CPUQuota=` value. Percentages only, and values
|
||||||
|
|
@ -390,11 +382,59 @@ mod tests {
|
||||||
assert!(!limits(Some("400%"), None).is_empty());
|
assert!(!limits(Some("400%"), None).is_empty());
|
||||||
}
|
}
|
||||||
|
|
||||||
|
const TRUNCATED: &str =
|
||||||
|
"{\n \"sock\": { \"cpu_quota\": \"50%\", \"memory_max\": \"1G\" },\n \"iris\": { \"mem";
|
||||||
|
|
||||||
|
fn corrupt_file() -> (tempfile::TempDir, PathBuf) {
|
||||||
|
let dir = tempfile::tempdir().expect("tempdir");
|
||||||
|
let path = dir.path().join(RESOURCE_LIMITS_FILE);
|
||||||
|
std::fs::write(&path, TRUNCATED).expect("seed");
|
||||||
|
(dir, path)
|
||||||
|
}
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
fn malformed_json_reads_as_empty_map() {
|
fn effective_of_a_corrupt_file_is_an_error() {
|
||||||
let parsed: BTreeMap<String, AgentLimits> =
|
let (_dir, path) = corrupt_file();
|
||||||
serde_json::from_str("{ not json").unwrap_or_default();
|
let err = effective(&path, "sock", HIVE_CPU, HIVE_MEM)
|
||||||
assert!(parsed.is_empty());
|
.expect_err("a corrupt file must not resolve to hive defaults");
|
||||||
|
assert_eq!(err.kind(), std::io::ErrorKind::InvalidData);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn set_limits_leaves_a_corrupt_file_untouched() {
|
||||||
|
let (_dir, path) = corrupt_file();
|
||||||
|
let err = set_limits_at(&path, "ruth", &limits(Some("400%"), None))
|
||||||
|
.expect_err("a corrupt file must not be overwritten");
|
||||||
|
assert_eq!(err.kind(), std::io::ErrorKind::InvalidData);
|
||||||
|
assert_eq!(std::fs::read(&path).expect("read"), TRUNCATED.as_bytes());
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn missing_file_resolves_to_hive_defaults() {
|
||||||
|
let dir = tempfile::tempdir().expect("tempdir");
|
||||||
|
let path = dir.path().join(RESOURCE_LIMITS_FILE);
|
||||||
|
let (cpu, mem) = effective(&path, "sock", HIVE_CPU, HIVE_MEM).expect("missing file");
|
||||||
|
assert_eq!((cpu.as_str(), mem.as_str()), (HIVE_CPU, HIVE_MEM));
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn set_limits_keeps_other_agents_and_leaves_no_temp_file() {
|
||||||
|
let dir = tempfile::tempdir().expect("tempdir");
|
||||||
|
let path = dir.path().join(RESOURCE_LIMITS_FILE);
|
||||||
|
set_limits_at(&path, "sock", &limits(None, Some("8G"))).expect("set on missing file");
|
||||||
|
set_limits_at(&path, "iris", &limits(Some("50%"), None)).expect("set");
|
||||||
|
let (cpu, mem) = effective(&path, "sock", HIVE_CPU, HIVE_MEM).expect("read");
|
||||||
|
assert_eq!((cpu.as_str(), mem.as_str()), (HIVE_CPU, "8G"));
|
||||||
|
let (cpu, _) = effective(&path, "iris", HIVE_CPU, HIVE_MEM).expect("read");
|
||||||
|
assert_eq!(cpu, "50%");
|
||||||
|
let entries: Vec<_> = std::fs::read_dir(dir.path())
|
||||||
|
.expect("read_dir")
|
||||||
|
.map(|e| e.expect("entry").file_name())
|
||||||
|
.collect();
|
||||||
|
assert_eq!(
|
||||||
|
entries,
|
||||||
|
vec![std::ffi::OsString::from(RESOURCE_LIMITS_FILE)]
|
||||||
|
);
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Absent fields must deserialize to `None`, not fail — an entry
|
/// Absent fields must deserialize to `None`, not fail — an entry
|
||||||
|
|
|
||||||
|
|
@ -88,7 +88,7 @@ fn set_groups_at(path: &Path, name: &str, groups: &[String]) -> anyhow::Result<(
|
||||||
if !groups.is_empty() {
|
if !groups.is_empty() {
|
||||||
validate_groups(groups)?;
|
validate_groups(groups)?;
|
||||||
}
|
}
|
||||||
let mut current = super::read_map(path)?;
|
let mut current: BTreeMap<String, Vec<String>> = super::read_map(path)?;
|
||||||
if groups.is_empty() {
|
if groups.is_empty() {
|
||||||
current.remove(name);
|
current.remove(name);
|
||||||
} else {
|
} else {
|
||||||
|
|
@ -104,7 +104,7 @@ pub fn remove_agent(name: &str) -> std::io::Result<()> {
|
||||||
}
|
}
|
||||||
|
|
||||||
fn remove_agent_at(path: &Path, name: &str) -> std::io::Result<()> {
|
fn remove_agent_at(path: &Path, name: &str) -> std::io::Result<()> {
|
||||||
let mut current = super::read_map(path)?;
|
let mut current: BTreeMap<String, Vec<String>> = super::read_map(path)?;
|
||||||
if current.remove(name).is_some() {
|
if current.remove(name).is_some() {
|
||||||
super::write_map(path, ¤t)?;
|
super::write_map(path, ¤t)?;
|
||||||
}
|
}
|
||||||
|
|
@ -151,7 +151,8 @@ mod tests {
|
||||||
let path = dir.path().join(TOOL_GROUPS_FILE);
|
let path = dir.path().join(TOOL_GROUPS_FILE);
|
||||||
set_groups_at(&path, "alice", &["inbox".to_owned()]).expect("set on missing file");
|
set_groups_at(&path, "alice", &["inbox".to_owned()]).expect("set on missing file");
|
||||||
set_groups_at(&path, "ruth", &["messaging".to_owned()]).expect("set");
|
set_groups_at(&path, "ruth", &["messaging".to_owned()]).expect("set");
|
||||||
let map = crate::agent_config::read_map(&path).expect("read");
|
let map: BTreeMap<String, Vec<String>> =
|
||||||
|
crate::agent_config::read_map(&path).expect("read");
|
||||||
assert_eq!(map["alice"], vec!["inbox".to_owned()]);
|
assert_eq!(map["alice"], vec!["inbox".to_owned()]);
|
||||||
assert_eq!(map["ruth"], vec!["messaging".to_owned()]);
|
assert_eq!(map["ruth"], vec!["messaging".to_owned()]);
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -17,7 +17,7 @@
|
||||||
//! `lifecycle::set_nspawn_flags`.
|
//! `lifecycle::set_nspawn_flags`.
|
||||||
|
|
||||||
use std::collections::BTreeSet;
|
use std::collections::BTreeSet;
|
||||||
use std::path::PathBuf;
|
use std::path::{Path, PathBuf};
|
||||||
|
|
||||||
use serde::Deserialize;
|
use serde::Deserialize;
|
||||||
|
|
||||||
|
|
@ -28,7 +28,7 @@ pub fn topology_path() -> PathBuf {
|
||||||
crate::paths::meta_root().join(TOPOLOGY_FILE)
|
crate::paths::meta_root().join(TOPOLOGY_FILE)
|
||||||
}
|
}
|
||||||
|
|
||||||
/// On-disk shapes [`read`] accepts. The array is what [`write()`] emits; the
|
/// On-disk shapes [`read_from`] accepts. The array is what [`reconcile`] writes; the
|
||||||
/// map is the legacy `name → parent | null` format, kept readable so a
|
/// map is the legacy `name → parent | null` format, kept readable so a
|
||||||
/// hive that upgrades across this change keeps its roster instead of
|
/// hive that upgrades across this change keeps its roster instead of
|
||||||
/// blanking it until the next `reconcile` pass — and a blank roster is not
|
/// blanking it until the next `reconcile` pass — and a blank roster is not
|
||||||
|
|
@ -42,33 +42,51 @@ enum OnDisk {
|
||||||
WithParents(std::collections::BTreeMap<String, Option<String>>),
|
WithParents(std::collections::BTreeMap<String, Option<String>>),
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Snapshot of the agent roster. Read on every `container_view::build_all`
|
/// A missing file is the empty roster (a fresh install that hasn't run
|
||||||
/// and every `render_flake` call. The file is small (one line per agent),
|
/// through `meta::sync_agents` yet).
|
||||||
/// so we re-read rather than caching — keeps the source of truth on disk.
|
impl Default for OnDisk {
|
||||||
///
|
fn default() -> Self {
|
||||||
/// Returns an empty set when the file is absent or unparsable. Safe
|
Self::Roster(BTreeSet::new())
|
||||||
/// degradation for fresh installs that haven't run through
|
|
||||||
/// `meta::sync_agents` yet.
|
|
||||||
#[must_use]
|
|
||||||
pub fn read() -> BTreeSet<String> {
|
|
||||||
let path = topology_path();
|
|
||||||
let Ok(raw) = std::fs::read_to_string(&path) else {
|
|
||||||
return BTreeSet::new();
|
|
||||||
};
|
|
||||||
match serde_json::from_str::<OnDisk>(&raw) {
|
|
||||||
Ok(OnDisk::Roster(names)) => names,
|
|
||||||
Ok(OnDisk::WithParents(map)) => map.into_keys().collect(),
|
|
||||||
Err(_) => BTreeSet::new(),
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// Snapshot of the agent roster in the file at `path` (normally
|
||||||
|
/// [`topology_path`]). The file is small (one line per agent), so it is
|
||||||
|
/// re-read rather than cached — keeps the source of truth on disk.
|
||||||
|
///
|
||||||
|
/// An absent file is the empty set. A file that exists but can't be read
|
||||||
|
/// or parsed is an error, so [`reconcile`] never writes over it.
|
||||||
|
fn read_from(path: &Path) -> std::io::Result<BTreeSet<String>> {
|
||||||
|
Ok(match super::read_map::<OnDisk>(path)? {
|
||||||
|
OnDisk::Roster(names) => names,
|
||||||
|
OnDisk::WithParents(map) => map.into_keys().collect(),
|
||||||
|
})
|
||||||
|
}
|
||||||
|
|
||||||
/// Every agent the roster knows about, in name order. This is the set
|
/// Every agent the roster knows about, in name order. This is the set
|
||||||
/// a [`hive_sh4re::permissions::Capability::ManageRootAgent`] holder gets
|
/// a [`hive_sh4re::permissions::Capability::ManageRootAgent`] holder gets
|
||||||
/// bind-mounted, and it is deliberately unfiltered: that capability means
|
/// bind-mounted, and it is deliberately unfiltered: that capability means
|
||||||
/// "may manage any agent", so the set is all of them.
|
/// "may manage any agent", so the set is all of them.
|
||||||
|
///
|
||||||
|
/// An unreadable roster is logged and yields no agents, so the holder
|
||||||
|
/// starts without cross-agent mounts rather than failing to start.
|
||||||
#[must_use]
|
#[must_use]
|
||||||
pub fn all_agents() -> Vec<String> {
|
pub fn all_agents() -> Vec<String> {
|
||||||
all_agents_in(&read())
|
all_agents_at(&topology_path())
|
||||||
|
}
|
||||||
|
|
||||||
|
fn all_agents_at(path: &Path) -> Vec<String> {
|
||||||
|
match read_from(path) {
|
||||||
|
Ok(topo) => all_agents_in(&topo),
|
||||||
|
Err(e) => {
|
||||||
|
tracing::error!(
|
||||||
|
path = %path.display(),
|
||||||
|
error = ?e,
|
||||||
|
"topology unreadable — no cross-agent mounts"
|
||||||
|
);
|
||||||
|
Vec::new()
|
||||||
|
}
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Pure form of [`all_agents`] for unit tests.
|
/// Pure form of [`all_agents`] for unit tests.
|
||||||
|
|
@ -77,33 +95,27 @@ pub fn all_agents_in(topo: &BTreeSet<String>) -> Vec<String> {
|
||||||
topo.iter().cloned().collect()
|
topo.iter().cloned().collect()
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Persist the roster. Sorted JSON output (`BTreeSet` iterates in key
|
|
||||||
/// order) keeps git diffs minimal across re-writes. Best-effort —
|
|
||||||
/// returns an `io::Error` so callers can decide whether a failure
|
|
||||||
/// should abort their op (`sync_agents`) or just log.
|
|
||||||
pub fn write(topology: &BTreeSet<String>) -> std::io::Result<()> {
|
|
||||||
let path = topology_path();
|
|
||||||
if let Some(parent) = path.parent() {
|
|
||||||
std::fs::create_dir_all(parent)?;
|
|
||||||
}
|
|
||||||
let text = serde_json::to_string_pretty(topology)
|
|
||||||
.map_err(|e| std::io::Error::new(std::io::ErrorKind::InvalidData, e))?;
|
|
||||||
std::fs::write(&path, format!("{text}\n"))
|
|
||||||
}
|
|
||||||
|
|
||||||
/// Reconcile `topology.json` against the current agent set. Adds any agent
|
/// Reconcile `topology.json` against the current agent set. Adds any agent
|
||||||
/// missing from the file; removes entries for agents no longer present.
|
/// missing from the file; removes entries for agents no longer present.
|
||||||
/// Returns true when the file changed and should be re-committed by the
|
/// Returns true when the file changed and should be re-committed by the
|
||||||
/// caller.
|
/// caller. The write is atomic, as a sorted array (`BTreeSet` iterates in
|
||||||
|
/// key order) to keep git diffs minimal.
|
||||||
///
|
///
|
||||||
/// `pending` lists agents that have a provisioned proposed config repo
|
/// `pending` lists agents that have a provisioned proposed config repo
|
||||||
/// but no container yet (provisioned, not yet spawned). They are KEPT
|
/// but no container yet (provisioned, not yet spawned). They are KEPT
|
||||||
/// (not dropped) so an agent that exists on disk but has never booted is
|
/// (not dropped) so an agent that exists on disk but has never booted is
|
||||||
/// still a name the hive knows about.
|
/// still a name the hive knows about.
|
||||||
|
///
|
||||||
|
/// Fails without writing when the existing file can't be read or parsed:
|
||||||
|
/// rebuilding it from the live set alone would drop every pending name.
|
||||||
pub fn reconcile(agent_names: &[String], pending: &[String]) -> std::io::Result<bool> {
|
pub fn reconcile(agent_names: &[String], pending: &[String]) -> std::io::Result<bool> {
|
||||||
let (next, changed) = apply_reconcile(&read(), agent_names, pending);
|
reconcile_at(&topology_path(), agent_names, pending)
|
||||||
|
}
|
||||||
|
|
||||||
|
fn reconcile_at(path: &Path, agent_names: &[String], pending: &[String]) -> std::io::Result<bool> {
|
||||||
|
let (next, changed) = apply_reconcile(&read_from(path)?, agent_names, pending);
|
||||||
if changed {
|
if changed {
|
||||||
write(&next)?;
|
super::write_map(path, &next)?;
|
||||||
}
|
}
|
||||||
Ok(changed)
|
Ok(changed)
|
||||||
}
|
}
|
||||||
|
|
@ -138,7 +150,10 @@ pub fn apply_reconcile(
|
||||||
|
|
||||||
#[cfg(test)]
|
#[cfg(test)]
|
||||||
mod tests {
|
mod tests {
|
||||||
use super::{BTreeSet, OnDisk, all_agents_in, apply_reconcile};
|
use super::{
|
||||||
|
BTreeSet, OnDisk, TOPOLOGY_FILE, all_agents_at, all_agents_in, apply_reconcile, read_from,
|
||||||
|
reconcile_at,
|
||||||
|
};
|
||||||
|
|
||||||
fn roster_three() -> BTreeSet<String> {
|
fn roster_three() -> BTreeSet<String> {
|
||||||
["alice", "bob", "carol"]
|
["alice", "bob", "carol"]
|
||||||
|
|
@ -228,4 +243,44 @@ mod tests {
|
||||||
};
|
};
|
||||||
assert_eq!(names, roster_three());
|
assert_eq!(names, roster_three());
|
||||||
}
|
}
|
||||||
|
|
||||||
|
const TRUNCATED: &str = "[\n \"alice\",\n \"bob\",\n \"do";
|
||||||
|
|
||||||
|
fn corrupt_file() -> (tempfile::TempDir, std::path::PathBuf) {
|
||||||
|
let dir = tempfile::tempdir().expect("tempdir");
|
||||||
|
let path = dir.path().join(TOPOLOGY_FILE);
|
||||||
|
std::fs::write(&path, TRUNCATED).expect("seed");
|
||||||
|
(dir, path)
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn reconcile_leaves_a_corrupt_file_untouched() {
|
||||||
|
let (_dir, path) = corrupt_file();
|
||||||
|
let live = vec!["alice".to_owned(), "bob".to_owned()];
|
||||||
|
let err =
|
||||||
|
reconcile_at(&path, &live, &[]).expect_err("a corrupt file must not be rewritten");
|
||||||
|
assert_eq!(err.kind(), std::io::ErrorKind::InvalidData);
|
||||||
|
assert_eq!(std::fs::read(&path).expect("read"), TRUNCATED.as_bytes());
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn all_agents_of_a_corrupt_file_is_empty() {
|
||||||
|
let (_dir, path) = corrupt_file();
|
||||||
|
assert!(all_agents_at(&path).is_empty());
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn reconcile_round_trips_through_a_missing_file_and_leaves_no_temp_file() {
|
||||||
|
let dir = tempfile::tempdir().expect("tempdir");
|
||||||
|
let path = dir.path().join(TOPOLOGY_FILE);
|
||||||
|
assert!(read_from(&path).expect("missing file").is_empty());
|
||||||
|
let live = vec!["alice".to_owned(), "bob".to_owned(), "carol".to_owned()];
|
||||||
|
assert!(reconcile_at(&path, &live, &[]).expect("reconcile"));
|
||||||
|
assert_eq!(read_from(&path).expect("read"), roster_three());
|
||||||
|
let entries: Vec<_> = std::fs::read_dir(dir.path())
|
||||||
|
.expect("read_dir")
|
||||||
|
.map(|e| e.expect("entry").file_name())
|
||||||
|
.collect();
|
||||||
|
assert_eq!(entries, vec![std::ffi::OsString::from(TOPOLOGY_FILE)]);
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -156,8 +156,17 @@ pub async fn build_all(hive: &crate::coordinator::HiveEnv) -> Vec<ContainerView>
|
||||||
let raw = lifecycle::list().await.unwrap_or_default();
|
let raw = lifecycle::list().await.unwrap_or_default();
|
||||||
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.
|
// SSE scan. An unreadable file shows every row at the hive defaults;
|
||||||
let limits = crate::resource_limits::read();
|
// the error log is the only signal, as `ContainerView` has no field
|
||||||
|
// to carry it.
|
||||||
|
let limits = crate::resource_limits::read().unwrap_or_else(|e| {
|
||||||
|
tracing::error!(
|
||||||
|
path = %crate::resource_limits::resource_limits_path().display(),
|
||||||
|
error = ?e,
|
||||||
|
"resource limits unreadable — container rows show hive defaults"
|
||||||
|
);
|
||||||
|
std::collections::BTreeMap::new()
|
||||||
|
});
|
||||||
let mut out = Vec::new();
|
let mut out = Vec::new();
|
||||||
for c in &raw {
|
for c in &raw {
|
||||||
let Some(logical) = c.strip_prefix(AGENT_PREFIX) else {
|
let Some(logical) = c.strip_prefix(AGENT_PREFIX) else {
|
||||||
|
|
|
||||||
|
|
@ -24,24 +24,77 @@ pub async fn write_dropins(name: &str, hive: &HiveEnv, paths: &AgentPaths) -> Re
|
||||||
validate(name)?;
|
validate(name)?;
|
||||||
let container = container_name(name);
|
let container = container_name(name);
|
||||||
set_nspawn_flags(&container, &paths.agent, &paths.claude, &paths.notes).await?;
|
set_nspawn_flags(&container, &paths.agent, &paths.claude, &paths.notes).await?;
|
||||||
let (cpu_quota, memory_max) =
|
if let Some((cpu_quota, memory_max)) = limits_to_apply(
|
||||||
crate::resource_limits::effective(name, &hive.agent_cpu_quota, &hive.agent_memory_max);
|
name,
|
||||||
set_resource_limits(
|
&crate::resource_limits::resource_limits_path(),
|
||||||
&container,
|
&limits_dropin_path(&container),
|
||||||
&cpu_quota,
|
hive,
|
||||||
&memory_max,
|
) {
|
||||||
hive.agent_cpu_weight,
|
set_resource_limits(
|
||||||
hive.agent_io_weight,
|
&container,
|
||||||
)
|
&cpu_quota,
|
||||||
.await?;
|
&memory_max,
|
||||||
|
hive.agent_cpu_weight,
|
||||||
|
hive.agent_io_weight,
|
||||||
|
)
|
||||||
|
.await?;
|
||||||
|
}
|
||||||
systemd_daemon_reload().await
|
systemd_daemon_reload().await
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// The limits drop-in hive-priv's `write_resource_limits` writes for
|
||||||
|
/// `container`. Must match that path: a mismatch reads as "no drop-in
|
||||||
|
/// yet" in [`limits_to_apply`].
|
||||||
|
fn limits_dropin_path(container: &str) -> PathBuf {
|
||||||
|
PathBuf::from(format!(
|
||||||
|
"/run/systemd/system/container@{container}.service.d/hyperhive-limits.conf"
|
||||||
|
))
|
||||||
|
}
|
||||||
|
|
||||||
|
/// The `(CPUQuota, MemoryMax)` to write into `agent_name`'s limits
|
||||||
|
/// drop-in, or `None` to leave the drop-in at `dropin_path` as it is.
|
||||||
|
///
|
||||||
|
/// When the overrides file at `limits_path` can't be read, the error is
|
||||||
|
/// logged and the agent still starts. The last-applied drop-in is kept,
|
||||||
|
/// because the hive defaults can be looser than the agent's own override.
|
||||||
|
/// With no drop-in yet, the agent gets the hive defaults: without one it
|
||||||
|
/// would run uncapped.
|
||||||
|
fn limits_to_apply(
|
||||||
|
agent_name: &str,
|
||||||
|
limits_path: &Path,
|
||||||
|
dropin_path: &Path,
|
||||||
|
hive: &HiveEnv,
|
||||||
|
) -> Option<(String, String)> {
|
||||||
|
let hive_cpu = &hive.agent_cpu_quota;
|
||||||
|
let hive_mem = &hive.agent_memory_max;
|
||||||
|
let e = match crate::resource_limits::effective(limits_path, agent_name, hive_cpu, hive_mem) {
|
||||||
|
Ok(limits) => return Some(limits),
|
||||||
|
Err(e) => e,
|
||||||
|
};
|
||||||
|
if dropin_path.exists() {
|
||||||
|
tracing::error!(
|
||||||
|
agent = %agent_name,
|
||||||
|
path = %limits_path.display(),
|
||||||
|
error = ?e,
|
||||||
|
"resource limits unreadable — keeping the last-applied limits drop-in"
|
||||||
|
);
|
||||||
|
None
|
||||||
|
} else {
|
||||||
|
tracing::error!(
|
||||||
|
agent = %agent_name,
|
||||||
|
path = %limits_path.display(),
|
||||||
|
error = ?e,
|
||||||
|
"resource limits unreadable and no limits drop-in yet — applying hive defaults"
|
||||||
|
);
|
||||||
|
Some((hive_cpu.clone(), hive_mem.clone()))
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
/// Write a systemd drop-in for `container@<container>.service` that applies
|
/// Write a systemd drop-in for `container@<container>.service` that applies
|
||||||
/// the agent's effective resource caps — its per-agent overrides from
|
/// the agent's effective resource caps — its per-agent overrides from
|
||||||
/// `meta/resource-limits.json` where set, the hive-wide defaults
|
/// `meta/resource-limits.json` where set, the hive-wide defaults
|
||||||
/// otherwise. Goes under `/run/systemd/system/...` so it's ephemeral
|
/// otherwise. Goes under `/run/systemd/system/...` so it's ephemeral
|
||||||
/// (regenerated on every spawn / rebuild).
|
/// (regenerated on every spawn / rebuild, except as [`limits_to_apply`] says).
|
||||||
///
|
///
|
||||||
/// The weights are hive-wide (`services.hyperhive.agentCpuWeight` /
|
/// The weights are hive-wide (`services.hyperhive.agentCpuWeight` /
|
||||||
/// `agentIoWeight`) — unlike the caps they have no per-agent override in
|
/// `agentIoWeight`) — unlike the caps they have no per-agent override in
|
||||||
|
|
@ -426,8 +479,9 @@ async fn set_nspawn_flags(
|
||||||
mod tests {
|
mod tests {
|
||||||
use super::{
|
use super::{
|
||||||
BindMount, QUEUE_CLIENT_ID_CREDENTIAL, QUEUE_SECRET_CREDENTIAL, bind_child_agent_dirs,
|
BindMount, QUEUE_CLIENT_ID_CREDENTIAL, QUEUE_SECRET_CREDENTIAL, bind_child_agent_dirs,
|
||||||
holds_manage_root_agent, queue_agent_credentials,
|
holds_manage_root_agent, limits_to_apply, queue_agent_credentials,
|
||||||
};
|
};
|
||||||
|
use crate::coordinator::HiveEnv;
|
||||||
|
|
||||||
fn child_binds() -> Vec<BindMount> {
|
fn child_binds() -> Vec<BindMount> {
|
||||||
let mut binds = Vec::new();
|
let mut binds = Vec::new();
|
||||||
|
|
@ -561,4 +615,48 @@ mod tests {
|
||||||
assert!(holds_manage_root_agent("ruth", &path));
|
assert!(holds_manage_root_agent("ruth", &path));
|
||||||
assert!(!holds_manage_root_agent("alice", &path));
|
assert!(!holds_manage_root_agent("alice", &path));
|
||||||
}
|
}
|
||||||
|
|
||||||
|
const TRUNCATED_LIMITS: &str = "{\n \"sock\": { \"cpu_quota\": \"50%\", \"mem";
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn corrupt_limits_keep_an_existing_dropin() {
|
||||||
|
let dir = tempfile::tempdir().expect("tempdir");
|
||||||
|
let limits = dir.path().join("resource-limits.json");
|
||||||
|
let dropin = dir.path().join("hyperhive-limits.conf");
|
||||||
|
std::fs::write(&limits, TRUNCATED_LIMITS).expect("seed limits");
|
||||||
|
std::fs::write(&dropin, "MemoryMax=1G\n").expect("seed drop-in");
|
||||||
|
assert_eq!(
|
||||||
|
limits_to_apply("sock", &limits, &dropin, &HiveEnv::default()),
|
||||||
|
None
|
||||||
|
);
|
||||||
|
assert_eq!(
|
||||||
|
std::fs::read(&limits).expect("read"),
|
||||||
|
TRUNCATED_LIMITS.as_bytes()
|
||||||
|
);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn corrupt_limits_without_a_dropin_apply_hive_defaults() {
|
||||||
|
let dir = tempfile::tempdir().expect("tempdir");
|
||||||
|
let limits = dir.path().join("resource-limits.json");
|
||||||
|
std::fs::write(&limits, TRUNCATED_LIMITS).expect("seed limits");
|
||||||
|
let applied = limits_to_apply(
|
||||||
|
"sock",
|
||||||
|
&limits,
|
||||||
|
&dir.path().join("absent.conf"),
|
||||||
|
&HiveEnv::default(),
|
||||||
|
);
|
||||||
|
assert_eq!(applied, Some(("200%".to_owned(), "4G".to_owned())));
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn readable_limits_are_applied_over_an_existing_dropin() {
|
||||||
|
let dir = tempfile::tempdir().expect("tempdir");
|
||||||
|
let limits = dir.path().join("resource-limits.json");
|
||||||
|
let dropin = dir.path().join("hyperhive-limits.conf");
|
||||||
|
std::fs::write(&limits, r#"{"sock": {"memory_max": "1G"}}"#).expect("seed limits");
|
||||||
|
std::fs::write(&dropin, "MemoryMax=8G\n").expect("seed drop-in");
|
||||||
|
let applied = limits_to_apply("sock", &limits, &dropin, &HiveEnv::default());
|
||||||
|
assert_eq!(applied, Some(("200%".to_owned(), "1G".to_owned())));
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
|
||||||
|
|
@ -1318,10 +1318,12 @@ where
|
||||||
);
|
);
|
||||||
// Propagated, never read as empty: an agent with no tool-groups entry
|
// Propagated, never read as empty: an agent with no tool-groups entry
|
||||||
// renders `toolGroups = null` and gets `AGENT_DEFAULT`, which fails
|
// renders `toolGroups = null` and gets `AGENT_DEFAULT`, which fails
|
||||||
// open for any agent whose explicit entry is narrower.
|
// open for any agent whose explicit entry is narrower. Likewise an
|
||||||
|
// agent with no resource-limits entry gets the hive `memoryMaxBytes`,
|
||||||
|
// a larger heap ceiling than a tighter override.
|
||||||
let tool_groups_map = crate::tool_groups::read()?;
|
let tool_groups_map = crate::tool_groups::read()?;
|
||||||
let capabilities_map = crate::capabilities::read()?;
|
let capabilities_map = crate::capabilities::read()?;
|
||||||
let resource_limits_map = crate::resource_limits::read();
|
let resource_limits_map = crate::resource_limits::read()?;
|
||||||
for spec in agents {
|
for spec in agents {
|
||||||
// Emit `toolGroups = "group1,group2"` when the operator has
|
// Emit `toolGroups = "group1,group2"` when the operator has
|
||||||
// explicitly configured groups for this agent. Absent entry = null
|
// explicitly configured groups for this agent. Absent entry = null
|
||||||
|
|
|
||||||
|
|
@ -506,10 +506,11 @@ async fn handle_set_resource_limits(
|
||||||
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(
|
||||||
|
&crate::resource_limits::resource_limits_path(),
|
||||||
name.as_str(),
|
name.as_str(),
|
||||||
&hive.agent_cpu_quota,
|
&hive.agent_cpu_quota,
|
||||||
&hive.agent_memory_max,
|
&hive.agent_memory_max,
|
||||||
);
|
)?;
|
||||||
Ok(HostResponse::messages(vec![format!(
|
Ok(HostResponse::messages(vec![format!(
|
||||||
"{name}: CPUQuota={cpu} MemoryMax={mem} — the cgroup cap itself is live now (restart \
|
"{name}: CPUQuota={cpu} MemoryMax={mem} — the cgroup cap itself is live now (restart \
|
||||||
the container if it's running and needs the new cap immediately), but the derived \
|
the container if it's running and needs the new cap immediately), but the derived \
|
||||||
|
|
|
||||||
Loading…
Reference in a new issue