diff --git a/hive-ag3nt/src/web_ui.rs b/hive-ag3nt/src/web_ui.rs index 39d0fdc5..4728269c 100644 --- a/hive-ag3nt/src/web_ui.rs +++ b/hive-ag3nt/src/web_ui.rs @@ -195,6 +195,24 @@ fn bind_unix(path: &Path) -> Result { use std::os::unix::fs::PermissionsExt; std::fs::set_permissions(path, std::fs::Permissions::from_mode(0o660)) .with_context(|| format!("set perms on {}", path.display()))?; + // Drop a `.bound` marker next to the socket so c0re's + // `agent_sockets::write` can filter the JSON map to only include + // agents whose harness has actually opted in to (and bound) the + // unix socket. Without this gating, gateway would `proxy_pass` + // to a non-existent socket for any sub-agent that hasn't flipped + // `hyperhive.web.useUnixSocket = true` yet — atlas's concern on + // PR #813. Best-effort: a failed write isn't fatal (the harness + // still binds + serves on the socket), it just means the + // gateway side keeps using the TCP upstream for one more sync. + if let Some(parent) = path.parent() { + let marker = parent.join(".bound"); + if let Err(e) = std::fs::write(&marker, b"") { + tracing::warn!( + marker = %marker.display(), error = %e, + "failed to write .bound marker — gateway may keep TCP upstream" + ); + } + } Ok(listener) } diff --git a/hive-c0re/src/agent_sockets.rs b/hive-c0re/src/agent_sockets.rs index 0c0b6c58..cbfff2fc 100644 --- a/hive-c0re/src/agent_sockets.rs +++ b/hive-c0re/src/agent_sockets.rs @@ -75,6 +75,16 @@ pub const AGENT_SOCKET_DIR: &str = "/run/hive-agent"; /// 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". Atlas's gate on +/// PR #813 — without it, the gateway would `proxy_pass` to a +/// non-existent socket for every sub-agent that hasn't flipped the +/// option yet. +pub const READY_MARKER: &str = ".bound"; + #[must_use] pub fn host_sockets_path() -> PathBuf { PathBuf::from(HOST_SOCKETS_PATH) @@ -103,18 +113,46 @@ pub fn socket_path_for(name: &str) -> PathBuf { /// (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. Atlas's gate on #813 — 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`. +/// /// `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()) +} + +/// 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 `.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 per-agent `.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 @@ -142,6 +180,41 @@ fn render(map: &BTreeMap) -> String { /// 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. +/// 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(coord: std::sync::Arc) { + let _ = coord; + 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"); + } + } + Err(e) => { + tracing::debug!(error = ?e, "agent_sockets poll: failed to list agents"); + } + } + } + }); +} + pub fn write(names: &[String]) -> Result<()> { let map = build_map(names); let body = render(&map); @@ -199,12 +272,13 @@ mod tests { // exercises the filter path — a literal `"hm1nd"` would pass // trivially if the constant ever changed and the filter // silently became a no-op (same pattern as #748 fix on - // agent_ports::build_map). + // agent_ports::build_map). 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(&names); + let map = build_map_with(&names, |_| true); assert!(!map.contains_key(MANAGER_NAME)); assert!(map.contains_key("iris")); assert!(map.contains_key("argus")); @@ -216,13 +290,13 @@ mod tests { // (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(&names); + 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(&[]); + let map = build_map_with:: bool>(&[], |_| true); assert!(map.is_empty()); } @@ -233,10 +307,41 @@ mod tests { // 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(&names); + let map = build_map_with(&names, |_| true); assert_eq!(map.len(), 1); } + #[test] + fn build_map_filters_by_ready_predicate() { + // The new gate: only ready agents (with `.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 (mara on PR #813: + // "agents can only access their own sockets" — 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/.bound")); + } + #[test] fn render_is_pretty_and_sorted() { let mut map = BTreeMap::new(); diff --git a/hive-c0re/src/main.rs b/hive-c0re/src/main.rs index 35ec3f96..757f4c33 100644 --- a/hive-c0re/src/main.rs +++ b/hive-c0re/src/main.rs @@ -12,9 +12,9 @@ use hive_sh4re::{HostRequest, HostResponse}; // explicit (any new daemon entry point reads off the next add). use hive_c0re::coordinator::Coordinator; use hive_c0re::{ - auto_update, broker, client, crash_watch, dashboard, dashboard_events, events_vacuum, forge, - manager_server, matrix, migrate, rebuild_queue, reminder_scheduler, scheduled_prompts_worker, - server, stats_vacuum, + agent_sockets, auto_update, broker, client, crash_watch, dashboard, dashboard_events, + events_vacuum, forge, manager_server, matrix, migrate, rebuild_queue, reminder_scheduler, + scheduled_prompts_worker, server, stats_vacuum, }; #[derive(Parser)] @@ -253,6 +253,15 @@ async fn cmd_serve( // when a previously-running container goes away without an // operator-initiated transient state. crash_watch::spawn(coord.clone()); + // Agent-sockets marker poll: re-fires `agent_sockets::write` + // every 10s so the JSON picks up newly-bound `.bound` markers + // (a sub-agent flipping `hyperhive.web.useUnixSocket = true`, + // rebuilding, then having its harness bind the socket) without + // needing an explicit hook on each container start. write() is + // idempotent so steady-state cost is one stat per agent per + // tick. closes atlas's #813 concern that agents in the JSON + // would 502 the gateway until they actually opt in. + agent_sockets::spawn_poll(coord.clone()); // Reminder scheduler: drains due reminders + handles // file_path payload persistence. See reminder_scheduler.rs. reminder_scheduler::spawn(coord.clone());