diff --git a/crates/socket-patch-core/src/api/client.rs b/crates/socket-patch-core/src/api/client.rs index 7d862c050..1bba6874d 100644 --- a/crates/socket-patch-core/src/api/client.rs +++ b/crates/socket-patch-core/src/api/client.rs @@ -13,8 +13,8 @@ use serde::Serialize; use crate::api::ranking::severity_order as get_severity_order; use crate::api::ranking::{cmp_batch_infos, cmp_search_results}; use crate::api::retry::{ - is_retryable_status, jitter_sample as retry_jitter, parse_retry_after, ApiRetry, - ApiRetryPolicy, ApiTimeouts, RetryHooks, + is_retryable_status, is_retryable_transport, jitter_sample as retry_jitter, parse_retry_after, + ApiRetry, ApiRetryPolicy, ApiTimeouts, RetryHooks, }; use crate::api::types::*; use crate::api::vendor_prefetch::VendorPrefetch; @@ -487,9 +487,42 @@ impl ApiClient { let max = retry.policy.max_retries; let mut retries = 0u32; loop { - let resp = build().send().await.map_err(|e| { - ApiError::Network(format!("Network error: {}", network_error_detail(&e))) - })?; + let resp = match build().send().await { + Ok(resp) => resp, + Err(e) => { + let network = || { + ApiError::Network(format!("Network error: {}", network_error_detail(&e))) + }; + // A connection dropped before the request went out + // shares the 429 / 503 budget and window; every other + // transport error is final at once. + if !is_retryable_transport(&e) || retries >= max { + return Err(network()); + } + let next = retries + 1; + let Some(delay) = retry.policy.delay( + next, + None, + retry_jitter(retry.hooks.jitter_seed, label, next), + ) else { + return Err(network()); + }; + if !retry.reserve(delay) { + debug_log(&format!( + "{label} failed to connect; not retrying: the run's {} s retry window has closed", + retry.policy.retry_window.as_secs() + )); + return Err(network()); + } + debug_log(&format!( + "{label} failed to connect ({}); retry {next}/{max} in {delay:?}", + network_error_detail(&e) + )); + (retry.hooks.sleep)(delay).await; + retries = next; + continue; + } + }; let status = resp.status(); if !is_retryable_status(status) { return Ok(Sent::Response(resp)); diff --git a/crates/socket-patch-core/src/api/retry.rs b/crates/socket-patch-core/src/api/retry.rs index 643214539..84af87241 100644 --- a/crates/socket-patch-core/src/api/retry.rs +++ b/crates/socket-patch-core/src/api/retry.rs @@ -27,11 +27,19 @@ //! overlap rather than add up), while a throttled run as a whole adds at //! most about the window's length of waiting. //! +//! A connection the peer drops while it is being established (reset, +//! aborted, or closed mid-TLS-handshake — [`is_retryable_transport`]) is +//! retried under the same budget and window, with the exponential backoff: +//! the request was never sent, and a load balancer or proxy resetting a +//! fresh connection is a transient blip, not an answer. +//! //! Nothing else is retried here: other 4xx (401/403 keep driving the proxy -//! fallback), other 5xx and transport errors surface on the first answer, -//! exactly as before. Retries happen inside each request's future, so the -//! callers' ordered folds (`ordered_concurrent`) see the same sequence of -//! results a clean run produces. +//! fallback), other 5xx and every other transport error (refused, timed +//! out, DNS, TLS certificate, or a failure after the request went out) +//! surface on the first answer, exactly as before. Retries happen inside +//! each request's future, so the callers' ordered folds +//! (`ordered_concurrent`) see the same sequence of results a clean run +//! produces. use std::future::Future; use std::pin::Pin; @@ -186,6 +194,47 @@ pub fn is_retryable_status(status: StatusCode) -> bool { ) } +/// Is this transport error one the retry loop may repeat? Only a +/// connection dropped while it was being ESTABLISHED (TCP connect or the +/// TLS handshake): reset, aborted, or closed mid-handshake. No byte of the +/// request has been sent at that point, so repeating it is safe for a POST +/// too. Everything else stays final on the first failure: a refused +/// connection (nothing listening — tests and offline runs rely on it +/// failing fast), a connect or read timeout (a stall is not repeated, see +/// [`ApiTimeouts`]), DNS and certificate errors, and any failure after the +/// request went out. +pub fn is_retryable_transport(error: &reqwest::Error) -> bool { + error.is_connect() && !error.is_timeout() && is_dropped_connection(error) +} + +/// Whether `error`'s cause chain holds an I/O error saying the peer +/// dropped the connection (reset, aborted, or EOF mid-handshake). The +/// connector wraps the OS error in an `Other` I/O error whose `source()` +/// skips it, so a wrapped error is looked through with `get_ref`. +fn is_dropped_connection(error: &(dyn std::error::Error + 'static)) -> bool { + let mut source = Some(error); + while let Some(cause) = source { + if let Some(io) = cause.downcast_ref::() { + if matches!( + io.kind(), + std::io::ErrorKind::ConnectionReset + | std::io::ErrorKind::ConnectionAborted + | std::io::ErrorKind::UnexpectedEof + ) { + return true; + } + if io + .get_ref() + .is_some_and(|inner| is_dropped_connection(inner)) + { + return true; + } + } + source = cause.source(); + } + false +} + /// A `Retry-After` header as a wait from `now_unix_secs`: delta-seconds /// (`Retry-After: 7`) or an HTTP-date (`Retry-After: Fri, 27 Mar 2026 /// 19:12:42 GMT`; a date already past waits zero). `None` when absent or diff --git a/crates/socket-patch-core/tests/api_retry_e2e.rs b/crates/socket-patch-core/tests/api_retry_e2e.rs index 26325f281..c14604cc5 100644 --- a/crates/socket-patch-core/tests/api_retry_e2e.rs +++ b/crates/socket-patch-core/tests/api_retry_e2e.rs @@ -785,3 +785,103 @@ async fn persistent_proxy_batch_throttle_errors_without_per_package_fallback() { // `expect(4)` / `expect(0)` are verified when `server` drops. } } + +/// A TLS endpoint whose peer resets every connection mid-handshake (after +/// reading the ClientHello) — what a load balancer dropping fresh +/// connections looks like (`client error (Connect): Connection reset by +/// peer`). Returns its `https://` base URL and the accepted-connection +/// count. +async fn resetting_endpoint() -> (String, Arc) { + use std::sync::atomic::{AtomicUsize, Ordering}; + use tokio::io::AsyncReadExt as _; + let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); + let addr = listener.local_addr().unwrap(); + let accepted = Arc::new(AtomicUsize::new(0)); + let count = Arc::clone(&accepted); + tokio::spawn(async move { + while let Ok((mut stream, _)) = listener.accept().await { + count.fetch_add(1, Ordering::SeqCst); + let mut hello = [0u8; 5]; + let _ = stream.read_exact(&mut hello).await; + // Linger 0: closing sends RST, not FIN. + let _ = stream.set_zero_linger(); + drop(stream); + } + }); + (format!("https://{addr}"), accepted) +} + +/// A connection reset while it is being established is retried under the +/// 429 / 503 budget (the request never went out, so this holds for the +/// batch POST too), then surfaces as the same `Network error` as before. +#[tokio::test] +async fn a_connection_reset_mid_handshake_is_retried_then_reported() { + use std::sync::atomic::Ordering; + let p = purl(70); + + let (url, accepted) = resetting_endpoint().await; + let (c, log) = client(&url, ApiRetryPolicy::default()); + let err = c + .search_patches_by_package(&p) + .await + .expect_err("every connection is reset"); + assert!(matches!(err, ApiError::Network(_)), "{err:?}"); + assert!(err.to_string().starts_with("Network error: "), "{err}"); + assert_eq!( + accepted.load(Ordering::SeqCst), + 4, + "one attempt + 3 retries" + ); + let w = waits(&log); + assert_eq!(w.len(), 3); + // The exponential steps (500 ms, 1 s, 2 s), each jittered into its + // upper half. + for (wait, step) in w.iter().zip([500u64, 1000, 2000]) { + let step = Duration::from_millis(step); + assert!(*wait >= step / 2 && *wait < step, "{wait:?} for {step:?}"); + } + + let (url, accepted) = resetting_endpoint().await; + let (hooks, log) = virtual_clock(0); + let c = ApiClient::new(options(&url, true)).with_api_retry(ApiRetryPolicy::default(), hooks); + let err = c + .search_patches_batch(std::slice::from_ref(&p)) + .await + .expect_err("every connection is reset"); + assert!(matches!(err, ApiError::Network(_)), "{err:?}"); + assert_eq!( + accepted.load(Ordering::SeqCst), + 4, + "the batch POST retries too" + ); + assert_eq!(waits(&log).len(), 3); + + // Retries off: one connection, no wait, the same error. + let (url, accepted) = resetting_endpoint().await; + let (c, log) = client(&url, ApiRetryPolicy::none()); + let err = c.search_patches_by_package(&p).await.expect_err("reset"); + assert!(matches!(err, ApiError::Network(_)), "{err:?}"); + assert_eq!(accepted.load(Ordering::SeqCst), 1); + assert!(waits(&log).is_empty()); +} + +/// A refused connection (nothing listening) is final at once: offline +/// runs and the tests that point the CLI at a dead port fail fast, as +/// before. +#[tokio::test] +async fn a_refused_connection_is_not_retried() { + let port = { + let l = std::net::TcpListener::bind("127.0.0.1:0").unwrap(); + l.local_addr().unwrap().port() + }; + let (c, log) = client( + &format!("http://127.0.0.1:{port}"), + ApiRetryPolicy::default(), + ); + let err = c + .search_patches_by_package(&purl(71)) + .await + .expect_err("nothing listens there"); + assert!(matches!(err, ApiError::Network(_)), "{err:?}"); + assert!(waits(&log).is_empty()); +}