Watch
0
0
Fork
You've already forked hyperhive
0
hyperhive/hive-sock-client/src/lib.rs
atlas 7b1fe5f9d3 hive-sock-client, web proxy, HTTP clients: bound connect and response waits
hive-sock-client: each attempt now bounds connect (5s), write (10s) and
the wait for the response (60s by default). The response bound is per
call through the new `request_within`, which hive-agent's serve-loop
`Recv` uses with its 180s long-poll plus 30s headroom. A response
timeout is terminal rather than retried: the server holds the request,
so a retry re-sends something it may still act on and multiplies the
wait by the backoff schedule.

Outbound HTTP: the matrix login/whoami clients in swarm-controller and
hive-c0re's dashboard (5s connect, 30s request), the authelia-bridge
client (5s/30s; ensuring an identity runs an argon2 hash first) and the
ci-runner forge calls (5s/15s, config_pr_poll's forge budget) get a
connect_timeout and a request timeout. Timeout errors name the bound
that fired.

hive-agent's unix-socket extra web proxy bounds the connect (5s) and
the wait for the response head (30s, the http sibling's budget); the
body read stays unbounded.

Refs #4723
2026-09-27 03:46:53 +02:00

515 lines
18 KiB
Rust

//! JSON-line-over-unix-socket client, shared by every daemon that talks to
//! a hyperhive socket.
//!
//! The wire protocol is identical everywhere — connect, write one line of
//! JSON, read one line of JSON back — so this crate is generic over the
//! request and response types and knows nothing about either protocol. The
//! host-served control socket and the harness's in-agent socket both use
//! it with their own wire-type crates.
//!
//! Two axes of behaviour are real and stay caller-selectable; everything
//! else is shared:
//!
//! - **retry**: [`Retry::RideOutRestart`] for callers with no natural
//! retry of their own, [`Retry::None`] for callers already inside a poll
//! loop whose interval *is* the retry.
//! - **response**: [`request`] decodes it, [`notify`] drains and discards
//! it.
//!
//! Every attempt is bounded: connect, write and waiting for the response
//! each have a deadline, so a peer that accepts and never answers fails the
//! call instead of hanging it. The response deadline is per call —
//! [`request_within`] takes one for verbs that legitimately park
//! server-side (the broker's `Recv` long-poll); everything else gets
//! `DEFAULT_RESPONSE_TIMEOUT`.
//!
//! Whether a failure propagates or is logged and swallowed is the caller's
//! choice and stays at the call site — it is not a property of the
//! transport.
use std::path::Path;
use std::time::Duration;
use anyhow::{Context, Result, anyhow};
use serde::Serialize;
use serde::de::DeserializeOwned;
use tokio::io::{AsyncBufReadExt, AsyncWriteExt, BufReader};
use tokio::net::UnixStream;
/// Backoff schedule for [`Retry::RideOutRestart`]. Five entries → up to 5
/// retries on top of the initial attempt; total wall-clock cap =
/// 2+4+8+16+30 = 60s. Sized to ride out a service restart (systemd usually
/// has the unix socket back inside ~5s) without the caller having to
/// handle the transient itself.
const RIDE_OUT_RESTART_BACKOFFS_MS: &[u64] = &[2_000, 4_000, 8_000, 16_000, 30_000];
/// Bound on connecting. The listener is on the same host, so a connect
/// still pending after this long is not going to complete.
const CONNECT_TIMEOUT: Duration = Duration::from_secs(5);
/// Bound on writing (and flushing / half-closing) the request line. A peer
/// that has not drained one JSON line in this long has stopped reading.
const WRITE_TIMEOUT: Duration = Duration::from_secs(10);
/// Response deadline for every call that does not pass its own. Sized for
/// the slowest non-polling verb: `CreateRepo`, which makes several
/// sequential forge API calls in hive-c0re before it answers.
const DEFAULT_RESPONSE_TIMEOUT: Duration = Duration::from_mins(1);
/// What to do when a connect or I/O attempt fails.
///
/// Deliberately two named policies rather than a configurable schedule:
/// only two behaviours exist in the tree, and naming them keeps the
/// *reason* for each choice readable at the call site.
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub enum Retry {
/// Fail on the first failure. For callers inside a poll loop, where
/// the poll interval already is the retry — a second, in-request
/// backoff would stack sleeps on top of it and delay the rest of the
/// batch.
None,
/// Back off on `RIDE_OUT_RESTART_BACKOFFS_MS` (~60s total). For
/// callers with no natural retry of their own, where a surfaced
/// transient costs more than the wait.
RideOutRestart,
}
impl Retry {
/// Sleep schedule between attempts; its length is the retry budget.
fn backoffs(self) -> &'static [u64] {
match self {
Self::None => &[],
Self::RideOutRestart => RIDE_OUT_RESTART_BACKOFFS_MS,
}
}
}
/// Send `req` over `socket` and decode the single-line JSON response.
///
/// # Errors
///
/// Returns an error if `req` cannot be serialised, if the socket is
/// unreachable after the retry budget is spent, if the server closes
/// without responding or does not respond within
/// `DEFAULT_RESPONSE_TIMEOUT`, or if the response does not deserialise
/// into `Resp`.
pub async fn request<Req, Resp>(socket: &Path, req: &Req, retry: Retry) -> Result<Resp>
where
Req: Serialize + ?Sized,
Resp: DeserializeOwned,
{
request_within(socket, req, retry, DEFAULT_RESPONSE_TIMEOUT).await
}
/// Same as [`request`], but waits up to `response_timeout` for the
/// response instead of `DEFAULT_RESPONSE_TIMEOUT`.
///
/// For verbs the server deliberately holds open, such as a `Recv` that
/// long-polls: the deadline has to exceed the wait the request asks for,
/// or every idle poll reads as a stuck server.
///
/// # Errors
///
/// Same as [`request`], with `response_timeout` as the response deadline.
pub async fn request_within<Req, Resp>(
socket: &Path,
req: &Req,
retry: Retry,
response_timeout: Duration,
) -> Result<Resp>
where
Req: Serialize + ?Sized,
Resp: DeserializeOwned,
{
let payload = encode(req)?;
let (line, _) = with_retry(socket, &payload, Mode::Decode, retry, response_timeout).await?;
decode(socket, &line)
}
/// Same as [`request`], but also reports how many retries it took past the
/// initial attempt (0 = succeeded first try).
///
/// MCP tool handlers use this to append a one-line hint to the tool result
/// when retries happened, so claude reads the earlier socket flake as a
/// transient rather than a content error worth an LLM-level retry.
///
/// # Errors
///
/// Same as [`request`].
pub async fn request_retried<Req, Resp>(
socket: &Path,
req: &Req,
retry: Retry,
) -> Result<(Resp, u32)>
where
Req: Serialize + ?Sized,
Resp: DeserializeOwned,
{
let payload = encode(req)?;
let (line, retries) = with_retry(
socket,
&payload,
Mode::Decode,
retry,
DEFAULT_RESPONSE_TIMEOUT,
)
.await?;
Ok((decode(socket, &line)?, retries))
}
/// Send `req` over `socket`, half-close, and drain the response line
/// without decoding it.
///
/// For fire-and-forget ops whose reply carries nothing the caller acts on.
/// The drain is not optional: without it the server's write-back lands on
/// a closed socket.
///
/// # Errors
///
/// Returns an error if `req` cannot be serialised, or if the socket is
/// unreachable / the write fails after the retry budget is spent. A
/// failure to read the discarded response is not an error — the request
/// was already delivered.
pub async fn notify<Req>(socket: &Path, req: &Req, retry: Retry) -> Result<()>
where
Req: Serialize + ?Sized,
{
let payload = encode(req)?;
with_retry(
socket,
&payload,
Mode::Drain,
retry,
DEFAULT_RESPONSE_TIMEOUT,
)
.await?;
Ok(())
}
/// What to do with the server's response line.
#[derive(Clone, Copy)]
enum Mode {
/// Flush the write half, read the response, hand it back to be parsed.
Decode,
/// Half-close the write half, read-and-discard the response.
Drain,
}
/// Serialise `req` into one newline-terminated JSON line.
///
/// Kept outside the retry loop on purpose: a serialisation failure is
/// deterministic, so retrying it would only reproduce it.
fn encode<Req>(req: &Req) -> Result<Vec<u8>>
where
Req: Serialize + ?Sized,
{
let mut payload = serde_json::to_string(req).context("serialise request")?;
payload.push('\n');
Ok(payload.into_bytes())
}
/// Decode one response line into `Resp`.
fn decode<Resp: DeserializeOwned>(socket: &Path, line: &str) -> Result<Resp> {
serde_json::from_str(line.trim())
.with_context(|| format!("decode response from {}", socket.display()))
}
/// The server accepted the request and did not answer within the response
/// deadline.
///
/// Not retried: the server is up and holding the request, so another
/// attempt would re-send something it may still act on, and would multiply
/// the wait by the whole backoff schedule.
#[derive(Debug)]
struct ResponseTimeout {
socket: std::path::PathBuf,
after: Duration,
}
impl std::fmt::Display for ResponseTimeout {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(
f,
"waiting for a response from {} timed out after {:?}",
self.socket.display(),
self.after
)
}
}
impl std::error::Error for ResponseTimeout {}
/// Run [`try_once`] until it succeeds, the retry budget is spent, or the
/// server times out on a delivered request, returning the response line and
/// the number of retries it took.
async fn with_retry(
socket: &Path,
payload: &[u8],
mode: Mode,
retry: Retry,
response_timeout: Duration,
) -> Result<(String, u32)> {
let backoffs = retry.backoffs();
let mut attempt: usize = 0;
loop {
match try_once(socket, payload, mode, response_timeout).await {
Ok(line) => return Ok((line, u32::try_from(attempt).unwrap_or(u32::MAX))),
Err(e) if e.is::<ResponseTimeout>() => return Err(e),
Err(e) => {
let Some(&sleep_ms) = backoffs.get(attempt) else {
return Err(exhausted(e, backoffs));
};
tracing::warn!(
attempt = attempt + 1,
sleep_ms,
socket = %socket.display(),
error = %e,
"hive socket attempt failed; retrying"
);
tokio::time::sleep(Duration::from_millis(sleep_ms)).await;
attempt += 1;
}
}
}
}
/// Annotate the final failure with how long the retry schedule tried, so a
/// surfaced error says whether it was one shot or a full ride-out.
fn exhausted(err: anyhow::Error, backoffs: &[u64]) -> anyhow::Error {
if backoffs.is_empty() {
return err;
}
let total_s = backoffs.iter().sum::<u64>() / 1_000;
err.context(format!(
"gave up after {} retries over ~{total_s}s",
backoffs.len()
))
}
/// One connect / write / read cycle, each step under its own deadline.
/// Every error it returns except [`ResponseTimeout`] is retryable by
/// construction — the deterministic failures (serialise, deserialise)
/// happen outside the retry loop.
async fn try_once(
socket: &Path,
payload: &[u8],
mode: Mode,
response_timeout: Duration,
) -> Result<String> {
let connected = tokio::time::timeout(CONNECT_TIMEOUT, UnixStream::connect(socket))
.await
.map_err(|_| {
anyhow!(
"connect to {} timed out after {CONNECT_TIMEOUT:?}",
socket.display()
)
})?;
let stream = connected.map_err(|e| {
// A refused or missing socket usually means the listener is
// mid-restart (operator redeploy, harness restart) — it is
// recreated on boot and a retrying caller rides it out. When the
// error does surface (budget spent, or a fail-fast caller) the
// hint marks it as a likely transient rather than a hard failure
// worth escalating. The path stays in the error either way: a
// permission problem that reads as "is the daemon running?" sends
// the operator to fix the wrong thing.
let restarting = matches!(
e.kind(),
std::io::ErrorKind::ConnectionRefused | std::io::ErrorKind::NotFound
);
let err = anyhow::Error::new(e).context(format!("connect to {}", socket.display()));
if restarting {
err.context("the listener may be restarting (e.g. an operator redeploy)")
} else {
err
}
})?;
let (read, mut write) = stream.into_split();
let send = async {
write
.write_all(payload)
.await
.with_context(|| format!("write to {}", socket.display()))?;
match mode {
Mode::Decode => write
.flush()
.await
.with_context(|| format!("flush {}", socket.display())),
Mode::Drain => write
.shutdown()
.await
.with_context(|| format!("shutdown write to {}", socket.display())),
}
};
tokio::time::timeout(WRITE_TIMEOUT, send)
.await
.map_err(|_| {
anyhow!(
"write to {} timed out after {WRITE_TIMEOUT:?}",
socket.display()
)
})??;
let mut reader = BufReader::new(read);
let mut line = String::new();
let response = tokio::time::timeout(response_timeout, reader.read_line(&mut line)).await;
match mode {
Mode::Decode => {
let read_bytes = response
.map_err(|_| {
anyhow::Error::new(ResponseTimeout {
socket: socket.to_path_buf(),
after: response_timeout,
})
})?
.with_context(|| format!("read from {}", socket.display()))?;
if read_bytes == 0 || line.is_empty() {
return Err(anyhow!(
"{} closed the connection without responding",
socket.display()
));
}
}
Mode::Drain => {
// Discarded, including an error or a timeout reading it: the
// request is already delivered and the caller acts on nothing
// here.
let _ = response;
line.clear();
}
}
Ok(line)
}
#[cfg(test)]
mod tests {
use std::path::PathBuf;
use std::time::{Duration, Instant};
use tokio::io::{AsyncBufReadExt, AsyncWriteExt, BufReader};
use tokio::net::UnixListener;
use super::{Retry, notify, request, request_within};
/// A fresh socket path under the test's temp dir, unique per test.
fn socket_path(name: &str) -> PathBuf {
let path = std::env::temp_dir().join(format!(
"hive-sock-client-{}-{name}.sock",
std::process::id()
));
let _ = std::fs::remove_file(&path);
path
}
/// A peer that accepts and never answers fails the request within the
/// response deadline, names the bound in the error, and is not retried
/// even under `RideOutRestart` (whose first backoff alone is 2s).
#[tokio::test]
async fn silent_peer_times_out_within_the_response_deadline() {
let path = socket_path("silent");
let listener = UnixListener::bind(&path).expect("bind test socket");
let server = tokio::spawn(async move {
let mut held = Vec::new();
while let Ok((stream, _)) = listener.accept().await {
held.push(stream);
}
});
let started = Instant::now();
let outcome = tokio::time::timeout(
Duration::from_secs(10),
request_within::<(), serde_json::Value>(
&path,
&(),
Retry::RideOutRestart,
Duration::from_millis(200),
),
)
.await;
server.abort();
let _ = std::fs::remove_file(&path);
let err = outcome
.expect("request hung past its response deadline")
.expect_err("a peer that never answers must fail the request");
let msg = format!("{err:#}");
assert!(msg.contains("timed out after 200ms"), "{msg}");
assert!(
started.elapsed() < Duration::from_secs(2),
"timed-out request was retried: {:?}",
started.elapsed()
);
}
/// The control for the test above: the same listener shape, answering,
/// completes under the same deadline — so the failure there is the
/// peer's silence, not the harness.
#[tokio::test]
async fn answering_peer_completes_within_the_response_deadline() {
let path = socket_path("answering");
let listener = UnixListener::bind(&path).expect("bind test socket");
let server = tokio::spawn(async move {
let (stream, _) = listener.accept().await.expect("accept");
let (read, mut write) = stream.into_split();
let mut line = String::new();
BufReader::new(read)
.read_line(&mut line)
.await
.expect("read request");
write.write_all(b"\"ok\"\n").await.expect("write response");
});
let resp = request_within::<(), serde_json::Value>(
&path,
&(),
Retry::None,
Duration::from_millis(200),
)
.await;
server.abort();
let _ = std::fs::remove_file(&path);
assert_eq!(resp.expect("answering peer"), serde_json::json!("ok"));
}
/// A connect to a non-existent socket path (ENOENT → `NotFound`) is
/// annotated with both the socket path and the "may be restarting"
/// hint, so a surfaced tool error reads as the expected transient.
#[tokio::test]
async fn missing_socket_names_the_path_and_hints_restart() {
let bogus = std::path::Path::new("/nonexistent/hive/mcp.sock");
let err = request::<(), serde_json::Value>(bogus, &(), Retry::None)
.await
.expect_err("connect to a non-existent socket must fail");
let msg = format!("{err:#}");
assert!(msg.contains("restarting"), "missing restart hint: {msg}");
assert!(msg.contains("connect to"), "missing connect context: {msg}");
}
/// `Retry::None` returns promptly rather than sleeping the ride-out
/// schedule — its callers are poll ticks that need the rest of the
/// batch to still run this cycle.
#[tokio::test]
async fn no_retry_fails_fast() {
let bogus = std::path::Path::new("/nonexistent/hive/agent.sock");
let started = std::time::Instant::now();
notify(bogus, &(), Retry::None)
.await
.expect_err("connect to a non-existent socket must fail");
assert!(
started.elapsed() < std::time::Duration::from_secs(1),
"Retry::None slept: {:?}",
started.elapsed()
);
}
/// The ride-out schedule is a full minute of patience — the value the
/// harness's tool callers depend on to not surface a restart.
#[test]
fn ride_out_restart_budget_is_a_minute() {
let total_ms: u64 = Retry::RideOutRestart.backoffs().iter().sum();
assert_eq!(total_ms, 60_000);
assert!(Retry::None.backoffs().is_empty());
}
}