Skip to content

Commit efb631f

Browse files
committed
Retry patch API connects reset mid-handshake
The JSON retry loop repeated only 429/503; any transport error was final on the first failure. CI e2e legs talk to the production patch API, and a load balancer resetting a fresh connection ("client error (Connect): Connection reset by peer") failed the whole leg, e.g. the pnpm safety e2e on #563. A connection reset, aborted or closed while being established is now retried under the existing budget, backoff and run-wide window. No byte of the request has gone out at that point, so it is safe for the batch POST too. Refused connections, timeouts, DNS and TLS certificate errors, and failures after the request was sent stay final at once, so dead-port tests and stalled reads behave as before. Co-Authored-By: Claude Opus 5.5 (1M context) <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_01F1QpyLyaVEx88TT3UwW1L1
1 parent 045d7ec commit efb631f

3 files changed

Lines changed: 191 additions & 9 deletions

File tree

‎crates/socket-patch-core/src/api/client.rs‎

Lines changed: 38 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -13,8 +13,8 @@ use serde::Serialize;
1313
use crate::api::ranking::severity_order as get_severity_order;
1414
use crate::api::ranking::{cmp_batch_infos, cmp_search_results};
1515
use crate::api::retry::{
16-
is_retryable_status, jitter_sample as retry_jitter, parse_retry_after, ApiRetry,
17-
ApiRetryPolicy, ApiTimeouts, RetryHooks,
16+
is_retryable_status, is_retryable_transport, jitter_sample as retry_jitter, parse_retry_after,
17+
ApiRetry, ApiRetryPolicy, ApiTimeouts, RetryHooks,
1818
};
1919
use crate::api::types::*;
2020
use crate::api::vendor_prefetch::VendorPrefetch;
@@ -487,9 +487,42 @@ impl ApiClient {
487487
let max = retry.policy.max_retries;
488488
let mut retries = 0u32;
489489
loop {
490-
let resp = build().send().await.map_err(|e| {
491-
ApiError::Network(format!("Network error: {}", network_error_detail(&e)))
492-
})?;
490+
let resp = match build().send().await {
491+
Ok(resp) => resp,
492+
Err(e) => {
493+
let network = || {
494+
ApiError::Network(format!("Network error: {}", network_error_detail(&e)))
495+
};
496+
// A connection dropped before the request went out
497+
// shares the 429 / 503 budget and window; every other
498+
// transport error is final at once.
499+
if !is_retryable_transport(&e) || retries >= max {
500+
return Err(network());
501+
}
502+
let next = retries + 1;
503+
let Some(delay) = retry.policy.delay(
504+
next,
505+
None,
506+
retry_jitter(retry.hooks.jitter_seed, label, next),
507+
) else {
508+
return Err(network());
509+
};
510+
if !retry.reserve(delay) {
511+
debug_log(&format!(
512+
"{label} failed to connect; not retrying: the run's {} s retry window has closed",
513+
retry.policy.retry_window.as_secs()
514+
));
515+
return Err(network());
516+
}
517+
debug_log(&format!(
518+
"{label} failed to connect ({}); retry {next}/{max} in {delay:?}",
519+
network_error_detail(&e)
520+
));
521+
(retry.hooks.sleep)(delay).await;
522+
retries = next;
523+
continue;
524+
}
525+
};
493526
let status = resp.status();
494527
if !is_retryable_status(status) {
495528
return Ok(Sent::Response(resp));

‎crates/socket-patch-core/src/api/retry.rs‎

Lines changed: 53 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -27,11 +27,19 @@
2727
//! overlap rather than add up), while a throttled run as a whole adds at
2828
//! most about the window's length of waiting.
2929
//!
30+
//! A connection the peer drops while it is being established (reset,
31+
//! aborted, or closed mid-TLS-handshake — [`is_retryable_transport`]) is
32+
//! retried under the same budget and window, with the exponential backoff:
33+
//! the request was never sent, and a load balancer or proxy resetting a
34+
//! fresh connection is a transient blip, not an answer.
35+
//!
3036
//! Nothing else is retried here: other 4xx (401/403 keep driving the proxy
31-
//! fallback), other 5xx and transport errors surface on the first answer,
32-
//! exactly as before. Retries happen inside each request's future, so the
33-
//! callers' ordered folds (`ordered_concurrent`) see the same sequence of
34-
//! results a clean run produces.
37+
//! fallback), other 5xx and every other transport error (refused, timed
38+
//! out, DNS, TLS certificate, or a failure after the request went out)
39+
//! surface on the first answer, exactly as before. Retries happen inside
40+
//! each request's future, so the callers' ordered folds
41+
//! (`ordered_concurrent`) see the same sequence of results a clean run
42+
//! produces.
3543
3644
use std::future::Future;
3745
use std::pin::Pin;
@@ -186,6 +194,47 @@ pub fn is_retryable_status(status: StatusCode) -> bool {
186194
)
187195
}
188196

197+
/// Is this transport error one the retry loop may repeat? Only a
198+
/// connection dropped while it was being ESTABLISHED (TCP connect or the
199+
/// TLS handshake): reset, aborted, or closed mid-handshake. No byte of the
200+
/// request has been sent at that point, so repeating it is safe for a POST
201+
/// too. Everything else stays final on the first failure: a refused
202+
/// connection (nothing listening — tests and offline runs rely on it
203+
/// failing fast), a connect or read timeout (a stall is not repeated, see
204+
/// [`ApiTimeouts`]), DNS and certificate errors, and any failure after the
205+
/// request went out.
206+
pub fn is_retryable_transport(error: &reqwest::Error) -> bool {
207+
error.is_connect() && !error.is_timeout() && is_dropped_connection(error)
208+
}
209+
210+
/// Whether `error`'s cause chain holds an I/O error saying the peer
211+
/// dropped the connection (reset, aborted, or EOF mid-handshake). The
212+
/// connector wraps the OS error in an `Other` I/O error whose `source()`
213+
/// skips it, so a wrapped error is looked through with `get_ref`.
214+
fn is_dropped_connection(error: &(dyn std::error::Error + 'static)) -> bool {
215+
let mut source = Some(error);
216+
while let Some(cause) = source {
217+
if let Some(io) = cause.downcast_ref::<std::io::Error>() {
218+
if matches!(
219+
io.kind(),
220+
std::io::ErrorKind::ConnectionReset
221+
| std::io::ErrorKind::ConnectionAborted
222+
| std::io::ErrorKind::UnexpectedEof
223+
) {
224+
return true;
225+
}
226+
if io
227+
.get_ref()
228+
.is_some_and(|inner| is_dropped_connection(inner))
229+
{
230+
return true;
231+
}
232+
}
233+
source = cause.source();
234+
}
235+
false
236+
}
237+
189238
/// A `Retry-After` header as a wait from `now_unix_secs`: delta-seconds
190239
/// (`Retry-After: 7`) or an HTTP-date (`Retry-After: Fri, 27 Mar 2026
191240
/// 19:12:42 GMT`; a date already past waits zero). `None` when absent or

‎crates/socket-patch-core/tests/api_retry_e2e.rs‎

Lines changed: 100 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -785,3 +785,103 @@ async fn persistent_proxy_batch_throttle_errors_without_per_package_fallback() {
785785
// `expect(4)` / `expect(0)` are verified when `server` drops.
786786
}
787787
}
788+
789+
/// A TLS endpoint whose peer resets every connection mid-handshake (after
790+
/// reading the ClientHello) — what a load balancer dropping fresh
791+
/// connections looks like (`client error (Connect): Connection reset by
792+
/// peer`). Returns its `https://` base URL and the accepted-connection
793+
/// count.
794+
async fn resetting_endpoint() -> (String, Arc<std::sync::atomic::AtomicUsize>) {
795+
use std::sync::atomic::{AtomicUsize, Ordering};
796+
use tokio::io::AsyncReadExt as _;
797+
let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
798+
let addr = listener.local_addr().unwrap();
799+
let accepted = Arc::new(AtomicUsize::new(0));
800+
let count = Arc::clone(&accepted);
801+
tokio::spawn(async move {
802+
while let Ok((mut stream, _)) = listener.accept().await {
803+
count.fetch_add(1, Ordering::SeqCst);
804+
let mut hello = [0u8; 5];
805+
let _ = stream.read_exact(&mut hello).await;
806+
// Linger 0: closing sends RST, not FIN.
807+
let _ = stream.set_zero_linger();
808+
drop(stream);
809+
}
810+
});
811+
(format!("https://{addr}"), accepted)
812+
}
813+
814+
/// A connection reset while it is being established is retried under the
815+
/// 429 / 503 budget (the request never went out, so this holds for the
816+
/// batch POST too), then surfaces as the same `Network error` as before.
817+
#[tokio::test]
818+
async fn a_connection_reset_mid_handshake_is_retried_then_reported() {
819+
use std::sync::atomic::Ordering;
820+
let p = purl(70);
821+
822+
let (url, accepted) = resetting_endpoint().await;
823+
let (c, log) = client(&url, ApiRetryPolicy::default());
824+
let err = c
825+
.search_patches_by_package(&p)
826+
.await
827+
.expect_err("every connection is reset");
828+
assert!(matches!(err, ApiError::Network(_)), "{err:?}");
829+
assert!(err.to_string().starts_with("Network error: "), "{err}");
830+
assert_eq!(
831+
accepted.load(Ordering::SeqCst),
832+
4,
833+
"one attempt + 3 retries"
834+
);
835+
let w = waits(&log);
836+
assert_eq!(w.len(), 3);
837+
// The exponential steps (500 ms, 1 s, 2 s), each jittered into its
838+
// upper half.
839+
for (wait, step) in w.iter().zip([500u64, 1000, 2000]) {
840+
let step = Duration::from_millis(step);
841+
assert!(*wait >= step / 2 && *wait < step, "{wait:?} for {step:?}");
842+
}
843+
844+
let (url, accepted) = resetting_endpoint().await;
845+
let (hooks, log) = virtual_clock(0);
846+
let c = ApiClient::new(options(&url, true)).with_api_retry(ApiRetryPolicy::default(), hooks);
847+
let err = c
848+
.search_patches_batch(std::slice::from_ref(&p))
849+
.await
850+
.expect_err("every connection is reset");
851+
assert!(matches!(err, ApiError::Network(_)), "{err:?}");
852+
assert_eq!(
853+
accepted.load(Ordering::SeqCst),
854+
4,
855+
"the batch POST retries too"
856+
);
857+
assert_eq!(waits(&log).len(), 3);
858+
859+
// Retries off: one connection, no wait, the same error.
860+
let (url, accepted) = resetting_endpoint().await;
861+
let (c, log) = client(&url, ApiRetryPolicy::none());
862+
let err = c.search_patches_by_package(&p).await.expect_err("reset");
863+
assert!(matches!(err, ApiError::Network(_)), "{err:?}");
864+
assert_eq!(accepted.load(Ordering::SeqCst), 1);
865+
assert!(waits(&log).is_empty());
866+
}
867+
868+
/// A refused connection (nothing listening) is final at once: offline
869+
/// runs and the tests that point the CLI at a dead port fail fast, as
870+
/// before.
871+
#[tokio::test]
872+
async fn a_refused_connection_is_not_retried() {
873+
let port = {
874+
let l = std::net::TcpListener::bind("127.0.0.1:0").unwrap();
875+
l.local_addr().unwrap().port()
876+
};
877+
let (c, log) = client(
878+
&format!("http://127.0.0.1:{port}"),
879+
ApiRetryPolicy::default(),
880+
);
881+
let err = c
882+
.search_patches_by_package(&purl(71))
883+
.await
884+
.expect_err("nothing listens there");
885+
assert!(matches!(err, ApiError::Network(_)), "{err:?}");
886+
assert!(waits(&log).is_empty());
887+
}

0 commit comments

Comments
 (0)