hyperhive/hive-c0re/src/server.rs
atlas a6dc980700 feat: per-agent CPU and memory limits
The hive applies one `agentCpuQuota` / `agentMemoryMax` to every
container. That's the right default and the wrong ceiling: a build-heavy
agent needs headroom the other twelve don't, and raising the hive-wide
value to suit it hands that headroom to everyone.

Adds a per-agent override, persisted host-side and resolved per-field
against the hive defaults.

Follows the existing `meta/*.json` pattern (`capabilities.json`,
`tool-groups.json`): a host-side map read by `hive-c0re`, staged and
committed in the meta repo so every change lands in the audit trail.

```json
{ "sock": { "cpu_quota": "400%", "memory_max": "8G" } }
```

Fallback is **per field**, not per agent: an entry with only
`memory_max` leaves that agent on the hive-wide CPU quota. Absent file,
absent agent and absent field all resolve to the hive default, so the
feature is inert until someone opts an agent in.

Unlike the other meta files this one is **not** injected into the
container — a limit is something done *to* an agent, not something it
reads about itself.

```
hivectl agents set-limits sock --cpu-quota 400% --memory-max 8G
hivectl agents set-limits sock --reset
```

Values are validated before they're persisted: they go into a systemd
drop-in verbatim, and a typo there makes the unit fail to *start* —
turning a fat-fingered quota into a container that won't come back.

The command is declarative: each call replaces the agent's whole entry.
That makes a forgotten flag a silent revert, so a bare `set-limits
<name>` is rejected at the clap layer and clearing needs an explicit
`--reset`.

`ContainerView` gains `cpu_quota` / `memory_max`, both always populated:
there's no "unset" state to render, only "same as everyone else". They
reflect what the drop-in *says* — what the next start will enforce — not
a live cgroup reading.

The write goes through `meta::commit_resource_limits` rather than the
bare setter, so it's staged and committed under `META_LOCK`. Writing
without committing would leave the meta working tree dirty for the next
`prepare_deploy` to trip over.

Docs: `persistence.md` (the new meta file, and why it isn't injected),
`tools/hivectl.md` (the prose guide), `tools/hivectl-cli.md`
(regenerated clap dump).

Closes: internal/requests issue 25
2026-07-26 14:15:05 +02:00

1111 lines
45 KiB
Rust

use std::path::Path;
use std::sync::Arc;
use anyhow::{Context, Result};
use hive_host_sock::{HostRequest, HostResponse, LifecycleScope};
use hive_priv_sock::{InfraAction, InfraContainer};
use tokio::io::{AsyncBufReadExt, AsyncWriteExt, BufReader};
use tokio::net::{UnixListener, UnixStream};
use crate::actions;
use crate::coordinator::Coordinator;
use crate::lifecycle;
pub async fn serve(socket: &Path, coord: Arc<Coordinator>) -> Result<()> {
// Prefer a socket passed by systemd socket-activation (LISTEN_FDS).
// When running under a `.socket` unit, systemd has already created,
// bound, and chmod-ed the socket for us — we just accept on it.
// Fall back to the traditional bind path when not socket-activated
// (direct invocation, dev, tests).
let listener = {
let mut listenfd = listenfd::ListenFd::from_env();
if let Some(std_listener) = listenfd
.take_unix_listener(0)
.context("take socket-activated unix listener")?
{
std_listener.set_nonblocking(true)?;
UnixListener::from_std(std_listener)
.context("convert socket-activated listener to tokio")?
} else {
// Standalone: create parent dir, remove any stale socket, bind.
if let Some(parent) = socket.parent() {
std::fs::create_dir_all(parent)
.with_context(|| format!("create socket parent {}", parent.display()))?;
}
if socket.exists() {
std::fs::remove_file(socket).context("remove stale socket")?;
}
UnixListener::bind(socket)
.with_context(|| format!("bind admin socket {}", socket.display()))?
}
};
tracing::info!(socket = %socket.display(), hyperhive_flake = %coord.hyperhive_flake, "hive-c0re admin listening");
loop {
let (stream, _) = listener.accept().await.context("accept connection")?;
let coord = coord.clone();
tokio::spawn(async move {
if let Err(e) = handle(stream, coord).await {
tracing::warn!(error = ?e, "connection failed");
}
});
}
}
async fn handle(stream: UnixStream, coord: Arc<Coordinator>) -> Result<()> {
let (read, mut write) = stream.into_split();
let mut reader = BufReader::new(read);
let mut line = String::new();
loop {
line.clear();
let n = reader.read_line(&mut line).await?;
if n == 0 {
return Ok(());
}
let resp = match serde_json::from_str::<HostRequest>(line.trim()) {
Ok(req) => dispatch(&req, coord.clone()).await,
Err(e) => HostResponse::error(format!("parse error: {e}")),
};
let mut payload = serde_json::to_string(&resp)?;
payload.push('\n');
write.write_all(payload.as_bytes()).await?;
write.flush().await?;
}
}
#[allow(
clippy::too_many_lines,
reason = "flat one-arm-per-HostRequest-variant router; each arm just \
delegates to a handler. Splitting the match would scatter the \
wire-command routing without shrinking it."
)]
async fn dispatch(req: &HostRequest, coord: Arc<Coordinator>) -> HostResponse {
let result: anyhow::Result<HostResponse> = async {
Ok(match req {
HostRequest::Spawn { name } => handle_spawn(&coord, name.as_str()).await?,
HostRequest::RequestSpawn { name } => {
tracing::info!(%name, "request_spawn");
let id = coord.approvals.submit_kind(
name.as_str(),
hive_sh4re::ApprovalKind::Spawn,
"",
None,
"operator",
None,
)?;
tracing::info!(%id, %name, "spawn approval queued");
HostResponse::success()
}
HostRequest::Kill { name } => submit_single(&coord, name.as_str(), Verb::Kill).await,
HostRequest::Restart { name } => {
submit_single(&coord, name.as_str(), Verb::Restart).await
}
HostRequest::SetPaused { name, paused } => {
handle_set_paused(&coord, name, *paused).await
}
HostRequest::RestartAll => handle_restart_all(&coord).await?,
HostRequest::RestartScoped { scope, graceful } => {
handle_restart_scoped(&coord, scope, *graceful).await?
}
HostRequest::Stop { scope, graceful } => {
// Resolve the scope to explicit container names at the entry
// point, then operate on names — never pass the bare "all
// agents" flag deeper (it'd force every consumer, incl. the
// graceful-stop queue, to re-expand it).
let agents = scoped_agents(scope).await?;
let infra = scoped_infra(scope);
// On a broad stop, remember which agents were actually
// running so a later broad `start` restores only those
// (not every configured container). A targeted `--agent`
// stop must not redefine the restore set.
if is_broad_scope(scope) {
let mut running = Vec::new();
for a in &agents {
if lifecycle::is_running(a).await {
running.push(a.clone());
}
}
coord.set_last_stopped_running(running);
}
handle_stop(&coord, &agents, &infra, *graceful).await?
}
HostRequest::Start { scope } => {
let mut agents = scoped_agents(scope).await?;
// A broad start restores only the set recorded at the
// last broad stop, if any. No record (cold "bring the
// hive up", or a daemon restart since the stop) → start
// all. Targeted `--agent` start is never filtered.
if is_broad_scope(scope)
&& let Some(prev) = coord.take_last_stopped_running()
{
agents.retain(|a| prev.contains(a));
}
let infra = scoped_infra(scope);
handle_start(&coord, &agents, &infra).await?
}
HostRequest::Destroy { name, purge } => {
actions::destroy(&coord, name.as_str(), *purge).await?;
HostResponse::success()
}
HostRequest::Rebuild { name } => {
submit_single(&coord, name.as_str(), Verb::Rebuild).await
}
HostRequest::QueueDag { id } => {
// A multi-step op is one DAG now (no fan-out children to gather).
let dags = coord
.job_queue
.snapshot()
.into_iter()
.filter(|d| d.id == *id)
.collect();
HostResponse::dags(dags)
}
HostRequest::List => HostResponse::list(lifecycle::list().await?),
HostRequest::AgentStatus => handle_agent_status(&coord).await,
// The hive domain + per-surface public URLs are injected into
// c0re's service env by hive-c0re.nix; surface them so the
// operator CLI can fill in this hive's own identity (the
// federation peer-config block) and open the web surfaces
// (`hivectl open`).
HostRequest::Urls => HostResponse::urls(hive_urls()),
HostRequest::Pending => HostResponse::pending(coord.approvals.pending()?),
HostRequest::Approve { id } => {
actions::approve(coord.clone(), *id).await?;
HostResponse::success()
}
HostRequest::Deny { id } => {
actions::deny(&coord, *id, None)?;
HostResponse::success()
}
HostRequest::SetParent { child, new_parent } => {
tracing::info!(%child, ?new_parent, "set_parent");
// `reparent_with_notify` wraps `topology::set_parent`
// with the three notification messages + the
// ContainerView rescan. Idempotent same-parent calls
// skip both the messages and the disk write per the
// topology fast-path.
coord
.reparent_with_notify(
child.as_str(),
new_parent.as_ref().map(hive_types::Ident::as_str),
)
.await
.map_err(anyhow::Error::msg)?;
HostResponse::success()
}
HostRequest::SetResourceLimits {
name,
cpu_quota,
memory_max,
} => {
handle_set_resource_limits(
&coord,
name,
cpu_quota.as_deref(),
memory_max.as_deref(),
)
.await?
}
HostRequest::MatrixCreateUser { name, password } => {
handle_matrix_create_user(name, password.as_deref()).await?
}
HostRequest::MatrixSyncAdmin => handle_matrix_sync_admin().await?,
HostRequest::MatrixPromoteUser { name } => {
handle_matrix_promote_user(name.as_str()).await?
}
HostRequest::MatrixResetPassword { name } => {
handle_matrix_reset_password(name.as_str()).await?
}
HostRequest::MatrixInvite { user, room } => {
handle_matrix_invite(user, room.as_deref()).await?
}
HostRequest::ForgeCreateUser { name, password } => {
handle_forge_create_user(name, password.as_deref()).await?
}
HostRequest::ReconcileConfigStatus { agent, verbose } => {
crate::forge::reconcile_config_status(agent.as_str(), *verbose).await?
}
HostRequest::ReconcileConfigApply { agent, direction } => {
crate::forge::reconcile_config_apply(agent.as_str(), *direction).await?
}
HostRequest::GatewayCreateUser { username, password } => {
HostResponse::messages(vec![crate::gateway_nginx::create_user(username, password)?])
}
HostRequest::GatewayDeleteUser { username } => {
HostResponse::messages(vec![crate::gateway_nginx::delete_user(username)?])
}
HostRequest::GatewayListUsers => {
HostResponse::messages(crate::gateway_nginx::list_users()?)
}
HostRequest::SetAgentGithubToken { agent, token } => {
handle_set_agent_github_token(agent.as_str(), token).await?
}
HostRequest::QuotaEnable => handle_quota_enable().await?,
HostRequest::QuotaLimit { name, limit } => {
handle_quota_limit(name.as_str(), *limit).await?
}
HostRequest::QuotaShow { name } => {
handle_quota_show(name.as_ref().map(hive_types::Ident::as_str)).await?
}
HostRequest::UpgradeSubvolume { name } => {
handle_upgrade_subvolume(name.as_str()).await?
}
HostRequest::SnapshotSubvolume { name, label } => {
handle_snapshot_subvolume(name.as_str(), label).await?
}
HostRequest::DeleteSnapshot { name, label } => {
handle_delete_snapshot(name.as_str(), label).await?
}
HostRequest::SendSnapshot {
name,
label,
parent,
dest,
} => handle_send_snapshot(name.as_str(), label, parent.as_deref(), dest).await?,
})
}
.await;
match result {
Ok(r) => r,
Err(e) => HostResponse::error(format!("{e:#}")),
}
}
/// Create + start the container for `name`, rolling back socket
/// registration and notifying the manager on failure.
async fn handle_spawn(coord: &Arc<Coordinator>, name: &str) -> Result<HostResponse> {
tracing::info!(%name, "spawn");
let agent_dir = crate::paths::agent_runtime_dir(name);
let hive = coord.hive_env();
let paths = Coordinator::agent_paths(name, agent_dir);
// lifecycle::spawn creates the runtime dir internally before start.
// MCP listener registration is event-driven: bind immediately on
// success so the harness can connect on its first turn without
// waiting for any poll interval.
match lifecycle::spawn(name, &hive, &paths).await {
Ok(()) => {
if let Err(e) = coord.power.set(name, crate::power::Wanted::Up) {
tracing::warn!(%name, error = ?e, "agent_power: set wanted=up failed");
}
// Bind the MCP listener now that the container is starting up.
// The harness connects to this socket on its first turn.
coord.register_agent(name)?;
coord.notify_manager(&hive_sh4re::HelperEvent::Spawned {
agent: name.to_owned(),
ok: true,
note: None,
});
// Update tmpfiles.d so the new agent's dirs survive a reboot.
tokio::spawn(lifecycle::sync_tmpfiles());
}
Err(e) => {
// Spawn failed: register_agent was never called, so there is
// nothing to unregister. Notify the manager and propagate.
coord.notify_manager(&hive_sh4re::HelperEvent::Spawned {
agent: name.to_owned(),
ok: false,
note: Some(format!("{e:#}")),
});
return Err(e);
}
}
Ok(HostResponse::success())
}
/// `hivectl pause|resume` / the dashboard toggle: write or remove the
/// agent's pause marker.
///
/// Deliberately not a lifecycle DAG. There's no container operation to
/// sequence — it's one marker file, and the harness picks it up on its
/// next poll — so queueing it would only add latency and a lease. That
/// also means it works on a stopped agent: the marker is sticky, so the
/// agent comes up paused.
async fn handle_set_paused(
coord: &std::sync::Arc<Coordinator>,
name: &hive_types::Ident,
paused: bool,
) -> HostResponse {
if let Err(e) = Coordinator::set_paused(name, paused) {
return HostResponse::error(format!("set paused={paused} for {name}: {e}"));
}
tracing::info!(%name, paused, "agent pause marker updated");
// Refresh the dashboard's view so the paused badge flips without
// waiting for the next periodic rescan.
coord.rescan_containers_and_emit().await;
HostResponse::success()
}
/// Collect per-agent status rows for `hivectl status` and the dashboard.
async fn handle_agent_status(coord: &Arc<Coordinator>) -> HostResponse {
let rows = crate::container_view::build_all(&coord.hive_env())
.await
.into_iter()
.map(|v| hive_sh4re::AgentStatusRow {
name: v.name,
running: v.running,
needs_update: v.needs_update,
needs_login: v.needs_login,
deployed_sha: v.deployed_sha,
// Reminders are agent-local now; c0re has no cross-agent
// visibility into pending counts anymore. Stubbed
// to 0 rather than deleting the wire field outright — leaves
// `hivectl status`/the dashboard column intact syntactically,
// just always empty, until iris's frontend follow-up decides
// whether to drop the column entirely.
pending_reminders: 0,
parent: v.parent,
paused: v.paused,
})
.collect();
HostResponse::agent_statuses(rows)
}
// ---------------------------------------------------------------------------
// Matrix provisioning handlers
//
// The `hivectl matrix` subcommands used to run these in-process, which forced
// the standalone CLI to link the whole daemon crate (matrix-sdk, reqwest, …).
// They now run daemon-side over the host socket: the daemon already holds the
// register + admin tokens and the matrix creds dir. Each op returns the
// operator-facing lines hivectl used to `println!` in `HostResponse::messages`
// for the client to print verbatim.
// ---------------------------------------------------------------------------
/// Shared reqwest client for the matrix admin HTTP calls (30s timeout,
/// mirroring the old in-CLI client).
fn matrix_http_client() -> Result<reqwest::Client> {
reqwest::Client::builder()
.timeout(std::time::Duration::from_secs(30))
.build()
.context("build reqwest client")
}
/// True when `name` has a state dir under the agents root, i.e. it's a
/// managed agent rather than a bare (operator/human) matrix account.
fn agent_exists(name: &hive_types::Ident) -> Result<bool> {
crate::paths::agent_state_dir(name)
.try_exists()
.with_context(|| format!("check agent state dir for {name}"))
}
/// Validate + persist an agent's CPU/memory overrides, then re-apply the
/// drop-in so the change lands without waiting for a rebuild.
///
/// Validation is here rather than only in `hivectl` because the values
/// are written verbatim into the systemd drop-in: a malformed
/// `CPUQuota=` makes systemd reject the unit, and the container stops
/// starting. Every client (CLI, dashboard, anything later) goes through
/// this path, so the guard belongs on this side of the socket.
///
/// `None`/`None` removes the agent's entry, returning it to the
/// hive-wide defaults.
async fn handle_set_resource_limits(
coord: &Arc<Coordinator>,
name: &hive_types::Ident,
cpu_quota: Option<&str>,
memory_max: Option<&str>,
) -> Result<HostResponse> {
if let Some(value) = cpu_quota {
crate::resource_limits::validate_cpu_quota(value).map_err(anyhow::Error::msg)?;
}
if let Some(value) = memory_max {
crate::resource_limits::validate_memory_max(value).map_err(anyhow::Error::msg)?;
}
tracing::info!(%name, ?cpu_quota, ?memory_max, "set_resource_limits");
let limits = crate::resource_limits::AgentLimits {
cpu_quota: cpu_quota.map(ToOwned::to_owned),
memory_max: memory_max.map(ToOwned::to_owned),
};
// Goes through `meta::commit_resource_limits`, not the bare
// `resource_limits::set_limits`: the write has to be staged +
// committed under `META_LOCK` or it leaves the meta working tree
// dirty for the next `prepare_deploy` / `sync_agents` to trip over.
crate::meta::commit_resource_limits(name.as_str(), &limits).await?;
// Re-apply the drop-in straight away — same three lines as the job
// queue's `WriteDropin` node. Without this the new values would sit
// in the JSON until the agent's next spawn or rebuild.
let agent_dir = crate::paths::agent_runtime_dir(name.as_str());
let hive = coord.hive_env();
let paths = Coordinator::agent_paths(name.as_str(), agent_dir);
crate::lifecycle::write_dropins(name.as_str(), &hive, &paths).await?;
let (cpu, mem) = crate::resource_limits::effective(
name.as_str(),
&hive.agent_cpu_quota,
&hive.agent_memory_max,
);
Ok(HostResponse::messages(vec![format!(
"{name}: CPUQuota={cpu} MemoryMax={mem} (restart the container if it is running \
and the new caps need to take effect immediately)"
)]))
}
/// Guard: matrix provisioning needs the homeserver container running.
async fn require_matrix_present() -> Result<()> {
if crate::matrix::is_present().await {
return Ok(());
}
anyhow::bail!(
"hive-matrix container not running — start it (services.hyperhive.matrix.enable = true) before provisioning matrix users"
)
}
async fn handle_matrix_create_user(
name: &hive_types::Ident,
password: Option<&str>,
) -> Result<HostResponse> {
require_matrix_present().await?;
let register_token =
crate::matrix::ensure_register_token().context("read matrix register token")?;
let client = matrix_http_client()?;
let mut out = Vec::new();
if agent_exists(name)? {
if password.is_some() {
// Agents auth by access_token, never by password — the
// boot-sweep provisioning path doesn't accept one. Refuse
// rather than silently dropping it.
anyhow::bail!(
"matrix create-user: a password is for non-agent (operator) accounts only; '{name}' is an agent which authenticates via access_token"
);
}
crate::matrix::ensure_user_for(&client, name.as_str(), &register_token)
.await
.with_context(|| format!("matrix create-user {name}"))?;
let path = Coordinator::agent_notes_dir(name).join("matrix-token");
out.push(format!("matrix: provisioned agent user '{name}'"));
out.push(format!("token persisted at: {}", path.display()));
} else {
let effective_password = match password {
Some(p) => p.to_owned(),
None => crate::matrix::random_password().context("generate random matrix password")?,
};
let token = crate::matrix::provision_user_token(
&client,
name.as_str(),
&register_token,
&effective_password,
)
.await
.with_context(|| format!("matrix create-user {name}"))?;
out.push(format!(
"matrix: provisioned user '{name}' (not an agent — token not persisted)"
));
out.push(format!("token: {token}"));
if password.is_some() {
out.push(
"password: set as supplied — use it to log into a matrix web client".to_owned(),
);
} else {
out.push(
"password: random throwaway (not surfaced — pass --password or --password-stdin to set one you can use)".to_owned(),
);
}
}
Ok(HostResponse::messages(out))
}
async fn handle_forge_create_user(
name: &hive_types::Ident,
password: Option<&str>,
) -> Result<HostResponse> {
if !crate::forge::is_present().await {
anyhow::bail!(
"hive-forge container not running — wait for hive-c0re to start it before provisioning forge users"
);
}
let mut out = Vec::new();
if agent_exists(name)? {
if password.is_some() {
// Agents authenticate by API token, never by password — refuse
// rather than silently dropping a supplied one.
anyhow::bail!(
"forge create-user: a password is for non-agent (operator) accounts only; '{name}' is an agent which authenticates via API token"
);
}
crate::forge::ensure_user_for(name.as_str())
.await
.with_context(|| format!("forge create-user {name}"))?;
let path = Coordinator::agent_notes_dir(name).join("forge-token");
out.push(format!("forge: provisioned agent user '{name}'"));
out.push(format!("token persisted at: {}", path.display()));
} else {
let token = crate::forge::provision_user_token(name.as_str(), password)
.await
.with_context(|| format!("forge create-user {name}"))?;
out.push(format!(
"forge: provisioned user '{name}' (not an agent — token not persisted)"
));
out.push(format!("token: {token}"));
if password.is_some() {
out.push("password: set as supplied — use it to log into the forge web UI".to_owned());
} else {
out.push(
"password: random throwaway (not surfaced — pass --password or --password-stdin to set one you can use)".to_owned(),
);
}
}
Ok(HostResponse::messages(out))
}
async fn handle_set_agent_github_token(agent: &str, token: &str) -> Result<HostResponse> {
crate::priv_client::write_agent_github_token(agent, token)
.await
.with_context(|| format!("write github-token for agent {agent}"))?;
Ok(HostResponse::messages(vec![format!(
"wrote github-token for agent '{agent}' \
(read live by the gh wrapper / git credential helper — no rebuild needed)"
)]))
}
async fn handle_quota_enable() -> Result<HostResponse> {
crate::priv_client::ensure_btrfs_quota()
.await
.context("enable btrfs qgroup accounting")?;
Ok(HostResponse::messages(vec![
"btrfs qgroup accounting enabled on the agent-state filesystem.".to_owned(),
"(usage may read 0 until btrfs finishes its background rescan)".to_owned(),
]))
}
async fn handle_quota_limit(name: &str, limit: Option<u64>) -> Result<HostResponse> {
crate::priv_client::set_subvolume_quota(name, limit)
.await
.with_context(|| format!("set quota for {name}"))?;
// Bare success: the client prints the human-readable confirmation from
// the value it sent (it holds the `human_bytes` formatter).
Ok(HostResponse::success())
}
async fn handle_quota_show(name: Option<&str>) -> Result<HostResponse> {
let agents: Vec<String> = match name {
Some(n) => vec![n.to_owned()],
None => Coordinator::kept_state_names()
.into_iter()
.map(hive_types::Ident::into_string)
.collect(),
};
let mut rows = Vec::with_capacity(agents.len());
for agent in &agents {
match crate::priv_client::read_subvolume_usage(agent).await {
Ok((referenced, exclusive)) => rows.push(hive_host_sock::QuotaRow {
agent: agent.clone(),
referenced: Some(referenced),
exclusive: Some(exclusive),
note: None,
}),
Err(e) => {
let msg = format!("{e:#}");
// btrfs-progs prints "ERROR: ... quota not enabled" to stderr
// when qgroups are off; short-circuit the whole sweep with the
// enable hint (case-insensitive fragment match — the wording
// varies across btrfs-progs versions).
if msg.to_ascii_lowercase().contains("quota not enabled") {
return Ok(HostResponse::error(
"btrfs quota not enabled — run `hivectl quota enable` first",
));
}
// A plain-dir agent (no subvolume) has no qgroup; note it
// inline and keep going rather than aborting the whole sweep.
rows.push(hive_host_sock::QuotaRow {
agent: agent.clone(),
referenced: None,
exclusive: None,
note: Some(format!("no qgroup data — plain dir or: {msg}")),
});
}
}
}
Ok(HostResponse::quota(rows))
}
async fn handle_upgrade_subvolume(name: &str) -> Result<HostResponse> {
crate::priv_client::upgrade_agent_subvolume(name)
.await
.with_context(|| format!("upgrade {name} state subvolume"))?;
Ok(HostResponse::success())
}
async fn handle_snapshot_subvolume(name: &str, label: &str) -> Result<HostResponse> {
let path = crate::priv_client::snapshot_agent_subvolume(name, label)
.await
.with_context(|| format!("snapshot {name} state subvolume (label {label:?})"))?;
Ok(HostResponse::messages(vec![path]))
}
async fn handle_delete_snapshot(name: &str, label: &str) -> Result<HostResponse> {
crate::priv_client::delete_agent_snapshot(name, label)
.await
.with_context(|| format!("delete {name} snapshot (label {label:?})"))?;
Ok(HostResponse::success())
}
async fn handle_send_snapshot(
name: &str,
label: &str,
parent: Option<&str>,
dest: &str,
) -> Result<HostResponse> {
let path = crate::priv_client::send_agent_snapshot_to_file(name, label, parent, dest)
.await
.with_context(|| format!("send {name} snapshot (label {label:?}) to file {dest:?}"))?;
Ok(HostResponse::messages(vec![path]))
}
async fn handle_matrix_sync_admin() -> Result<HostResponse> {
require_matrix_present().await?;
let register_token =
crate::matrix::ensure_register_token().context("read matrix register token")?;
let client = matrix_http_client()?;
crate::matrix::ensure_admin_user(&client, &register_token)
.await
.context("matrix sync-admin")?;
let path = crate::matrix::admin_token_path();
Ok(HostResponse::messages(vec![
format!(
"matrix: hive admin user '@{}' provisioned",
crate::matrix::HIVE_ADMIN_LOCALPART
),
format!("token persisted at: {}", path.display()),
]))
}
async fn handle_matrix_promote_user(name: &str) -> Result<HostResponse> {
require_matrix_present().await?;
let admin_token = crate::matrix::read_admin_token()?;
let client = matrix_http_client()?;
let server_name = crate::matrix::discover_server_name(&client)
.await
.context("discover matrix server_name")?;
crate::matrix::promote_user_to_admin(&client, &admin_token, name, &server_name)
.await
.with_context(|| format!("matrix promote-user {name}"))?;
Ok(HostResponse::messages(vec![format!(
"matrix: promoted @{name}:{server_name} to admin"
)]))
}
async fn handle_matrix_invite(user: &str, room: Option<&str>) -> Result<HostResponse> {
require_matrix_present().await?;
let admin_token = crate::matrix::read_admin_token()?;
let client = matrix_http_client()?;
let server_name = crate::matrix::discover_server_name(&client)
.await
.context("discover matrix server_name")?;
let room_id = crate::matrix::invite_user(&client, &admin_token, user, room, &server_name)
.await
.with_context(|| format!("matrix invite {user}"))?;
let target = if user.starts_with('@') {
user.to_owned()
} else {
format!("@{user}:{server_name}")
};
Ok(HostResponse::messages(vec![format!(
"matrix: invited {target} to {room_id}"
)]))
}
async fn handle_matrix_reset_password(name: &str) -> Result<HostResponse> {
require_matrix_present().await?;
let admin_token = crate::matrix::read_admin_token()?;
let client = matrix_http_client()?;
let server_name = crate::matrix::discover_server_name(&client)
.await
.context("discover matrix server_name")?;
crate::matrix::reset_user_password(&client, &admin_token, name, &server_name)
.await
.with_context(|| format!("matrix reset-password {name}"))?;
// Password is persisted by reset_user_password.
let pw_path = crate::paths::matrix_creds_dir().join(format!("{name}-password"));
Ok(HostResponse::messages(vec![
format!("matrix: password for @{name}:{server_name} reset"),
format!("password persisted at: {}", pw_path.display()),
format!("next: hivectl matrix create-user {name} # mints a fresh access token"),
]))
}
/// Single-agent queue verbs the admin socket exposes. Each submits the
/// matching DAG (persisting the `wanted` intent, serializing on the
/// agent's lease, with the transient/crash-watch suppression the old
/// direct lifecycle calls lacked) and returns the DAG id for the
/// client's wait loop.
#[derive(Clone, Copy)]
enum Verb {
/// Stop DAG (`wanted = Offline`; Reconcile kills + unregisters +
/// fires `Killed`).
Kill,
/// Restart DAG (`wanted = Up`; mechanical stop + reconcile-start).
Restart,
/// Rebuild DAG — the Swap tail owns the manager `Rebuilt` events +
/// kick, so the CLI path can't drift from the dashboard's.
Rebuild,
}
async fn submit_single(coord: &Arc<Coordinator>, name: &str, verb: Verb) -> HostResponse {
use crate::job_queue::{Source, submit};
let id = match verb {
Verb::Kill => {
tracing::info!(%name, "kill");
submit::stop(
coord,
name,
Source::Manual,
"manual kill via hivectl".to_owned(),
)
.await
}
Verb::Restart => {
tracing::info!(%name, "restart");
submit::restart(
coord,
name,
Source::Manual,
"manual restart via hivectl".to_owned(),
)
.await
}
Verb::Rebuild => {
tracing::info!(%name, "rebuild");
submit::rebuild(
coord,
name,
Source::Manual,
"manual rebuild via hivectl".to_owned(),
)
}
};
HostResponse::queued(vec![id])
}
/// Restart every container in **one** DAG — a per-agent restart subgraph
/// each, running concurrently on their own leases (so unrelated agents'
/// restarts overlap while nothing races an in-flight rebuild). Returns
/// once queued; per-node progress surfaces on the single DAG.
async fn handle_restart_all(coord: &Arc<Coordinator>) -> Result<HostResponse> {
tracing::info!("restart-all");
let containers = lifecycle::list().await?;
let agents: Vec<String> = containers
.iter()
.filter_map(|a| a.strip_prefix(lifecycle::AGENT_PREFIX).map(str::to_owned))
.collect();
let queued = if agents.is_empty() {
Vec::new()
} else {
vec![
crate::job_queue::submit::restart_many(
coord,
&agents,
false,
crate::job_queue::Source::Manual,
"manual restart via hivectl restart-all".to_owned(),
)
.await,
]
};
let mut resp = HostResponse::list(agents);
resp.queued_dags = Some(queued);
Ok(resp)
}
/// Stop the given `agents` (resolved logical names) then `infra` containers
/// (`hivectl stop`). Agents go down before infra so they're not mid-request
/// against a forge/matrix that's already gone. Per-target failures are
/// aggregated rather than aborting on the first error, mirroring
/// `handle_restart_all`. Callers resolve the [`LifecycleScope`] to these
/// explicit name lists up front — this never sees the "all" flag.
///
/// Every agent rides the job queue: a `graceful` stop submits the
/// quiesce DAG (signal → drain → reconcile-stop; all drains overlap),
/// a hard stop a plain stop DAG — both persist `wanted = Offline` and
/// serialize on the agent's lease so nothing races an in-flight
/// rebuild. The response carries the DAG ids so `hivectl` can wait
/// with per-node progress. Infra containers have no harness / lease
/// and stay direct + synchronous.
async fn handle_stop(
coord: &Arc<Coordinator>,
agents: &[String],
infra: &[InfraContainer],
graceful: bool,
) -> Result<HostResponse> {
tracing::info!(?agents, ?infra, graceful, "stop");
let mut ok_items: Vec<String> = Vec::new();
let mut errors: Vec<String> = Vec::new();
let mut queued: Vec<u64> = Vec::new();
// One DAG for all targeted agents — a per-agent stop subgraph each
// (`SetWanted(Offline) → [Signal → Drain →] Reconcile`), independent
// roots that run concurrently on their own leases. A hive-wide
// `hivectl stop` is now a single DAG, not N.
if !agents.is_empty() {
let reason = if graceful {
"manual via hivectl graceful stop"
} else {
"manual via hivectl stop"
};
queued.push(
crate::job_queue::submit::stop_many(
coord,
agents,
graceful,
crate::job_queue::Source::Manual,
reason.to_owned(),
)
.await,
);
ok_items.extend(agents.iter().cloned());
}
// Agents go down before infra so they're not mid-request against a
// forge/matrix that's already gone. Hard stops are quick kills —
// await their DAGs (bounded) before touching infra. Graceful stops
// keep the immediate return (drains take minutes and the
// agents-then-infra race pre-existed there).
if !graceful && !infra.is_empty() {
await_dags(coord, &queued, std::time::Duration::from_mins(2)).await;
}
for &container in infra {
let name = container.unit_name();
match crate::priv_client::control_infra_container(container, InfraAction::Stop).await {
Ok(()) => ok_items.push(name.to_owned()),
Err(e) => {
tracing::warn!(%name, error = ?e, "stop: infra stop failed");
errors.push(format!("{name}: {e:#}"));
}
}
}
let mut resp = finish_lifecycle(ok_items, &errors);
resp.queued_dags = Some(queued);
Ok(resp)
}
/// Best-effort server-side wait for a set of DAGs to settle terminal,
/// bounded by `timeout` — used to preserve ordering invariants inside
/// one request (agent stops before infra stops) without trusting the
/// client to wait.
async fn await_dags(coord: &Arc<Coordinator>, ids: &[u64], timeout: std::time::Duration) {
let deadline = std::time::Instant::now() + timeout;
loop {
let snap = coord.job_queue.snapshot();
// A DAG has settled when it's either gone from the snapshot (fully
// `Done` DAGs drop out) or still present but with every node terminal
// (a `Failed`/`Cancelled` DAG lingers). It's pending only while it has
// a non-terminal node.
let pending = ids.iter().any(|id| {
snap.iter()
.any(|d| d.id == *id && d.nodes.iter().any(|n| !n.state.is_terminal()))
});
if !pending {
return;
}
if std::time::Instant::now() >= deadline {
tracing::warn!(?ids, "await_dags: timed out; proceeding");
return;
}
tokio::time::sleep(std::time::Duration::from_millis(250)).await;
}
}
/// Start the given `infra` containers then `agents` (`hivectl start`) — the
/// inverse of [`handle_stop`]. Infra comes up before agents so the agents
/// find forge/matrix/gateway ready. Per-target failures aggregated. Callers
/// resolve the [`LifecycleScope`] to these explicit name lists up front.
async fn handle_start(
coord: &Arc<Coordinator>,
agents: &[String],
infra: &[InfraContainer],
) -> Result<HostResponse> {
tracing::info!(?agents, ?infra, "start");
let mut ok_items: Vec<String> = Vec::new();
let mut errors: Vec<String> = Vec::new();
for &container in infra {
let name = container.unit_name();
match crate::priv_client::control_infra_container(container, InfraAction::Start).await {
Ok(()) => ok_items.push(name.to_owned()),
Err(e) => {
tracing::warn!(%name, error = ?e, "start: infra start failed");
errors.push(format!("{name}: {e:#}"));
}
}
}
// One DAG for all targeted agents — a per-agent start subgraph each
// (`SetWanted(Up) → Reconcile`, or a rebuild-then-start for a stale
// rev), independent roots that run concurrently on their own leases. A
// hive-wide `hivectl start` is now a single DAG, not N. Through the
// queue: persists `wanted = Up`, per-agent stale-rev upgrade to a full
// rebuild, serializes on each agent's lease. The id rides back for
// hivectl's wait loop.
let mut queued: Vec<u64> = Vec::new();
if !agents.is_empty() {
queued.push(
crate::job_queue::submit::start_many(
coord,
agents,
crate::job_queue::Source::Manual,
"manual via hivectl start".to_owned(),
)
.await,
);
ok_items.extend(agents.iter().cloned());
}
let mut resp = finish_lifecycle(ok_items, &errors);
resp.queued_dags = Some(queued);
Ok(resp)
}
/// Restart containers hive-wide (`hivectl restart`) — the DAG-based
/// sibling of [`handle_stop`]/[`handle_start`]. All targeted agents ride
/// **one** DAG (a per-agent restart subgraph each: `SetWanted → [Signal →
/// Drain →] StopForUpdate → Reconcile`, independent roots that run
/// concurrently on their own leases), not N separate DAGs — a hive-wide
/// restart is one job. `graceful` prepends signal→drain per agent. No
/// client-side stop-then-start composition, so a dropped `hivectl`
/// connection never strands an agent. Infra containers have no lease/DAG
/// and restart synchronously (stop then start).
async fn handle_restart_scoped(
coord: &Arc<Coordinator>,
scope: &LifecycleScope,
graceful: bool,
) -> Result<HostResponse> {
tracing::info!(?scope, graceful, "restart");
let agents = scoped_agents(scope).await?;
let infra = scoped_infra(scope);
let mut ok_items: Vec<String> = Vec::new();
let mut errors: Vec<String> = Vec::new();
let mut queued: Vec<u64> = Vec::new();
// One DAG for all targeted agents — a per-agent restart subgraph each
// (`SetWanted → [Signal → Drain →] StopForUpdate → Reconcile`),
// independent roots that run concurrently on their own leases. A
// hive-wide `hivectl restart` is now a single DAG, not N. No
// client-side stop-then-start composition — the whole restart survives
// a dropped connection because the DAG owns it.
if !agents.is_empty() {
queued.push(
crate::job_queue::submit::restart_many(
coord,
&agents,
graceful,
crate::job_queue::Source::Manual,
if graceful {
"manual via hivectl restart --graceful".to_owned()
} else {
"manual restart via hivectl restart".to_owned()
},
)
.await,
);
ok_items.extend(agents.iter().cloned());
}
for &container in &infra {
let name = container.unit_name();
let res = async {
crate::priv_client::control_infra_container(container, InfraAction::Stop).await?;
crate::priv_client::control_infra_container(container, InfraAction::Start).await
}
.await;
match res {
Ok(()) => ok_items.push(name.to_owned()),
Err(e) => {
tracing::warn!(%name, error = ?e, "restart: infra restart failed");
errors.push(format!("{name}: {e:#}"));
}
}
}
let mut resp = finish_lifecycle(ok_items, &errors);
resp.queued_dags = Some(queued);
Ok(resp)
}
/// Resolve which sub-agent logical names a scope targets: every live
/// container (from `lifecycle::list`) when `agents` is set or the scope is
/// "everything", plus any explicit `agent_names`. Returns de-duplicated
/// logical names with the `h-` container prefix stripped.
/// A scope that targets *every* agent rather than an explicit
/// `--agent <name>` list: either the `agents` flag or a bare
/// "everything" scope. Broad scopes are the ones whose stop/start pair
/// drives the previously-running restore set (see `handle` Stop/Start
/// arms); a targeted `--agent` stop/start must not redefine it.
fn is_broad_scope(scope: &LifecycleScope) -> bool {
scope.agents || scope.is_everything()
}
/// Assemble this hive's domain + browser-facing web URLs from c0re's
/// service env (injected by hive-c0re.nix). Each field is `None` when its
/// surface isn't browser-reachable (domain unset, forge not behind the
/// gateway, matrix GUI off), so the CLI can hint precisely instead of
/// opening a dead link. Scheme matches the existing `HIVE_FORGE_PUBLIC_URL`
/// convention (gateway terminates TLS, so https).
fn hive_urls() -> hive_host_sock::HiveUrls {
// Treat an empty env value as unset everywhere — an empty domain would
// otherwise render `swarm.peers."" = …` (invalid nix) and `https:///`.
let env = |k: &str| std::env::var(k).ok().filter(|v| !v.is_empty());
let domain = env("HYPERHIVE_HIVE_DOMAIN");
hive_host_sock::HiveUrls {
home: domain.as_ref().map(|d| format!("https://{d}/")),
forge: env("HIVE_FORGE_PUBLIC_URL"),
matrix: env("HIVE_MATRIX_PUBLIC_URL"),
domain,
}
}
async fn scoped_agents(scope: &LifecycleScope) -> Result<Vec<String>> {
use std::collections::BTreeSet;
let mut set: BTreeSet<String> = BTreeSet::new();
if is_broad_scope(scope) {
for c in lifecycle::list().await? {
let logical = c
.strip_prefix(lifecycle::AGENT_PREFIX)
.unwrap_or(&c)
.to_owned();
set.insert(logical);
}
}
for n in &scope.agent_names {
set.insert(n.clone());
}
Ok(set.into_iter().collect())
}
/// Resolve which infra containers a scope targets. An "everything" scope
/// (no flags set) selects all controllable infra; otherwise each set flag
/// maps to its [`InfraContainer`]. Fixed order for deterministic output.
fn scoped_infra(scope: &LifecycleScope) -> Vec<InfraContainer> {
let everything = scope.is_everything();
let mut out = Vec::new();
if everything || scope.ci {
out.push(InfraContainer::Ci);
}
if everything || scope.forge {
out.push(InfraContainer::Forge);
}
if everything || scope.gateway {
out.push(InfraContainer::Gateway);
}
if everything || scope.matrix {
out.push(InfraContainer::Matrix);
}
out
}
/// Build the aggregated lifecycle response: `ok` with the touched names when
/// every target succeeded, otherwise `ok: false` with the joined errors and
/// the partial success list (matches `handle_restart_all`).
fn finish_lifecycle(ok_items: Vec<String>, errors: &[String]) -> HostResponse {
if errors.is_empty() {
HostResponse::list(ok_items)
} else {
HostResponse {
ok: false,
error: Some(errors.join("; ")),
agents: Some(ok_items),
..HostResponse::default()
}
}
}