Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
43 changes: 38 additions & 5 deletions crates/socket-patch-core/src/api/client.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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));
Expand Down
57 changes: 53 additions & 4 deletions crates/socket-patch-core/src/api/retry.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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::<std::io::Error>() {
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
Expand Down
100 changes: 100 additions & 0 deletions crates/socket-patch-core/tests/api_retry_e2e.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<std::sync::atomic::AtomicUsize>) {
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());
}
Loading