diff --git a/hive-agent/src/main.rs b/hive-agent/src/main.rs index 064f79cd..6756ff90 100644 --- a/hive-agent/src/main.rs +++ b/hive-agent/src/main.rs @@ -48,6 +48,15 @@ const DEFAULT_SOCKET: &str = "/run/hive/mcp.sock"; /// so a hive-c0re restart is worth waiting out rather than surfacing. const CONTROL_SOCKET_RETRY: Retry = Retry::RideOutRestart; +/// How long the serve loop's `Recv` asks the broker to park when the inbox +/// is empty. +const RECV_WAIT: Duration = Duration::from_mins(3); + +/// Response deadline for that `Recv`: the park itself plus headroom for the +/// broker to answer once it ends, so an idle poll never reads as a stuck +/// hive-c0re. +const RECV_RESPONSE_TIMEOUT: Duration = RECV_WAIT.saturating_add(Duration::from_secs(30)); + /// Default web UI port — used when `HIVE_PORT` env is unset. const DEFAULT_WEB_PORT: u16 = 8042; @@ -389,13 +398,14 @@ impl Surface for AgentSurface { } async fn recv_next(socket: &Path) -> RecvOutcome { - let recv: Result = hive_sock_client::request( + let recv: Result = hive_sock_client::request_within( socket, &Request::Recv { - wait_seconds: Some(180), + wait_seconds: Some(RECV_WAIT.as_secs()), max: None, }, CONTROL_SOCKET_RETRY, + RECV_RESPONSE_TIMEOUT, ) .await; match recv { diff --git a/hive-agent/src/web_ui/proxy.rs b/hive-agent/src/web_ui/proxy.rs index 047cb079..1bdb5e45 100644 --- a/hive-agent/src/web_ui/proxy.rs +++ b/hive-agent/src/web_ui/proxy.rs @@ -40,6 +40,16 @@ const HOP_BY_HOP: [&str; 6] = [ "upgrade", ]; +/// Bound on dialing a `unix:` upstream. The socket is local, so a connect +/// still pending after this long is not going to complete. +const UNIX_CONNECT_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(5); + +/// Bound on a `unix:` upstream answering with its response head — the same +/// budget the `http(s)://` path gives a whole request. Only the head is +/// bounded: once the upstream has answered, a long body is the upstream +/// working, not stuck. +const UNIX_RESPONSE_HEAD_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(30); + /// Nest every proxy declared in `HIVE_EXTRA_WEB_PROXIES` (a JSON object /// `{"": ""}`) under `/extra//` on `app`. Absent / /// blank / invalid JSON is a no-op (logged). Called from [`super::serve`] @@ -213,7 +223,12 @@ async fn forward_over_unix_socket( mut headers: HeaderMap, body: Bytes, ) -> anyhow::Result { - let stream = tokio::net::UnixStream::connect(sock_path).await?; + let stream = tokio::time::timeout( + UNIX_CONNECT_TIMEOUT, + tokio::net::UnixStream::connect(sock_path), + ) + .await + .map_err(|_| anyhow::anyhow!("connect timed out after {UNIX_CONNECT_TIMEOUT:?}"))??; let io = TokioIo::new(stream); let (mut sender, conn) = hyper::client::conn::http1::handshake(io).await?; @@ -241,7 +256,13 @@ async fn forward_over_unix_socket( } let req = req_builder.body(Full::new(body))?; - let upstream_resp = sender.send_request(req).await?; + let upstream_resp = tokio::time::timeout(UNIX_RESPONSE_HEAD_TIMEOUT, sender.send_request(req)) + .await + .map_err(|_| { + anyhow::anyhow!( + "waiting for the response head timed out after {UNIX_RESPONSE_HEAD_TIMEOUT:?}" + ) + })??; let status = StatusCode::from_u16(upstream_resp.status().as_u16())?; let mut resp_headers = HeaderMap::new(); for (name, value) in upstream_resp.headers() { diff --git a/hive-c0re/src/dashboard/matrix_accounts.rs b/hive-c0re/src/dashboard/matrix_accounts.rs index 7a5f3717..eefcc6ad 100644 --- a/hive-c0re/src/dashboard/matrix_accounts.rs +++ b/hive-c0re/src/dashboard/matrix_accounts.rs @@ -369,6 +369,33 @@ pub(super) async fn get_github_account(Query(q): Query) -> R axum::Json(GithubAccountStatus { present }).into_response() } +/// Bound on reaching the homeserver. +const HTTP_CONNECT_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(5); +/// Bound on one whole homeserver round trip, body included. A password +/// login makes the homeserver hash the password before it answers, so this +/// is looser than a plain API call needs. +const HTTP_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(30); + +/// A client with both homeserver bounds applied. +fn http_client() -> Result { + reqwest::Client::builder() + .connect_timeout(HTTP_CONNECT_TIMEOUT) + .timeout(HTTP_TIMEOUT) + .build() + .map_err(|e| format!("build HTTP client: {e}")) +} + +/// `what` failed with `e`; a timeout names the bound that fired. +fn http_error(what: &str, e: &reqwest::Error) -> String { + if e.is_connect() && e.is_timeout() { + format!("{what}: connect timed out after {HTTP_CONNECT_TIMEOUT:?}") + } else if e.is_timeout() { + format!("{what}: timed out after {HTTP_TIMEOUT:?}") + } else { + format!("{what}: {e}") + } +} + /// POST `m.login.password` to `/_matrix/client/v3/login`. /// Returns `(access_token, user_id)`. async fn matrix_password_login( @@ -383,17 +410,17 @@ async fn matrix_password_login( "password": password, "initial_device_display_name": "hyperhive", }); - let resp = reqwest::Client::new() + let resp = http_client()? .post(&url) .json(&body) .send() .await - .map_err(|e| format!("POST /login: {e}"))?; + .map_err(|e| http_error("POST /login", &e))?; let status = resp.status(); let json: serde_json::Value = resp .json() .await - .map_err(|e| format!("parse /login response: {e}"))?; + .map_err(|e| http_error("parse /login response", &e))?; if !status.is_success() { let err = json .get("error") @@ -416,17 +443,17 @@ async fn matrix_password_login( /// to validate it and recover the `user_id`. async fn matrix_whoami(homeserver: &str, token: &str) -> Result { let url = format!("{homeserver}/_matrix/client/v3/account/whoami"); - let resp = reqwest::Client::new() + let resp = http_client()? .get(&url) .bearer_auth(token) .send() .await - .map_err(|e| format!("GET /whoami: {e}"))?; + .map_err(|e| http_error("GET /whoami", &e))?; let status = resp.status(); let json: serde_json::Value = resp .json() .await - .map_err(|e| format!("parse /whoami response: {e}"))?; + .map_err(|e| http_error("parse /whoami response", &e))?; if !status.is_success() { let err = json .get("error") diff --git a/hive-c0re/src/forge/ci_runner.rs b/hive-c0re/src/forge/ci_runner.rs index 9c7631cb..f7eef7fb 100644 --- a/hive-c0re/src/forge/ci_runner.rs +++ b/hive-c0re/src/forge/ci_runner.rs @@ -126,13 +126,46 @@ fn runner_address_matches(raw: &str, configured: &str) -> bool { } } +/// Bound on reaching the forge. +const HTTP_CONNECT_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(5); +/// Bound on one whole forge admin call, body included; the same budget +/// `config_pr_poll` gives this forge. +const HTTP_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(15); + +/// A client with both forge bounds applied. +fn http_client() -> Result { + reqwest::Client::builder() + .connect_timeout(HTTP_CONNECT_TIMEOUT) + .timeout(HTTP_TIMEOUT) + .build() + .context("build HTTP client") +} + +/// Context for `what` failing with `e`; a timeout names the bound that fired. +fn http_context(what: &str, e: &reqwest::Error) -> String { + if e.is_connect() && e.is_timeout() { + format!("{what}: connect timed out after {HTTP_CONNECT_TIMEOUT:?}") + } else if e.is_timeout() { + format!("{what}: timed out after {HTTP_TIMEOUT:?}") + } else { + what.to_owned() + } +} + /// `GET /admin/runners/{id}` — `true` iff the runner still exists on the forge /// (HTTP 200). A 404 (deleted from the admin panel) or any other status means /// re-registration is needed. A transport error (forge unreachable) is treated /// as "keep the existing creds" so a network blip never wipes a valid runner. async fn runner_valid(core_token: &str, id: u64) -> bool { let url = format!("{}/api/v1/admin/runners/{id}", forge_http_base()); - match reqwest::Client::new() + let client = match http_client() { + Ok(c) => c, + Err(e) => { + tracing::warn!(error = ?e, "ci runner: build HTTP client failed; keeping existing creds"); + return true; + } + }; + match client .get(&url) .header("Authorization", format!("token {core_token}")) .send() @@ -140,7 +173,8 @@ async fn runner_valid(core_token: &str, id: u64) -> bool { { Ok(resp) => resp.status().is_success(), Err(e) => { - tracing::warn!(error = ?e, "ci runner: validation request failed; keeping existing creds"); + let what = http_context("validation request failed", &e); + tracing::warn!(error = ?e, "ci runner: {what}; keeping existing creds"); true } } @@ -155,17 +189,20 @@ async fn fetch_registration_token(core_token: &str) -> Result { "{}/api/v1/admin/runners/registration-token", forge_http_base() ); - let resp = reqwest::Client::new() + let resp = http_client()? .get(&url) .header("Authorization", format!("token {core_token}")) .send() .await - .context("GET admin/runners/registration-token")?; + .map_err(|e| { + let what = http_context("GET admin/runners/registration-token", &e); + anyhow::Error::new(e).context(what) + })?; let status = resp.status(); - let json: Value = resp - .json() - .await - .context("parse registration-token response")?; + let json: Value = resp.json().await.map_err(|e| { + let what = http_context("parse registration-token response", &e); + anyhow::Error::new(e).context(what) + })?; if !status.is_success() { anyhow::bail!("registration-token HTTP {status}: {json}"); } diff --git a/hive-sock-client/src/lib.rs b/hive-sock-client/src/lib.rs index 3ae62432..ade26ef9 100644 --- a/hive-sock-client/src/lib.rs +++ b/hive-sock-client/src/lib.rs @@ -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(socket: &Path, req: &Req, retry: Retry) -> Result 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( + socket: &Path, + req: &Req, + retry: Retry, + response_timeout: Duration, +) -> Result +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(socket: &Path, line: &str) -> Result { + 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::() => 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 { - let stream = UnixStream::connect(socket).await.map_err(|e| { +async fn try_once( + socket: &Path, + payload: &[u8], + mode: Mode, + response_timeout: Duration, +) -> Result { + 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 { })?; 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 { } } 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 { #[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" diff --git a/swarm-controller/src/auth.rs b/swarm-controller/src/auth.rs index 931f9719..db34118a 100644 --- a/swarm-controller/src/auth.rs +++ b/swarm-controller/src/auth.rs @@ -23,6 +23,13 @@ use anyhow::{Context, Result}; use swarm_authelia_bridge_sock::{BridgeRequest, BridgeResponse}; +/// Bound on reaching the bridge. +const HTTP_CONNECT_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(5); +/// Bound on one whole bridge call, body included. Ensuring an identity has +/// the bridge shell out to `authelia` for an argon2 hash before it answers, +/// so this is looser than a plain API call needs. +const HTTP_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(30); + /// A configured connection to `swarm-authelia-bridge`. #[derive(Clone)] pub struct AuthBridge { @@ -53,8 +60,13 @@ impl AuthBridge { that identity must exist first" ) })?; + let http = reqwest::Client::builder() + .connect_timeout(HTTP_CONNECT_TIMEOUT) + .timeout(HTTP_TIMEOUT) + .build() + .context("building the swarm-authelia-bridge HTTP client")?; Ok(Some(Self { - http: reqwest::Client::new(), + http, base_url, queue_cfg, })) @@ -109,7 +121,18 @@ impl AuthBridge { .json(req) .send() .await - .context("calling swarm-authelia-bridge")?; + .map_err(|e| { + let what = if e.is_connect() && e.is_timeout() { + format!( + "calling swarm-authelia-bridge: connect timed out after {HTTP_CONNECT_TIMEOUT:?}" + ) + } else if e.is_timeout() { + format!("calling swarm-authelia-bridge: timed out after {HTTP_TIMEOUT:?}") + } else { + "calling swarm-authelia-bridge".to_owned() + }; + anyhow::Error::new(e).context(what) + })?; let status = response.status(); let body = response.text().await.unwrap_or_default(); diff --git a/swarm-controller/src/matrix_account.rs b/swarm-controller/src/matrix_account.rs index 59ebb0da..1b9c68c8 100644 --- a/swarm-controller/src/matrix_account.rs +++ b/swarm-controller/src/matrix_account.rs @@ -346,6 +346,33 @@ async fn resolve_credential( } } +/// Bound on reaching the homeserver. +const HTTP_CONNECT_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(5); +/// Bound on one whole homeserver round trip, body included. A password +/// login makes the homeserver hash the password before it answers, so this +/// is looser than a plain API call needs. +const HTTP_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(30); + +/// A client with both homeserver bounds applied. +fn http_client() -> Result { + reqwest::Client::builder() + .connect_timeout(HTTP_CONNECT_TIMEOUT) + .timeout(HTTP_TIMEOUT) + .build() + .map_err(|e| format!("build HTTP client: {e}")) +} + +/// `what` failed with `e`; a timeout names the bound that fired. +fn http_error(what: &str, e: &reqwest::Error) -> String { + if e.is_connect() && e.is_timeout() { + format!("{what}: connect timed out after {HTTP_CONNECT_TIMEOUT:?}") + } else if e.is_timeout() { + format!("{what}: timed out after {HTTP_TIMEOUT:?}") + } else { + format!("{what}: {e}") + } +} + /// POST `m.login.password` to `/_matrix/client/v3/login`. /// Returns `(access_token, user_id)`. /// @@ -366,17 +393,17 @@ async fn matrix_password_login( "password": password, "initial_device_display_name": "hyperhive", }); - let resp = reqwest::Client::new() + let resp = http_client()? .post(&url) .json(&body) .send() .await - .map_err(|e| format!("POST /login: {e}"))?; + .map_err(|e| http_error("POST /login", &e))?; let status = resp.status(); let json: serde_json::Value = resp .json() .await - .map_err(|e| format!("parse /login response: {e}"))?; + .map_err(|e| http_error("parse /login response", &e))?; if !status.is_success() { let err = json .get("error")