Watch
0
0
Fork
You've already forked hyperhive
0

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
This commit is contained in:
atlas 2026-09-26 00:17:59 +02:00 • committed by mara
commit 386741d38f
2 changed files with 110 additions and 36 deletions

View file

@ -63,9 +63,10 @@ pub fn start(agent: &str, socket_path: &Path, coord: Arc<Coordinator>) -> Result
let path = socket_path.to_path_buf(); let path = socket_path.to_path_buf();
let handle = tokio::spawn(async move { let handle = tokio::spawn(async move {
loop { let err = accept_until_fatal(
match listener.accept().await { &agent,
Ok((stream, _)) => { || async { listener.accept().await.map(|(stream, _)| stream) },
|stream| {
let agent = agent.clone(); let agent = agent.clone();
let coord = coord.clone(); let coord = coord.clone();
tokio::spawn(async move { tokio::spawn(async move {
@ -73,13 +74,10 @@ pub fn start(agent: &str, socket_path: &Path, coord: Arc<Coordinator>) -> Result
tracing::warn!(error = ?e, "agent connection failed"); tracing::warn!(error = ?e, "agent connection failed");
} }
}); });
} },
Err(e) => { )
tracing::warn!(error = ?e, "agent listener accept failed; exiting"); .await;
return; exit_on_dead_listener(&agent, &err);
}
}
}
}); });
Ok(AgentSocket { path, handle }) Ok(AgentSocket { path, handle })
} }
@ -107,9 +105,10 @@ pub fn start_manager(coord: Arc<Coordinator>) -> Result<()> {
tracing::info!(socket = %socket.display(), "manager socket listening"); tracing::info!(socket = %socket.display(), "manager socket listening");
tokio::spawn(async move { tokio::spawn(async move {
loop { let err = accept_until_fatal(
match listener.accept().await { MANAGER_AGENT,
Ok((stream, _)) => { || async { listener.accept().await.map(|(stream, _)| stream) },
|stream| {
let coord = coord.clone(); let coord = coord.clone();
tokio::spawn(async move { tokio::spawn(async move {
// Pure transport: serve as `ruth`, no privilege grant. // Pure transport: serve as `ruth`, no privilege grant.
@ -117,17 +116,71 @@ pub fn start_manager(coord: Arc<Coordinator>) -> Result<()> {
tracing::warn!(error = ?e, "manager connection failed"); tracing::warn!(error = ?e, "manager connection failed");
} }
}); });
} },
Err(e) => { )
tracing::warn!(error = ?e, "manager listener accept failed"); .await;
return; exit_on_dead_listener(MANAGER_AGENT, &err);
}
}
}
}); });
Ok(()) 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<A, F>(
agent: &str,
mut accept: A,
mut on_conn: impl FnMut(UnixStream),
) -> std::io::Error
where
A: FnMut() -> F,
F: std::future::Future<Output = std::io::Result<UnixStream>>,
{
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;
}
}
}
}
}
/// 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 /// Longest request line accepted on the agent and manager sockets, newline
/// included. The largest request a client sends is an `OperatorMsg` from the /// included. The largest request a client sends is an `OperatorMsg` from the
/// agent web UI's `/send` (`hive-agent/src/web_ui/actions.rs:21`), whose body /// agent web UI's `/send` (`hive-agent/src/web_ui/actions.rs:21`), whose body
@ -893,4 +946,23 @@ mod tests {
); );
writer.await.expect("join"); 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));
}
} }

View file

@ -10,9 +10,11 @@
//! - `run_reconcile` calls `register_agent` immediately after `start_with_fallback`. //! - `run_reconcile` calls `register_agent` immediately after `start_with_fallback`.
//! - `kill`/`destroy` paths call `unregister_agent`. //! - `kill`/`destroy` paths call `unregister_agent`.
//! //!
//! No recurring poll is needed because c0re owns the listener lifecycle — //! No recurring poll is needed because c0re owns the listener lifecycle. An
//! a listener can only disappear when c0re itself restarts, which is exactly //! accept loop retries every error except one that leaves its fd unusable,
//! the case `sync_on_start` covers. //! 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; use std::sync::Arc;