use std::path::Path; use std::sync::Arc; use anyhow::{Context, Result}; use hive_sh4re::priv_proto::{InfraAction, InfraContainer}; use hive_sh4re::{HostRequest, HostResponse, LifecycleScope}; 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) -> 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) -> 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::(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?; } } async fn dispatch(req: &HostRequest, coord: Arc) -> HostResponse { let result: anyhow::Result = async { Ok(match req { HostRequest::Spawn { name } => handle_spawn(&coord, name).await?, HostRequest::RequestSpawn { name } => { tracing::info!(%name, "request_spawn"); let id = coord.approvals.submit_kind( name, hive_sh4re::ApprovalKind::Spawn, "", None, "operator", )?; tracing::info!(%id, %name, "spawn approval queued"); HostResponse::success() } HostRequest::Kill { name } => submit_single(&coord, name, Verb::Kill), HostRequest::Restart { name } => submit_single(&coord, name, Verb::Restart), HostRequest::RestartAll => handle_restart_all(&coord).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, *purge).await?; HostResponse::success() } HostRequest::Rebuild { name } => submit_single(&coord, name, Verb::Rebuild), HostRequest::QueueDag { id } => { // The polled DAG first, then its live fan-out children. let dags = coord .job_queue .snapshot() .into_iter() .filter(|d| d.id == *id || d.parent_id == Some(*id)) .collect(); HostResponse::dags(dags) } HostRequest::List => HostResponse::list(lifecycle::list().await?), HostRequest::AgentStatus => { let rows = crate::container_view::build_all(&coord) .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, pending_reminders: v.pending_reminders, parent: v.parent, }) .collect(); HostResponse::agent_statuses(rows) } // 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).await?; 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, new_parent.as_deref()) .await .map_err(anyhow::Error::msg)?; HostResponse::success() } }) } .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, name: &str) -> Result { tracing::info!(%name, "spawn"); let agent_dir = coord.ensure_runtime(name)?; let hive = coord.hive_env(); let paths = Coordinator::agent_paths(name, agent_dir); 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"); } coord.notify_manager(&hive_sh4re::HelperEvent::Spawned { agent: name.to_owned(), ok: true, note: None, sha: None, }); } Err(e) => { // Roll back socket registration if container creation failed. coord.unregister_agent(name); coord.notify_manager(&hive_sh4re::HelperEvent::Spawned { agent: name.to_owned(), ok: false, note: Some(format!("{e:#}")), sha: None, }); return Err(e); } } Ok(HostResponse::success()) } /// 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, } fn submit_single(coord: &Arc, 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(), ) } Verb::Restart => { tracing::info!(%name, "restart"); submit::restart( coord, name, Source::Manual, "manual restart via hivectl".to_owned(), ) } 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 by submitting one restart DAG per agent — /// each serializes on its own lease, so unrelated agents' restarts /// overlap while nothing races an in-flight rebuild. Returns once all /// are queued; per-agent results surface on the queue. async fn handle_restart_all(coord: &Arc) -> Result { tracing::info!("restart-all"); let agents = lifecycle::list().await?; let mut ok_agents: Vec = Vec::new(); for agent in &agents { let Some(logical) = agent.strip_prefix(lifecycle::AGENT_PREFIX) else { continue; }; crate::job_queue::submit::restart( coord, logical, crate::job_queue::Source::Manual, "manual restart via hivectl restart-all".to_owned(), ); ok_agents.push(logical.to_owned()); } Ok(HostResponse::list(ok_agents)) } /// 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, agents: &[String], infra: &[InfraContainer], graceful: bool, ) -> Result { tracing::info!(?agents, ?infra, graceful, "stop"); let mut ok_items: Vec = Vec::new(); let mut errors: Vec = Vec::new(); let mut queued: Vec = Vec::new(); for agent in agents { let reason = if graceful { "manual via hivectl graceful stop" } else { "manual via hivectl stop" }; let id = if graceful { crate::job_queue::submit::graceful_stop( coord, agent, crate::job_queue::Source::Manual, reason.to_owned(), ) } else { crate::job_queue::submit::stop( coord, agent, crate::job_queue::Source::Manual, reason.to_owned(), ) }; queued.push(id); ok_items.push(agent.clone()); } 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) } /// 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, agents: &[String], infra: &[InfraContainer], ) -> Result { tracing::info!(?agents, ?infra, "start"); let mut ok_items: Vec = Vec::new(); let mut errors: Vec = 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:#}")); } } } let mut queued: Vec = Vec::new(); for agent in agents { // Through the queue: persists `wanted = Up`, upgrades a // stale-rev start to a full rebuild, and serializes on the // agent's lease. Ids ride back for hivectl's wait loop. queued.push(crate::job_queue::submit::start( coord, agent, crate::job_queue::Source::Manual, "manual via hivectl start".to_owned(), )); ok_items.push(agent.clone()); } 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 ` 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_sh4re::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_sh4re::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> { use std::collections::BTreeSet; let mut set: BTreeSet = 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 { 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, 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() } } }