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
This commit is contained in:
atlas 2026-09-26 18:29:48 +02:00 • committed by mara
commit 7b1fe5f9d3
7 changed files with 412 additions and 60 deletions

View file

@ -16,6 +16,13 @@
//! - **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.
@ -36,6 +43,19 @@ use tokio::net::UnixStream;
/// 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:
@ -70,16 +90,40 @@ impl Retry {
///
/// 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 if the response does not deserialise into
/// `Resp`.
/// 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_retried(socket, req, retry)
.await
.map(|(resp, _)| resp)
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
@ -102,10 +146,15 @@ where
Resp: DeserializeOwned,
{
let payload = encode(req)?;
let (line, retries) = with_retry(socket, &payload, Mode::Decode, retry).await?;
let resp = serde_json::from_str(line.trim())
.with_context(|| format!("decode response from {}", socket.display()))?;
Ok((resp, retries))
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
@ -126,7 +175,14 @@ where
Req: Serialize + ?Sized,
{
let payload = encode(req)?;
with_retry(socket, &payload, Mode::Drain, retry).await?;
with_retry(
socket,
&payload,
Mode::Drain,
retry,
DEFAULT_RESPONSE_TIMEOUT,
)
.await?;
Ok(())
}
@ -152,19 +208,53 @@ where
Ok(payload.into_bytes())
}
/// Run [`try_once`] until it succeeds or the retry budget is spent,
/// returning the response line and the number of retries it took.
/// 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).await {
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));
@ -196,11 +286,25 @@ fn exhausted(err: anyhow::Error, backoffs: &[u64]) -> anyhow::Error {
))
}
/// One connect / write / read cycle. Every error it returns is retryable
/// by construction — the deterministic failures (serialise, deserialise)
/// 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) -> Result<String> {
let stream = UnixStream::connect(socket).await.map_err(|e| {
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
@ -222,28 +326,43 @@ async fn try_once(socket: &Path, payload: &[u8], mode: Mode) -> Result<String> {
})?;
let (read, mut write) = stream.into_split();
write
.write_all(payload)
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
.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()))?,
}
.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 = reader
.read_line(&mut line)
.await
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!(
@ -253,9 +372,10 @@ async fn try_once(socket: &Path, payload: &[u8], mode: Mode) -> Result<String> {
}
}
Mode::Drain => {
// Discarded, including any error reading it: the request is
// already delivered and the caller acts on nothing here.
let _ = reader.read_line(&mut line).await;
// 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();
}
}
@ -264,7 +384,94 @@ async fn try_once(socket: &Path, payload: &[u8], mode: Mode) -> Result<String> {
#[cfg(test)]
mod tests {
use super::{Retry, notify, request};
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"