293 lines
9.1 KiB
Rust
293 lines
9.1 KiB
Rust
//! 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_sh4re::priv_proto::{
|
|
BindMount, JournalQuery, NetworkIsolation, PRIV_SOCK, PrivEvent, PrivRequest, PrivResponse,
|
|
PrivStream,
|
|
};
|
|
use tokio::io::{AsyncBufReadExt, AsyncWriteExt, BufReader};
|
|
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<PrivResponse> {
|
|
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::<PrivEvent>(&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<PrivResponse> {
|
|
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::<PrivEvent>(&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?)
|
|
}
|
|
|
|
pub async fn update_container(name: &str) -> Result<(String, String)> {
|
|
check(
|
|
call(&PrivRequest::UpdateContainer {
|
|
name: name.to_owned(),
|
|
stream: false,
|
|
})
|
|
.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?)
|
|
}
|
|
|
|
pub async fn create_container(name: &str) -> Result<(String, String)> {
|
|
check(
|
|
call(&PrivRequest::CreateContainer {
|
|
name: name.to_owned(),
|
|
stream: false,
|
|
})
|
|
.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<String> {
|
|
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: Option<NetworkIsolation>,
|
|
) -> Result<()> {
|
|
ok(call(&PrivRequest::WriteNspawnFlags {
|
|
container: container.to_owned(),
|
|
binds: binds.to_vec(),
|
|
isolation,
|
|
})
|
|
.await?)
|
|
}
|
|
|
|
pub async fn write_resource_limits(
|
|
container: &str,
|
|
memory_max: &str,
|
|
cpu_quota: &str,
|
|
) -> Result<()> {
|
|
ok(call(&PrivRequest::WriteResourceLimits {
|
|
container: container.to_owned(),
|
|
memory_max: memory_max.to_owned(),
|
|
cpu_quota: cpu_quota.to_owned(),
|
|
})
|
|
.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?)
|
|
}
|
|
|
|
pub async fn chown_socket_dir(agent_name: &str, uid: u32, gid: u32) -> Result<()> {
|
|
ok(call(&PrivRequest::ChownSocketDir {
|
|
agent_name: agent_name.to_owned(),
|
|
uid,
|
|
gid,
|
|
})
|
|
.await?)
|
|
}
|
|
|
|
pub async fn chmod_socket_dir(agent_name: &str, mode: u32) -> Result<()> {
|
|
ok(call(&PrivRequest::ChmodSocketDir {
|
|
agent_name: agent_name.to_owned(),
|
|
mode,
|
|
})
|
|
.await?)
|
|
}
|
|
|
|
/// Run `forgejo admin <args>` 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<String> = args.iter().map(|s| (*s).to_owned()).collect();
|
|
check(call(&PrivRequest::RunForgeAdmin { args: owned }).await?)
|
|
}
|
|
|
|
/// Write the Forgejo access token for `agent_name` to
|
|
/// `<agent_state_root>/<agent_name>/state/forge-token` via hive-priv
|
|
/// (running as root). The file is written 0600 and chowned to the agent
|
|
/// user so it is readable from inside the agent container.
|
|
pub async fn write_agent_forge_token(agent_name: &str, token: &str) -> Result<()> {
|
|
ok(call(&PrivRequest::WriteAgentForgeToken {
|
|
agent_name: agent_name.to_owned(),
|
|
token: token.to_owned(),
|
|
})
|
|
.await?)
|
|
}
|
|
|
|
/// Write the Matrix access token for `agent_name` to
|
|
/// `<agent_state_root>/<agent_name>/state/matrix-token` via hive-priv
|
|
/// (running as root). The file is written 0600 and chowned to the agent
|
|
/// user so it is readable from inside the agent container.
|
|
pub async fn write_agent_matrix_token(agent_name: &str, token: &str) -> Result<()> {
|
|
ok(call(&PrivRequest::WriteAgentMatrixToken {
|
|
agent_name: agent_name.to_owned(),
|
|
token: token.to_owned(),
|
|
})
|
|
.await?)
|
|
}
|
|
|
|
/// Restart `hive-matrix-daemon.service` inside an agent container via
|
|
/// `systemctl --machine=h-<agent_name> restart hive-matrix-daemon.service`.
|
|
/// Non-fatal: callers should handle errors gracefully — if the container is
|
|
/// not running the restart will fail (the unit starts naturally on next boot).
|
|
pub async fn restart_matrix_daemon(agent_name: &str) -> Result<()> {
|
|
ok(call(&PrivRequest::RestartMatrixDaemon {
|
|
agent_name: agent_name.to_owned(),
|
|
})
|
|
.await?)
|
|
}
|
|
|
|
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(())
|
|
}
|