Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
b1243f149f | ||
|
|
ae3ecc1de2 | ||
|
|
949bab62d3 | ||
|
|
dce97f4826 | ||
|
|
14c2c6d4a5 | ||
|
|
f32fba0238 | ||
|
|
9d1f5ebe76 | ||
|
|
9cd408de8b |
7 changed files with 176 additions and 3 deletions
|
|
@ -967,6 +967,9 @@ pub async fn destroy(coord: &Arc<Coordinator>, name: &str, purge: bool) -> Resul
|
||||||
// roster, so any schedule that still targets the just-destroyed agent
|
// roster, so any schedule that still targets the just-destroyed agent
|
||||||
// now drops that ghost column live (no page reload needed).
|
// now drops that ghost column live (no page reload needed).
|
||||||
coord.emit_schedules_snapshot();
|
coord.emit_schedules_snapshot();
|
||||||
|
// Update tmpfiles.d to remove the destroyed agent's dirs from the
|
||||||
|
// boot-time pre-creation list. Best-effort: failure is logged only.
|
||||||
|
tokio::spawn(lifecycle::sync_tmpfiles());
|
||||||
Ok(())
|
Ok(())
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -722,6 +722,33 @@ pub async fn list() -> Result<Vec<String>> {
|
||||||
.collect())
|
.collect())
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// Sync `/etc/tmpfiles.d/hyperhive-agents.conf` with the currently-known
|
||||||
|
/// agent set (from `nixos-container list`). Strips the `h-` prefix to get
|
||||||
|
/// logical names. Best-effort: errors are logged but never propagated — a
|
||||||
|
/// failed tmpfiles write shouldn't block a spawn or destroy.
|
||||||
|
///
|
||||||
|
/// Called at hive-c0re startup and after each spawn / destroy so the file
|
||||||
|
/// always reflects the live agent set. `systemd-tmpfiles-setup.service`
|
||||||
|
/// reads the file at boot (before any container units start), pre-creating
|
||||||
|
/// bind-mount source dirs so container@h-* units don't race hive-c0re.
|
||||||
|
pub async fn sync_tmpfiles() {
|
||||||
|
let agents = match list().await {
|
||||||
|
Ok(containers) => containers
|
||||||
|
.into_iter()
|
||||||
|
.filter_map(|c| c.strip_prefix(AGENT_PREFIX).map(str::to_owned))
|
||||||
|
.collect::<Vec<_>>(),
|
||||||
|
Err(e) => {
|
||||||
|
tracing::warn!(error = ?e, "sync_tmpfiles: list failed; skipping");
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
};
|
||||||
|
if let Err(e) = crate::priv_client::sync_agent_tmpfiles(&agents).await {
|
||||||
|
tracing::warn!(error = ?e, "sync_tmpfiles: priv call failed");
|
||||||
|
} else {
|
||||||
|
tracing::debug!(count = agents.len(), "sync_tmpfiles: ok");
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
/// Build the per-line callback for `create_container_streaming` /
|
/// Build the per-line callback for `create_container_streaming` /
|
||||||
/// `update_container_streaming`. Both ops share identical dispatch logic
|
/// `update_container_streaming`. Both ops share identical dispatch logic
|
||||||
/// (stdout → info + `append_stdout`, stderr → warn + `append_stderr`); this
|
/// (stdout → info + `append_stdout`, stderr → warn + `append_stderr`); this
|
||||||
|
|
|
||||||
|
|
@ -283,6 +283,10 @@ async fn cmd_serve(
|
||||||
if let Err(e) = auto_update::ensure_root_agent(&coord).await {
|
if let Err(e) = auto_update::ensure_root_agent(&coord).await {
|
||||||
tracing::warn!(error = ?e, "auto-spawn root agent failed");
|
tracing::warn!(error = ?e, "auto-spawn root agent failed");
|
||||||
}
|
}
|
||||||
|
// Sync /etc/tmpfiles.d/hyperhive-agents.conf so agent runtime dirs are
|
||||||
|
// pre-declared for the next boot. Best-effort background task — a failure
|
||||||
|
// here must not block hive-c0re startup. See lifecycle::sync_tmpfiles.
|
||||||
|
tokio::spawn(hive_c0re::lifecycle::sync_tmpfiles());
|
||||||
// Auto-update in the background — don't block service start.
|
// Auto-update in the background — don't block service start.
|
||||||
// Sub-agent rebuilds can take tens of seconds; we want the admin
|
// Sub-agent rebuilds can take tens of seconds; we want the admin
|
||||||
// socket up immediately.
|
// socket up immediately.
|
||||||
|
|
|
||||||
|
|
@ -389,6 +389,21 @@ pub async fn upgrade_agent_subvolume(agent_name: &str) -> Result<()> {
|
||||||
.await?)
|
.await?)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// Write `/etc/tmpfiles.d/hyperhive-agents.conf` for `agents` (logical names,
|
||||||
|
/// e.g. `"atlas"`) and immediately apply it with `systemd-tmpfiles --create`.
|
||||||
|
/// See [`PrivRequest::SyncAgentTmpfiles`] for the full semantics.
|
||||||
|
///
|
||||||
|
/// # Errors
|
||||||
|
///
|
||||||
|
/// Returns an error if the priv socket call fails, if any agent name is
|
||||||
|
/// invalid, or if `systemd-tmpfiles --create` exits non-zero.
|
||||||
|
pub async fn sync_agent_tmpfiles(agents: &[String]) -> Result<()> {
|
||||||
|
ok(call(&PrivRequest::SyncAgentTmpfiles {
|
||||||
|
agents: agents.to_vec(),
|
||||||
|
})
|
||||||
|
.await?)
|
||||||
|
}
|
||||||
|
|
||||||
/// Parse `(referenced, exclusive)` bytes from `btrfs qgroup show -f --raw`
|
/// Parse `(referenced, exclusive)` bytes from `btrfs qgroup show -f --raw`
|
||||||
/// output (a qgroup row is `<id-with-slash> <rfer> <excl> …`).
|
/// output (a qgroup row is `<id-with-slash> <rfer> <excl> …`).
|
||||||
///
|
///
|
||||||
|
|
|
||||||
|
|
@ -216,6 +216,8 @@ async fn handle_spawn(coord: &Arc<Coordinator>, name: &str) -> Result<HostRespon
|
||||||
note: None,
|
note: None,
|
||||||
sha: None,
|
sha: None,
|
||||||
});
|
});
|
||||||
|
// Update tmpfiles.d so the new agent's dirs survive a reboot.
|
||||||
|
tokio::spawn(lifecycle::sync_tmpfiles());
|
||||||
}
|
}
|
||||||
Err(e) => {
|
Err(e) => {
|
||||||
// Roll back socket registration if container creation failed.
|
// Roll back socket registration if container creation failed.
|
||||||
|
|
|
||||||
|
|
@ -33,6 +33,10 @@ use tokio::process::Command;
|
||||||
/// Root of the per-agent unix-socket dirs on the host.
|
/// Root of the per-agent unix-socket dirs on the host.
|
||||||
const SOCKET_DIR_ROOT: &str = "/run/hive-agent";
|
const SOCKET_DIR_ROOT: &str = "/run/hive-agent";
|
||||||
|
|
||||||
|
/// Root of the per-agent MCP socket dirs on the host.
|
||||||
|
/// Matches `coordinator::AGENT_RUNTIME_ROOT` in hive-c0re.
|
||||||
|
const AGENT_RUNTIME_ROOT: &str = "/run/hyperhive/agents";
|
||||||
|
|
||||||
#[tokio::main]
|
#[tokio::main]
|
||||||
async fn main() -> Result<()> {
|
async fn main() -> Result<()> {
|
||||||
tracing_subscriber::fmt()
|
tracing_subscriber::fmt()
|
||||||
|
|
@ -162,7 +166,16 @@ async fn exec(req: PrivRequest, writer: &mut OwnedWriteHalf) -> Result<(String,
|
||||||
match req {
|
match req {
|
||||||
PrivRequest::StartContainer { ref name } => {
|
PrivRequest::StartContainer { ref name } => {
|
||||||
validate_container_name(name)?;
|
validate_container_name(name)?;
|
||||||
container_run(&["start", &container_system_name(name)]).await
|
let machine = container_system_name(name);
|
||||||
|
// Clear any start-limit lockout left by earlier failures so a
|
||||||
|
// now-correct start isn't blocked. nixos-container start does not
|
||||||
|
// do this itself. Best-effort: if the unit doesn't exist yet
|
||||||
|
// (first-time create) reset-failed is a no-op and we proceed.
|
||||||
|
let _ = Command::new("systemctl")
|
||||||
|
.args(["reset-failed", &format!("container@{machine}.service")])
|
||||||
|
.status()
|
||||||
|
.await;
|
||||||
|
container_run(&["start", &machine]).await
|
||||||
}
|
}
|
||||||
|
|
||||||
PrivRequest::StopContainer { ref name } => {
|
PrivRequest::StopContainer { ref name } => {
|
||||||
|
|
@ -317,6 +330,8 @@ async fn exec(req: PrivRequest, writer: &mut OwnedWriteHalf) -> Result<(String,
|
||||||
validate_agent_name(agent_name)?;
|
validate_agent_name(agent_name)?;
|
||||||
upgrade_agent_subvolume(agent_name).await
|
upgrade_agent_subvolume(agent_name).await
|
||||||
}
|
}
|
||||||
|
|
||||||
|
PrivRequest::SyncAgentTmpfiles { ref agents } => sync_agent_tmpfiles(agents).await,
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|
@ -409,17 +424,40 @@ fn chmod_socket_dir(agent_name: &str, mode: u32) -> Result<(String, String)> {
|
||||||
}
|
}
|
||||||
|
|
||||||
/// `WriteResourceLimits` — drop a systemd `MemoryMax`/`CPUQuota`
|
/// `WriteResourceLimits` — drop a systemd `MemoryMax`/`CPUQuota`
|
||||||
/// override into the container service's drop-in dir.
|
/// override into the container service's drop-in dir, together with a
|
||||||
|
/// `ConditionPathIsDirectory=` guard on the agent's MCP runtime dir.
|
||||||
|
///
|
||||||
|
/// The condition causes systemd to *skip* (not *fail*) the unit when the
|
||||||
|
/// bind-mount source dir is absent — result is `condition`, which does not
|
||||||
|
/// increment the start-limit counter. This is belt-and-braces on top of
|
||||||
|
/// the tmpfiles.d entries written by `SyncAgentTmpfiles`: in the unlikely
|
||||||
|
/// event the dir is missing at start time, the unit idles rather than
|
||||||
|
/// restart-looping into `start-limit-hit`.
|
||||||
fn write_resource_limits(
|
fn write_resource_limits(
|
||||||
container: &str,
|
container: &str,
|
||||||
memory_max: &str,
|
memory_max: &str,
|
||||||
cpu_quota: &str,
|
cpu_quota: &str,
|
||||||
) -> Result<(String, String)> {
|
) -> Result<(String, String)> {
|
||||||
validate_container_system_name(container)?;
|
validate_container_system_name(container)?;
|
||||||
|
// Derive the logical agent name (strip h- prefix) to form the runtime
|
||||||
|
// dir path. Falls back to the full container name for infra containers
|
||||||
|
// that don't use the h- prefix.
|
||||||
|
let logical = container.strip_prefix(AGENT_PREFIX).unwrap_or(container);
|
||||||
|
let runtime_dir = format!("{AGENT_RUNTIME_ROOT}/{logical}");
|
||||||
let dir = format!("/run/systemd/system/container@{container}.service.d");
|
let dir = format!("/run/systemd/system/container@{container}.service.d");
|
||||||
std::fs::create_dir_all(&dir).with_context(|| format!("create {dir}"))?;
|
std::fs::create_dir_all(&dir).with_context(|| format!("create {dir}"))?;
|
||||||
let path = format!("{dir}/hyperhive-limits.conf");
|
let path = format!("{dir}/hyperhive-limits.conf");
|
||||||
let content = format!("[Service]\nMemoryMax={memory_max}\nCPUQuota={cpu_quota}\n");
|
// [Unit] section: condition checked at start time — skips (not fails)
|
||||||
|
// the unit when the MCP socket dir is absent, avoiding restart loops.
|
||||||
|
// [Service] section: resource caps.
|
||||||
|
let content = format!(
|
||||||
|
"[Unit]\n\
|
||||||
|
ConditionPathIsDirectory={runtime_dir}\n\
|
||||||
|
\n\
|
||||||
|
[Service]\n\
|
||||||
|
MemoryMax={memory_max}\n\
|
||||||
|
CPUQuota={cpu_quota}\n"
|
||||||
|
);
|
||||||
std::fs::write(&path, content).with_context(|| format!("write {path}"))?;
|
std::fs::write(&path, content).with_context(|| format!("write {path}"))?;
|
||||||
Ok((String::new(), String::new()))
|
Ok((String::new(), String::new()))
|
||||||
}
|
}
|
||||||
|
|
@ -1484,3 +1522,72 @@ fn write_bridge_dns_marker(container: &str, isolation: Option<&NetworkIsolation>
|
||||||
}
|
}
|
||||||
Ok(())
|
Ok(())
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// `SyncAgentTmpfiles` — write `/etc/tmpfiles.d/hyperhive-agents.conf` for
|
||||||
|
/// the given agent set and immediately apply it with `systemd-tmpfiles --create`.
|
||||||
|
///
|
||||||
|
/// Each call atomically replaces the file with entries for all current agents,
|
||||||
|
/// then creates any missing dirs on the running host. The file survives reboots
|
||||||
|
/// and is read by `systemd-tmpfiles-setup.service` (runs in `sysinit.target`,
|
||||||
|
/// before any container units can start), so bind-mount source dirs are always
|
||||||
|
/// pre-created regardless of whether hive-c0re has reached `ensure_runtime`.
|
||||||
|
///
|
||||||
|
/// Directories written per agent:
|
||||||
|
/// - `/run/hyperhive/agents/<name>` (MCP socket dir, bind-mounted into container
|
||||||
|
/// as `/run/hive`)
|
||||||
|
/// - `/run/hive-agent/<name>` (web socket dir, bind-mounted into container)
|
||||||
|
const TMPFILES_PATH: &str = "/etc/tmpfiles.d/hyperhive-agents.conf";
|
||||||
|
|
||||||
|
async fn sync_agent_tmpfiles(agents: &[String]) -> Result<(String, String)> {
|
||||||
|
use std::fmt::Write as _;
|
||||||
|
for name in agents {
|
||||||
|
validate_agent_name(name)?;
|
||||||
|
}
|
||||||
|
|
||||||
|
// Build tmpfiles.d content. Root dirs first, then per-agent.
|
||||||
|
let mut content =
|
||||||
|
String::from("# managed by hive-c0re — do not edit (regenerated on spawn/destroy)\n");
|
||||||
|
// Parent dirs — created with permissive mode so hive-c0re can make subdirs.
|
||||||
|
// /run/hyperhive itself is also a RuntimeDirectory of hive-c0re.service; the
|
||||||
|
// tmpfiles.d entry here ensures it exists before hive-c0re starts (boot race).
|
||||||
|
content.push_str("d /run/hyperhive 0750 hive-core hive-core -\n");
|
||||||
|
writeln!(content, "d {AGENT_RUNTIME_ROOT} 0755 hive-core hive-core -").ok();
|
||||||
|
writeln!(content, "d {SOCKET_DIR_ROOT} 0755 root root -").ok();
|
||||||
|
// Per-agent dirs.
|
||||||
|
for name in agents {
|
||||||
|
writeln!(
|
||||||
|
content,
|
||||||
|
"d {AGENT_RUNTIME_ROOT}/{name} 0755 hive-core hive-core -"
|
||||||
|
)
|
||||||
|
.ok();
|
||||||
|
// 0777: agent harness (non-root uid) must bind sockets here.
|
||||||
|
// `d` adjusts mode/owner on existing dirs; world-writable matches
|
||||||
|
// the chmod_socket_dir(0o777) fallback so a runtime re-sync doesn't
|
||||||
|
// break a live agent's socket dir. host_config's chown_socket_dir
|
||||||
|
// tightens ownership afterwards when the agent uid is available.
|
||||||
|
writeln!(content, "d {SOCKET_DIR_ROOT}/{name} 0777 root root -").ok();
|
||||||
|
}
|
||||||
|
|
||||||
|
// Atomic write: write to a tmp file then rename so a concurrent reader
|
||||||
|
// always sees a complete file.
|
||||||
|
let tmp = format!("{TMPFILES_PATH}.tmp");
|
||||||
|
std::fs::write(&tmp, &content).with_context(|| format!("write {tmp}"))?;
|
||||||
|
std::fs::rename(&tmp, TMPFILES_PATH)
|
||||||
|
.with_context(|| format!("rename {TMPFILES_PATH}.tmp -> {TMPFILES_PATH}"))?;
|
||||||
|
tracing::info!(agents = agents.len(), "tmpfiles.d: wrote {TMPFILES_PATH}");
|
||||||
|
|
||||||
|
// Apply immediately so dirs exist on the running host, not just after next boot.
|
||||||
|
let out = Command::new("systemd-tmpfiles")
|
||||||
|
.args(["--create", TMPFILES_PATH])
|
||||||
|
.output()
|
||||||
|
.await
|
||||||
|
.context("systemd-tmpfiles --create")?;
|
||||||
|
if !out.status.success() {
|
||||||
|
let stderr = String::from_utf8_lossy(&out.stderr).trim().to_owned();
|
||||||
|
anyhow::bail!(
|
||||||
|
"systemd-tmpfiles --create failed ({}): {stderr}",
|
||||||
|
out.status
|
||||||
|
);
|
||||||
|
}
|
||||||
|
Ok((String::new(), String::new()))
|
||||||
|
}
|
||||||
|
|
|
||||||
|
|
@ -512,6 +512,21 @@ pub enum PrivRequest {
|
||||||
/// Logical agent name (validated by `validate_agent_name`).
|
/// Logical agent name (validated by `validate_agent_name`).
|
||||||
agent_name: String,
|
agent_name: String,
|
||||||
},
|
},
|
||||||
|
|
||||||
|
/// Write `/etc/tmpfiles.d/hyperhive-agents.conf` for the given agent set
|
||||||
|
/// and immediately apply it with `systemd-tmpfiles --create`. Each entry
|
||||||
|
/// declares the per-agent runtime dirs (`/run/hyperhive/agents/<name>` and
|
||||||
|
/// `/run/hive-agent/<name>`) so systemd recreates them at every boot before
|
||||||
|
/// any container units start — preventing bind-mount source missing errors
|
||||||
|
/// when container@h-* units race hive-c0re after a reboot.
|
||||||
|
///
|
||||||
|
/// Called at hive-c0re startup and after every agent spawn / destroy.
|
||||||
|
/// Agents are logical names (validated by `validate_agent_name`).
|
||||||
|
SyncAgentTmpfiles {
|
||||||
|
/// Logical agent names (e.g. `"atlas"`, `"ruth"`). hive-priv validates
|
||||||
|
/// each name before writing any path component derived from it.
|
||||||
|
agents: Vec<String>,
|
||||||
|
},
|
||||||
}
|
}
|
||||||
|
|
||||||
/// Response from the privileged helper.
|
/// Response from the privileged helper.
|
||||||
|
|
|
||||||
Loading…
Reference in a new issue