//! Minimal privileged helper for hive-c0re. //! //! Runs as root. Exposes a narrow unix socket at `/run/hive/priv.sock` //! that accepts `PrivRequest` JSON lines and executes only the //! operations that genuinely require root. All coordination logic, //! broker, HTTP, and scheduling stay in the unprivileged hive-c0re //! process. //! //! **Security model**: every request is validated against a strict //! container-name allowlist before any filesystem or process operation. //! Only containers whose names match the hive convention (`h-*`, //! the manager container, or known sibling service containers) are //! accepted. Every variant maps to a single known operation — no //! arbitrary command pass-through. //! //! **Socket activation**: when systemd passes the listener socket via //! `LISTEN_FDS=1` + `LISTEN_PID=`, the inherited fd 3 is used //! instead of binding a fresh socket. use std::os::fd::{AsRawFd as _, FromRawFd as _, OwnedFd, RawFd}; use std::path::{Path, PathBuf}; use anyhow::{Context as _, Result, anyhow, bail}; use hive_priv_sock::{ AGENT_PREFIX, AGENT_RUNTIME_ROOT, AGENT_STATE_ROOT, BindMount, CredentialMount, InfraAction, InfraContainer, JournalQuery, META_DIR, MIGRATE_STAGING_ROOT, NetworkIsolation, PAUSED_MARKER_FILE, PRIV_SOCK, PrivEvent, PrivRequest, PrivResponse, PrivStream, PrivStreamLine, SIBLING_CONTAINERS, }; use tokio::io::{AsyncWriteExt, BufReader}; use tokio::net::unix::OwnedWriteHalf; use tokio::net::{UnixListener, UnixStream}; use tokio::process::Command; /// Root of the per-agent unix-socket dirs on the host. const SOCKET_DIR_ROOT: &str = "/run/hive-agent"; #[tokio::main] async fn main() -> Result<()> { tracing_subscriber::fmt() .with_env_filter( tracing_subscriber::EnvFilter::try_from_default_env().unwrap_or_else(|_| "info".into()), ) // This is a systemd-managed daemon — stdout always goes to journald, // never a human terminal, and journald doesn't strip ANSI escapes: // they land in victorialogs as raw byte-array spam otherwise. .with_ansi(false) .init(); let listener = socket_listener()?; if let Err(e) = remove_legacy_tmpfiles() { tracing::warn!(error = %format!("{e:#}"), "legacy tmpfiles cleanup failed"); } tracing::info!("hive-priv listening"); loop { match listener.accept().await { Ok((stream, _)) => { tokio::spawn(handle(stream)); } Err(e) => { tracing::error!(error = %e, "accept failed"); } } } } fn socket_listener() -> Result { // hive-priv is ALWAYS socket-activated by the `hive-priv.socket` unit // (fd 3 via LISTEN_FDS). There is intentionally no self-bind fallback, // so dev and prod take the same path; see docs/trust-boundary/boundary.md. let listen_fds: Option = std::env::var("LISTEN_FDS") .ok() .and_then(|s| s.parse().ok()); let listen_pid: Option = std::env::var("LISTEN_PID") .ok() .and_then(|s| s.parse().ok()); let activated = matches!(listen_fds, Some(n) if n >= 1) && listen_pid == Some(std::process::id()); if !activated { bail!( "hive-priv requires systemd socket activation (expected LISTEN_FDS>=1 + \ LISTEN_PID= for {PRIV_SOCK}); run it via the hive-priv.socket unit, \ not directly" ); } // SAFETY: systemd has passed us a ready UnixListener on fd 3. let std_listener = unsafe { use std::os::unix::io::FromRawFd; std::os::unix::net::UnixListener::from_raw_fd(3) }; std_listener .set_nonblocking(true) .context("set socket non-blocking")?; let listener = tokio::net::UnixListener::from_std(std_listener).context("wrap systemd socket")?; tracing::info!("using systemd-activated socket"); Ok(listener) } /// Ancillary-data buffer sized and aligned for one `SCM_RIGHTS` message. /// /// `CMSG_SPACE` is not a `const fn`, so the size is a literal with room /// to spare (24 bytes are needed for a single descriptor on x86-64). The /// union member gives the `cmsghdr` alignment `CMSG_FIRSTHDR` requires — /// a bare `[u8; N]` is only byte-aligned and would be undefined behaviour /// to walk. #[repr(C)] union CmsgSpace { _align: libc::cmsghdr, bytes: [u8; 32], } /// One `recvmsg` into `buf`, returning the bytes read plus any file /// descriptors that rode along as `SCM_RIGHTS`. /// /// Why not a plain read: ancillary data is attached to a *specific* /// `recvmsg` call, so a buffered line reader cannot surface it — it /// reads the bytes and silently drops the descriptor. /// /// `MSG_CMSG_CLOEXEC` is not optional: without it a received descriptor /// is inherited by every `btrfs` / `nixos-container` child this helper /// later spawns. fn recv_with_fds(sock: RawFd, buf: &mut [u8]) -> std::io::Result<(usize, Vec)> { const FD_SIZE: usize = std::mem::size_of::(); let mut iov = libc::iovec { iov_base: buf.as_mut_ptr().cast(), iov_len: buf.len(), }; let mut cmsg = CmsgSpace { bytes: [0; 32] }; // SAFETY: msghdr is a plain C struct with no invalid bit patterns; // every field we care about is set immediately below. let mut msg: libc::msghdr = unsafe { std::mem::zeroed() }; msg.msg_iov = &raw mut iov; msg.msg_iovlen = 1; msg.msg_control = std::ptr::addr_of_mut!(cmsg.bytes).cast(); msg.msg_controllen = 32; // SAFETY: `msg` points at a live iovec covering `buf` and a live, // correctly aligned control buffer of the length we just declared. let n = unsafe { libc::recvmsg(sock, &raw mut msg, libc::MSG_CMSG_CLOEXEC) }; if n < 0 { return Err(std::io::Error::last_os_error()); } // Take ownership of every descriptor the kernel attached, even ones // this protocol never expects: an `OwnedFd` we drop is closed, an // fd we fail to claim is leaked for the lifetime of the process. let mut fds = Vec::new(); // SAFETY: `msg` was just filled in by a successful `recvmsg`. let mut cmsgp = unsafe { libc::CMSG_FIRSTHDR(&raw const msg) }; while !cmsgp.is_null() { // SAFETY: CMSG_FIRSTHDR / CMSG_NXTHDR only ever return a pointer // to a complete header inside the control buffer. let hdr = unsafe { std::ptr::read_unaligned(cmsgp) }; if hdr.cmsg_level == libc::SOL_SOCKET && hdr.cmsg_type == libc::SCM_RIGHTS { // SAFETY: same, and CMSG_LEN(0) is the header's own length. let payload = hdr.cmsg_len as usize - unsafe { libc::CMSG_LEN(0) } as usize; let count = payload / FD_SIZE; // SAFETY: CMSG_DATA points at `payload` bytes of descriptors. let data = unsafe { libc::CMSG_DATA(cmsgp) }; for i in 0..count { // Copied out byte-wise rather than read through a // `*const RawFd`: the control buffer is only guaranteed // `cmsghdr`-aligned, so casting to a more strictly // aligned pointer would be unsound even where it happens // to work. let mut raw = [0u8; FD_SIZE]; // SAFETY: i < count, so this reads inside the payload. unsafe { std::ptr::copy_nonoverlapping(data.add(i * FD_SIZE), raw.as_mut_ptr(), FD_SIZE); } // SAFETY: the kernel just created this descriptor for // us — we are its only owner. fds.push(unsafe { OwnedFd::from_raw_fd(RawFd::from_ne_bytes(raw)) }); } } // SAFETY: `cmsgp` came from this same message. cmsgp = unsafe { libc::CMSG_NXTHDR(&raw const msg, cmsgp) }; } // `n >= 0` was checked above, so the conversion cannot fail; going // through `try_from` keeps it a cast-free, lint-clean widening. let read = usize::try_from(n).unwrap_or_default(); Ok((read, fds)) } /// Reads newline-delimited requests off one connection, pairing each /// with the descriptor that arrived with it. /// /// The pairing is deliberately trivial, because the protocol is: /// `hive-sock-client` connects per request, so a connection carries one /// line and at most one descriptor. The loop below still handles several /// sequential requests (the server always has), but it refuses to guess /// — a second descriptor arriving before its line is a protocol error, /// not something to queue and hope about. struct Requests<'a> { sock: &'a UnixStream, buf: Vec, fd: Option, } impl Requests<'_> { /// Next complete request line and its descriptor, or `None` at EOF. async fn next(&mut self) -> Result)>> { loop { if let Some(nl) = self.buf.iter().position(|&b| b == b'\n') { let line: Vec = self.buf.drain(..=nl).take(nl).collect(); let line = String::from_utf8(line).context("request line was not valid UTF-8")?; return Ok(Some((line, self.fd.take()))); } let mut chunk = [0u8; 8192]; let raw = self.sock.as_raw_fd(); let (n, fds) = self .sock .async_io(tokio::io::Interest::READABLE, || { recv_with_fds(raw, &mut chunk) }) .await .context("recvmsg on the priv socket")?; for fd in fds { if self.fd.replace(fd).is_some() { bail!("more than one file descriptor passed for a single request"); } } if n == 0 { if !self.buf.is_empty() { bail!("connection closed mid-request ({} bytes)", self.buf.len()); } return Ok(None); } self.buf.extend_from_slice(&chunk[..n]); } } } async fn handle(stream: UnixStream) { let (reader, mut writer) = stream.into_split(); let mut requests = Requests { sock: reader.as_ref(), buf: Vec::new(), fd: None, }; loop { let (line, fd) = match requests.next().await { Ok(Some(req)) => req, Ok(None) => break, Err(e) => { tracing::warn!(error = %format!("{e:#}"), "reading request failed"); break; } }; let resp = dispatch(&line, fd, &mut writer).await; // Write the terminal PrivResponse as a PrivEvent::Done. Wire-identical // to a bare PrivResponse (untagged), so old hive-c0re callers that // deserialise directly to PrivResponse continue to work. let event = PrivEvent::Done(resp); let mut json = serde_json::to_string(&event).unwrap_or_else(|e| { format!("{{\"ok\":false,\"stdout\":\"\",\"stderr\":\"\",\"error\":\"serialise failed: {e}\"}}") }); json.push('\n'); if let Err(e) = writer.write_all(json.as_bytes()).await { tracing::warn!(error = %e, "write response failed"); break; } } } /// Reject a request whose descriptor and operation disagree, in either /// direction. /// /// No guessing when the caller didn't say: an op that streams into a /// passed descriptor cannot invent one, and an op that takes none must /// not silently accept one. Returning the `Err` here drops the /// `OwnedFd`, which closes it. fn check_fd_agreement(req: &PrivRequest, fd: Option<&OwnedFd>) -> Result<()> { let wants_fd = matches!(req, PrivRequest::SendAgentSnapshotToFd { .. }); match (wants_fd, fd.is_some()) { (true, false) => bail!("this operation requires a passed file descriptor, none arrived"), (false, true) => bail!("this operation does not take a passed file descriptor"), _ => Ok(()), } } async fn dispatch(line: &str, fd: Option, writer: &mut OwnedWriteHalf) -> PrivResponse { match run(line, fd, writer).await { Ok((stdout, stderr)) => PrivResponse { ok: true, stdout, stderr, error: None, }, Err(e) => PrivResponse { ok: false, stdout: String::new(), stderr: String::new(), error: Some(format!("{e:#}")), }, } } /// Parse one request line, check it agrees with the descriptor that /// arrived with it, and execute it. /// /// Split out of [`dispatch`] so the three failure modes collapse into one /// `Result` instead of three nested matches building the same struct. async fn run( line: &str, fd: Option, writer: &mut OwnedWriteHalf, ) -> Result<(String, String)> { let req = serde_json::from_str::(line).context("parse request")?; check_fd_agreement(&req, fd.as_ref())?; exec(req, fd, writer).await } /// Write one `PrivEvent::Line` to the client. Best-effort: a write /// failure is logged but doesn't abort the running subprocess. async fn write_line_event(writer: &mut OwnedWriteHalf, stream: PrivStream, data: &str) { let event = PrivEvent::Line(PrivStreamLine { stream, data: data.to_owned(), }); if let Ok(mut json) = serde_json::to_string(&event) { json.push('\n'); if let Err(e) = writer.write_all(json.as_bytes()).await { tracing::warn!(error = %e, "write_line_event: write failed"); } } } /// Execute a validated `PrivRequest`. Returns `(stdout, stderr)` on success. /// For streaming ops (`CreateContainer`/`UpdateContainer` with `stream: true`) /// output lines are forwarded to `writer` as `PrivEvent::Line` messages and /// the returned strings are empty. /// /// `fd` is the descriptor that arrived with this request, already checked /// against the operation by [`check_fd_agreement`]: `Some` exactly for /// the variants that stream into a caller-supplied descriptor, `None` /// for every other operation. // One match arm per priv op — a flat 1:1 dispatch table. Every arm either // delegates directly or validates then delegates; an op whose handling is // more than that gets its own named function instead, so the match's length // tracks the op count, not complexity. #[allow(clippy::too_many_lines)] async fn exec( req: PrivRequest, fd: Option, writer: &mut OwnedWriteHalf, ) -> Result<(String, String)> { match req { PrivRequest::StartContainer { ref name } => start_container(name).await, PrivRequest::StopContainer { ref name } => stop_container(name).await, PrivRequest::KillContainer { ref name } => kill_container(name).await, PrivRequest::UpdateContainer { ref name, stream } => { container_flake_action("update", name, stream, writer).await } PrivRequest::CreateContainer { ref name, stream } => { container_flake_action("create", name, stream, writer).await } PrivRequest::DestroyContainer { ref name } => destroy_container(name).await, PrivRequest::ListContainers => container_run(&["list"]).await, PrivRequest::ReadContainerJournal { ref container, ref query, } => exec_read_container_journal(container, query).await, PrivRequest::WriteNspawnFlags { ref container, ref binds, ref isolation, ref load_credentials, } => handle_write_nspawn_flags(container, binds, isolation, load_credentials), PrivRequest::EnsureAgentSocketDir { ref name } => ensure_agent_socket_dir(name), PrivRequest::WriteResourceLimits { ref container, ref memory_max, ref cpu_quota, cpu_weight, io_weight, } => write_resource_limits(container, memory_max, cpu_quota, cpu_weight, io_weight), PrivRequest::RemoveServiceDropin { ref container } => remove_service_dropin(container), PrivRequest::DaemonReload => daemon_reload().await, PrivRequest::ReloadGatewayNginx => sync_gateway_nginx().await, PrivRequest::RunForgeAdmin { ref args } => exec_forge_admin(args).await, PrivRequest::SetAgentPaused { ref agent_name, paused, } => exec_set_agent_paused(agent_name, paused), PrivRequest::RegisterCiRunner { ref token } => register_ci_runner(token).await, PrivRequest::ControlInfraContainer { container, action } => { control_infra_container(container, action).await } PrivRequest::EnsureAgentSubvolume { ref agent_name } => { exec_ensure_agent_subvolume(agent_name).await } PrivRequest::DeleteAgentSubvolume { ref agent_name } => { exec_delete_agent_subvolume(agent_name).await } PrivRequest::EnsureBtrfsQuota => ensure_btrfs_quota().await, PrivRequest::ReadSubvolumeUsage { ref agent_name } => { exec_read_subvolume_usage(agent_name).await } PrivRequest::SetSubvolumeQuota { ref agent_name, limit_bytes, } => exec_set_subvolume_quota(agent_name, limit_bytes).await, PrivRequest::UpgradeAgentSubvolume { ref agent_name } => { exec_upgrade_agent_subvolume(agent_name).await } PrivRequest::SnapshotAgentSubvolume { ref agent_name, ref snapshot_name, } => exec_snapshot_agent_subvolume(agent_name, snapshot_name).await, PrivRequest::DeleteAgentSnapshot { ref agent_name, ref snapshot_name, } => exec_delete_agent_snapshot(agent_name, snapshot_name).await, PrivRequest::SendAgentSnapshotToFile { ref agent_name, ref snapshot_name, ref parent_snapshot_name, ref dest_file_name, } => { exec_send_agent_snapshot_to_file( agent_name, snapshot_name, parent_snapshot_name.as_deref(), dest_file_name, ) .await } PrivRequest::SendAgentSnapshotToFd { ref agent_name, ref snapshot_name, ref parent_snapshot_name, } => { exec_send_agent_snapshot_to_fd( agent_name, snapshot_name, parent_snapshot_name.as_deref(), fd, ) .await } PrivRequest::SyncAgentTmpfiles => { remove_legacy_tmpfiles()?; Ok((String::new(), String::new())) } } } /// `StartContainer`: 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 the start proceeds regardless. async fn start_container(name: &str) -> Result<(String, String)> { validate_agent_name(name)?; let machine = container_system_name(name); let _ = Command::new("systemctl") .args(["reset-failed", &format!("container@{machine}.service")]) .status() .await; container_run(&["start", &machine]).await } /// `RunForgeAdmin`: every arg must pass [`validate_forge_admin_arg`] before /// the admin CLI ever sees it. async fn exec_forge_admin(args: &[String]) -> Result<(String, String)> { for arg in args { validate_forge_admin_arg(arg)?; } run_forge_admin(args).await } /// `StopContainer`. async fn stop_container(name: &str) -> Result<(String, String)> { validate_agent_name(name)?; stop_and_release(&container_system_name(name)).await } /// `KillContainer`: `nixos-container` has no kill verb, so this uses /// `machinectl` to send `SIGKILL` to every process in the container — the /// right semantics for a forced shutdown after a graceful stop has already /// been attempted. async fn kill_container(name: &str) -> Result<(String, String)> { validate_agent_name(name)?; let machine = container_system_name(name); machinectl_run(&["kill", &machine, "--signal=SIGKILL"]).await } /// `DestroyContainer`. async fn destroy_container(name: &str) -> Result<(String, String)> { validate_agent_name(name)?; container_run(&["destroy", &container_system_name(name)]).await } /// `ReadContainerJournal`. async fn exec_read_container_journal( container: &str, query: &JournalQuery, ) -> Result<(String, String)> { validate_container_system_name(container)?; read_container_journal(container, query).await } /// `SetAgentPaused`. fn exec_set_agent_paused(agent_name: &str, paused: bool) -> Result<(String, String)> { validate_agent_name(agent_name)?; set_agent_paused(agent_name, paused) } /// `EnsureAgentSubvolume`. async fn exec_ensure_agent_subvolume(agent_name: &str) -> Result<(String, String)> { validate_agent_name(agent_name)?; ensure_agent_subvolume(agent_name).await } /// `DeleteAgentSubvolume`. async fn exec_delete_agent_subvolume(agent_name: &str) -> Result<(String, String)> { validate_agent_name(agent_name)?; delete_agent_subvolume(agent_name).await } /// `ReadSubvolumeUsage`. async fn exec_read_subvolume_usage(agent_name: &str) -> Result<(String, String)> { validate_agent_name(agent_name)?; read_subvolume_usage(agent_name).await } /// `SetSubvolumeQuota`. async fn exec_set_subvolume_quota( agent_name: &str, limit_bytes: Option, ) -> Result<(String, String)> { validate_agent_name(agent_name)?; set_subvolume_quota(agent_name, limit_bytes).await } /// `UpgradeAgentSubvolume`. async fn exec_upgrade_agent_subvolume(agent_name: &str) -> Result<(String, String)> { validate_agent_name(agent_name)?; upgrade_agent_subvolume(agent_name).await } /// `SnapshotAgentSubvolume`. async fn exec_snapshot_agent_subvolume( agent_name: &str, snapshot_name: &str, ) -> Result<(String, String)> { validate_agent_name(agent_name)?; validate_snapshot_name(snapshot_name)?; snapshot_agent_subvolume(agent_name, snapshot_name).await } /// `DeleteAgentSnapshot`. async fn exec_delete_agent_snapshot( agent_name: &str, snapshot_name: &str, ) -> Result<(String, String)> { validate_agent_name(agent_name)?; validate_snapshot_name(snapshot_name)?; delete_agent_snapshot(agent_name, snapshot_name).await } /// `SendAgentSnapshotToFile`. async fn exec_send_agent_snapshot_to_file( agent_name: &str, snapshot_name: &str, parent_snapshot_name: Option<&str>, dest_file_name: &str, ) -> Result<(String, String)> { validate_agent_name(agent_name)?; validate_snapshot_name(snapshot_name)?; if let Some(parent) = parent_snapshot_name { validate_snapshot_name(parent)?; } validate_credential_name(dest_file_name)?; send_agent_snapshot_to_file( agent_name, snapshot_name, parent_snapshot_name, dest_file_name, ) .await } /// `SendAgentSnapshotToFd`. async fn exec_send_agent_snapshot_to_fd( agent_name: &str, snapshot_name: &str, parent_snapshot_name: Option<&str>, fd: Option, ) -> Result<(String, String)> { validate_agent_name(agent_name)?; validate_snapshot_name(snapshot_name)?; if let Some(parent) = parent_snapshot_name { validate_snapshot_name(parent)?; } let dest = fd.context("no descriptor to stream into")?; send_agent_snapshot_to_fd(agent_name, snapshot_name, parent_snapshot_name, dest).await } /// Shared body for `CreateContainer` / `UpdateContainer`: validate the /// name, build the toplevel ourselves, and apply it (streaming line /// events to `writer` when `stream` is set). /// /// **Both verbs build explicitly now — not just `update`.** The first cut /// of this fix only rewrote `update`, reasoning that `create` was already /// safe: it wraps its whole action in an exclusive `flock` before calling /// `nixos-container`'s own `buildFlake()`, and once `update` stopped /// writing to `buildFlake()`'s shared `.tmp` out-link, concurrent /// `create`s were the only remaining writers — mutually excluded by that /// lock. True, but it leaves `create`'s safety resting on an internal /// implementation detail of a script we don't own (its current locking /// behavior, which could change upstream without notice) instead of on /// something we control. Building here for both verbs removes `buildFlake()` /// from the picture entirely — there's no shared `.tmp` left to race on, /// so there's nothing left to reason about staying in sync with. /// /// **`update` no longer calls `nixos-container update` at all.** That /// action's own version-compat probe runs unconditionally before /// `--system-path` is ever honored, dying on every agent's update. /// [`swap_container_profile`] replicates exactly what `update`'s own /// action does *past* that probe (confirmed against its source): `nix-env /// --set` the per-container profile, then `systemctl reload` if the /// container is running. `create` is untouched — it isn't the failing /// verb. async fn container_flake_action( verb: &str, name: &str, stream: bool, writer: &mut OwnedWriteHalf, ) -> Result<(String, String)> { validate_agent_name(name)?; // The build is the multi-minute phase of this operation — give it the // same live-line treatment `container_run_streaming` gives // `nixos-container` itself when the caller asked for it. Without // this, moving the build out of the streamed `nixos-container` call // (which is the whole point of this fix) would silently regress every // UI that shows build progress: nothing until the longest phase // finishes, then everything at once. let toplevel = nix_build_toplevel(name, stream.then_some(&mut *writer)).await?; let system_name = container_system_name(name); if verb == "update" { return swap_container_profile(&system_name, &toplevel, stream, writer).await; } let args = [verb, &system_name, "--system-path", &toplevel]; if stream { container_run_streaming(&args, writer).await } else { container_run(&args).await } } /// Apply a prebuilt toplevel to an existing container directly, without /// going through `nixos-container update` (see `container_flake_action`'s /// doc comment for why). Mirrors `nixos-container.pl`'s own `update` /// action verbatim, past its version probe: point the per-container `nix- /// env` profile at the new toplevel, then reload the container unit *if* /// it's currently running — the container's own next start already reads /// from the profile, so a stopped container needs nothing further. /// /// Both steps are near-instant in the happy case, but still stream- /// forwarded like the rest of this operation: a failure here (e.g. a /// wedged `systemctl reload`) is exactly when the caller most wants the /// live line in the build log, not just a summary error afterward. async fn swap_container_profile( system_name: &str, toplevel: &str, stream: bool, writer: &mut OwnedWriteHalf, ) -> Result<(String, String)> { let profile = format!("/nix/var/nix/profiles/per-container/{system_name}/system"); let set_out = Command::new("nix-env") .args(["-p", &profile, "--set", toplevel]) .output() .await .context("invoke nix-env --set")?; let mut stdout = String::from_utf8_lossy(&set_out.stdout).into_owned(); let mut stderr = String::from_utf8_lossy(&set_out.stderr).into_owned(); log_and_forward(&stdout, &stderr, stream, writer).await; if !set_out.status.success() { bail!( "nix-env -p {profile} --set failed ({}): {}", set_out.status, stderr.trim() ); } let unit = format!("container@{system_name}"); let state_out = Command::new("systemctl") .args(["show", "--property=ActiveState", "--value", &unit]) .output() .await .context("query container ActiveState")?; let active = String::from_utf8_lossy(&state_out.stdout).trim() == "active"; if active { let reload_out = Command::new("systemctl") .args(["reload", &unit]) .output() .await .context("invoke systemctl reload")?; let r_stdout = String::from_utf8_lossy(&reload_out.stdout).into_owned(); let r_stderr = String::from_utf8_lossy(&reload_out.stderr).into_owned(); log_and_forward(&r_stdout, &r_stderr, stream, writer).await; if !reload_out.status.success() { bail!( "systemctl reload {unit} failed ({}): {}", reload_out.status, r_stderr.trim() ); } stdout.push_str(&r_stdout); stderr.push_str(&r_stderr); } Ok((stdout, stderr)) } /// Log a command's captured stdout/stderr the same way every other /// `nixos-container`-adjacent shellout in this file does, and — when /// `stream` is set — forward each line to the client as a live /// [`PrivEvent::Line`] too, so a caller watching the build log sees these /// lines exactly like any other step's, not just a summary on failure. async fn log_and_forward(stdout: &str, stderr: &str, stream: bool, writer: &mut OwnedWriteHalf) { for line in stdout.lines() { tracing::info!(target: "nixos-container", "{line}"); if stream { write_line_event(writer, PrivStream::Stdout, line).await; } } for line in stderr.lines() { tracing::warn!(target: "nixos-container", "{line}"); if stream { write_line_event(writer, PrivStream::Stderr, line).await; } } } /// The explicit `nixosConfigurations..config.system.build.toplevel` /// flake attr path — same construction `hive-c0re`'s own /// `lifecycle::prebuild_toplevel` uses, kept here as a pure function so /// the exact string shape is unit-tested without needing to run `nix`. fn toplevel_attr(name: &str) -> String { format!("{META_DIR}#nixosConfigurations.{name}.config.system.build.toplevel") } /// Wall-clock bound on [`nix_build_toplevel`]. Four times CI's observed /// cold-cache `nix flake check` (~15 min, `.forgejo/workflows/ci.yml`), which /// compiles this workspace's crates, the part of an agent toplevel no binary /// cache holds; the rebuild path has usually warmed this exact attr already /// (`hive-c0re`'s `prebuild_toplevel`). A stalled download is bounded by nix's /// own `stalled-download-timeout`; this covers the rest, such as a wedged /// nix-daemon holding the caller's build slot forever. const NIX_BUILD_TIMEOUT: std::time::Duration = std::time::Duration::from_hours(1); /// How a [`run_bounded`] child ended. enum BoundedRun { Exited { status: std::process::ExitStatus, stdout: String, stderr: String, }, /// Still running at the limit; its whole process group was killed. TimedOut, } /// Run `cmd` to completion or until `limit`, draining both pipes: stdout is /// captured silently, stderr is logged and, when `writer` is set, forwarded /// live as [`PrivEvent::Line`]s. /// /// The child leads its own process group and a timeout `SIGKILL`s the whole /// group, so a grandchild still holding the pipes (or the daemon connection) /// dies with it instead of outliving the error. async fn run_bounded( mut cmd: Command, limit: std::time::Duration, mut writer: Option<&mut OwnedWriteHalf>, ) -> std::io::Result { use tokio::io::AsyncBufReadExt as _; let mut child = cmd .stdout(std::process::Stdio::piped()) .stderr(std::process::Stdio::piped()) .process_group(0) .spawn()?; // Group leader, so its pid is the pgid. let pgid = child.id().and_then(|pid| i32::try_from(pid).ok()); let stdout = child.stdout.take().expect("stdout piped"); let stderr = child.stderr.take().expect("stderr piped"); let run = async { let mut stdout_lines = BufReader::new(stdout).lines(); let mut stderr_lines = BufReader::new(stderr).lines(); let mut stdout_buf = String::new(); let mut stderr_buf = String::new(); // ⚠️ Both pipes are drained concurrently even though only one is // streamed: reading stderr alone would let stdout fill its pipe // buffer and deadlock the child on a build with enough stdout output // to fill it. // // ⚠️ Each stream's EOF is tracked separately rather than breaking on // the first `None`: `next_line()` on a closed stream returns // `Ok(None)` immediately and forever, so a loop that keeps polling a // finished stream spins hot until the other one ends too. let mut stdout_done = false; let mut stderr_done = false; while !(stdout_done && stderr_done) { tokio::select! { line = stdout_lines.next_line(), if !stdout_done => { match line { // Captured, never streamed — this is the store path. Ok(Some(l)) => { stdout_buf.push_str(&l); stdout_buf.push('\n'); } Ok(None) => stdout_done = true, Err(e) => { tracing::warn!(error = %e, "nix build stdout read error"); stdout_done = true; } } } line = stderr_lines.next_line(), if !stderr_done => { match line { // Streamed as it arrives, so the dashboard and // `journalctl -f` show build progress live. Ok(Some(l)) => { tracing::info!(target: "nix-build-toplevel", "{l}"); if let Some(w) = writer.as_deref_mut() { write_line_event(w, PrivStream::Stderr, &l).await; } if !stderr_buf.is_empty() { stderr_buf.push('\n'); } stderr_buf.push_str(&l); } Ok(None) => stderr_done = true, Err(e) => { tracing::warn!(error = %e, "nix build stderr read error"); stderr_done = true; } } } } } let status = child.wait().await?; Ok::<_, std::io::Error>(BoundedRun::Exited { status, stdout: stdout_buf, stderr: stderr_buf, }) }; if let Ok(exited) = tokio::time::timeout(limit, run).await { return exited; } if let Some(pgid) = pgid { // SAFETY: plain `kill(2)`; `-pgid` targets the child's own process // group, which `process_group(0)` made it lead. unsafe { libc::kill(-pgid, libc::SIGKILL); } } child.wait().await?; Ok(BoundedRun::TimedOut) } /// Build `nixosConfigurations..config.system.build.toplevel` and /// return the resulting store path, so `create`/`update` can hand /// `nixos-container` an explicit `--system-path` instead of letting its /// own `buildFlake()` build to a racy shared out-link. See "Container /// toplevel builds" in this crate's README for the full story — the /// concurrency bug this closes (the "agent container gets closure of /// other agent" mystery bug) and why stdout/stderr are drained /// concurrently but handled asymmetrically (stderr streamed live, /// stdout captured and required to be exactly one line). /// /// `writer` is `None` for the non-streaming call shape (`stream: false`); /// stderr still logs to journald either way, just without the /// `PrivEvent::Line` forwarding. Bounded by [`NIX_BUILD_TIMEOUT`]. async fn nix_build_toplevel(name: &str, writer: Option<&mut OwnedWriteHalf>) -> Result { let attr = toplevel_attr(name); let mut cmd = Command::new("nix"); cmd.args([ "--extra-experimental-features", "nix-command flakes", "build", "--no-link", "--print-out-paths", &attr, ]); let (status, stdout_buf, stderr_buf) = match run_bounded(cmd, NIX_BUILD_TIMEOUT, writer) .await .with_context(|| format!("run nix build {attr}"))? { BoundedRun::Exited { status, stdout, stderr, } => (status, stdout, stderr), BoundedRun::TimedOut => bail!( "nix build {attr} timed out after {}s; killed it", NIX_BUILD_TIMEOUT.as_secs() ), }; // ⚠️ Success is decided by the exit status, not by "we parsed a // path" — a build can print to stdout and still fail. if !status.success() { bail!( "nix build {attr} failed ({status}): {}", stderr_buf.lines().last().unwrap_or("").trim() ); } single_output_path(&stdout_buf) .map(str::to_owned) .map_err(|count| { anyhow!("nix build {attr} produced {count} output path(s), expected exactly 1: {stdout_buf:?}") }) } /// Parse `nix build --print-out-paths`' stdout down to the single output /// path this function's caller expects. `--print-out-paths` prints one /// line *per output*, not one line total (`nix build --no-link /// --print-out-paths nixpkgs#openssl` prints two: `…-bin`, `…-man`); /// `config.system.build.toplevel` is single-output today, so this is one /// line in practice — but a bare whole-buffer `.trim()` would silently /// hand a multi-line string on to `--system-path` the day that ever /// changes, the same corrupted-argument failure this function exists to /// avoid. Trims and drops empty lines *before* counting, so a lone /// `"\n"` (or trailing whitespace on the real line) can't be mistaken /// for a present-but-blank path — see the unit tests below for the exact /// table this closes. `Err` carries the surviving line count, for the /// caller's error message. fn single_output_path(stdout: &str) -> Result<&str, usize> { let lines: Vec<&str> = stdout .lines() .map(str::trim) .filter(|l| !l.is_empty()) .collect(); match lines[..] { [path] => Ok(path), _ => Err(lines.len()), } } /// `WriteNspawnFlags` — validate the container + every bind path + every /// credential entry, then write the container's nspawn flag overrides. fn handle_write_nspawn_flags( container: &str, binds: &[BindMount], isolation: &NetworkIsolation, load_credentials: &[CredentialMount], ) -> Result<(String, String)> { validate_container_system_name(container)?; for bind in binds { validate_bind_path(&bind.host_path)?; validate_bind_path(&bind.container_path)?; } for cred in load_credentials { validate_credential_name(&cred.name)?; // Same path rules as binds (absolute, no colon/newline/quote/null): // the colon ban is essential since `--load-credential=name:path` // uses `:` as the name/path separator. validate_bind_path(&cred.host_path)?; } write_nspawn_flags(container, binds, isolation, load_credentials)?; Ok((String::new(), String::new())) } /// A btrfs snapshot label must start with `hive-` — this doubles as an /// allow-list: only names hivectl itself constructs (or an operator who /// knows the convention) can reach the `btrfs subvolume snapshot`/`delete` /// shellouts, so an arbitrary caller can't use the snapshot ops to probe or /// churn unrelated paths under `AGENT_STATE_ROOT`. Beyond the prefix, the /// same charset restriction as [`validate_credential_name`] applies (it's /// interpolated straight into a filesystem path). fn validate_snapshot_name(name: &str) -> Result<()> { if !name.starts_with("hive-") { bail!("invalid snapshot label {name:?}: must start with \"hive-\""); } validate_credential_name(name) } /// A systemd credential id must be a short token — restrict to /// `[A-Za-z0-9_-]` (no `.`) so it can't inject extra `--load-credential` /// argv or break the `name:path` shape. `.` is deliberately excluded, not /// just a bare `..`: this name gets interpolated into filesystem paths /// (snapshot labels) and there's no legitimate need for a dot in either a /// systemd credential id or a `hive-`-prefixed snapshot label — we're /// defining this token format from scratch, so keep it maximally strict /// rather than allow-then-patch each traversal-adjacent character /// (mara: "we are making up the rules here, lets go strict"). fn validate_credential_name(name: &str) -> Result<()> { if name.is_empty() || !name .bytes() .all(|b| b.is_ascii_alphanumeric() || matches!(b, b'_' | b'-')) { bail!("invalid credential name {name:?}: must be non-empty [A-Za-z0-9_-]"); } Ok(()) } /// `RemoveServiceDropin` — remove the container service's drop-in dir /// if present (idempotent). fn remove_service_dropin(container: &str) -> Result<(String, String)> { validate_container_system_name(container)?; let dir = format!("/run/systemd/system/container@{container}.service.d"); if Path::new(&dir).exists() { std::fs::remove_dir_all(&dir).with_context(|| format!("remove {dir}"))?; } Ok((String::new(), String::new())) } /// `WriteResourceLimits` — drop the systemd resource settings into the /// container service's drop-in dir, together with a /// `ConditionPathIsDirectory=` guard on the agent's MCP runtime dir. /// /// Two different kinds of setting land in the same file. `MemoryMax=` / /// `CPUQuota=` are hard caps that throttle even on an idle host; /// `CPUWeight=` / `IOWeight=` are cgroup v2 relative shares that only /// decide who yields *under contention*. A weight of `None` means "not /// configured" and omits the line, so a hive-c0re built before the weights /// existed keeps producing the old two-line drop-in. /// /// 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 /// hive-c0re creating the dir in its start preamble: if it is missing at /// start time, the unit idles rather than restart-looping into /// `start-limit-hit`. fn write_resource_limits( container: &str, memory_max: &str, cpu_quota: &str, cpu_weight: Option, io_weight: Option, ) -> Result<(String, String)> { 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"); std::fs::create_dir_all(&dir).with_context(|| format!("create {dir}"))?; let path = format!("{dir}/hyperhive-limits.conf"); let content = limits_dropin_body(&runtime_dir, memory_max, cpu_quota, cpu_weight, io_weight); // Published whole: a `daemon-reload` between truncate and write would // otherwise load the unit with its limits and start guard missing. publish_file(Path::new(&path), content.as_bytes(), 0o644, None)?; Ok((String::new(), String::new())) } /// How long a window the start-limit counts over, and how many starts it /// allows inside it. /// /// `container@.service` sets `Restart=on-failure` and **no** start limit, so /// systemd's defaults apply: 5 starts per 10s, `RestartSec` 100ms. That /// makes the bound depend on *how fast* a container dies — one that fails /// instantly trips the limit in under a second, one that takes longer than /// ~2s never trips it and restarts forever. Whether an agent gets bounded /// is not meant to be a function of its failure speed. /// /// The window has to exceed the worst-case time to burn the burst, or the /// counter ages out between attempts and the limit is again unreachable: /// `TimeoutStartSec` is 1min, so `BURST` slow failures plus their backoff /// can span several minutes. 10min covers that with room. /// /// Giving up is cheap here **because it is not terminal** — hive-c0re's /// reconcile sweep retries later, and `reset-failed` (see `StartContainer`) /// clears the latch first. That is what makes a tight burst safe. const START_LIMIT_INTERVAL_SEC: u32 = 600; /// One start plus two retries — the operator's ruling was "retry once or /// twice", with the reconcile sweep as the slow path after that. const START_LIMIT_BURST: u32 = 3; /// Backoff between those retries. The 100ms default is for processes that /// respawn instantly; a container that just failed to boot gains nothing /// from being retried a tenth of a second later. const RESTART_SEC: u32 = 5; /// Render the body of `hyperhive-limits.conf`. /// /// `[Unit]`: the condition is checked at start time — it skips (not fails) /// the unit when the MCP socket dir is absent, avoiding restart loops — /// plus the bounded start limit (see the constants above; `StartLimit*` are /// `[Unit]` settings since systemd 229 and are silently ignored under /// `[Service]`). /// `[Service]`: the restart backoff, the hard caps, then the relative /// weights. A weight of `None` means "not configured" and omits its line /// entirely, so a request from a hive-c0re built before the weights /// existed — or one whose nix option is `null` — renders no weight lines. fn limits_dropin_body( runtime_dir: &str, memory_max: &str, cpu_quota: &str, cpu_weight: Option, io_weight: Option, ) -> String { // Built as two possibly-empty lines rather than pushed onto the // string: `format!` appended to a `String` trips clippy::pedantic's // `format_push_string`, and a `write!` would need an unwrap. let cpu_weight_line = cpu_weight.map_or_else(String::new, |w| format!("CPUWeight={w}\n")); let io_weight_line = io_weight.map_or_else(String::new, |w| format!("IOWeight={w}\n")); format!( "[Unit]\n\ ConditionPathIsDirectory={runtime_dir}\n\ StartLimitIntervalSec={START_LIMIT_INTERVAL_SEC}\n\ StartLimitBurst={START_LIMIT_BURST}\n\ \n\ [Service]\n\ RestartSec={RESTART_SEC}\n\ MemoryMax={memory_max}\n\ CPUQuota={cpu_quota}\n\ {cpu_weight_line}{io_weight_line}" ) } /// `DaemonReload` — `systemctl daemon-reload` on the host. async fn daemon_reload() -> Result<(String, String)> { let out = Command::new("systemctl") .arg("daemon-reload") .output() .await .context("invoke systemctl daemon-reload")?; if !out.status.success() { bail!( "systemctl daemon-reload failed ({}): {}", out.status, String::from_utf8_lossy(&out.stderr).trim() ); } Ok((String::new(), String::new())) } /// Host path to the hive-ci runner's persisted registration credentials. /// /// Paired with `hive-c0re`'s `forge::ci_runner::RUNNER_FILE`, which reads the /// same file to decide whether a runner is registered and whether it still /// names the configured forge host. Deliberately duplicated rather than shared: /// `hive-priv` is the minimal root helper and does not depend on `hive-c0re`. const RUNNER_CREDENTIALS: &str = "/var/lib/nixos-containers/hive-ci/var/lib/gitea-runner/hive/.runner"; /// Delete the runner's persisted credentials so upstream's `ExecStartPre` takes /// its **absence** branch on the next start. /// /// Absence is the state we want, so `NotFound` is success. Anything else — a /// permission error above all — is NOT swallowed: it means the file is still /// there, the restart will take upstream's already-registered branch, and the /// caller would return `Ok` for a registration that never happened. That is the /// same shape as a precondition that "passes" because it could not read the file /// it was checking, and it is worth failing loudly to avoid. /// /// Split from [`register_ci_runner`] purely so this rule is testable without a /// container or a `systemctl`. fn clear_runner_credentials(path: &str) -> Result<()> { match std::fs::remove_file(path) { Ok(()) => Ok(()), Err(e) if e.kind() == std::io::ErrorKind::NotFound => Ok(()), Err(e) => { Err(anyhow::Error::new(e).context(format!("remove stale runner credentials {path}"))) } } } /// `RegisterCiRunner` — write the runner registration token to the host-side /// `/run/hive-ci/runner-token` env-file, then restart the in-container runner /// so it re-registers. The forge admin token never enters the container; only /// the registration token c0re passes here is written, and it lands on a host /// path bind-mounted read-only into hive-ci. async fn register_ci_runner(token: &str) -> Result<(String, String)> { use std::io::Write as _; use std::os::unix::fs::{OpenOptionsExt as _, PermissionsExt as _}; // Reject anything that could corrupt the `KEY=VALUE` env-file or smuggle a // second line — a forge registration token is an opaque single-line string. if token.is_empty() || token.contains(['\n', '\r', '\0']) { bail!("ci runner registration token empty or contains control characters"); } let token_path = "/run/hive-ci/runner-token"; // In-place truncate+write of the existing inode (mirrors the prefetch's // `echo > $FILE`), NOT a temp+rename: nspawn pins this file's inode into // hive-ci at container start, so a rename would leave the running runner // reading the old content. Format + perms match the tmpfiles seed and the // prefetch: `TOKEN=`, mode 0600, root-owned. Created 0600 when the // seed is missing, so the token is never on disk at a wider mode. std::fs::OpenOptions::new() .write(true) .create(true) .truncate(true) .mode(0o600) .open(token_path) .and_then(|mut f| f.write_all(format!("TOKEN={token}\n").as_bytes())) .with_context(|| format!("write {token_path}"))?; std::fs::set_permissions(token_path, std::fs::Permissions::from_mode(0o600)) .with_context(|| format!("chmod {token_path}"))?; // Remove the persisted credentials BEFORE restarting, or the restart is a // no-op as far as registration goes. // // Upstream's `ExecStartPre` only re-registers when `.runner` is absent, the // labels changed, or the *registration token hash* changed — never when the // instance URL changed. c0re only calls this helper once it has already // decided the existing credentials are absent or stale // (`forge::ci_runner::ensure_ci_runner_registered` returns early otherwise), // so by the time we are here a re-registration is exactly what is wanted and // deleting the file is the narrow way to guarantee it happens. // // Writing a fresh token is NOT sufficient on its own: whether the hash // changes depends on whether the forge mints a new registration token per // request or hands back a stable one, which is Forgejo's behaviour to // choose and change. Gating our remediation on the absence branch — the one // upstream evaluates unconditionally — makes that question moot instead of // load-bearing. // // See [`clear_runner_credentials`] for why absence is the branch we aim at // and why only `NotFound` counts as success. clear_runner_credentials(RUNNER_CREDENTIALS)?; // Restart the in-container runner so it reads the new token and registers. let out = Command::new("systemctl") .args(["--machine=hive-ci", "restart", "gitea-runner-hive.service"]) .output() .await .context("systemctl restart gitea-runner-hive.service in hive-ci")?; if !out.status.success() { bail!( "systemctl restart gitea-runner-hive.service in hive-ci exited {}: {}", out.status, String::from_utf8_lossy(&out.stderr).trim() ); } Ok(( String::from_utf8_lossy(&out.stdout).into_owned(), String::from_utf8_lossy(&out.stderr).into_owned(), )) } /// `ControlInfraContainer` — start/stop a hive infrastructure service via /// `systemctl `. The [`InfraContainer`] enum is the allowlist: /// serde already rejected any unknown / unsafe name (hive-c0re has no /// variant, so a stop can't sever the daemon socket) at deserialisation, so /// no root-side `.contains()` check is needed here. Serves the hive-wide /// `hivectl stop`/`start` flow and the dashboard's operator-only infra /// panel. No agent-facing path exists. /// /// ⚠️ The unit is derived from the variant, never sent by the caller — /// which is what keeps this from being a general `systemctl` pass-through. /// It is not always `container@.service`: the gateway resolves to the /// host's `nginx.service`. async fn control_infra_container( container: InfraContainer, action: InfraAction, ) -> Result<(String, String)> { let verb = action.systemctl_verb(); let unit = container.service_unit(); let out = Command::new("systemctl") .args([verb, &unit]) .output() .await .with_context(|| format!("systemctl {verb} {unit}"))?; if !out.status.success() { bail!( "systemctl {verb} {unit} exited {}: {}", out.status, String::from_utf8_lossy(&out.stderr).trim() ); } tracing::info!(target: "infra-control", "{verb} {unit}"); Ok(( String::from_utf8_lossy(&out.stdout).into_owned(), String::from_utf8_lossy(&out.stderr).into_owned(), )) } /// A leaf safe to join onto an agent-owned directory: one plain component, so /// it can neither climb out of the directory nor name the directory itself. fn ensure_plain_filename(who: &str, filename: &str) -> Result<()> { if filename.is_empty() || filename == "." || filename == ".." || filename.contains('/') { bail!("{who}: refusing non-plain filename {filename:?}"); } Ok(()) } /// Name the content is written under before being renamed onto `filename`, /// unique per call so two concurrent writes of one file never share a temp. /// /// The leading dot and the `.partial` suffix are load-bearing: a reader /// globbing on a published name's prefix, or systemd reading `*.conf` out of /// `tmpfiles.d` and drop-in dirs, cannot pick up a half-written temp. fn partial_name(filename: &str) -> String { use std::sync::atomic::{AtomicU64, Ordering}; static SEQ: AtomicU64 = AtomicU64::new(0); let n = SEQ.fetch_add(1, Ordering::Relaxed); format!(".{filename}.{}-{n}.partial", std::process::id()) } /// Open `path` as a directory fd for [`StagedFile`] to work relative to. fn open_dir(path: &Path) -> Result { use std::os::unix::fs::OpenOptionsExt as _; std::fs::OpenOptions::new() .read(true) .custom_flags(libc::O_DIRECTORY) .open(path) .with_context(|| format!("open directory {}", path.display())) } /// A file written under a [`partial_name`] in its destination directory and /// then renamed onto its final name, so a reader sees either the old file or /// the complete new one, never a partial or wrong-mode one. Mode, and owner /// when the caller `fchown`s [`file`](Self::file), are on the inode before the /// rename makes it visible. /// /// Every step is relative to the `dir` fd, so it is safe in a directory an /// agent or container controls: the temp is created `O_EXCL|O_NOFOLLOW`, which /// refuses any entry already planted at its name, and `rename` replaces /// whatever sits at the final name without following it. /// /// Dropped without [`publish`](Self::publish), it unlinks the temp, so a /// failure between the two leaves the old file untouched and no residue. struct StagedFile<'a> { dir: &'a std::fs::File, tmp: std::ffi::CString, dest: std::ffi::CString, dest_path: PathBuf, file: std::fs::File, published: bool, } impl<'a> StagedFile<'a> { /// Create the temp at `mode` and write `content` to it. The mode is set /// again on the fd after the write, so the umask cannot change it. fn write( dir: &'a std::fs::File, dir_path: &Path, filename: &str, content: &[u8], mode: u32, ) -> Result { use std::io::Write as _; use std::os::unix::fs::PermissionsExt as _; ensure_plain_filename("StagedFile::write", filename)?; let dest_path = dir_path.join(filename); let tmp_name = partial_name(filename); let tmp = std::ffi::CString::new(tmp_name.as_str()) .with_context(|| format!("NUL in filename {filename:?}"))?; let dest = std::ffi::CString::new(filename) .with_context(|| format!("NUL in filename {filename:?}"))?; // SAFETY: `dir` is an open directory fd and `tmp` is NUL-terminated; // the result is checked before it is wrapped. let fd = unsafe { libc::openat( dir.as_raw_fd(), tmp.as_ptr(), libc::O_WRONLY | libc::O_CREAT | libc::O_EXCL | libc::O_NOFOLLOW | libc::O_CLOEXEC, mode, ) }; if fd < 0 { return Err(std::io::Error::last_os_error()) .with_context(|| format!("create {}", dir_path.join(&tmp_name).display())); } // SAFETY: `fd` is a fresh, valid fd that nothing else owns. let file = std::fs::File::from(unsafe { OwnedFd::from_raw_fd(fd) }); let mut staged = Self { dir, tmp, dest, dest_path, file, published: false, }; staged .file .write_all(content) .with_context(|| format!("write {}", staged.dest_path.display()))?; staged .file .set_permissions(std::fs::Permissions::from_mode(mode)) .with_context(|| format!("chmod {mode:o} {}", staged.dest_path.display()))?; Ok(staged) } /// The temp's open fd, for an `fchown` that must land before publishing. fn file(&self) -> &std::fs::File { &self.file } /// Flush the content, rename the temp onto the final name, and flush the /// directory: without the last step a crash can lose the rename itself. fn publish(mut self) -> Result<()> { self.file .sync_all() .with_context(|| format!("fsync {}", self.dest_path.display()))?; let fd = self.dir.as_raw_fd(); // SAFETY: `fd` is an open directory fd; both names are NUL-terminated. if unsafe { libc::renameat(fd, self.tmp.as_ptr(), fd, self.dest.as_ptr()) } != 0 { return Err(std::io::Error::last_os_error()) .with_context(|| format!("publishing {}", self.dest_path.display())); } self.published = true; self.dir .sync_all() .with_context(|| format!("fsync directory of {}", self.dest_path.display())) } } impl Drop for StagedFile<'_> { fn drop(&mut self) { if self.published { return; } // SAFETY: `self.dir` is an open directory fd; `self.tmp` is // NUL-terminated. if unsafe { libc::unlinkat(self.dir.as_raw_fd(), self.tmp.as_ptr(), 0) } != 0 { tracing::warn!( dest = %self.dest_path.display(), error = %std::io::Error::last_os_error(), "failed to remove unpublished temp" ); } } } /// Publish `content` at `path` through a [`StagedFile`], at `mode` and, when /// given, owned by `owner` (`(uid, gid)`). fn publish_file(path: &Path, content: &[u8], mode: u32, owner: Option<(u32, u32)>) -> Result<()> { let (Some(dir_path), Some(filename)) = (path.parent(), path.file_name().and_then(|n| n.to_str())) else { bail!( "publish_file: {} has no directory or file name", path.display() ); }; let dir = open_dir(dir_path)?; let staged = StagedFile::write(&dir, dir_path, filename, content, mode)?; if let Some((uid, gid)) = owner { std::os::unix::fs::fchown(staged.file(), Some(uid), Some(gid)) .with_context(|| format!("chown {}", path.display()))?; } staged.publish() } /// Create or remove an agent's pause marker under its harness dir. The /// marker is written empty and chowned to the harness dir's owner (the /// agent), matching how the harness itself would have created it. /// /// Both directions are idempotent: re-pausing truncates the existing empty /// marker rather than failing, and a `NotFound` on removal is the /// already-resumed case, not an error. fn set_agent_paused(agent_name: &str, paused: bool) -> Result<(String, String)> { let harness_dir = PathBuf::from(AGENT_STATE_ROOT) .join(agent_name) .join("harness"); if paused { return write_agent_dir_file(agent_name, &harness_dir, PAUSED_MARKER_FILE, ""); } remove_marker_in(&harness_dir, PAUSED_MARKER_FILE)?; tracing::info!(agent = %agent_name, "cleared pause marker"); Ok((String::new(), String::new())) } /// Unlink `dir/filename`, treating "already gone" as success. /// /// `remove_file` unlinks the leaf itself and never follows a symlink, so an /// agent-planted link at the marker path cannot redirect this root unlink /// at another file — the same threat [`StagedFile`] closes on the create side. fn remove_marker_in(dir: &Path, filename: &str) -> Result<()> { let path = dir.join(filename); match std::fs::remove_file(&path) { Ok(()) => Ok(()), Err(e) if e.kind() == std::io::ErrorKind::NotFound => Ok(()), Err(e) => Err(e).with_context(|| format!("remove {}", path.display())), } } /// Write `content` to `dir/filename` as root, chowning the result to `dir`'s /// owner so the agent process can read it back. Its caller is the pause /// marker, in `harness/`, a directory owned by the agent, which is precisely /// why it needs hive-priv at all. /// /// The file is published through a [`StagedFile`], so a reader woken by its /// appearance cannot catch it empty or half-written: several of these paths /// have a `systemd.path` unit watching them, and one of those triggers on the /// file existing rather than changing. fn write_agent_dir_file( agent_name: &str, dir: &Path, filename: &str, content: &str, ) -> Result<(String, String)> { use std::os::fd::AsRawFd as _; use std::os::unix::fs::MetadataExt as _; // Checked before `create_dir_all`, so a bad leaf does no filesystem work. ensure_plain_filename("write_agent_dir_file", filename)?; let state_dir = dir.to_path_buf(); // NOTE: `create_dir_all` is normally a no-op — lifecycle creates and chowns // the state dir during spawn. On the rare edge where the dir doesn't exist // yet (container being provisioned for the first time), the newly created dir // is root:root. The `stat state_dir` chown below will then see uid=0 and // leave the file root-owned (0600). The agent won't be able to read it until // its lifecycle completes. std::fs::create_dir_all(&state_dir) .with_context(|| format!("create state dir {}", state_dir.display()))?; // Security-critical: the (state/-owning) agent may plant a symlink or // hardlink in this dir, and `StagedFile` neither follows nor reuses one, // so this root-privileged create/write/chmod/chown can't be redirected at // an arbitrary file. let path = state_dir.join(filename); let dir_fd = open_dir(&state_dir)?; let staged = StagedFile::write(&dir_fd, &state_dir, filename, content.as_bytes(), 0o600)?; // Chown to the state dir's owner so the agent process can read the file. // fchown on the temp's fd — TOCTOU-immune (the inode the write hit, never a // swapped path). If stat fails (e.g. dir just created, owner is root), the // file stays root-owned and 0600 — still unreadable by others, just not // agent-readable. Log a warning so operators can diagnose. match dir_fd.metadata() { Ok(meta) => { // SAFETY: the temp's fd is open and owned by `staged` for the whole // call; `fchown` only mutates that inode's uid/gid. let rc = unsafe { libc::fchown(staged.file().as_raw_fd(), meta.uid(), meta.gid()) }; if rc != 0 { let e = std::io::Error::last_os_error(); tracing::warn!( agent = %agent_name, path = %path.display(), error = %e, "write_agent_dir_file: fchown failed" ); } } Err(e) => { tracing::warn!( agent = %agent_name, error = %e, "write_agent_dir_file: stat state_dir failed, leaving file root-owned" ); } } // After the chown, never before: the file must never be visible under its // final name while still root-owned. A failed publish drops `staged`, // which removes the temp and the credential in it. staged.publish()?; tracing::info!( agent = %agent_name, dir = %state_dir.display(), file = %filename, "wrote agent file" ); Ok((String::new(), String::new())) } /// btrfs superblock magic, as reported by `statfs(2)`'s `f_type`. const BTRFS_SUPER_MAGIC: i64 = 0x9123_683E; /// Inode number of a btrfs subvolume root (`BTRFS_FIRST_FREE_OBJECTID`). /// Every subvolume's top directory has this inode; plain directories do /// not, so `statfs == btrfs && st_ino == 256` reliably identifies a /// subvolume root. const BTRFS_SUBVOL_ROOT_INO: u64 = 256; /// Whether `path` lives on a btrfs filesystem (via `statfs(2)`). fn is_on_btrfs(path: &Path) -> Result { use std::os::unix::ffi::OsStrExt as _; let c_path = std::ffi::CString::new(path.as_os_str().as_bytes()) .with_context(|| format!("path {} has an interior null byte", path.display()))?; // SAFETY: `c_path` is a valid NUL-terminated C string that outlives the // call; `statfs` only writes into the zero-initialised `buf`. let mut buf: libc::statfs = unsafe { std::mem::zeroed() }; let rc = unsafe { libc::statfs(c_path.as_ptr(), &raw mut buf) }; if rc != 0 { return Err(std::io::Error::last_os_error()) .with_context(|| format!("statfs {}", path.display())); } Ok(buf.f_type == BTRFS_SUPER_MAGIC) } /// Whether `path` is the root of a btrfs subvolume (on btrfs and inode 256). fn is_btrfs_subvolume(path: &Path) -> bool { use std::os::unix::fs::MetadataExt as _; let on_btrfs = is_on_btrfs(path).unwrap_or(false); let ino_match = std::fs::metadata(path).is_ok_and(|m| m.ino() == BTRFS_SUBVOL_ROOT_INO); on_btrfs && ino_match } /// `EnsureAgentSubvolume` — make the agent's state root a btrfs subvolume /// when the FS supports it. Idempotent + progressive: no-op when the root /// already exists or the FS isn't btrfs. See the wire doc on the variant. async fn ensure_agent_subvolume(agent_name: &str) -> Result<(String, String)> { use std::os::unix::fs::MetadataExt as _; let root = PathBuf::from(AGENT_STATE_ROOT); let agent_root = root.join(agent_name); // Progressive: existing agents (plain dir OR already a subvol) are left // untouched — never auto-migrated. if agent_root.exists() { return Ok((String::new(), String::new())); } // Only btrfs supports subvolumes; on anything else hive-c0re's normal // `create_dir_all` makes a plain dir (the pre-subvolume behaviour). The // parent must exist for both statfs and `btrfs subvolume create`. std::fs::create_dir_all(&root) .with_context(|| format!("create agents root {}", root.display()))?; if !is_on_btrfs(&root)? { return Ok((String::new(), String::new())); } let out = Command::new("btrfs") .args(["subvolume", "create"]) .arg(&agent_root) .output() .await .with_context(|| format!("spawn btrfs subvolume create {}", agent_root.display()))?; if !out.status.success() { bail!( "btrfs subvolume create {} failed: {}", agent_root.display(), String::from_utf8_lossy(&out.stderr).trim() ); } // The subvol root is created root-owned, but hive-c0re (the `hive-core` // user) must be able to mkdir state/ claude/ harness/ inside it — exactly // as it would in a plain dir. Chown it to AGENT_STATE_ROOT's owner // (hive-core). This MUST succeed: a root-owned subvol would make the // downstream dir creation fail with a confusing permission error, and the // c0re-side exists-check would then skip re-running this op on retry, // wedging the agent. So on any failure roll the subvol back and bail — the // create path surfaces a clear error and a retry starts clean. let chown_result = std::fs::metadata(&root) .with_context(|| format!("stat agents root {} for ownership", root.display())) .and_then(|meta| { std::os::unix::fs::chown(&agent_root, Some(meta.uid()), Some(meta.gid())).with_context( || format!("chown subvol {} to agents-root owner", agent_root.display()), ) }); if let Err(e) = chown_result { // Best-effort rollback so we never leave a root-owned subvol behind. let _ = Command::new("btrfs") .args(["subvolume", "delete"]) .arg(&agent_root) .output() .await; return Err(e.context(format!( "rolled back subvolume {} after chown failed", agent_root.display() ))); } tracing::info!(agent = %agent_name, path = %agent_root.display(), "created agent state subvolume"); Ok((String::new(), String::new())) } /// `DeleteAgentSubvolume` — delete the agent's state root iff it's a btrfs /// subvolume (purge path only). No-op for plain dirs / missing paths; /// hive-c0re's own `remove_dir_all` covers those. See the wire doc. async fn delete_agent_subvolume(agent_name: &str) -> Result<(String, String)> { let agent_root = PathBuf::from(AGENT_STATE_ROOT).join(agent_name); if !is_btrfs_subvolume(&agent_root) { return Ok((String::new(), String::new())); } let out = Command::new("btrfs") .args(["subvolume", "delete"]) .arg(&agent_root) .output() .await .with_context(|| format!("spawn btrfs subvolume delete {}", agent_root.display()))?; if !out.status.success() { bail!( "btrfs subvolume delete {} failed: {}", agent_root.display(), String::from_utf8_lossy(&out.stderr).trim() ); } tracing::info!(agent = %agent_name, path = %agent_root.display(), "deleted agent state subvolume"); Ok((String::new(), String::new())) } /// `EnsureBtrfsQuota` — enable btrfs qgroup accounting on the filesystem /// holding `AGENT_STATE_ROOT`. Idempotent + statfs-gated (no-op off btrfs). /// Operator opt-in only; see the wire doc. async fn ensure_btrfs_quota() -> Result<(String, String)> { let root = PathBuf::from(AGENT_STATE_ROOT); std::fs::create_dir_all(&root) .with_context(|| format!("create agents root {}", root.display()))?; if !is_on_btrfs(&root)? { // Non-btrfs host: quota/qgroups don't apply. No-op success so the // operator-facing verb degrades cleanly. return Ok(( String::new(), "filesystem is not btrfs — quota not applicable".to_owned(), )); } let out = Command::new("btrfs") .args(["quota", "enable"]) .arg(&root) .output() .await .with_context(|| format!("spawn btrfs quota enable {}", root.display()))?; if !out.status.success() { bail!( "btrfs quota enable {} failed: {}", root.display(), String::from_utf8_lossy(&out.stderr).trim() ); } tracing::info!(path = %root.display(), "enabled btrfs qgroup accounting"); Ok((String::new(), String::new())) } /// `ReadSubvolumeUsage` — return an agent subvolume's qgroup rows /// (`btrfs qgroup show -f --raw <…/agent_name>`) verbatim in stdout for /// hive-c0re to parse. `-f` lists the qgroups impacting the given path, /// excluding ancestral qgroups (per btrfs-qgroup-show(8)) — so it scopes /// to this subvolume and never mixes in other agents'. hive-c0re then /// selects the level-0 (`0/`) leaf row. See the wire doc. async fn read_subvolume_usage(agent_name: &str) -> Result<(String, String)> { let agent_root = PathBuf::from(AGENT_STATE_ROOT).join(agent_name); let out = Command::new("btrfs") .args(["qgroup", "show", "-f", "--raw"]) .arg(&agent_root) .output() .await .with_context(|| format!("spawn btrfs qgroup show {}", agent_root.display()))?; if !out.status.success() { // The common failure is "quota not enabled" — pass the stderr // through so hive-c0re can surface it gracefully. bail!( "btrfs qgroup show {} failed: {}", agent_root.display(), String::from_utf8_lossy(&out.stderr).trim() ); } Ok(( String::from_utf8_lossy(&out.stdout).into_owned(), String::new(), )) } /// Best-effort removal of a leftover migration path from a prior aborted /// upgrade: try `btrfs subvolume delete` (in case it's a half-created /// subvolume) then a plain recursive remove. Both failures are ignored — /// the path may simply not exist. async fn cleanup_stale_path(path: &Path) { if path.exists() { let _ = Command::new("btrfs") .args(["subvolume", "delete"]) .arg(path) .output() .await; let _ = std::fs::remove_dir_all(path); } } /// Stage a populated subvolume at `tmp` mirroring `agent_root`: create the /// subvolume, copy `agent_root`'s contents into it preserving /// ownership/permissions/xattrs, then match the subvolume root's owner + mode /// to the original. On any failure the partially-staged `tmp` is cleaned up /// (so the caller can bail with the original dir still untouched). async fn stage_upgrade_subvolume(agent_root: &Path, tmp: &Path) -> Result<()> { use std::os::unix::fs::{MetadataExt as _, PermissionsExt as _}; // Fresh subvolume to receive the copy. let out = Command::new("btrfs") .args(["subvolume", "create"]) .arg(tmp) .output() .await .with_context(|| format!("spawn btrfs subvolume create {}", tmp.display()))?; if !out.status.success() { bail!( "btrfs subvolume create {} failed: {}", tmp.display(), String::from_utf8_lossy(&out.stderr).trim() ); } // Copy contents preserving everything (`-a` = --preserve=all → mode, // ownership, timestamps, links, xattrs); reflink for fast CoW clones on the // same btrfs. `/.` copies the directory's contents (incl. dotfiles) // into the subvolume rather than nesting it. let copy = Command::new("cp") .arg("-a") .arg("--reflink=auto") .arg(format!("{}/.", agent_root.display())) .arg(tmp) .output() .await .with_context(|| format!("spawn cp into {}", tmp.display()))?; if !copy.status.success() { cleanup_stale_path(tmp).await; bail!( "copy {} -> {} failed (original left untouched): {}", agent_root.display(), tmp.display(), String::from_utf8_lossy(©.stderr).trim() ); } // Match the new subvolume root's ownership + mode to the original dir. // `cp -a /.` copies the *contents* but the subvolume root keeps its // create-time root ownership, so set it explicitly — the swapped-in // subvolume must be indistinguishable from the original to hive-c0re. let apply = std::fs::metadata(agent_root) .with_context(|| format!("stat {} for ownership", agent_root.display())) .and_then(|m| { std::os::unix::fs::chown(tmp, Some(m.uid()), Some(m.gid())) .with_context(|| format!("chown {} to match original", tmp.display()))?; std::fs::set_permissions(tmp, std::fs::Permissions::from_mode(m.mode())) .with_context(|| format!("chmod {} to match original", tmp.display()))?; Ok(()) }); if let Err(e) = apply { cleanup_stale_path(tmp).await; return Err(e.context("upgrade aborted before swap; original left untouched")); } Ok(()) } /// `UpgradeAgentSubvolume` — convert an existing plain-dir agent state root /// into a btrfs subvolume in place. Operator opt-in; the caller (hivectl) /// stops the agent first and restarts it after. See the wire doc. /// /// Migration: stage a sibling subvolume mirroring the dir /// ([`stage_upgrade_subvolume`]), then rename the original aside and the /// subvolume into place, then remove the original. Any failure before the /// rename-swap leaves the original dir untouched. async fn upgrade_agent_subvolume(agent_name: &str) -> Result<(String, String)> { let root = PathBuf::from(AGENT_STATE_ROOT); let agent_root = root.join(agent_name); if !agent_root.exists() { // A crash between the two swap renames (original → `..old` // succeeded, `..migrating` → agent_root did not) leaves the // agent root missing but the original data intact under `..old`. // Point at the recovery rather than a bare "nothing to do" so the // operator isn't left guessing where the data went. let old = root.join(format!(".{agent_name}.old")); if old.exists() { bail!( "no state dir at {agent} — but {old} holds the original data from an \ interrupted upgrade (host crashed mid-swap). Restore it with \ `mv {old} {agent}`, then re-run the upgrade.", agent = agent_root.display(), old = old.display(), ); } bail!( "no state dir to upgrade at {} — nothing to do", agent_root.display() ); } // Idempotent: already a subvolume → nothing to do. if is_btrfs_subvolume(&agent_root) { return Ok((String::new(), String::new())); } if !is_on_btrfs(&root)? { bail!( "{} is not on btrfs — subvolumes are unsupported, cannot upgrade", root.display() ); } // Sibling temp paths on the same filesystem (so the copy can reflink and // the swap renames are atomic). Leading dots keep them out of the agent // namespace (`validate_agent_name` rejects dot-prefixed names). let tmp = root.join(format!(".{agent_name}.migrating")); let old = root.join(format!(".{agent_name}.old")); // Clear any debris from a previously interrupted run before starting. cleanup_stale_path(&tmp).await; cleanup_stale_path(&old).await; stage_upgrade_subvolume(&agent_root, &tmp).await?; // Swap. `rename` is atomic within a filesystem. The window between the // two renames is the only unsafe point: a crash there leaves the agent // root missing but both `.old` (original) and the new subvolume present // — recoverable by hand, hence the loud logging. if let Err(e) = std::fs::rename(&agent_root, &old) { cleanup_stale_path(&tmp).await; return Err(anyhow::Error::new(e).context(format!( "rename {} -> {} failed; original left untouched", agent_root.display(), old.display() ))); } if let Err(e) = std::fs::rename(&tmp, &agent_root) { // Restore the original from its renamed-aside copy. let restored = std::fs::rename(&old, &agent_root).is_ok(); cleanup_stale_path(&tmp).await; return Err(anyhow::Error::new(e).context(format!( "rename {} -> {} failed; original {}", tmp.display(), agent_root.display(), if restored { "restored" } else { "COULD NOT BE RESTORED — manual recovery needed" } ))); } // 5. Success: drop the original (a plain dir) and report. if let Err(e) = std::fs::remove_dir_all(&old) { // The migration succeeded; a leftover `.old` is cosmetic. Warn only. tracing::warn!( agent = %agent_name, path = %old.display(), "upgraded subvolume but failed to remove old dir: {e}" ); } tracing::info!( agent = %agent_name, path = %agent_root.display(), "upgraded agent state dir to btrfs subvolume" ); Ok(( format!("upgraded {} to a btrfs subvolume", agent_root.display()), String::new(), )) } /// Derive a snapshot's path from the agent name + label: a dot-prefixed /// sibling of the agent's state root so it can never collide with a real /// agent directory (`validate_agent_name` rejects dot-prefixed names). fn snapshot_path(agent_name: &str, snapshot_name: &str) -> PathBuf { PathBuf::from(AGENT_STATE_ROOT).join(format!(".{agent_name}.snapshot.{snapshot_name}")) } /// `SnapshotAgentSubvolume` — create a read-only btrfs snapshot of an /// agent's state subvolume, for `btrfs send` to stream from during /// inter-hive migration. See the wire doc. async fn snapshot_agent_subvolume( agent_name: &str, snapshot_name: &str, ) -> Result<(String, String)> { let agent_root = PathBuf::from(AGENT_STATE_ROOT).join(agent_name); if !is_btrfs_subvolume(&agent_root) { bail!( "{} is not a btrfs subvolume — nothing to snapshot (run `hivectl agent subvol upgrade` first)", agent_root.display() ); } let snap = snapshot_path(agent_name, snapshot_name); if snap.exists() { bail!( "snapshot {} already exists — delete it first or pick a different name", snap.display() ); } let out = Command::new("btrfs") .args(["subvolume", "snapshot", "-r"]) .arg(&agent_root) .arg(&snap) .output() .await .with_context(|| { format!( "spawn btrfs subvolume snapshot -r {} {}", agent_root.display(), snap.display() ) })?; if !out.status.success() { bail!( "btrfs subvolume snapshot -r {} {} failed: {}", agent_root.display(), snap.display(), String::from_utf8_lossy(&out.stderr).trim() ); } tracing::info!( agent = %agent_name, snapshot = %snap.display(), "created read-only agent state snapshot" ); Ok((snap.display().to_string(), String::new())) } /// `DeleteAgentSnapshot` — delete a previously-created read-only snapshot. /// No-op if the path doesn't exist. See the wire doc. async fn delete_agent_snapshot(agent_name: &str, snapshot_name: &str) -> Result<(String, String)> { let snap = snapshot_path(agent_name, snapshot_name); if !snap.exists() { return Ok((String::new(), String::new())); } let out = Command::new("btrfs") .args(["subvolume", "delete"]) .arg(&snap) .output() .await .with_context(|| format!("spawn btrfs subvolume delete {}", snap.display()))?; if !out.status.success() { bail!( "btrfs subvolume delete {} failed: {}", snap.display(), String::from_utf8_lossy(&out.stderr).trim() ); } tracing::info!(agent = %agent_name, snapshot = %snap.display(), "deleted agent state snapshot"); Ok((String::new(), String::new())) } /// Both paths a `SendAgentSnapshotToFile` request reads from — the /// snapshot itself, and its optional incremental parent — must exist /// before we touch `dest`. Pulled out on its own so that invariant is /// checkable without a `dest` path, a `MIGRATE_STAGING_ROOT`, or `btrfs` /// at all: a caller cannot create the output file by accident while /// still holding a reference to this check. fn check_send_snapshot_preconditions(snap: &Path, parent_path: Option<&Path>) -> Result<()> { if !snap.exists() { bail!( "snapshot {} does not exist — create it with `subvol snapshot create` first", snap.display() ); } if let Some(parent_path) = parent_path && !parent_path.exists() { bail!( "parent snapshot {} does not exist — pick an existing parent or omit it for a full send", parent_path.display() ); } Ok(()) } /// Validate preconditions, then atomically create `dest` for a snapshot /// export. The two are combined in one function on purpose: keeping /// `check_send_snapshot_preconditions` ahead of the `create_new` call is /// the entire fix, and folding them together here means that ordering /// can't drift back apart in a later edit to `send_agent_snapshot_to_file`. /// No `btrfs` involved, so this is unit-testable against plain files. fn open_export_dest(snap: &Path, parent_path: Option<&Path>, dest: &Path) -> Result { check_send_snapshot_preconditions(snap, parent_path)?; // `create_new` (O_CREAT|O_EXCL) makes the no-overwrite guarantee atomic // instead of a check-then-create race against a concurrent request. match std::fs::File::options() .write(true) .create_new(true) .open(dest) { Ok(f) => Ok(f), Err(e) if e.kind() == std::io::ErrorKind::AlreadyExists => bail!( "{} already exists — pick a different destination or remove it first \ (send never overwrites an existing export)", dest.display() ), Err(e) => Err(e).with_context(|| format!("create {}", dest.display())), } } /// Deletes `path` when dropped, unless [`disarm`](Self::disarm)ed first. /// Guards `dest` between the moment `send_agent_snapshot_to_file` creates /// it and the moment the export fully succeeds, so any `?`/`bail!` on /// that path — including ones a future edit adds — removes the partial /// file instead of leaving a zero-byte export that then makes every /// retry against the same `--dest` fail with "already exists". struct PartialExportGuard<'a> { path: &'a Path, armed: bool, } impl<'a> PartialExportGuard<'a> { fn new(path: &'a Path) -> Self { Self { path, armed: true } } fn disarm(mut self) { self.armed = false; } } impl Drop for PartialExportGuard<'_> { fn drop(&mut self) { if !self.armed { return; } // Best-effort: warn (don't fail the whole call over it) if removal // itself fails, so a stuck partial file that later masquerades as a // completed export is at least visible in the log. if let Err(rm_err) = std::fs::remove_file(self.path) { tracing::warn!( dest = %self.path.display(), error = %rm_err, "failed to remove partial export after failure — \ next attempt at this dest will hit the already-exists guard" ); } } } /// `SendAgentSnapshotToFile` — stream a read-only snapshot (optionally /// incremental against `parent_name`) to a file under /// `MIGRATE_STAGING_ROOT` via `btrfs send`. Local-file half of the /// inter-hive migration transport; see `PrivRequest::SendAgentSnapshotToFile` /// for the cross-hive follow-up. async fn send_agent_snapshot_to_file( agent_name: &str, snapshot_name: &str, parent_name: Option<&str>, dest_file_name: &str, ) -> Result<(String, String)> { let snap = snapshot_path(agent_name, snapshot_name); let parent_path = parent_name.map(|parent| snapshot_path(agent_name, parent)); std::fs::create_dir_all(MIGRATE_STAGING_ROOT) .with_context(|| format!("create {MIGRATE_STAGING_ROOT}"))?; let dest = Path::new(MIGRATE_STAGING_ROOT).join(dest_file_name); // Preconditions are checked *inside* this call, before `dest` is // created — see `open_export_dest`'s doc comment for why that ordering // matters and is kept in one place. let dest_file = open_export_dest(&snap, parent_path.as_deref(), &dest)?; let cleanup = PartialExportGuard::new(&dest); let mut cmd = Command::new("btrfs"); cmd.arg("send"); if let Some(parent_path) = &parent_path { cmd.arg("-p").arg(parent_path); } cmd.arg(&snap); cmd.stdout(std::process::Stdio::from(dest_file)); cmd.stderr(std::process::Stdio::piped()); let out = cmd .spawn() .with_context(|| format!("spawn btrfs send {}", snap.display()))? .wait_with_output() .await .with_context(|| format!("wait on btrfs send {}", snap.display()))?; if !out.status.success() { bail!( "btrfs send {} failed: {}", snap.display(), String::from_utf8_lossy(&out.stderr).trim() ); } cleanup.disarm(); tracing::info!( agent = %agent_name, snapshot = %snap.display(), dest = %dest.display(), parent = ?parent_name, "exported agent snapshot to file" ); Ok((dest.display().to_string(), String::new())) } /// `SendAgentSnapshotToFd` — stream a read-only snapshot (optionally /// incremental against `parent_name`) straight into a descriptor the /// caller passed us. /// /// The network half of the inter-hive migration transport, arranged so /// this helper never learns there *is* a network: hive-c0re connects to /// the peer's snapshot store, writes the header itself, and hands the /// connected socket over. We only ever see "a thing to write bytes into", /// which keeps a root process out of any address, protocol or trust /// decision — and keeps everyone out of the data path once `btrfs send` /// starts, which matters at multi-gigabyte sizes. async fn send_agent_snapshot_to_fd( agent_name: &str, snapshot_name: &str, parent_name: Option<&str>, dest: OwnedFd, ) -> Result<(String, String)> { let snap = snapshot_path(agent_name, snapshot_name); if !snap.exists() { bail!( "snapshot {} does not exist — create it with `subvol snapshot create` first", snap.display() ); } let mut cmd = Command::new("btrfs"); cmd.arg("send"); if let Some(parent) = parent_name { let parent_path = snapshot_path(agent_name, parent); if !parent_path.exists() { bail!( "parent snapshot {} does not exist — pick an existing parent or omit it for a full send", parent_path.display() ); } cmd.arg("-p").arg(&parent_path); } cmd.arg(&snap); cmd.stdout(std::process::Stdio::from(dest)); cmd.stderr(std::process::Stdio::piped()); let out = cmd .spawn() .with_context(|| format!("spawn btrfs send {}", snap.display()))? .wait_with_output() .await .with_context(|| format!("wait on btrfs send {}", snap.display()))?; if !out.status.success() { // Nothing to clean up: the destination isn't ours. A partial // stream is the receiving end's problem, and `btrfs receive` // refuses to commit an incomplete subvolume anyway. bail!( "btrfs send {} failed: {}", snap.display(), String::from_utf8_lossy(&out.stderr).trim() ); } tracing::info!( agent = %agent_name, snapshot = %snap.display(), parent = ?parent_name, "streamed agent snapshot into a passed descriptor" ); Ok((String::new(), String::new())) } /// `SetSubvolumeQuota` — set or clear a qgroup size limit on an agent /// subvolume (`btrfs qgroup limit <…/agent_name>`). See the /// wire doc. async fn set_subvolume_quota( agent_name: &str, limit_bytes: Option, ) -> Result<(String, String)> { let agent_root = PathBuf::from(AGENT_STATE_ROOT).join(agent_name); let limit = limit_bytes.map_or_else(|| "none".to_owned(), |n| n.to_string()); let out = Command::new("btrfs") .args(["qgroup", "limit", &limit]) .arg(&agent_root) .output() .await .with_context(|| format!("spawn btrfs qgroup limit {}", agent_root.display()))?; if !out.status.success() { bail!( "btrfs qgroup limit {limit} {} failed: {}", agent_root.display(), String::from_utf8_lossy(&out.stderr).trim() ); } tracing::info!(agent = %agent_name, %limit, "set agent subvolume quota"); Ok((String::new(), String::new())) } /// Validate a single argument destined for `forgejo admin`. Rejects null /// bytes, newlines and carriage returns (which could corrupt the /// subprocess args list or log output). /// /// Shell metacharacters are **not** rejected, and do not need to be: the /// command is spawned directly, with no shell to interpret them. fn validate_forge_admin_arg(arg: &str) -> Result<()> { if arg.bytes().any(|b| b == 0 || b == b'\n' || b == b'\r') { bail!("forge admin arg {arg:?} contains null byte or newline"); } Ok(()) } /// Redact a line before it hits the (root-readable, but still /// unnecessarily exposed) host journal. /// /// Two independent rules, because the previous single rule failed open. /// It matched only the substring "password", chosen to be robust against /// forgejo *rewording* its password line — and the leak arrived from the /// other axis entirely: a **different kind of secret** on a differently /// worded line. `forgejo admin user generate-access-token` prints /// `Access token was successfully created: <40 hex>`, which contains no /// "password" and went to the journal verbatim for every agent ever /// provisioned. /// /// So the second rule matches on **shape, not vocabulary**: a long /// unbroken run of secret-alphabet characters. A new secret type is then /// caught by default rather than by someone remembering to add a keyword. /// /// ⚠️ This deliberately over-matches. A nix store hash is also a long /// opaque run and will redact its line. That is the correct direction to /// be wrong in: the cost of a false positive is one less log line, and /// the cost of a false negative is a live credential sitting in the host /// journal, readable by anything with access to it. fn redact_secret_line(line: &str) -> std::borrow::Cow<'_, str> { if line.to_ascii_lowercase().contains("password") { return std::borrow::Cow::Borrowed("[redacted: line mentions a password]"); } if contains_secret_shaped_run(line) { return std::borrow::Cow::Borrowed("[redacted: line contains a secret-shaped token]"); } std::borrow::Cow::Borrowed(line) } /// True when the line contains an unbroken run of at least 32 characters /// from the hex / base64url alphabet. 32 sits below forgejo's 40-hex /// access token and above the ordinary words and path segments that /// appear in `forgejo admin` output. /// /// ⚠️ Scans for a RUN, not for a whitespace-delimited word. An earlier /// version split on whitespace and required the whole word to match, /// which a secret with punctuation glued to it defeats: `","` and /// `"[]"` both fail an all-chars check on the word while still /// containing the credential in full. Whitespace is not what delimits a /// secret — the alphabet is (thanks @argus for catching it). fn contains_secret_shaped_run(line: &str) -> bool { const MIN: usize = 32; let mut run = 0usize; for b in line.bytes() { if b.is_ascii_alphanumeric() || matches!(b, b'+' | b'/' | b'=' | b'_' | b'-') { run += 1; if run >= MIN { return true; } } else { run = 0; } } false } /// Name a `forgejo admin` invocation for an error message without /// reproducing its arguments: keep the leading verb path, stop at the /// first flag. `["user", "change-password", "--username", "iris", /// "--password", "…"]` becomes `forgejo admin user change-password`. /// /// Same allowlist as hive-c0re's `describe_forge_admin` /// (`hive-c0re/src/forge/mod.rs`), for the same reason: the verbs are a /// closed set this crate chooses itself, but hive-c0re passes a live /// `--password` value as an argument on some calls, and argument values /// are never safe to assume closed. Redacting the value after /// `--password` instead would repeat the bug `redact_secret_line`'s own /// doc comment describes — a denylist of one flag name fails open the /// moment a differently named secret-bearing flag is added. fn describe_forge_admin(args: &[String]) -> String { let verbs: Vec<&str> = args .iter() .take_while(|a| !a.starts_with('-')) .map(String::as_str) .collect(); if verbs.is_empty() { "forgejo admin".to_owned() } else { format!("forgejo admin {}", verbs.join(" ")) } } /// Run `forgejo admin ` inside the `hive-forge` container as the /// `forgejo` unix user. Requires root (for nsenter into the container's /// namespaces). Returns `(stdout, stderr)`. async fn run_forge_admin(args: &[String]) -> Result<(String, String)> { let mut cmd_args: Vec<&str> = vec![ "run", "hive-forge", "--", "runuser", "-u", "forgejo", "--", "forgejo", "--work-path", "/var/lib/forgejo", "admin", ]; for a in args { cmd_args.push(a.as_str()); } let out = Command::new("nixos-container") .args(&cmd_args) .output() .await .context("invoke nixos-container run hive-forge -- forgejo admin")?; let stdout = String::from_utf8_lossy(&out.stdout).into_owned(); let stderr = String::from_utf8_lossy(&out.stderr).into_owned(); // stdout at DEBUG, not INFO: on the success path this stream carries // the *product* of the command (the freshly minted token, the created // user's details) and nothing an operator needs at default verbosity. // Redaction stays on as the second layer — the level decides who sees // it, the redactor decides what it says, and neither alone is enough. for line in stdout.lines() { tracing::debug!(target: "forgejo-admin", "{}", redact_secret_line(line)); } for line in stderr.lines() { tracing::warn!(target: "forgejo-admin", "{}", redact_secret_line(line)); } if !out.status.success() { // Redact here too. The error string is propagated to the caller and // ends up logged; a partial-failure stderr can carry the same // material stdout would have. And name the invocation via // `describe_forge_admin` rather than `args.join(" ")`: some callers // pass a live `--password` value as an argument, so reproducing the // raw vector here is the same "two of three sites" gap that makes // these leaks survive a fix — the log was redacted, but the error // was not. let safe_stderr: String = stderr .lines() .map(|l| redact_secret_line(l).into_owned()) .collect::>() .join("; "); bail!( "{} failed ({}): {}", describe_forge_admin(args), out.status, safe_stderr.trim() ); } Ok((stdout, stderr)) } /// Invoke `nixos-container` with the given args, log output to journald. async fn container_run(args: &[&str]) -> Result<(String, String)> { let out = Command::new("nixos-container") .args(args) .output() .await .context("invoke nixos-container")?; let stdout = String::from_utf8_lossy(&out.stdout).into_owned(); let stderr = String::from_utf8_lossy(&out.stderr).into_owned(); // `list` is a read-only enumeration called on the hot path (dashboard // rescan, forge + boot sweeps) — its stdout is the return value, not // progress, so logging every container name on every call floods the // journal. Log stdout only for the mutating ops, where each line is // genuine progress. stderr is always logged (errors matter regardless). if args.first() != Some(&"list") { for line in stdout.lines() { tracing::info!(target: "nixos-container", "{line}"); } } for line in stderr.lines() { tracing::warn!(target: "nixos-container", "{line}"); } if !out.status.success() { bail!( "nixos-container {} failed ({}): {}", args.join(" "), out.status, stderr.trim() ); } Ok((stdout, stderr)) } /// Invoke `machinectl` with the given args, log output to journald. /// Used for operations that nixos-container doesn't expose (e.g. sending /// signals to running containers). async fn machinectl_run(args: &[&str]) -> Result<(String, String)> { let out = Command::new("machinectl") .args(args) .output() .await .context("invoke machinectl")?; let stdout = String::from_utf8_lossy(&out.stdout).into_owned(); let stderr = String::from_utf8_lossy(&out.stderr).into_owned(); for line in stdout.lines() { tracing::info!(target: "machinectl", "{line}"); } for line in stderr.lines() { tracing::warn!(target: "machinectl", "{line}"); } if !out.status.success() { bail!( "machinectl {} failed ({}): {}", args.join(" "), out.status, stderr.trim() ); } Ok((stdout, stderr)) } /// How long to wait for machined to drop a machine's registration after a /// shutdown has been asked for. Generous: a container with slow-stopping /// units legitimately takes a while, and escalating early would SIGKILL a /// shutdown that was going to finish on its own. const NAME_RELEASE_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(20); /// Poll cadence while waiting for the registration to go away. const NAME_RELEASE_POLL: std::time::Duration = std::time::Duration::from_millis(500); /// Cap on a single registration probe. The probe is a D-Bus round trip to /// machined; if machined itself is wedged the call would otherwise sit on /// the D-Bus method timeout, which is far longer than the whole stop /// sequence should take. const MACHINE_PROBE_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(5); /// True when machined still holds a registration for `machine`. /// /// This asks machined the same question the registration itself answers — /// `machinectl show` resolves the name through `GetMachine`, the very lookup /// that makes a later boot fail with `Failed to register machine: already /// exists`. So a positive answer here is exactly the condition that breaks /// the next start, not a proxy for it. /// /// Deliberately NOT the container's systemd unit state: the unit can be /// `inactive` while the registration is still held, and that gap is the /// whole bug this probe exists to catch. /// /// Errors and timeouts answer "still registered". Being wrong that way costs /// a redundant SIGKILL to something already gone; being wrong the other way /// hands back a stop that silently leaked the name. async fn machine_registered(machine: &str) -> bool { let probe = Command::new("machinectl") .args(["show", machine, "--property=Name"]) // Don't leave a probe behind when the timeout below fires. .kill_on_drop(true) .output(); match tokio::time::timeout(MACHINE_PROBE_TIMEOUT, probe).await { Ok(Ok(out)) => out.status.success(), Ok(Err(e)) => { tracing::warn!(%machine, error = %e, "machinectl show failed to run; assuming still registered"); true } Err(_) => { tracing::warn!(%machine, "machinectl show timed out; assuming still registered"); true } } } /// Poll [`machine_registered`] until the name is free or the timeout expires. /// Returns true once the registration is gone. async fn wait_name_released(machine: &str) -> bool { let deadline = tokio::time::Instant::now() + NAME_RELEASE_TIMEOUT; loop { if !machine_registered(machine).await { return true; } if tokio::time::Instant::now() >= deadline { return false; } tokio::time::sleep(NAME_RELEASE_POLL).await; } } /// Stop `machine` and only report success once machined has actually released /// the name. /// /// `nixos-container stop` exiting 0 does not mean the machine is gone. A /// process that sits in the machine's cgroup without being a child of the /// container's init — a shell exec'd in from outside, say — never receives /// the shutdown's SIGTERM if it has been stopped with SIGSTOP, so the /// registration outlives the "successful" stop. Every later start of that /// container then fails with `Failed to register machine: already exists`, /// and machined re-persists the stale record across its own restart, so /// there is no cleaning it up after the fact. The only place to catch it is /// here, in the stop. /// /// So: ask for the stop, wait for the name, and fail loudly if it is still /// held — a caller that is told the stop worked will go on to start the /// container and hit the confusing registration error instead of this one. /// /// This used to escalate to `machinectl kill --signal=SIGKILL` here. Dropped /// per real incident data: the case this guards is a container genuinely /// wedged (e.g. a root-login process survived the stop), and SIGKILL /// doesn't recover that in practice — only a host-level reboot has. /// Pretending a kill attempt handled it hides a condition that needs a human /// to look at the host, so this now just reports the failure instead of /// quietly (and ineffectually) trying to force it. /// /// The verify-then-fail lives in the helper rather than at a call site so /// that every stop gets it: dashboard, reconcile, destroy, cold-start /// fallback. The start path already distrusts its own exit code the same way; /// this is the missing half of that pair. async fn stop_and_release(machine: &str) -> Result<(String, String)> { let stop = container_run(&["stop", machine]).await; if wait_name_released(machine).await { // Happy path, and also the path where a stop that reported failure // nonetheless brought the machine down. Either way the caller gets // the original result untouched. return stop; } tracing::error!( %machine, "stop finished but machined still holds the registration — the \ container is likely wedged (a process outside its init tree \ survived the stop); this needs a host-level look, not another \ stop attempt" ); bail!( "stop {machine}: machined still holds the machine name after {}s. \ This container is likely wedged and starting it again will fail to \ register — needs host-level intervention (a reboot has been the \ only reliable fix in practice).", NAME_RELEASE_TIMEOUT.as_secs() ) } /// Invoke `nixos-container` with the given args and forward output lines /// to the caller as `PrivEvent::Line` messages in real time, logging each /// line to journald as it arrives. Returns `(String::new(), String::new())` /// on success (all output was streamed); the error string includes stderr /// tail on failure. async fn container_run_streaming( args: &[&str], writer: &mut OwnedWriteHalf, ) -> Result<(String, String)> { use tokio::io::AsyncBufReadExt as _; use tokio::process::Command; let mut child = Command::new("nixos-container") .args(args) .stdout(std::process::Stdio::piped()) .stderr(std::process::Stdio::piped()) .spawn() .context("invoke nixos-container (streaming)")?; let stdout = child.stdout.take().expect("stdout piped"); let stderr = child.stderr.take().expect("stderr piped"); let mut stdout_lines = BufReader::new(stdout).lines(); let mut stderr_lines = BufReader::new(stderr).lines(); // Collect stderr for the error message; stream both to the client. let mut stderr_buf = String::new(); // Drive stdout and stderr concurrently. `tokio::select!` interleaves // them without bias — both streams drain at roughly the same rate // as the subprocess produces output. loop { tokio::select! { line = stdout_lines.next_line() => { match line { Ok(Some(l)) => { tracing::info!(target: "nixos-container", "{l}"); write_line_event(writer, PrivStream::Stdout, &l).await; } Ok(None) => break, Err(e) => { tracing::warn!(error = %e, "nixos-container stdout read error"); break; } } } line = stderr_lines.next_line() => { match line { Ok(Some(l)) => { tracing::warn!(target: "nixos-container", "{l}"); write_line_event(writer, PrivStream::Stderr, &l).await; if !stderr_buf.is_empty() { stderr_buf.push('\n'); } stderr_buf.push_str(&l); } Ok(None) => {} Err(e) => { tracing::warn!(error = %e, "nixos-container stderr read error"); } } } } } // Drain any remaining stderr after stdout closed. while let Ok(Some(l)) = stderr_lines.next_line().await { tracing::warn!(target: "nixos-container", "{l}"); write_line_event(writer, PrivStream::Stderr, &l).await; if !stderr_buf.is_empty() { stderr_buf.push('\n'); } stderr_buf.push_str(&l); } let status = child.wait().await.context("wait nixos-container")?; if !status.success() { // Only the last stderr line is embedded — the full stderr was // already forwarded line-by-line as PrivEvent::Line messages and // is captured in build_logs.sqlite by the caller. Keeping the // error message short avoids bloating the anyhow chain. bail!( "nixos-container {} failed ({}): {}", args.join(" "), status, stderr_buf.lines().last().unwrap_or("").trim() ); } Ok((String::new(), String::new())) } /// Read a container's journal as root via `journalctl -M`. Returns /// `(stdout, stderr)`. Unlike `container_run` a non-zero exit is *not* a /// hard error — journalctl's own diagnostic (folded into `stderr` with /// the exit status) is what the caller surfaces to the operator, so the /// helper never bails. async fn read_container_journal(container: &str, query: &JournalQuery) -> Result<(String, String)> { let mut args: Vec = vec![ "-M".to_owned(), container.to_owned(), "--no-pager".to_owned(), format!("--output={}", query.output.as_journalctl()), "-n".to_owned(), query.lines.to_string(), ]; if query.boot { args.push("-b".to_owned()); } if let Some(u) = &query.unit { args.push("-u".to_owned()); args.push(u.clone()); } if let Some(p) = &query.priority { args.push("-p".to_owned()); args.push(p.clone()); } // `--grep=`/`--since=`/`--until=` use the `=`-joined form so a value // can never be parsed as a separate journalctl flag. if let Some(g) = &query.grep { args.push(format!("--grep={g}")); } if let Some(s) = &query.since { args.push(format!("--since={s}")); } if let Some(u) = &query.until { args.push(format!("--until={u}")); } let out = Command::new("journalctl") .args(&args) .output() .await .context("invoke journalctl -M")?; let stdout = String::from_utf8_lossy(&out.stdout).into_owned(); let stderr = if out.status.success() { String::from_utf8_lossy(&out.stderr).into_owned() } else { format!( "journalctl -M {container} exited {}: {}", out.status, String::from_utf8_lossy(&out.stderr).trim() ) }; Ok((stdout, stderr)) } /// Synchronise the host's nginx unit after an `agents.conf` write. /// /// Queries `ActiveState` and dispatches: /// - `active` → `systemctl reload nginx` (SIGHUP, zero-downtime) /// - `failed` → `systemctl reset-failed nginx` + `systemctl start nginx` /// - otherwise → `systemctl start nginx` /// /// ⚠️ `nginx` is hard-coded on purpose — see `PrivRequest::ReloadGatewayNginx`. /// The unit name is the scope of this verb: nginx is a host unit, so no /// namespace bounds it and the literal is the only thing standing between /// "reload the gateway" and "reload anything". /// /// Returns `(String::new(), String::new())` on success so it fits the /// `exec` return type directly. async fn sync_gateway_nginx() -> Result<(String, String)> { let state_out = Command::new("systemctl") .args(["show", "--property=ActiveState", "--value", "nginx"]) .output() .await .context("query gateway nginx ActiveState")?; if !state_out.status.success() { tracing::warn!( exit_code = ?state_out.status.code(), stderr = %String::from_utf8_lossy(&state_out.stderr).trim(), "systemctl show ActiveState exited non-zero — gateway nginx may be down" ); } let state = String::from_utf8_lossy(&state_out.stdout).trim().to_owned(); // State-aware dispatch: reload when running; reset+start after // start-limit failure; plain start when inactive or unknown. match state.as_str() { "active" => { let out = Command::new("systemctl") .args(["reload", "nginx"]) .output() .await .context("reload gateway nginx")?; if !out.status.success() { bail!( "gateway nginx reload failed ({}): {}", out.status, String::from_utf8_lossy(&out.stderr).trim() ); } } "failed" => { // Clear start-limit so the next start can proceed. let _ = Command::new("systemctl") .args(["reset-failed", "nginx"]) .status() .await; let out = Command::new("systemctl") .args(["start", "nginx"]) .output() .await .context("start gateway nginx after reset-failed")?; if !out.status.success() { bail!( "gateway nginx start (after reset-failed) failed ({}): {}", out.status, String::from_utf8_lossy(&out.stderr).trim() ); } } _ => { // inactive, activating, deactivating, unknown — just start. let out = Command::new("systemctl") .args(["start", "nginx"]) .output() .await .context("start gateway nginx")?; if !out.status.success() { bail!( "gateway nginx start failed (state={state:?}) ({}): {}", out.status, String::from_utf8_lossy(&out.stderr).trim() ); } } } Ok((String::new(), String::new())) } /// Return the system container name for a logical agent name. /// All agents (including the manager) use the `h-` prefix. /// /// **This prefix is what confines the lifecycle verbs to agent containers** — /// not their name validation, which only checks characters. A caller cannot /// name `hive-forge` and reach it: it becomes `h-hive-forge`. Infra containers /// are reached through [`PrivRequest::ControlInfraContainer`] and the /// [`InfraContainer`] enum instead, which is why no lifecycle verb here takes /// a sibling name. fn container_system_name(name: &str) -> String { format!("{AGENT_PREFIX}{name}") } /// Validate a logical agent name (the name hive-c0re uses internally, /// before the `h-` container prefix is applied). fn validate_agent_name(name: &str) -> Result<()> { validate_name_chars(name)?; Ok(()) } /// Validate a system-level container name (already has `h-` prefix for /// all agents including the manager, or is a sibling service name). fn validate_container_system_name(name: &str) -> Result<()> { if SIBLING_CONTAINERS.contains(&name) { return Ok(()); } if let Some(suffix) = name.strip_prefix(AGENT_PREFIX) { validate_name_chars(suffix)?; return Ok(()); } bail!("container name {name:?} is not managed by hive"); } fn validate_name_chars(name: &str) -> Result<()> { if name.is_empty() || !name .chars() .all(|c| c.is_ascii_lowercase() || c.is_ascii_digit() || c == '-') { bail!("invalid name {name:?}: must be non-empty lowercase ascii + digits + hyphens"); } Ok(()) } /// Validate a bind-mount path: must be absolute, non-empty, and contain /// no newlines, null bytes, double-quotes, or colons. /// /// The first three would break the `EXTRA_NSPAWN_FLAGS="..."` conf line /// format. The colon is the one that matters most and had been left out /// of this list: `--bind=SRC:DST` is colon-separated, so a path carrying /// one does not corrupt the line — it silently becomes a *different /// mount* than the caller asked for. fn validate_bind_path(path: &str) -> Result<()> { if path.is_empty() || !path.starts_with('/') || path .bytes() .any(|b| b == 0 || b == b'\n' || b == b'"' || b == b':') { bail!( "invalid bind path {path:?}: must be an absolute path with no colons, newlines, null bytes, or double-quotes" ); } Ok(()) } /// `--tmpfs=/.git` for every bound git repo, hiding its metadata /// from inside the container. /// /// Two kinds of mount qualify, for one reason: **the agent is given a /// working tree, never a repository.** `/knowledge` is the hive's shared /// docs — whose `.git/config` has held a credential the host-side worker /// embedded — and `/agents//config` is a config repo, an agent's own /// or a parent's read-only view of a child's. In both cases `.git` carries /// every branch and the full history of a document whose *currently /// deployed* value is the only thing a reader may act on, and an abandoned /// branch is indistinguishable from a live one. /// /// An overlay rather than an exported copy: there is no second tree to /// keep in sync, so nothing can go stale, and no code path has to remember /// to refresh it. /// /// ⚠️ Ordering matters — these must be appended **after** the `--bind` /// flags so nspawn mounts them on top of the already-mounted trees. /// Config mounts are matched by shape, not by a name list: the set is /// dynamic, growing with each child bound into a parent. fn git_overlay_flags(binds: &[BindMount]) -> Vec { binds .iter() .map(|b| b.container_path.as_str()) .filter(|p| *p == "/knowledge" || (p.starts_with("/agents/") && p.ends_with("/config"))) .map(|p| format!("--tmpfs={p}/.git")) .collect() } /// Update `/etc/nixos-containers/.conf`: strip old network vars /// (`PRIVATE_NETWORK`, `HOST_ADDRESS*`, `LOCAL_ADDRESS*`, `HOST_BRIDGE`), /// write the current network-isolation settings, then append /// `EXTRA_NSPAWN_FLAGS`. Always writes `PRIVATE_NETWORK=1` + veth /// wiring — isolation is the only mode, so there is no branch that /// leaves a container on the host's network namespace. fn write_nspawn_flags( container: &str, binds: &[BindMount], isolation: &NetworkIsolation, load_credentials: &[CredentialMount], ) -> Result<()> { use std::fmt::Write as _; use std::os::unix::fs::MetadataExt as _; let path = format!("/etc/nixos-containers/{container}.conf"); let original = std::fs::read_to_string(&path).with_context(|| format!("read {path}"))?; let lines: Vec<&str> = original .lines() .filter(|line| { let t = line.trim_start(); !t.starts_with("EXTRA_NSPAWN_FLAGS=") && !t.starts_with("PRIVATE_NETWORK=") && !t.starts_with("HOST_ADDRESS=") && !t.starts_with("LOCAL_ADDRESS=") && !t.starts_with("HOST_ADDRESS6=") && !t.starts_with("LOCAL_ADDRESS6=") && !t.starts_with("HOST_BRIDGE=") }) .collect(); let mut out = lines.join("\n"); if !out.is_empty() { out.push('\n'); } { let iso = isolation; out.push_str("PRIVATE_NETWORK=1\n"); // HOST_ADDRESS = the bridge gateway IP. nixos-container's // container-side setup only installs a default route // (`ip route add default via $HOST_ADDRESS`) when HOST_ADDRESS is // non-empty; leaving it blank gave the container an address but no // route off the bridge subnet (no internet, no api.anthropic.com). // In bridge mode (HOST_BRIDGE set) the host-side address/route // setup is skipped, so this only affects the container's route — // exactly what we want. let _ = writeln!(out, "HOST_ADDRESS={}", iso.gateway_ip); // LOCAL_ADDRESS is intentionally empty: containers take their IP by // DHCP from the bridge dnsmasq pool. (Why HOST_ADDRESS is still // written: the comment directly above.) out.push_str("LOCAL_ADDRESS=\n"); out.push_str("HOST_ADDRESS6=\n"); out.push_str("LOCAL_ADDRESS6=\n"); let _ = writeln!(out, "HOST_BRIDGE={}", iso.bridge); } let mut flags: Vec = binds .iter() .map(|b| { let flag = if b.read_only { "--bind-ro" } else { "--bind" }; format!("{flag}={}:{}", b.host_path, b.container_path) }) .collect(); flags.extend(git_overlay_flags(binds)); // Credential forwarding: nspawn loads each host secret into the // container's credential store under ``; inner units inherit it // via `LoadCredential=`. Validated (name charset + bind-path // rules) in handle_write_nspawn_flags above. for cred in load_credentials { flags.push(format!( "--load-credential={}:{}", cred.name, cred.host_path )); } let flags_joined = flags.join(" "); let _ = writeln!(out, "EXTRA_NSPAWN_FLAGS=\"{flags_joined}\""); // Keeps the owner and mode `nixos-container` created the file with; this // helper only rewrites its lines. let meta = std::fs::metadata(&path).with_context(|| format!("stat {path}"))?; publish_file( Path::new(&path), out.as_bytes(), meta.mode() & 0o7777, Some((meta.uid(), meta.gid())), )?; // DNS marker for the in-container resolver oneshot: the oneshot only // rewrites the container's resolv.conf when this marker exists, and the // marker carries the gateway IP so the container need not re-derive it. // Always written — every container is isolated, so there is no mode in // which the marker should be absent. Why the rewrite is needed at all, // and which unit does it: `docs/networking/network.md` § *How the isolated // container gets its resolver*. write_bridge_dns_marker(container, isolation)?; Ok(()) } /// Filename of the in-container DNS marker, in the container's own `/etc`. const BRIDGE_DNS_MARKER: &str = "hyperhive-bridge-dns"; /// Write the bridge-DNS marker the `hyperhive-isolated-dns` oneshot keys /// off. The marker file contains just the gateway IP. Always written: /// every container is isolated, so there is no host-netns case that /// wants the marker absent. fn write_bridge_dns_marker(container: &str, isolation: &NetworkIsolation) -> Result<()> { let rootfs = PathBuf::from(format!("/var/lib/nixos-containers/{container}")); write_bridge_dns_marker_in(&rootfs, &isolation.gateway_ip) } /// Write `/etc/hyperhive-bridge-dns` without following a symlink. /// /// Root inside the container owns its `/etc`, and this write runs as host /// root, where an absolute symlink planted in the container resolves against /// the host's `/`. So this opens `etc` with `O_DIRECTORY|O_NOFOLLOW` (a /// symlinked `etc` fails with `ELOOP`) and publishes the marker through a /// [`StagedFile`] relative to that fd: whatever the container planted at the /// leaf, symlink or FIFO or device node, is replaced by the rename, never /// written through or opened. /// /// On a fresh install the container's `/etc` may not exist yet (nixos-container /// materialises the rootfs on the first start), so this creates it first; the /// dir and the marker persist through that start. Mode 0644: the gateway IP /// is not a secret. fn write_bridge_dns_marker_in(rootfs: &Path, gateway_ip: &str) -> Result<()> { use std::os::unix::fs::OpenOptionsExt as _; let etc_path = rootfs.join("etc"); // The rootfs dir is the container's `/`, which the container cannot // replace, so creating it by path is safe. std::fs::create_dir_all(rootfs) .with_context(|| format!("create container rootfs {}", rootfs.display()))?; let root = std::fs::OpenOptions::new() .read(true) .custom_flags(libc::O_DIRECTORY | libc::O_NOFOLLOW) .open(rootfs) .with_context(|| format!("open container rootfs {}", rootfs.display()))?; // SAFETY: `root` is an open directory fd and the name a NUL-terminated // literal. if unsafe { libc::mkdirat(root.as_raw_fd(), c"etc".as_ptr(), 0o755) } != 0 { let e = std::io::Error::last_os_error(); if e.kind() != std::io::ErrorKind::AlreadyExists { return Err(e).with_context(|| format!("create {}", etc_path.display())); } } // SAFETY: as above; this checks the result before wrapping it. let etc_fd = unsafe { libc::openat( root.as_raw_fd(), c"etc".as_ptr(), libc::O_RDONLY | libc::O_DIRECTORY | libc::O_NOFOLLOW | libc::O_CLOEXEC, ) }; if etc_fd < 0 { return Err(std::io::Error::last_os_error()) .with_context(|| format!("open (no-follow) {}", etc_path.display())); } // SAFETY: `etc_fd` is a fresh, valid fd that nothing else owns. let etc = std::fs::File::from(unsafe { OwnedFd::from_raw_fd(etc_fd) }); StagedFile::write( &etc, &etc_path, BRIDGE_DNS_MARKER, format!("{gateway_ip}\n").as_bytes(), 0o644, )? .publish() } /// Written by older hive-priv builds. Nothing writes it now, but left on disk /// systemd-tmpfiles would still apply it at every boot. const LEGACY_TMPFILES_PATH: &str = "/etc/tmpfiles.d/hyperhive-agents.conf"; /// Unlink [`LEGACY_TMPFILES_PATH`]; already absent is success. Run at every /// hive-priv start and by the legacy `SyncAgentTmpfiles` request. fn remove_legacy_tmpfiles() -> Result<()> { match std::fs::remove_file(LEGACY_TMPFILES_PATH) { Ok(()) => { tracing::info!("removed legacy {LEGACY_TMPFILES_PATH}"); Ok(()) } Err(e) if e.kind() == std::io::ErrorKind::NotFound => Ok(()), Err(e) => Err(e).with_context(|| format!("remove {LEGACY_TMPFILES_PATH}")), } } /// `EnsureAgentSocketDir` — create `/run/hive-agent/` if missing. fn ensure_agent_socket_dir(name: &str) -> Result<(String, String)> { ensure_socket_dir_in(Path::new(SOCKET_DIR_ROOT), name)?; Ok((String::new(), String::new())) } /// Create `/` as a `0751` directory owned by the caller (root in /// production), or accept an existing directory untouched. /// `name` must be an agent ident; it is validated here, before any path is /// built from it. /// /// The result is an nspawn bind source, so a symlink here would bind a host /// path of the planter's choosing into the container. Every step is relative /// to an `O_DIRECTORY|O_NOFOLLOW` fd for `root`, and an existing entry is /// checked with `AT_SYMLINK_NOFOLLOW` and refused unless it is a directory. /// /// An existing directory keeps its owner and mode: the container's activation /// hands it to the agent user, and re-asserting root here would undo that on /// every start. fn ensure_socket_dir_in(root: &Path, name: &str) -> Result<()> { use std::os::unix::fs::OpenOptionsExt as _; validate_agent_name(name)?; let path = root.join(name); let c_name = std::ffi::CString::new(name) .with_context(|| format!("agent name {name:?} contains a NUL byte"))?; let root_dir = std::fs::OpenOptions::new() .read(true) .custom_flags(libc::O_DIRECTORY | libc::O_NOFOLLOW | libc::O_CLOEXEC) .open(root) .with_context(|| format!("open (no-follow) {}", root.display()))?; // SAFETY: `root_dir` is an open directory fd and `c_name` NUL-terminated. if unsafe { libc::mkdirat(root_dir.as_raw_fd(), c_name.as_ptr(), 0o751) } == 0 { // mkdirat's mode is masked by the umask; the mode is part of the // contract, so it is set explicitly. // SAFETY: as above. let rc = unsafe { libc::fchmodat( root_dir.as_raw_fd(), c_name.as_ptr(), 0o751, libc::AT_SYMLINK_NOFOLLOW, ) }; if rc != 0 { return Err(std::io::Error::last_os_error()) .with_context(|| format!("chmod 0751 {}", path.display())); } tracing::info!(path = %path.display(), "created agent socket dir"); return Ok(()); } let err = std::io::Error::last_os_error(); if err.kind() != std::io::ErrorKind::AlreadyExists { return Err(err).with_context(|| format!("create {}", path.display())); } // SAFETY: `stat` is plain old data; the zeroed value is only read after // `fstatat` has filled it in. let mut st: libc::stat = unsafe { std::mem::zeroed() }; // SAFETY: as for `mkdirat`; `st` is a valid, writable `stat`. let rc = unsafe { libc::fstatat( root_dir.as_raw_fd(), c_name.as_ptr(), &raw mut st, libc::AT_SYMLINK_NOFOLLOW, ) }; if rc != 0 { return Err(std::io::Error::last_os_error()) .with_context(|| format!("stat (no-follow) {}", path.display())); } if st.st_mode & libc::S_IFMT != libc::S_IFDIR { bail!( "{} exists but is not a directory; refusing it as a bind source", path.display() ); } Ok(()) } #[cfg(test)] mod tests { use super::{ BindMount, BoundedRun, OwnedFd, PAUSED_MARKER_FILE, PrivRequest, StagedFile, check_fd_agreement, clear_runner_credentials, contains_secret_shaped_run, describe_forge_admin, ensure_plain_filename, ensure_socket_dir_in, git_overlay_flags, limits_dropin_body, open_dir, open_export_dest, partial_name, publish_file, redact_secret_line, remove_marker_in, run_bounded, single_output_path, toplevel_attr, validate_credential_name, validate_snapshot_name, write_agent_dir_file, write_bridge_dns_marker_in, }; use std::path::PathBuf; use std::sync::atomic::{AtomicU32, Ordering}; fn bind(container_path: &str) -> BindMount { BindMount { host_path: "/var/lib/hyperhive/whatever".to_owned(), container_path: container_path.to_owned(), read_only: true, } } /// Pins the exact attr path we hand to `nix build` — it has to match /// `hive-c0re`'s own `lifecycle::prebuild_toplevel` construction, since /// that step's whole point is warming the store for this later build. /// A drifted attr path defeats the cache-warming silently — no error, /// just a slower `update`. #[test] fn toplevel_attr_matches_prebuild_toplevels_construction() { assert_eq!( toplevel_attr("atlas"), "/var/lib/hyperhive/meta#nixosConfigurations.atlas.config.system.build.toplevel" ); } /// `--print-out-paths`' exact per-shape table (measured against real /// `rustc` semantics, not just reasoned about) — a lone `"\n"` and a /// trailing-space path are the two shapes a bare `.lines().collect()` /// gets wrong, both catchable only by trimming + dropping empties /// before counting rather than after. #[test] fn single_output_path_rejects_blank_and_trims_whitespace() { assert_eq!(single_output_path(""), Err(0)); assert_eq!(single_output_path("\n"), Err(0)); assert_eq!(single_output_path("\n\n"), Err(0)); assert_eq!(single_output_path("/nix/store/abc\n"), Ok("/nix/store/abc")); assert_eq!( single_output_path("/nix/store/abc \n"), Ok("/nix/store/abc") ); assert_eq!( single_output_path("/nix/store/abc\r\n"), Ok("/nix/store/abc") ); assert_eq!( single_output_path("/nix/store/abc\n/nix/store/def\n"), Err(2) ); } /// Every bound git repo gets its `.git` overlaid — the knowledge tree /// and *each* config mount, an agent's own plus every child's. /// /// The child case is the one worth pinning: that set grows at runtime /// as agents gain children, so a rule written as a list of names would /// silently stop covering new ones. #[test] fn every_bound_git_repo_gets_its_dot_git_hidden() { let flags = git_overlay_flags(&[ bind("/knowledge"), bind("/agents/atlas/config"), bind("/agents/kiddo/config"), ]); assert_eq!( flags, [ "--tmpfs=/knowledge/.git", "--tmpfs=/agents/atlas/config/.git", "--tmpfs=/agents/kiddo/config/.git", ] ); } /// ...and nothing else does. A blanket "overlay .git on every bind" /// would mask a real `.git` under `state/`, where an agent legitimately /// keeps working clones of its own. #[test] fn non_repo_mounts_are_left_alone() { let flags = git_overlay_flags(&[ bind("/agents/atlas/state"), bind("/shared"), bind("/applied"), bind("/agents/atlas/config-notes"), ]); assert!(flags.is_empty(), "overlaid a non-repo mount: {flags:?}"); } /// A request that streams into a caller-supplied descriptor. fn fd_taking_request() -> PrivRequest { PrivRequest::SendAgentSnapshotToFd { agent_name: "atlas".to_owned(), snapshot_name: "hive-migrate".to_owned(), parent_snapshot_name: None, } } /// A real descriptor — `/dev/null` rather than a fake, so the drop /// that closes it on a rejection path is genuinely exercised. fn some_fd() -> OwnedFd { std::fs::File::open("/dev/null") .expect("open /dev/null") .into() } /// An op that streams into a passed descriptor cannot invent one: /// falling back to anything (a temp file, the response socket) would /// send an agent's state somewhere the caller never asked for. #[test] fn an_fd_taking_op_without_a_descriptor_is_rejected() { let err = check_fd_agreement(&fd_taking_request(), None) .expect_err("no descriptor arrived, so this must not proceed"); let msg = format!("{err:#}"); assert!(msg.contains("requires a passed file descriptor"), "{msg}"); } /// The mirror case: a descriptor sent alongside an op that takes /// none is a protocol error, not something to ignore. Returning the /// error drops the `OwnedFd`, which closes it — the alternative /// leaks one descriptor per stray request in a long-lived root /// process. #[test] fn a_descriptor_sent_to_an_op_that_takes_none_is_rejected() { let fd = some_fd(); let err = check_fd_agreement(&PrivRequest::DaemonReload, Some(&fd)) .expect_err("an unexpected descriptor must not be silently ignored"); let msg = format!("{err:#}"); assert!( msg.contains("does not take a passed file descriptor"), "{msg}" ); } /// Both agreeing combinations pass, so the check rejects mismatches /// rather than descriptors in general. #[test] fn agreeing_combinations_are_accepted() { let fd = some_fd(); check_fd_agreement(&fd_taking_request(), Some(&fd)) .expect("an fd-taking op with its descriptor is the normal case"); check_fd_agreement(&PrivRequest::DaemonReload, None) .expect("every ordinary request arrives without a descriptor"); } /// The whole rendered body, pinned literally — including the values, /// so changing the restart policy is a visible test edit rather than a /// silent one. /// /// Placement is the part most worth pinning: `StartLimit*` are `[Unit]` /// settings and systemd **silently ignores** them under `[Service]`, so /// a bound that moved sections would look configured and do nothing. /// An unset weight still emits no line at all, which is what keeps a /// hive-c0re older than that field from changing what lands on disk. #[test] fn dropin_body_is_pinned_exactly() { assert_eq!( limits_dropin_body("/run/hyperhive/agents/iris", "4G", "200%", None, None), "[Unit]\n\ ConditionPathIsDirectory=/run/hyperhive/agents/iris\n\ StartLimitIntervalSec=600\n\ StartLimitBurst=3\n\ \n\ [Service]\n\ RestartSec=5\n\ MemoryMax=4G\n\ CPUQuota=200%\n" ); } /// The bound is hive-wide **policy**, not a per-agent parameter: it is /// rendered from constants and no caller-supplied value can omit or /// alter it. This is the property that justifies keeping it out of the /// wire protocol — if it ever varies by request, that argument is gone. #[test] fn the_start_limit_is_present_whatever_the_caller_passes() { for (mem, cpu, cw, iw) in [ ("4G", "200%", None, None), ("512M", "50%", Some(10), Some(10)), ("infinity", "infinity", Some(10_000), None), ] { let body = limits_dropin_body("/rt/x", mem, cpu, cw, iw); let unit = body .split("[Service]") .next() .expect("the body always has a [Unit] section before [Service]"); assert!( unit.contains("StartLimitIntervalSec=600") && unit.contains("StartLimitBurst=3"), "start limit missing from [Unit] for ({mem}, {cpu}): {body}" ); assert!( body.contains("\nRestartSec=5\n"), "restart backoff missing for ({mem}, {cpu}): {body}" ); } } /// Weights are appended to the `[Service]` section, each omitted /// independently when `None`. #[test] fn weights_are_emitted_only_when_set() { let both = limits_dropin_body("/rt/x", "4G", "200%", Some(80), Some(80)); assert!( both.ends_with("CPUQuota=200%\nCPUWeight=80\nIOWeight=80\n"), "{both}" ); let cpu_only = limits_dropin_body("/rt/x", "4G", "200%", Some(80), None); assert!( cpu_only.ends_with("CPUQuota=200%\nCPUWeight=80\n"), "{cpu_only}" ); assert!(!cpu_only.contains("IOWeight"), "{cpu_only}"); let io_only = limits_dropin_body("/rt/x", "4G", "200%", None, Some(80)); assert!( io_only.ends_with("CPUQuota=200%\nIOWeight=80\n"), "{io_only}" ); assert!(!io_only.contains("CPUWeight"), "{io_only}"); } #[test] fn redacts_lines_mentioning_password_case_insensitively() { assert_eq!( redact_secret_line("New password: hunter2"), "[redacted: line mentions a password]" ); assert_eq!( redact_secret_line("PASSWORD=hunter2"), "[redacted: line mentions a password]" ); assert_eq!( redact_secret_line("User \"foo\" was successfully created."), "User \"foo\" was successfully created." ); } /// `run_forge_admin`'s `bail!` names the invocation through this /// function instead of `args.join(" ")` — the regression this test /// guards. A `--password` value never survives to the error string. #[test] fn describe_forge_admin_does_not_reproduce_a_password_argument() { let described = describe_forge_admin(&[ "user".to_owned(), "create".to_owned(), "--username".to_owned(), "x".to_owned(), "--password".to_owned(), "S3cret".to_owned(), ]); assert_eq!(described, "forgejo admin user create"); assert!(!described.contains("S3cret"), "{described}"); } /// The `--password=value` single-argument form is also excluded: it /// starts with `-`, so it never enters the kept verb path. #[test] fn describe_forge_admin_redacts_the_equals_sign_form_too() { let described = describe_forge_admin(&[ "user".to_owned(), "create".to_owned(), "--password=S3cret".to_owned(), ]); assert!(!described.contains("S3cret"), "{described}"); } /// Control: a non-secret argument (the verb path) still appears, so /// the message stays useful for diagnosing which call failed. #[test] fn describe_forge_admin_keeps_the_verb_path() { assert_eq!( describe_forge_admin(&["user".to_owned(), "change-password".to_owned()]), "forgejo admin user change-password" ); assert_eq!( describe_forge_admin(&["--help".to_owned()]), "forgejo admin" ); assert_eq!(describe_forge_admin(&[]), "forgejo admin"); } /// The regression this function exists for. The keyword rule passes /// this line straight through — it says nothing about a password — so /// only the shape rule catches it. #[test] fn redacts_access_token_line_which_mentions_no_password() { let line = "Access token was successfully created: 0123456789abcdef0123456789abcdef01234567"; assert!( !line.to_ascii_lowercase().contains("password"), "fixture must not contain the keyword, or it proves nothing" ); assert_eq!( redact_secret_line(line), "[redacted: line contains a secret-shaped token]" ); } #[test] fn secret_shape_boundaries() { // Ordinary forgejo-admin output survives: no run is long enough. assert_eq!( redact_secret_line("Command 'user' 'create' finished with no errors."), "Command 'user' 'create' finished with no errors." ); // 31 chars is below the floor, 32 is at it. assert!(!contains_secret_shaped_run(&"a".repeat(31))); assert!(contains_secret_shaped_run(&"a".repeat(32))); // Punctuation breaks the run — a sentence never trips it however long. assert!(!contains_secret_shaped_run( "this.is.a.very.long.dotted.identifier.but.not.a.secret" )); // base64url and hex alphabets both count. assert!(contains_secret_shaped_run( "ZGVhZGJlZWZkZWFkYmVlZmRlYWRiZWVmZGVhZA==" )); assert!(contains_secret_shaped_run( "aG93-dy_there-aG93dy1theresomething" )); } /// argus on the review: a whitespace-delimited check is defeated by /// punctuation glued to the secret — the punctuation joins the "word" /// and fails the alphabet test for the whole run, while the credential /// sits there in full. Scanning for a RUN rather than a WORD closes it. /// These are the shapes that used to slip through. #[test] fn secret_is_caught_with_punctuation_glued_to_it() { const TOK: &str = "0123456789abcdef0123456789abcdef01234567"; for line in [ format!("token: {TOK},"), format!("token: {TOK}."), format!("using [{TOK}] now"), format!("value=\"{TOK}\""), format!("(created {TOK})"), // no whitespace anywhere -- one glued blob format!("Bearer:{TOK};next"), ] { assert_eq!( redact_secret_line(&line), "[redacted: line contains a secret-shaped token]", "leaked through: {line}" ); } } /// Unique scratch dir per test, no external tempfile dep. fn scratch() -> PathBuf { static CTR: AtomicU32 = AtomicU32::new(0); let n = CTR.fetch_add(1, Ordering::Relaxed); let dir = std::env::temp_dir().join(format!( "hive-priv-nofollow-test-{}-{n}", std::process::id() )); std::fs::create_dir_all(&dir).unwrap(); dir } #[test] fn rejects_non_plain_filenames() { let dir = scratch(); let dir_fd = open_dir(&dir).unwrap(); for bad in ["", ".", "..", "a/b", "/etc/passwd", "../escape", "sub/tok"] { assert!( StagedFile::write(&dir_fd, &dir, bad, b"x", 0o600).is_err(), "must reject filename {bad:?}" ); // The wrapper's own check runs before its `create_dir_all`. // Pointing it at a directory that does not exist yet is what makes // this arm bite: the property worth pinning is that a bad leaf // fails BEFORE any root-privileged filesystem work. let absent = dir.join("never-created"); assert!( write_agent_dir_file("a", &absent, bad, "x").is_err(), "wrapper must reject filename {bad:?}" ); assert!( !absent.exists(), "rejected write created a directory for {bad:?}" ); } std::fs::remove_dir_all(&dir).ok(); } /// The temp name is the whole reason the rename is safe to watch: a /// `systemd.path` unit globbing `matrix-token*` would fire on a temp that /// merely suffixed the real name, on exactly the empty file the rename /// exists to hide. #[test] fn a_temp_cannot_match_a_glob_on_the_name_it_publishes() { for real in ["matrix-token", "matrix-token-alice", "forge-token"] { let tmp = partial_name(real); assert!( !tmp.starts_with(real), "{tmp:?} matches a `{real}*` path unit while half-written" ); assert!(!tmp.contains('/'), "temp must stay in the same directory"); } } /// A published file arrives complete and alone. The residue arm is the /// load-bearing one: a stranded temp holds the same secret, is owned by /// root, and nothing else in the system would ever remove it. #[test] fn publishing_leaves_only_the_finished_file() { use std::os::unix::fs::PermissionsExt as _; let dir = scratch(); let names = |d: &PathBuf| { let mut v: Vec = std::fs::read_dir(d) .unwrap() .map(|e| e.unwrap().file_name().to_string_lossy().into_owned()) .collect(); v.sort(); v }; assert!(names(&dir).is_empty(), "control: scratch starts empty"); let path = dir.join("matrix-token-alice"); for content in ["SECRET", "ROTATED"] { write_agent_dir_file("a", &dir, "matrix-token-alice", content).unwrap(); assert_eq!(names(&dir), ["matrix-token-alice"], "temp left behind"); assert_eq!(std::fs::read_to_string(&path).unwrap(), content); } let mode = std::fs::metadata(&path).unwrap().permissions().mode() & 0o777; assert_eq!(mode, 0o600, "credential must be 0600 once published"); std::fs::remove_dir_all(&dir).ok(); } /// A failed publish must not strand the temp. #[test] fn a_failed_publish_removes_the_temp() { let dir = scratch(); // A directory at the destination fails the `rename` (EISDIR) without // needing to drop privileges; the write and chown ahead of it still run. std::fs::create_dir(dir.join("forge-token")).unwrap(); let err = write_agent_dir_file("a", &dir, "forge-token", "SECRET") .expect_err("publishing onto a directory must fail"); assert!( format!("{err:#}").contains("publishing"), "must fail AT the publish -- failing earlier leaves no temp to \ clean up and makes the next assertion vacuous. got: {err:#}" ); assert_eq!( entries(&dir), ["forge-token"], "temp still holds the secret after a failed publish" ); std::fs::remove_dir_all(&dir).ok(); } /// Names in `dir`, sorted — the residue check every publish test ends on. fn entries(dir: &PathBuf) -> Vec { let mut v: Vec = std::fs::read_dir(dir) .unwrap() .map(|e| e.unwrap().file_name().to_string_lossy().into_owned()) .collect(); v.sort(); v } /// The agent plants a symlink where the token goes. The rename replaces /// the link; the root-privileged write never reaches its target. #[test] fn a_symlink_at_the_name_is_replaced_not_followed() { let dir = scratch(); let target = dir.join("target"); std::fs::write(&target, "original").unwrap(); std::os::unix::fs::symlink(&target, dir.join("forge-token")).unwrap(); write_agent_dir_file("a", &dir, "forge-token", "SECRET").unwrap(); assert_eq!( std::fs::read_to_string(&target).unwrap(), "original", "symlink target must be untouched" ); let meta = std::fs::symlink_metadata(dir.join("forge-token")).unwrap(); assert!(meta.file_type().is_file(), "the link must be replaced"); assert_eq!( std::fs::read_to_string(dir.join("forge-token")).unwrap(), "SECRET" ); std::fs::remove_dir_all(&dir).ok(); } /// Any failure between staging and publishing drops the `StagedFile`. /// Until the rename a reader sees only the old file, and afterwards the /// old file is still whole and the temp is gone. #[test] fn a_failed_write_leaves_the_old_file_intact() { let dir = scratch(); let dest = dir.join("hyperhive-limits.conf"); std::fs::write(&dest, "OLD").unwrap(); let dir_fd = open_dir(&dir).unwrap(); let staged = StagedFile::write(&dir_fd, &dir, "hyperhive-limits.conf", b"NEW", 0o644).unwrap(); assert_eq!(std::fs::read_to_string(&dest).unwrap(), "OLD"); assert_eq!( entries(&dir).len(), 2, "control: the temp exists while staged" ); drop(staged); assert_eq!(std::fs::read_to_string(&dest).unwrap(), "OLD"); assert_eq!(entries(&dir), ["hyperhive-limits.conf"], "temp left behind"); std::fs::remove_dir_all(&dir).ok(); } /// Two overlapping writes of one file each own a temp, so whichever /// publishes last lands whole — never the other's bytes, never a mix. #[test] fn overlapping_writes_of_one_file_never_share_a_temp() { let dir = scratch(); let dir_fd = open_dir(&dir).unwrap(); let name = "hyperhive-agents.conf"; let first = StagedFile::write(&dir_fd, &dir, name, b"roster A\n", 0o644).unwrap(); let second = StagedFile::write(&dir_fd, &dir, name, b"roster B, longer\n", 0o644).unwrap(); second.publish().unwrap(); first.publish().unwrap(); assert_eq!( std::fs::read_to_string(dir.join(name)).unwrap(), "roster A\n" ); assert_eq!(entries(&dir), [name], "temp left behind"); std::fs::remove_dir_all(&dir).ok(); } /// The mode lands exactly as asked, whatever the umask would have made /// of it (0666 under the usual 022 would otherwise come out 0644), and /// the requested owner is on the file. #[test] fn publish_file_applies_mode_and_owner() { use std::os::unix::fs::{MetadataExt as _, PermissionsExt as _}; let dir = scratch(); let path = dir.join("c.conf"); let owner = std::fs::metadata(&dir).unwrap(); publish_file(&path, b"x\n", 0o666, Some((owner.uid(), owner.gid()))).unwrap(); let meta = std::fs::metadata(&path).unwrap(); assert_eq!(meta.permissions().mode() & 0o7777, 0o666); assert_eq!((meta.uid(), meta.gid()), (owner.uid(), owner.gid())); publish_file(&path, b"y\n", 0o600, None).unwrap(); let meta = std::fs::metadata(&path).unwrap(); assert_eq!(meta.permissions().mode() & 0o7777, 0o600); assert_eq!(std::fs::read_to_string(&path).unwrap(), "y\n"); assert_eq!(entries(&dir), ["c.conf"], "temp left behind"); std::fs::remove_dir_all(&dir).ok(); } /// The container plants a symlink where the bridge-DNS marker goes; the /// rename replaces it without writing through it. #[test] fn bridge_dns_marker_replaces_symlink_leaf() { let dir = scratch(); let rootfs = dir.join("rootfs"); std::fs::create_dir_all(rootfs.join("etc")).unwrap(); let target = dir.join("target"); std::fs::write(&target, "original").unwrap(); let leaf = rootfs.join("etc/hyperhive-bridge-dns"); std::os::unix::fs::symlink(&target, &leaf).unwrap(); write_bridge_dns_marker_in(&rootfs, "10.0.0.1").unwrap(); assert_eq!( std::fs::read_to_string(&target).unwrap(), "original", "symlink target must be untouched" ); assert!( std::fs::symlink_metadata(&leaf) .unwrap() .file_type() .is_file() ); assert_eq!(std::fs::read_to_string(&leaf).unwrap(), "10.0.0.1\n"); std::fs::remove_dir_all(&dir).ok(); } /// The container replaces its whole `/etc` with a symlink, so the leaf /// name resolves inside a directory of its choosing. #[test] fn bridge_dns_marker_refuses_symlink_etc() { let dir = scratch(); let rootfs = dir.join("rootfs"); std::fs::create_dir_all(&rootfs).unwrap(); let elsewhere = dir.join("elsewhere"); std::fs::create_dir_all(&elsewhere).unwrap(); let target = elsewhere.join("hyperhive-bridge-dns"); std::fs::write(&target, "original").unwrap(); std::os::unix::fs::symlink(&elsewhere, rootfs.join("etc")).unwrap(); let res = write_bridge_dns_marker_in(&rootfs, "10.0.0.1"); assert!(res.is_err(), "O_NOFOLLOW must refuse a symlinked etc"); assert_eq!( std::fs::read_to_string(&target).unwrap(), "original", "file behind the symlinked etc must be untouched" ); std::fs::remove_dir_all(&dir).ok(); } /// A FIFO at the leaf is replaced, never opened, so it cannot block the /// helper. #[test] fn bridge_dns_marker_replaces_fifo_leaf() { use std::os::unix::ffi::OsStrExt as _; let dir = scratch(); let rootfs = dir.join("rootfs"); std::fs::create_dir_all(rootfs.join("etc")).unwrap(); let fifo = rootfs.join("etc/hyperhive-bridge-dns"); let c_fifo = std::ffi::CString::new(fifo.as_os_str().as_bytes()).unwrap(); // SAFETY: `c_fifo` is a valid NUL-terminated path. assert_eq!(unsafe { libc::mkfifo(c_fifo.as_ptr(), 0o644) }, 0); write_bridge_dns_marker_in(&rootfs, "10.0.0.1").unwrap(); assert!( std::fs::symlink_metadata(&fifo) .unwrap() .file_type() .is_file() ); assert_eq!(std::fs::read_to_string(&fifo).unwrap(), "10.0.0.1\n"); std::fs::remove_dir_all(&dir).ok(); } /// The control: a plain rootfs with no `etc` yet (a fresh install) gets /// the dir and a 0644 marker, and a second write replaces it. #[test] fn bridge_dns_marker_written_into_plain_etc() { use std::os::unix::fs::PermissionsExt as _; let dir = scratch(); let rootfs = dir.join("rootfs"); std::fs::create_dir_all(&rootfs).unwrap(); write_bridge_dns_marker_in(&rootfs, "10.0.0.1").unwrap(); let path = rootfs.join("etc/hyperhive-bridge-dns"); assert_eq!(std::fs::read_to_string(&path).unwrap(), "10.0.0.1\n"); let mode = std::fs::metadata(&path).unwrap().permissions().mode() & 0o777; assert_eq!(mode, 0o644, "marker must stay readable in the container"); write_bridge_dns_marker_in(&rootfs, "10.9.8.7").unwrap(); assert_eq!(std::fs::read_to_string(&path).unwrap(), "10.9.8.7\n"); std::fs::remove_dir_all(&dir).ok(); } /// Whether `pid` is gone, or only a zombie awaiting its reaper: either /// way it is no longer running. fn is_dead(pid: &str) -> bool { match std::fs::read_to_string(format!("/proc/{pid}/stat")) { Err(_) => true, // The state letter follows the parenthesised command name. Ok(stat) => stat .rsplit_once(") ") .is_some_and(|(_, rest)| rest.starts_with('Z')), } } /// A child still running at the limit comes back `TimedOut`, and the kill /// reaches its whole process group: the backgrounded grandchild dies too, /// not just the shell that started it. #[tokio::test] async fn a_build_past_its_limit_is_killed_with_its_group() { use std::time::{Duration, Instant}; let dir = scratch(); let pidfile = dir.join("grandchild"); let mut cmd = tokio::process::Command::new("sh"); cmd.arg("-c") .arg(format!("sleep 600 & echo $! > {}; wait", pidfile.display())); let started = Instant::now(); let run = run_bounded(cmd, Duration::from_secs(2), None) .await .unwrap(); assert!(matches!(run, BoundedRun::TimedOut), "must time out"); assert!( started.elapsed() < Duration::from_mins(1), "the limit did not bound the run" ); let pid = std::fs::read_to_string(&pidfile).unwrap(); let pid = pid.trim(); let deadline = Instant::now() + Duration::from_secs(5); while !is_dead(pid) && Instant::now() < deadline { tokio::time::sleep(Duration::from_millis(50)).await; } assert!(is_dead(pid), "grandchild {pid} outlived the timeout"); std::fs::remove_dir_all(&dir).ok(); } /// The control: a child that exits inside its limit is `Exited`, with its /// status and both streams kept apart. #[tokio::test] async fn a_build_inside_its_limit_returns_its_output() { let mut cmd = tokio::process::Command::new("sh"); cmd.arg("-c") .arg("echo progress >&2; echo /nix/store/abc-toplevel"); match run_bounded(cmd, std::time::Duration::from_mins(1), None) .await .unwrap() { BoundedRun::Exited { status, stdout, stderr, } => { assert!(status.success()); assert_eq!(stdout, "/nix/store/abc-toplevel\n"); assert_eq!(stderr, "progress"); } BoundedRun::TimedOut => panic!("a fast child must not time out"), } } /// A missing socket dir is created `0751`, owned by the caller (root in /// production), which is the mode the container's activation expects to /// hand over and the one dialers traverse. #[test] fn socket_dir_created_0751_when_absent() { use std::os::unix::fs::{MetadataExt as _, PermissionsExt as _}; let dir = scratch(); ensure_socket_dir_in(&dir, "probe").unwrap(); let meta = std::fs::symlink_metadata(dir.join("probe")).unwrap(); assert!(meta.is_dir()); assert_eq!(meta.permissions().mode() & 0o7777, 0o751); // SAFETY: `geteuid` has no preconditions. assert_eq!(meta.uid(), unsafe { libc::geteuid() }); std::fs::remove_dir_all(&dir).ok(); } /// An existing directory keeps its mode and contents: by the second start /// it belongs to the agent user, and a live socket may sit inside it. #[test] fn socket_dir_existing_directory_left_alone() { use std::os::unix::fs::PermissionsExt as _; let dir = scratch(); let sub = dir.join("probe"); std::fs::create_dir(&sub).unwrap(); std::fs::set_permissions(&sub, std::fs::Permissions::from_mode(0o700)).unwrap(); std::fs::write(sub.join("web.sock"), "").unwrap(); ensure_socket_dir_in(&dir, "probe").unwrap(); let mode = std::fs::metadata(&sub).unwrap().permissions().mode() & 0o7777; assert_eq!(mode, 0o700, "an existing dir's mode must not be reset"); assert!(sub.join("web.sock").exists()); std::fs::remove_dir_all(&dir).ok(); } /// A symlink at the leaf would make nspawn bind wherever it points; it is /// refused, whether it points at a directory or nowhere. #[test] fn socket_dir_refuses_symlink_leaf() { use std::os::unix::fs::PermissionsExt as _; let dir = scratch(); let elsewhere = dir.join("elsewhere"); std::fs::create_dir(&elsewhere).unwrap(); std::fs::set_permissions(&elsewhere, std::fs::Permissions::from_mode(0o700)).unwrap(); std::os::unix::fs::symlink(&elsewhere, dir.join("probe")).unwrap(); std::os::unix::fs::symlink(dir.join("missing"), dir.join("dangling")).unwrap(); assert!(ensure_socket_dir_in(&dir, "probe").is_err()); assert!(ensure_socket_dir_in(&dir, "dangling").is_err()); let mode = std::fs::metadata(&elsewhere).unwrap().permissions().mode() & 0o7777; assert_eq!(mode, 0o700, "the symlink target must be untouched"); assert!(!dir.join("missing").exists()); std::fs::remove_dir_all(&dir).ok(); } /// A regular file where the dir should be is refused, not bound. #[test] fn socket_dir_refuses_non_directory_leaf() { let dir = scratch(); std::fs::write(dir.join("probe"), "x").unwrap(); let err = ensure_socket_dir_in(&dir, "probe").unwrap_err(); assert!(format!("{err:#}").contains("not a directory"), "{err:#}"); std::fs::remove_dir_all(&dir).ok(); } /// A symlinked root would place the new dir in a tree of the planter's /// choosing; the `O_NOFOLLOW` open refuses it before anything is created. #[test] fn socket_dir_refuses_symlinked_root() { let dir = scratch(); let elsewhere = dir.join("elsewhere"); std::fs::create_dir(&elsewhere).unwrap(); let root = dir.join("root"); std::os::unix::fs::symlink(&elsewhere, &root).unwrap(); assert!(ensure_socket_dir_in(&root, "probe").is_err()); assert!(!elsewhere.join("probe").exists()); std::fs::remove_dir_all(&dir).ok(); } /// Names that are not a plain agent ident are refused by validation, /// before anything is created. Runs against a scratch parent, so the /// refusal cannot come from a missing `/run/hive-agent`; `.`/`..` exist /// as directories and `Atlas` would be created, so each bad name gets /// through only if validation is gone. Control: a plain name succeeds /// in the same parent. #[test] fn socket_dir_refuses_bogus_names() { let dir = scratch(); for bad in ["", ".", "..", "../etc", "a/b", "Atlas", "a b", "a\0b"] { let err = ensure_socket_dir_in(&dir, bad).unwrap_err(); assert!( format!("{err:#}").contains("invalid name"), "agent name {bad:?} must be refused by validation, got: {err:#}" ); } assert_eq!(std::fs::read_dir(&dir).unwrap().count(), 0); ensure_socket_dir_in(&dir, "atlas").unwrap(); assert!(dir.join("atlas").is_dir()); std::fs::remove_dir_all(&dir).ok(); } /// An older hive-c0re's `SyncAgentTmpfiles` still decodes, payload and /// all, so it gets the legacy cleanup rather than a parse error. Control: /// an unknown op does not decode. #[test] fn legacy_sync_agent_tmpfiles_request_still_decodes() { let legacy = r#"{"op":"sync_agent_tmpfiles","agents":[{"name":"atlas","uid":1000,"gid":994}]}"#; assert!(matches!( serde_json::from_str::(legacy), Ok(PrivRequest::SyncAgentTmpfiles) )); assert!(serde_json::from_str::(r#"{"op":"sync_agent_tmpfilez"}"#).is_err()); } /// The runner-credential clear, in all three states that matter. The /// PRESENCE arm is the load-bearing one: an implementation that did nothing /// at all would pass the "absent is fine" arm perfectly, and the whole point /// of the call is that the file is *gone* afterwards — upstream re-registers /// on absence and on nothing else we control. #[test] fn clearing_runner_credentials_removes_it_and_tolerates_absence() { let dir = scratch(); let path = dir.join(".runner"); let path_str = path.to_str().unwrap(); // Absent → success (this is the state we are aiming for). assert!(!path.exists()); clear_runner_credentials(path_str).unwrap(); // Present → success AND actually gone. Without this arm a no-op passes. std::fs::write(&path, r#"{"id":7,"address":"http://old.invalid"}"#).unwrap(); assert!( path.exists(), "control: the file must exist before the clear" ); clear_runner_credentials(path_str).unwrap(); assert!( !path.exists(), "stale credentials must be GONE, or the restart takes upstream's \ already-registered branch and registration silently never happens" ); std::fs::remove_dir_all(&dir).ok(); } /// A failure that is not `NotFound` must propagate, never read as success. /// A directory in the file's place makes `remove_file` fail with a non- /// `NotFound` error without needing to drop privileges in a test. #[test] fn clearing_runner_credentials_propagates_a_real_failure() { let dir = scratch(); let path = dir.join(".runner"); std::fs::create_dir(&path).unwrap(); let err = clear_runner_credentials(path.to_str().unwrap()) .expect_err("a non-NotFound failure must NOT be reported as success"); assert!( format!("{err:#}").contains("remove stale runner credentials"), "error must name what it failed to do, got: {err:#}" ); assert!( path.exists(), "nothing was removed, and the caller must know" ); std::fs::remove_dir_all(&dir).ok(); } /// The pause marker round-trips through the same root-only path the /// credential writes use, and BOTH directions are idempotent — the /// dashboard toggle and `hivectl pause|resume` fire blind, without /// reading the current state first. #[test] fn pause_marker_create_and_remove_are_idempotent() { let dir = scratch(); let path = dir.join(PAUSED_MARKER_FILE); for _ in 0..2 { write_agent_dir_file("a", &dir, PAUSED_MARKER_FILE, "").unwrap(); assert!(path.exists(), "marker must exist after pause"); assert_eq!(std::fs::read_to_string(&path).unwrap(), ""); } for _ in 0..2 { remove_marker_in(&dir, PAUSED_MARKER_FILE).unwrap(); assert!(!path.exists(), "marker must be gone after resume"); } std::fs::remove_dir_all(&dir).ok(); } /// A resume must never follow an agent-planted symlink at the marker /// path: this unlink runs as root, so following it would let an agent /// delete an arbitrary file on the host. #[test] fn resume_unlinks_the_symlink_not_its_target() { let dir = scratch(); let target = dir.join("target"); std::fs::write(&target, "original").unwrap(); let link = dir.join(PAUSED_MARKER_FILE); std::os::unix::fs::symlink(&target, &link).unwrap(); remove_marker_in(&dir, PAUSED_MARKER_FILE).unwrap(); // `exists()` follows the link, so it can't tell "link removed" from // "target removed, dangling link left" — stat the link itself. assert!( std::fs::symlink_metadata(&link).is_err(), "the link itself must be unlinked" ); assert_eq!( std::fs::read_to_string(&target).unwrap(), "original", "symlink target must survive the root unlink" ); std::fs::remove_dir_all(&dir).ok(); } /// The only gate on the root-privileged socket for a `--load-credential` /// name. The doc comment is explicit that `.` is excluded on purpose /// ("this name gets interpolated into filesystem paths") — a future /// "let's allow dots for version numbers" loosening must fail this, /// not just the `/` case below. #[test] fn validate_credential_name_rejects_dot_slash_and_empty_but_allows_the_charset() { assert!( validate_credential_name("").is_err(), "empty must be rejected" ); assert!(validate_credential_name("valid-name_123").is_ok()); assert!( validate_credential_name("bad.name").is_err(), "dot must be rejected — see the fn's doc comment" ); assert!(validate_credential_name("bad/name").is_err()); } /// A snapshot label must additionally carry the `hive-` prefix (it /// doubles as the allow-list gating the btrfs snapshot/delete /// shellouts) on top of [`validate_credential_name`]'s charset rule. #[test] fn validate_snapshot_name_requires_hive_prefix() { assert!(validate_snapshot_name("not-hive-prefixed").is_err()); assert!(validate_snapshot_name("hive-valid").is_ok()); } /// `ensure_plain_filename` is the gate between an agent-writable /// directory and a root-privileged file write. `.`/`..` would resolve /// to the directory itself or its parent; any `/` climbs into a path /// component the caller never named. #[test] fn ensure_plain_filename_rejects_dot_dotdot_and_any_slash() { assert!(ensure_plain_filename("test", ".").is_err()); assert!(ensure_plain_filename("test", "..").is_err()); assert!(ensure_plain_filename("test", "sub/dir").is_err()); assert!(ensure_plain_filename("test", "matrix-token").is_ok()); } /// Unique scratch dir per test, no external tempfile dep. Same pattern /// as `scratch()` above, kept separate so a rename of one doesn't /// collide with the other's directory names under parallel test runs. fn snapshot_export_scratch() -> PathBuf { static CTR: AtomicU32 = AtomicU32::new(0); let n = CTR.fetch_add(1, Ordering::Relaxed); let dir = std::env::temp_dir().join(format!( "hive-priv-snapshot-export-test-{}-{n}", std::process::id() )); std::fs::create_dir_all(&dir).unwrap(); dir } /// The fix: a missing parent must be caught before `dest` is created, /// not after. Fails against the pre-fix ordering (`dest` created via /// `create_new` first, parent checked second) — see the issue for the /// inverted run. #[test] fn missing_parent_leaves_no_file_at_dest() { let dir = snapshot_export_scratch(); let snap = dir.join("snap"); std::fs::write(&snap, b"").unwrap(); let parent = dir.join("parent-that-does-not-exist"); let dest = dir.join("dest"); let err = open_export_dest(&snap, Some(&parent), &dest).unwrap_err(); assert!( err.to_string().contains("does not exist"), "wrong error: {err}" ); assert!( !dest.exists(), "a failed export must not leave a file at dest" ); } /// A retry against the same `dest` after a precondition failure must /// not trip the no-overwrite guard — there is nothing there to /// overwrite, since the failed attempt never created `dest` (the test /// above). This is the behavior a "stale parent falls back to a full /// send, retried without -p" caller needs against a fixed `--dest`. #[test] fn retry_after_precondition_failure_succeeds() { let dir = snapshot_export_scratch(); let snap = dir.join("snap"); std::fs::write(&snap, b"").unwrap(); let parent = dir.join("parent-that-does-not-exist"); let dest = dir.join("dest"); assert!(open_export_dest(&snap, Some(&parent), &dest).is_err()); assert!(!dest.exists()); // Retry without the missing parent, same dest. let file = open_export_dest(&snap, None, &dest); assert!( file.is_ok(), "retry against a dest no previous attempt touched must succeed: {:?}", file.err() ); assert!(dest.exists()); } /// Control for the two tests above: an export that actually landed at /// `dest` must still refuse a second one. Without this, a bug that /// simply stopped creating `dest` at all (rather than fixing the /// ordering) would also make the retry test above pass. #[test] fn existing_dest_still_refuses_overwrite() { let dir = snapshot_export_scratch(); let snap = dir.join("snap"); std::fs::write(&snap, b"").unwrap(); let dest = dir.join("dest"); std::fs::write(&dest, b"already here").unwrap(); let err = open_export_dest(&snap, None, &dest).unwrap_err(); assert!( err.to_string().contains("already exists"), "wrong error: {err}" ); assert_eq!( std::fs::read(&dest).unwrap(), b"already here", "must not have touched the existing file" ); } }