//! Async client for the `hive-priv` privileged-helper socket. //! //! Exposes a standalone async function per operation. Each call opens a //! fresh connection to `/run/hive/priv.sock`, sends one JSON line, reads //! the response, and closes. Connection-per-call is intentional: priv //! calls are infrequent (once per rebuild step), so simplicity wins over //! a persistent connection. use anyhow::{Context as _, Result, bail}; use hive_priv_sock::{ BindMount, CredentialMount, InfraAction, InfraContainer, JournalQuery, NetworkIsolation, PRIV_SOCK, PrivEvent, PrivRequest, PrivResponse, PrivStream, }; use std::os::fd::{AsRawFd as _, OwnedFd, RawFd}; use tokio::io::{AsyncBufReadExt, AsyncWriteExt, BufReader, Interest}; use tokio::net::UnixStream; /// Send a single request to `hive-priv` and return the response. /// For streaming ops use `call_streaming` instead. pub async fn call(req: &PrivRequest) -> Result { let mut stream = UnixStream::connect(PRIV_SOCK) .await .context("connect to hive-priv socket")?; let line = serde_json::to_string(req).context("serialise PrivRequest")? + "\n"; stream .write_all(line.as_bytes()) .await .context("send request to hive-priv")?; stream.shutdown().await.context("shutdown write half")?; let mut resp_line = String::new(); BufReader::new(stream) .read_line(&mut resp_line) .await .context("read response from hive-priv")?; // New hive-priv sends `PrivEvent::Done(PrivResponse)` (untagged, wire- // identical to bare `PrivResponse`) — deserialise as `PrivEvent` to // handle both the old bare format and the new tagged format. match serde_json::from_str::(&resp_line).context("parse PrivResponse")? { PrivEvent::Done(resp) => Ok(resp), PrivEvent::Line(_) => bail!("unexpected stream line from non-streaming priv op"), } } /// 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; the union member supplies the `cmsghdr` /// alignment `CMSG_FIRSTHDR` requires. #[repr(C)] union CmsgSpace { _align: libc::cmsghdr, bytes: [u8; 32], } /// `sendmsg` `bytes` with `fd` attached as `SCM_RIGHTS`, returning how /// many bytes were accepted. /// /// The descriptor rides on this one call — ancillary data cannot be /// sent separately from payload — so the caller must not have written /// any of `bytes` beforehand. fn send_with_fd(sock: RawFd, bytes: &[u8], fd: RawFd) -> std::io::Result { const FD_SIZE: usize = std::mem::size_of::(); let mut iov = libc::iovec { iov_base: bytes.as_ptr().cast::().cast_mut(), iov_len: bytes.len(), }; let mut cmsg = CmsgSpace { bytes: [0; 32] }; // SAFETY: msghdr is a plain C struct with no invalid bit patterns; // every field we rely on 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(); // SAFETY: CMSG_SPACE is a pure size computation. let space = unsafe { libc::CMSG_SPACE(u32::try_from(FD_SIZE).unwrap_or(4)) }; msg.msg_controllen = space as _; // SAFETY: the control buffer is live, aligned, and long enough for // the header CMSG_SPACE just sized. let hdr = unsafe { libc::CMSG_FIRSTHDR(&raw const msg) }; if hdr.is_null() { return Err(std::io::Error::other( "control buffer too small for SCM_RIGHTS", )); } // SAFETY: CMSG_LEN is a pure size computation; `hdr` points into // our own buffer, and write_unaligned tolerates its alignment. unsafe { let len = libc::CMSG_LEN(u32::try_from(FD_SIZE).unwrap_or(4)); std::ptr::write_unaligned( hdr, libc::cmsghdr { cmsg_len: len as _, cmsg_level: libc::SOL_SOCKET, cmsg_type: libc::SCM_RIGHTS, }, ); // Copied in byte-wise: the control buffer is only cmsghdr- // aligned, so casting CMSG_DATA to a *mut RawFd would be // unsound even where it happens to work. std::ptr::copy_nonoverlapping( std::ptr::from_ref(&fd).cast::(), libc::CMSG_DATA(hdr), FD_SIZE, ); } // SAFETY: `msg` points at a live iovec over `bytes` and the control // buffer we just filled in. let n = unsafe { libc::sendmsg(sock, &raw const msg, 0) }; if n < 0 { return Err(std::io::Error::last_os_error()); } Ok(usize::try_from(n).unwrap_or_default()) } /// Send a request to `hive-priv` with an open file descriptor attached, /// and return the response. /// /// The helper receives the descriptor itself — not a path or an address /// — so it can act on something it was handed without being told what /// that thing is or how to reach it. Used for /// [`PrivRequest::SendAgentSnapshotToFd`], where the descriptor is a /// socket already connected to a peer hive's snapshot store. /// /// ⚠️ Takes the descriptor by value and closes it as soon as the kernel /// has it, *before* awaiting the response. That is not tidiness: a /// socket stays open until every copy of it is closed, so a caller /// holding one back would leave the receiving end waiting for an EOF /// that never comes — `btrfs receive` blocks, and this side reports /// success for a transfer the peer has not committed. Passing ownership /// makes that mistake unrepresentable. /// /// # Errors /// /// Fails if the socket is unreachable, the descriptor cannot be /// attached, or hive-priv answers with something other than a terminal /// event. pub async fn call_with_fd(req: &PrivRequest, fd: OwnedFd) -> Result { let mut stream = UnixStream::connect(PRIV_SOCK) .await .context("connect to hive-priv socket (fd-passing)")?; let line = serde_json::to_string(req).context("serialise PrivRequest")? + "\n"; let bytes = line.as_bytes(); let sock = stream.as_raw_fd(); let raw_fd = fd.as_raw_fd(); let sent = stream .async_io(Interest::WRITABLE, || send_with_fd(sock, bytes, raw_fd)) .await .context("send request + descriptor to hive-priv")?; // The kernel has duplicated the descriptor into hive-priv's queue, // so our copy has done its job. Close it now, before waiting on the // response: see the EOF note on this function. drop(fd); // A short sendmsg is legal; the descriptor went with the first // call, so the tail is an ordinary write. if sent < bytes.len() { stream .write_all(&bytes[sent..]) .await .context("send remainder of request to hive-priv")?; } stream.shutdown().await.context("shutdown write half")?; let mut resp_line = String::new(); BufReader::new(stream) .read_line(&mut resp_line) .await .context("read response from hive-priv")?; match serde_json::from_str::(&resp_line).context("parse PrivResponse")? { PrivEvent::Done(resp) => Ok(resp), PrivEvent::Line(_) => bail!("unexpected stream line from non-streaming priv op"), } } /// Send a streaming request to `hive-priv`, calling `on_line` for each /// `PrivEvent::Line` as it arrives, then returning the terminal /// `PrivResponse`. Used for long-running ops (`create` / `update`). pub async fn call_streaming( req: &PrivRequest, mut on_line: impl FnMut(PrivStream, &str), ) -> Result { let mut stream = UnixStream::connect(PRIV_SOCK) .await .context("connect to hive-priv socket (streaming)")?; let line = serde_json::to_string(req).context("serialise PrivRequest")? + "\n"; stream .write_all(line.as_bytes()) .await .context("send request to hive-priv")?; stream.shutdown().await.context("shutdown write half")?; let mut reader = BufReader::new(stream); loop { let mut event_line = String::new(); reader .read_line(&mut event_line) .await .context("read event from hive-priv")?; if event_line.is_empty() { bail!("hive-priv closed connection before sending Done event"); } match serde_json::from_str::(&event_line).context("parse PrivEvent")? { PrivEvent::Line(l) => on_line(l.stream, &l.data), PrivEvent::Done(resp) => return Ok(resp), } } } pub async fn start_container(name: &str) -> Result<()> { ok(call(&PrivRequest::StartContainer { name: name.to_owned(), }) .await?) } pub async fn stop_container(name: &str) -> Result<()> { ok(call(&PrivRequest::StopContainer { name: name.to_owned(), }) .await?) } pub async fn kill_container(name: &str) -> Result<()> { ok(call(&PrivRequest::KillContainer { name: name.to_owned(), }) .await?) } /// Streaming variant: forward stdout/stderr lines to `on_line` as they /// arrive. Returns `Ok(())` on success; the callback is responsible for /// appending lines to `build_logs` or otherwise capturing the output. pub async fn update_container_streaming( name: &str, on_line: impl FnMut(PrivStream, &str), ) -> Result<()> { ok(call_streaming( &PrivRequest::UpdateContainer { name: name.to_owned(), stream: true, }, on_line, ) .await?) } /// Streaming variant: forward stdout/stderr lines to `on_line` as they /// arrive. Returns `Ok(())` on success. pub async fn create_container_streaming( name: &str, on_line: impl FnMut(PrivStream, &str), ) -> Result<()> { ok(call_streaming( &PrivRequest::CreateContainer { name: name.to_owned(), stream: true, }, on_line, ) .await?) } pub async fn destroy_container(name: &str) -> Result<()> { ok(call(&PrivRequest::DestroyContainer { name: name.to_owned(), }) .await?) } pub async fn list_containers() -> Result { let (stdout, _) = check(call(&PrivRequest::ListContainers).await?)?; Ok(stdout) } /// Read a container's journal via the root helper (`journalctl -M`). /// Returns `(stdout, stderr)`; a non-zero journalctl exit is reported in /// `stderr` rather than as an `Err`, so callers can surface either. pub async fn read_container_journal( container: &str, query: JournalQuery, ) -> Result<(String, String)> { check( call(&PrivRequest::ReadContainerJournal { container: container.to_owned(), query, }) .await?, ) } pub async fn write_nspawn_flags( container: &str, binds: &[BindMount], isolation: NetworkIsolation, load_credentials: &[CredentialMount], ) -> Result<()> { ok(call(&PrivRequest::WriteNspawnFlags { container: container.to_owned(), binds: binds.to_vec(), isolation, load_credentials: load_credentials.to_vec(), }) .await?) } /// See [`PrivRequest::EnsureAgentSocketDir`]. pub async fn ensure_agent_socket_dir(name: &str) -> Result<()> { ok(call(&PrivRequest::EnsureAgentSocketDir { name: name.to_owned(), }) .await?) } pub async fn write_resource_limits( container: &str, memory_max: &str, cpu_quota: &str, cpu_weight: Option, io_weight: Option, ) -> Result<()> { ok(call(&PrivRequest::WriteResourceLimits { container: container.to_owned(), memory_max: memory_max.to_owned(), cpu_quota: cpu_quota.to_owned(), cpu_weight, io_weight, }) .await?) } pub async fn remove_service_dropin(container: &str) -> Result<()> { ok(call(&PrivRequest::RemoveServiceDropin { container: container.to_owned(), }) .await?) } pub async fn daemon_reload() -> Result<()> { ok(call(&PrivRequest::DaemonReload).await?) } pub async fn reload_gateway_nginx() -> Result<()> { ok(call(&PrivRequest::ReloadGatewayNginx).await?) } /// Run `forgejo admin ` inside the `hive-forge` container via /// hive-priv (which runs as root and can nsenter into the container). /// Returns `(stdout, stderr)` on success. pub async fn run_forge_admin(args: &[&str]) -> Result<(String, String)> { let owned: Vec = args.iter().map(|s| (*s).to_owned()).collect(); check(call(&PrivRequest::RunForgeAdmin { args: owned }).await?) } /// Create (`paused: true`) or remove (`paused: false`) the pause marker in /// `agent_name`'s harness dir via hive-priv (running as root). /// /// hive-c0re cannot do this itself: the harness dir is chowned to the agent /// user on the container's first boot and left mode 0755, so this process /// can stat the marker (that's what `Coordinator::is_paused` does) but gets /// `EACCES` on create *and* unlink. Both directions are idempotent. pub async fn set_agent_paused(agent_name: &str, paused: bool) -> Result<()> { ok(call(&PrivRequest::SetAgentPaused { agent_name: agent_name.to_owned(), paused, }) .await?) } /// Write a GitHub personal access token (PAT) for `agent_name` via hive-priv /// (running as root). Writes `/github-token` 0600, chowned to the agent /// user so the `gh` wrapper / git credential helper can read it from inside the /// container. Single account per agent — no account suffix. The token value is /// operator-supplied (for the agent's GitHub integration, `services.hyperhive.agent.github.enable`). /// /// # Errors /// /// Returns an error if the hive-priv call fails — the socket is unreachable, /// `agent_name` is rejected by the root-side validation, or the file /// write/chown fails. pub async fn write_agent_github_token(agent_name: &str, token: &str) -> Result<()> { ok(call(&PrivRequest::WriteAgentGithubToken { agent_name: agent_name.to_owned(), token: token.to_owned(), }) .await?) } /// Register the hive-ci Forgejo Actions runner: hand the freshly-minted /// registration token to hive-priv, which writes it to the host-side /// `/run/hive-ci/runner-token` env-file and restarts the in-container runner. /// The forge admin token stays in hive-c0re; only the registration token /// crosses to the (host-path) env-file the container bind-mounts read-only. /// /// # Errors /// Propagates the hive-priv call failure: the socket / IPC error, or the /// root-side error when the token is rejected (empty or control characters), /// the env-file write fails, or the runner restart exits non-zero. pub async fn register_ci_runner(token: &str) -> Result<()> { ok(call(&PrivRequest::RegisterCiRunner { token: token.to_owned(), }) .await?) } /// Start / stop a hive infrastructure service (`hive-ci`, `hive-gateway`, /// `hive-forge`, `hive-matrix`) on the host via `systemctl `, /// where the unit is derived root-side from the variant /// (`container@.service`, or `nginx.service` for the gateway). The /// [`InfraContainer`] enum is the allowlist — hive-priv needs no name /// re-validation. Used by the hive-wide `hivectl stop` / `start` flow and /// the dashboard's operator-only infra panel. No agent-facing path exists. pub async fn control_infra_container(container: InfraContainer, action: InfraAction) -> Result<()> { ok(call(&PrivRequest::ControlInfraContainer { container, action }).await?) } /// Ensure the agent's persistent state root is a btrfs subvolume when the /// host FS supports it (via hive-priv, which runs as root). Idempotent and /// progressive: a no-op when the root already exists or the FS isn't btrfs. /// Safe to call on every provision. pub async fn ensure_agent_subvolume(agent_name: &str) -> Result<()> { ok(call(&PrivRequest::EnsureAgentSubvolume { agent_name: agent_name.to_owned(), }) .await?) } /// Delete the agent's state root iff it is a btrfs subvolume (purge path /// only). No-op for plain dirs / missing paths — hive-c0re's own /// `remove_dir_all` handles those. pub async fn delete_agent_subvolume(agent_name: &str) -> Result<()> { ok(call(&PrivRequest::DeleteAgentSubvolume { agent_name: agent_name.to_owned(), }) .await?) } /// Enable btrfs qgroup accounting on the agent-state filesystem (operator /// opt-in; prerequisite for usage reads + quotas). Idempotent; no-op off /// btrfs. Via hive-priv (root). /// /// # Errors /// Returns an error if the hive-priv call fails or `btrfs quota enable` /// reports a non-zero exit. pub async fn ensure_btrfs_quota() -> Result<()> { ok(call(&PrivRequest::EnsureBtrfsQuota).await?) } /// Read an agent state subvolume's btrfs qgroup usage, returning /// `(referenced_bytes, exclusive_bytes)`. Via hive-priv (root). /// /// # Errors /// Returns an error if the hive-priv call fails, `btrfs qgroup show` exits /// non-zero (e.g. quota not enabled — the message propagates so the caller /// can surface it), or the output has no level-0 qgroup row to parse. pub async fn read_subvolume_usage(agent_name: &str) -> Result<(u64, u64)> { let (stdout, _) = check( call(&PrivRequest::ReadSubvolumeUsage { agent_name: agent_name.to_owned(), }) .await?, )?; parse_qgroup_usage(&stdout).with_context(|| format!("parse qgroup usage for {agent_name}")) } /// Set or clear a btrfs qgroup size limit on an agent's state subvolume. /// `limit_bytes = None` clears it. Requires quota enabled. Via hive-priv. /// /// # Errors /// Returns an error if the hive-priv call fails or `btrfs qgroup limit` /// reports a non-zero exit (e.g. quota not enabled). pub async fn set_subvolume_quota(agent_name: &str, limit_bytes: Option) -> Result<()> { ok(call(&PrivRequest::SetSubvolumeQuota { agent_name: agent_name.to_owned(), limit_bytes, }) .await?) } /// Convert an existing plain-dir agent state root into a btrfs subvolume in /// place (operator opt-in; via hive-priv as root). The caller must stop the /// agent first (so its state bind-mount is gone) and restart it after. /// Idempotent: a no-op when the root is already a subvolume. /// /// # Errors /// Returns an error if the hive-priv call fails, the state dir is missing, /// the FS isn't btrfs, or the migration (subvolume create / copy / swap) /// fails — in which case the original dir is left untouched. pub async fn upgrade_agent_subvolume(agent_name: &str) -> Result<()> { ok(call(&PrivRequest::UpgradeAgentSubvolume { agent_name: agent_name.to_owned(), }) .await?) } /// Create a read-only btrfs snapshot of an agent's state subvolume (via /// hive-priv as root) — the first step of `hivectl migrate`'s send/receive /// path. Returns the snapshot's absolute host path. Fails if the agent's /// state root isn't a subvolume yet, or a snapshot with the same /// `snapshot_name` already exists. /// /// # Errors /// Returns an error if the hive-priv call fails, the state dir isn't a /// btrfs subvolume, or the snapshot already exists. pub async fn snapshot_agent_subvolume(agent_name: &str, snapshot_name: &str) -> Result { let (stdout, _) = check( call(&PrivRequest::SnapshotAgentSubvolume { agent_name: agent_name.to_owned(), snapshot_name: snapshot_name.to_owned(), }) .await?, )?; Ok(stdout) } /// Delete a previously-created read-only agent-state snapshot (cleanup /// counterpart to [`snapshot_agent_subvolume`]). No-op if the snapshot /// doesn't exist. Via hive-priv as root. /// /// # Errors /// Returns an error if the hive-priv call fails or the underlying /// `btrfs subvolume delete` fails. pub async fn delete_agent_snapshot(agent_name: &str, snapshot_name: &str) -> Result<()> { ok(call(&PrivRequest::DeleteAgentSnapshot { agent_name: agent_name.to_owned(), snapshot_name: snapshot_name.to_owned(), }) .await?) } /// Stream a read-only agent snapshot to a local file via `btrfs send` /// (optionally incremental against `parent_snapshot_name`). Returns the /// full path of the written file under `MIGRATE_STAGING_ROOT`. Via /// hive-priv as root. See [`PrivRequest::SendAgentSnapshotToFile`]. /// /// # Errors /// Returns an error if the hive-priv call fails, the snapshot (or parent) /// doesn't exist, or the destination file already exists. pub async fn send_agent_snapshot_to_file( agent_name: &str, snapshot_name: &str, parent_snapshot_name: Option<&str>, dest_file_name: &str, ) -> Result { let (stdout, _) = check( call(&PrivRequest::SendAgentSnapshotToFile { agent_name: agent_name.to_owned(), snapshot_name: snapshot_name.to_owned(), parent_snapshot_name: parent_snapshot_name.map(str::to_owned), dest_file_name: dest_file_name.to_owned(), }) .await?, )?; Ok(stdout) } /// Parse `(referenced, exclusive)` bytes from `btrfs qgroup show -f --raw` /// output (a qgroup row is ` …`). /// /// `-f ` already restricts the listing to qgroups impacting that path /// (excluding ancestral qgroups — see btrfs-qgroup-show(8)), so it never /// mixes in other agents' subvolumes. Among the rows it returns we select /// the **level-0** qgroup (`0/`) — the subvolume's own automatic /// usage qgroup — rather than blindly taking the last line. Picking the /// `0/` leaf is unambiguous even if an operator has assigned the subvolume /// to a higher-level aggregate qgroup (`1/`, …) that `-F` would surface. fn parse_qgroup_usage(out: &str) -> Result<(u64, u64)> { for line in out.lines() { let cols: Vec<&str> = line.split_whitespace().collect(); if cols.len() >= 3 && cols[0].starts_with("0/") && let (Ok(rfer), Ok(excl)) = (cols[1].parse::(), cols[2].parse::()) { return Ok((rfer, excl)); } } bail!("no level-0 qgroup data row in `btrfs qgroup show` output: {out:?}") } fn check(resp: PrivResponse) -> Result<(String, String)> { if resp.ok { Ok((resp.stdout, resp.stderr)) } else { bail!( "{}", resp.error.as_deref().unwrap_or("hive-priv returned error") ) } } fn ok(resp: PrivResponse) -> Result<()> { check(resp)?; Ok(()) }