From 386741d38f3096faa2f76af7f4ca668b36a3d777 Mon Sep 17 00:00:00 2001 From: atlas Date: Sat, 26 Sep 2026 00:17:59 +0200 Subject: [PATCH] hive-c0re: keep agent and manager listeners alive across accept errors Both accept loops returned on the first accept() error. The AgentSocket stayed in Coordinator.agents, so mcp_sockets::sync_on_start skipped the agent, and only a Start node or a hive-c0re restart bound it again. The unit sets no LimitNOFILE, so a single EMFILE at the 1024-fd soft limit cut the agent off from send/recv/ack with one warn line as the only trace. Both loops now share accept_until_fatal. Connection errors (ECONNABORTED/ECONNRESET/ECONNREFUSED) retry at once, any other error retries after 1 s, following axum's serve loop that the dashboard already runs. Errors that leave the listening fd unusable (EBADF, EFAULT, EINVAL, ENOTSOCK, EOPNOTSUPP) log at error and exit the process: nothing re-binds a listener while hive-c0re runs, the manager listener has no re-bind path at all, and the unit's Restart=on-failure restart runs start_manager and sync_on_start, which re-bind every socket. mcp_sockets.rs now states that invariant instead of asserting that a listener can only disappear on restart. LimitNOFILE is left unset: per-agent fd use is a listener plus a Recv long-poll plus short-lived requests, and a higher limit would raise the memory ceiling of the per-line bound (bound x open connections). Closes #4721 --- hive-c0re/src/socket_server/mod.rs | 138 ++++++++++++++++++++------- hive-c0re/src/workers/mcp_sockets.rs | 8 +- 2 files changed, 110 insertions(+), 36 deletions(-) diff --git a/hive-c0re/src/socket_server/mod.rs b/hive-c0re/src/socket_server/mod.rs index aff0862f..052efee9 100644 --- a/hive-c0re/src/socket_server/mod.rs +++ b/hive-c0re/src/socket_server/mod.rs @@ -63,23 +63,21 @@ pub fn start(agent: &str, socket_path: &Path, coord: Arc) -> Result let path = socket_path.to_path_buf(); let handle = tokio::spawn(async move { - loop { - match listener.accept().await { - Ok((stream, _)) => { - let agent = agent.clone(); - let coord = coord.clone(); - tokio::spawn(async move { - if let Err(e) = serve(stream, agent, coord).await { - tracing::warn!(error = ?e, "agent connection failed"); - } - }); - } - Err(e) => { - tracing::warn!(error = ?e, "agent listener accept failed; exiting"); - return; - } - } - } + let err = accept_until_fatal( + &agent, + || async { listener.accept().await.map(|(stream, _)| stream) }, + |stream| { + let agent = agent.clone(); + let coord = coord.clone(); + tokio::spawn(async move { + if let Err(e) = serve(stream, agent, coord).await { + tracing::warn!(error = ?e, "agent connection failed"); + } + }); + }, + ) + .await; + exit_on_dead_listener(&agent, &err); }); Ok(AgentSocket { path, handle }) } @@ -107,25 +105,80 @@ pub fn start_manager(coord: Arc) -> Result<()> { tracing::info!(socket = %socket.display(), "manager socket listening"); tokio::spawn(async move { - loop { - match listener.accept().await { - Ok((stream, _)) => { - let coord = coord.clone(); - tokio::spawn(async move { - // Pure transport: serve as `ruth`, no privilege grant. - if let Err(e) = serve(stream, MANAGER_AGENT.to_owned(), coord).await { - tracing::warn!(error = ?e, "manager connection failed"); - } - }); - } - Err(e) => { - tracing::warn!(error = ?e, "manager listener accept failed"); - return; + let err = accept_until_fatal( + MANAGER_AGENT, + || async { listener.accept().await.map(|(stream, _)| stream) }, + |stream| { + let coord = coord.clone(); + tokio::spawn(async move { + // Pure transport: serve as `ruth`, no privilege grant. + if let Err(e) = serve(stream, MANAGER_AGENT.to_owned(), coord).await { + tracing::warn!(error = ?e, "manager connection failed"); + } + }); + }, + ) + .await; + exit_on_dead_listener(MANAGER_AGENT, &err); + }); + Ok(()) +} + +/// Pause after an accept error that is not about a single connection. EMFILE +/// and ENFILE leave the listener readable, so retrying at once would spin until +/// some fd is closed. Same policy as axum's `serve`, which the dashboard runs. +const ACCEPT_RETRY_DELAY: std::time::Duration = std::time::Duration::from_secs(1); + +/// Accept errors that mean the listening fd itself is unusable: retrying it +/// returns the same error forever. +fn accept_error_is_fatal(e: &std::io::Error) -> bool { + matches!( + e.raw_os_error(), + Some(libc::EBADF | libc::EFAULT | libc::EINVAL | libc::ENOTSOCK | libc::EOPNOTSUPP) + ) +} + +/// Hand each accepted connection to `on_conn`, retrying every accept error +/// except a fatal one, which is returned. `agent` labels the log lines. +async fn accept_until_fatal( + agent: &str, + mut accept: A, + mut on_conn: impl FnMut(UnixStream), +) -> std::io::Error +where + A: FnMut() -> F, + F: std::future::Future>, +{ + loop { + match accept().await { + Ok(stream) => on_conn(stream), + Err(e) if accept_error_is_fatal(&e) => return e, + Err(e) => { + tracing::warn!(%agent, error = ?e, "socket accept failed; retrying"); + if !matches!( + e.kind(), + std::io::ErrorKind::ConnectionAborted + | std::io::ErrorKind::ConnectionReset + | std::io::ErrorKind::ConnectionRefused + ) { + tokio::time::sleep(ACCEPT_RETRY_DELAY).await; } } } - }); - Ok(()) + } +} + +/// Exit on a listener that can no longer accept. A running hive-c0re re-binds +/// an agent socket only on that agent's next start, and the manager socket +/// never; on restart (`Restart = "on-failure"` on the unit) `start_manager` and +/// `mcp_sockets::sync_on_start` re-bind every one. +fn exit_on_dead_listener(agent: &str, e: &std::io::Error) -> ! { + tracing::error!( + %agent, + error = ?e, + "socket listener cannot accept; exiting so hive-c0re restarts and re-binds all sockets" + ); + std::process::exit(1) } /// Longest request line accepted on the agent and manager sockets, newline @@ -893,4 +946,23 @@ mod tests { ); writer.await.expect("join"); } + + #[tokio::test] + async fn accept_loop_survives_a_transient_error() { + let (client, _peer) = UnixStream::pair().expect("socketpair"); + let mut script = std::collections::VecDeque::from([ + Err(std::io::Error::from_raw_os_error(libc::EMFILE)), + Ok(client), + Err(std::io::Error::from_raw_os_error(libc::EBADF)), + ]); + let mut served = 0; + let err = accept_until_fatal( + "iris", + || std::future::ready(script.pop_front().expect("accept after a fatal error")), + |_stream| served += 1, + ) + .await; + assert_eq!(served, 1, "connection after EMFILE not served"); + assert_eq!(err.raw_os_error(), Some(libc::EBADF)); + } } diff --git a/hive-c0re/src/workers/mcp_sockets.rs b/hive-c0re/src/workers/mcp_sockets.rs index ee0f3de3..49d9b860 100644 --- a/hive-c0re/src/workers/mcp_sockets.rs +++ b/hive-c0re/src/workers/mcp_sockets.rs @@ -10,9 +10,11 @@ //! - `run_reconcile` calls `register_agent` immediately after `start_with_fallback`. //! - `kill`/`destroy` paths call `unregister_agent`. //! -//! No recurring poll is needed because c0re owns the listener lifecycle — -//! a listener can only disappear when c0re itself restarts, which is exactly -//! the case `sync_on_start` covers. +//! No recurring poll is needed because c0re owns the listener lifecycle. An +//! accept loop retries every error except one that leaves its fd unusable, +//! and on that one hive-c0re exits (`socket_server::exit_on_dead_listener`), +//! so a listener only disappears across a restart — exactly the case +//! `sync_on_start` covers. use std::sync::Arc;