//! `/var/lib/hyperhive/agent-sockets.json` writer. Sibling to //! `agent_ports.rs`; same atomic `.tmp` + `rename()` shape so //! the gateway's nginx worker never reads a partial file. Manager //! excluded from the map (manager UI routes via the dashboard //! upstream, not per-agent `/agent//`). //! //! Full mechanism — per-agent subdir bind-mount, `hyperhive-socket-bound` //! marker gate, gateway UDS upstream, transition vs `agent-ports.json`, //! 10s poll loop: `docs/gateway.md::Per-agent unix-socket upstream`. use std::collections::BTreeMap; use std::path::{Path, PathBuf}; use anyhow::{Context, Result}; use crate::lifecycle::MANAGER_NAME; const HOST_SOCKETS_PATH: &str = "/var/lib/hyperhive/agent-sockets.json"; /// Host-side parent directory holding per-agent socket subdirs. The /// gateway container bind-mounts this whole tree (read-only) so it /// can `proxy_pass` to any agent. Each agent's container bind-mounts /// only its own `/` subdir — agents can only access their own /// sockets. pub const AGENT_SOCKET_DIR: &str = "/run/hive-agent"; /// Socket filename inside each per-agent subdir. Fixed so the path /// derives entirely from `(AGENT_SOCKET_DIR, name)` — no second /// degree of freedom for callers to get wrong. pub const SOCKET_FILENAME: &str = "web.sock"; /// Marker file the harness drops next to the socket after a /// successful `bind_unix`. Presence = "this agent has opted in to /// `hyperhive.web.useUnixSocket = true` and its harness has bound /// the socket"; absence = "the harness is still on TCP, don't /// publish the unix upstream for this agent yet". Without this gate /// the gateway would `proxy_pass` to a non-existent socket for every /// sub-agent that hasn't flipped the option yet. /// /// Renamed from `.bound` (legacy) to match the `hyperhive-` prefix /// convention for all harness-written state files (#838). `build_map` /// checks both names during the transition window so existing containers /// don't lose gateway routing before their next rebuild. pub const READY_MARKER: &str = "hyperhive-socket-bound"; const READY_MARKER_LEGACY: &str = ".bound"; #[must_use] pub fn host_sockets_path() -> PathBuf { PathBuf::from(HOST_SOCKETS_PATH) } /// Per-agent socket subdir on the host. Lifecycle pre-creates this /// before container start so the bind-mount source exists; the /// harness binds the socket inside it as `web.sock`. #[must_use] pub fn agent_dir_for(name: &str) -> PathBuf { Path::new(AGENT_SOCKET_DIR).join(name) } /// Compute the deterministic socket path for an agent. Pure function /// of the agent name so the value matches whatever /// [`agent_sockets::write`] writes for that agent, and whatever the /// harness binds via `HIVE_WEB_SOCKET`. #[must_use] pub fn socket_path_for(name: &str) -> PathBuf { agent_dir_for(name).join(SOCKET_FILENAME) } /// Compute the agent-socket map for the given logical agent names. /// Sub-agents only — manager is filtered out at the call boundary /// for the same reason it's filtered from `agent_ports::build_map` /// (manager UI is routed via the c0re dashboard upstream, not via /// `/agent//`). /// /// Also filters by `READY_MARKER` presence: only agents whose /// harness has actually bound the unix socket (and dropped the /// marker) appear in the map. Without this, the gateway would /// `proxy_pass` to a non-existent socket for every sub-agent that /// hasn't yet flipped `hyperhive.web.useUnixSocket = true`. /// /// Accepts either the new `hyperhive-socket-bound` marker or the legacy /// `.bound` marker so existing containers keep their gateway routing /// through the transition window (before their next rebuild writes the /// new marker name). /// /// `BTreeMap` keeps the JSON output sorted by key so a re-emit /// without churn produces byte-identical output — same idempotency /// shape `agent_ports::write` relies on. #[must_use] pub fn build_map(names: &[String]) -> BTreeMap { build_map_with(names, |name| { ready_marker_for(name).exists() || agent_dir_for(name).join(READY_MARKER_LEGACY).exists() }) } /// Body of `build_map` with the ready-check parameterised. Tests /// pass a predicate they control (no real filesystem access). /// Production callers go through `build_map` which wires the /// predicate to the on-disk `hyperhive-socket-bound` (or legacy /// `.bound`) marker check. fn build_map_with(names: &[String], is_ready: F) -> BTreeMap where F: Fn(&str) -> bool, { names .iter() .filter(|n| n.as_str() != MANAGER_NAME) .filter(|n| is_ready(n)) .map(|n| (n.clone(), socket_path_for(n))) .collect() } /// Path to the `hyperhive-socket-bound` marker file the harness writes /// after a successful `bind_unix`. Lives next to `web.sock` in the /// per-agent subdir so it's covered by the same bind-mount and same /// per-agent isolation as the socket itself. #[must_use] pub fn ready_marker_for(name: &str) -> PathBuf { agent_dir_for(name).join(READY_MARKER) } /// Render the map as pretty-printed JSON. Pretty so a human peek at /// `cat /var/lib/hyperhive/agent-sockets.json` shows one row per agent /// — keeps the file readable without a separate jq step (mirrors /// `agent_ports::render`). fn render(map: &BTreeMap) -> String { // Serialize as strings (PathBuf → JSON string via the Display // impl). BTreeMap → serde_json::to_string_pretty preserves key // order, so the output is deterministic across calls with the // same agent set. let stringly: BTreeMap<&String, String> = map .iter() .map(|(k, v)| (k, v.display().to_string())) .collect(); serde_json::to_string_pretty(&stringly) .expect("BTreeMap<&String, String> is always serialisable") } /// Atomically write the JSON for `names` to /// `/var/lib/hyperhive/agent-sockets.json`. Writes via a sibling /// `.tmp` + rename so a crashing process never leaves a /// partial file behind that the gateway worker would fail to parse. /// /// Idempotent — if the rendered content matches what's already on /// disk, the write + rename are skipped so the file's mtime stays /// stable and inotify watchers in the gateway (or any future /// watchers) don't fire spurious reload events. Mirrors the /// `agent_ports::write` shape — keep them in lockstep. pub fn write(names: &[String]) -> Result<()> { let map = build_map(names); let body = render(&map); let path = host_sockets_path(); if std::fs::read_to_string(&path).ok().as_deref() == Some(&body) { return Ok(()); } if let Some(parent) = path.parent() { std::fs::create_dir_all(parent) .with_context(|| format!("create {}", parent.display()))?; } let tmp = path.with_extension("json.tmp"); std::fs::write(&tmp, &body) .with_context(|| format!("write {}", tmp.display()))?; std::fs::rename(&tmp, &path).with_context(|| { format!( "rename {} -> {} (atomic publish)", tmp.display(), path.display() ) })?; Ok(()) } /// Spawn the marker poll task. Periodically re-runs `write` so the /// JSON map picks up newly-bound sockets (an agent flipping /// `hyperhive.web.useUnixSocket = true`, rebuilding, then having its /// harness drop a fresh `.bound` marker) without needing an explicit /// hook on container start. `write` is idempotent (skips the rename /// when content unchanged) so the steady-state cost is one directory /// stat per agent per poll interval. /// /// Mirrors the spawn-loop shape used by `crash_watch`, /// `reminder_scheduler`, etc. — the existing background-task /// convention in `main.rs`. pub fn spawn_poll() { tokio::spawn(async move { let mut interval = tokio::time::interval(std::time::Duration::from_secs(10)); // First tick fires immediately; that's fine — meta::sync_agents // also writes on boot, this just catches up the window before // the next agent restart. loop { interval.tick().await; match crate::lifecycle::agents_for_meta_listing().await { Ok(agents) => { let names: Vec = agents.into_iter().map(|a| a.name).collect(); if let Err(e) = write(&names) { tracing::debug!(error = ?e, "agent_sockets poll write failed"); } // Regenerate the gateway nginx include whenever // socket readiness changes — the upstream // selection (UDS vs TCP) depends on .bound markers // which change independently of topology. Write is // idempotent; skips rename when nothing changed. if let Err(e) = crate::gateway_nginx::write(&names) { tracing::debug!(error = ?e, "gateway_nginx poll write failed"); } } Err(e) => { tracing::debug!(error = ?e, "agent_sockets poll: failed to list agents"); } } } }); } #[cfg(test)] mod tests { use super::*; #[test] fn socket_path_for_uses_subdir_layout() { // Per-agent subdir + fixed socket filename — see module-level // "Per-agent subdir layout" for why this isn't a flat // `.sock`. Pin both ends so a future move (e.g. to // `/run/hyperhive/sockets/`) requires updating both the // constant and the consumers. let p = socket_path_for("iris"); assert_eq!(p, Path::new("/run/hive-agent/iris/web.sock")); } #[test] fn agent_dir_for_is_socket_parent() { // `agent_dir_for` is what lifecycle bind-mounts per agent; // `socket_path_for` lives inside it. Keep them in lockstep so // a divergence (e.g. typo in one constant) surfaces here // rather than as a confusing nspawn bind-source-not-found at // container start. let dir = agent_dir_for("iris"); let sock = socket_path_for("iris"); assert_eq!(sock.parent(), Some(dir.as_path())); } #[test] fn build_map_filters_manager() { // Use `MANAGER_NAME` in the input so the assert actually // exercises the filter path — a literal `"hm1nd"` would pass // trivially if the constant ever changed and the filter // silently became a no-op. All-ready predicate bypasses the // marker check so we exercise the manager filter in isolation. let names: Vec = ["iris", MANAGER_NAME, "argus"] .iter() .map(|s| (*s).to_owned()) .collect(); let map = build_map_with(&names, |_| true); assert!(!map.contains_key(MANAGER_NAME)); assert!(map.contains_key("iris")); assert!(map.contains_key("argus")); } #[test] fn build_map_uses_socket_path_for() { // Map values agree with the helper so callers can use either // (build_map for the bulk write, socket_path_for for one-off // lookups) without divergence. let names = vec!["iris".to_owned()]; let map = build_map_with(&names, |_| true); assert_eq!(map.get("iris"), Some(&socket_path_for("iris"))); } #[test] fn build_map_handles_empty_input() { let map = build_map_with:: bool>(&[], |_| true); assert!(map.is_empty()); } #[test] fn build_map_dedupes_via_btreemap_key_collision() { // Duplicate inputs collapse via the map; no callsite passes // dups today, but guarding the invariant here means a future // bug doesn't surface as a corrupt JSON doc (two `"iris":` // keys). Mirrors agent_ports test. let names = vec!["iris".to_owned(), "iris".to_owned()]; let map = build_map_with(&names, |_| true); assert_eq!(map.len(), 1); } #[test] fn build_map_filters_by_ready_predicate() { // Only ready agents (with `hyperhive-socket-bound` marker) get // published. Pin the behaviour so a future refactor that drops // the filter surfaces here, not as a 502-spew in the gateway. let names: Vec = ["iris", "argus", "atlas"] .iter() .map(|s| (*s).to_owned()) .collect(); // Pretend only `atlas` has flipped + bound — the gate makes // sure only opted-in agents get a UDS upstream. let map = build_map_with(&names, |name| name == "atlas"); assert!(map.contains_key("atlas")); assert!(!map.contains_key("iris")); assert!(!map.contains_key("argus")); } #[test] fn ready_marker_path_is_sibling_of_socket() { // Marker lives in the same per-agent subdir as the socket so // the same bind-mount covers both; harness writes both inside // the container, host (and gateway via shared bind-mount) // sees both at the deterministic path. let marker = ready_marker_for("iris"); let socket = socket_path_for("iris"); assert_eq!(marker.parent(), socket.parent()); assert_eq!(marker, Path::new("/run/hive-agent/iris/hyperhive-socket-bound")); } #[test] fn render_is_pretty_and_sorted() { let mut map = BTreeMap::new(); map.insert("zeta".to_owned(), PathBuf::from("/run/hive-agent/zeta/web.sock")); map.insert("alpha".to_owned(), PathBuf::from("/run/hive-agent/alpha/web.sock")); let body = render(&map); // Pretty-print = newlines between keys + indentation. assert!(body.contains('\n')); // BTreeMap sorts → alpha before zeta in output. let alpha_pos = body.find("alpha").expect("alpha in output"); let zeta_pos = body.find("zeta").expect("zeta in output"); assert!( alpha_pos < zeta_pos, "sorted order broken:\n{body}" ); } #[test] fn render_emits_paths_as_strings() { // PathBuf-valued map serialises as plain JSON strings (not // some {"inner": "..."} wrapper). Pin the shape so the // gateway-side reader can deserialise into String values // without nested struct logic. let mut map = BTreeMap::new(); map.insert("iris".to_owned(), PathBuf::from("/run/hive-agent/iris/web.sock")); let body = render(&map); assert!(body.contains("\"iris\"")); assert!(body.contains("\"/run/hive-agent/iris/web.sock\"")); } }