diff --git a/crates/socket-patch-core/src/api/client.rs b/crates/socket-patch-core/src/api/client.rs index 642213af4..89ef48073 100644 --- a/crates/socket-patch-core/src/api/client.rs +++ b/crates/socket-patch-core/src/api/client.rs @@ -282,8 +282,8 @@ const PROXY_BATCH_PATH_CONCURRENCY: usize = 10; /// Retry policy for the vendoring service's package-reference POST and /// archive GET: `attempts` tries in total, exponential delays from `base` -/// with ±25% jitter, each capped at `max_delay` (a `Retry-After` in seconds -/// is honored under the same cap). Retried: transport errors (per-attempt +/// with ±25% jitter, each capped at `max_delay` (a `Retry-After`, in +/// seconds or as an HTTP-date, is honored under the same cap). Retried: transport errors (per-attempt /// timeouts and bodies cut off mid-transfer included) and HTTP 429, /// 500, 502, 503, 504. Never retried: auth (401/403), terminal misses /// (404/410), still-building (408), other 4xx, parse errors. @@ -343,32 +343,15 @@ impl VendorRetryPolicy { /// `service` fails closed — the existing miss policy). pub(crate) const VENDOR_BREAKER_THRESHOLD: u32 = 2; -/// A jitter sample in `[0, 1)` from std's randomly keyed hasher (no RNG -/// dependency; the quality needed here is "not synchronized"). -fn jitter_sample() -> f64 { - use std::hash::{BuildHasher as _, Hasher as _}; - let mut hasher = std::collections::hash_map::RandomState::new().build_hasher(); - hasher.write_u64(0); - (hasher.finish() >> 11) as f64 / (1u64 << 53) as f64 -} +/// The jitter key of the package-reference POST's retries (archive and +/// artifact downloads are keyed by their URL). +const VENDOR_REFERENCES_KEY: &str = "POST vendor package references"; /// Is this vendor-service HTTP status worth another attempt? fn vendor_status_retryable(status: StatusCode) -> bool { matches!(status.as_u16(), 429 | 500 | 502 | 503 | 504) } -/// A `Retry-After: ` header (the HTTP-date form is ignored). -fn retry_after_secs(headers: &HeaderMap) -> Option { - headers - .get(header::RETRY_AFTER)? - .to_str() - .ok()? - .trim() - .parse::() - .ok() - .map(Duration::from_secs) -} - /// One vendor-service attempt's failure: the error, and — when the failure /// is retryable — the server's `Retry-After` hint (`Some(None)` = retryable /// without a hint). @@ -1488,11 +1471,27 @@ impl ApiClient { } } - /// Pause before retry number `retry` (see [`VendorRetryPolicy`]). - async fn vendor_backoff(&self, retry: u32, retry_after: Option) { - let delay = self.vendor_retry.delay(retry, retry_after, jitter_sample()); + /// The retry hint for a vendor-service answer with this `status`: + /// `None` when the status isn't retried, else its `Retry-After` (either + /// form, read on the client's [`RetryHooks`] clock), if any. + fn vendor_retry_hint( + &self, + status: StatusCode, + headers: &HeaderMap, + ) -> Option> { + let hooks = &self.api_retry.hooks; + vendor_status_retryable(status).then(|| parse_retry_after(headers, (hooks.now_unix_secs)())) + } + + /// Pause before retry number `retry` of the request `key` (see + /// [`VendorRetryPolicy`]), on the client's [`RetryHooks`] jitter seed + /// and sleep. + async fn vendor_backoff(&self, key: &str, retry: u32, retry_after: Option) { + let hooks = &self.api_retry.hooks; + let jitter = retry_jitter(hooks.jitter_seed, key, retry); + let delay = self.vendor_retry.delay(retry, retry_after, jitter); debug_log(&format!("vendor service retry {retry} in {delay:?}")); - tokio::time::sleep(delay).await; + (hooks.sleep)(delay).await; } /// Step 1 of [`Self::fetch_vendor_package`], retried per the client's @@ -1549,7 +1548,8 @@ impl ApiClient { debug_log(&format!( "vendor package request attempt {attempt} failed: {e}" )); - self.vendor_backoff(attempt, retry_after).await; + self.vendor_backoff(VENDOR_REFERENCES_KEY, attempt, retry_after) + .await; attempt += 1; } Err((e, hint)) => return Err((e, hint.is_some())), @@ -1614,7 +1614,7 @@ impl ApiClient { } // 429 classifies as RateLimited but is still retried (the hint); // 401/403 carry no hint. - let hint = vendor_status_retryable(status).then(|| retry_after_secs(resp.headers())); + let hint = self.vendor_retry_hint(status, resp.headers()); if let Some(err) = classify_auth_error(status, !use_auth) { return Err((err, hint)); } @@ -1649,7 +1649,8 @@ impl ApiClient { loop { match self.download_vendor_archive_once(url).await { (ServeDownload::Failed(e), Some(retry_after)) if attempt < attempts => { - self.vendor_download_retry(attempt, &e, retry_after).await; + self.vendor_download_retry(url, attempt, &e, retry_after) + .await; attempt += 1; } (outcome, hint) => return (outcome, hint.is_some()), @@ -1661,6 +1662,7 @@ impl ApiClient { /// next one. async fn vendor_download_retry( &self, + url: &str, attempt: u32, e: &ApiError, retry_after: Option, @@ -1668,7 +1670,7 @@ impl ApiClient { debug_log(&format!( "vendor package download attempt {attempt} failed: {e}" )); - self.vendor_backoff(attempt, retry_after).await; + self.vendor_backoff(url, attempt, retry_after).await; } /// [`Self::download_vendor_archive_retrying`] without the flag. @@ -1733,8 +1735,7 @@ impl ApiClient { // caller's pending policy, never retried here. StatusCode::REQUEST_TIMEOUT => return (ServeDownload::Pending, None), _ => { - let hint = - vendor_status_retryable(status).then(|| retry_after_secs(resp.headers())); + let hint = self.vendor_retry_hint(status, resp.headers()); if let Some(err) = classify_auth_error(status, true) { return (ServeDownload::Failed(err), hint); } @@ -1812,7 +1813,7 @@ impl ApiClient { url: &str, deferred: DeferredAttempt, ) -> Result, ApiError> { - self.vendor_download_retry(1, &deferred.error, deferred.retry_after) + self.vendor_download_retry(url, 1, &deferred.error, deferred.retry_after) .await; artifact_download_result(self.download_vendor_archive_from(url, 2).await.0, url) } @@ -2808,7 +2809,7 @@ impl ApiClient { loop { match self.download_artifact_capped_once(url, max_bytes).await { (Err(_), Some(retry_after)) if attempt < attempts => { - self.vendor_backoff(attempt, retry_after).await; + self.vendor_backoff(url, attempt, retry_after).await; attempt += 1; } (outcome, _) => return outcome, @@ -2863,8 +2864,7 @@ impl ApiClient { ) } _ => { - let hint = - vendor_status_retryable(status).then(|| retry_after_secs(resp.headers())); + let hint = self.vendor_retry_hint(status, resp.headers()); let err = classify_auth_error(status, true).unwrap_or_else(|| { ApiError::Other(format!( "artifact download failed with status {}", @@ -5770,8 +5770,6 @@ mod vendor_retry_tests { p.delay(1, Some(Duration::from_secs(60)), 0.5), Duration::from_secs(4) ); - let j = jitter_sample(); - assert!((0.0..1.0).contains(&j)); assert_eq!(VendorRetryPolicy::none().attempts, 1); } @@ -5847,6 +5845,179 @@ mod vendor_retry_tests { assert_eq!(c.vendor_outage.load(Ordering::Relaxed), 0, "{status}"); } } + + /// The server's HTTP-date `Retry-After` and the clock reading 2 s + /// before it. + const RETRY_DATE: &str = "Fri, 27 Mar 2026 19:12:42 GMT"; + + /// A client on `policy` whose retry hooks read "now" as 2 s before + /// [`RETRY_DATE`] and record each backoff instead of sleeping. + fn hooked( + uri: &str, + policy: VendorRetryPolicy, + seed: u64, + ) -> (ApiClient, Arc>>) { + let now = crate::api::date::parse_timestamp_secs(RETRY_DATE).unwrap() - 2; + let slept = Arc::new(std::sync::Mutex::new(Vec::new())); + let record = Arc::clone(&slept); + let hooks = RetryHooks { + sleep: Arc::new(move |d| { + record.lock().unwrap().push(d); + Box::pin(async {}) + }), + now_unix_secs: Arc::new(move || now), + jitter_seed: seed, + ..RetryHooks::default() + }; + let c = client(uri, policy).with_api_retry(ApiRetryPolicy::default(), hooks); + (c, slept) + } + + /// A policy whose cap (10 s) leaves a 2 s `Retry-After` uncapped. + fn roomy() -> VendorRetryPolicy { + VendorRetryPolicy { + max_delay: Duration::from_secs(10), + ..fast() + } + } + + /// Answer `route` once with `status` and an HTTP-date `Retry-After`, + /// then with `then`. + async fn mount_once_then( + server: &MockServer, + verb: &str, + route: &str, + status: u16, + then: ResponseTemplate, + ) { + Mock::given(method(verb)) + .and(path(route)) + .respond_with(ResponseTemplate::new(status).insert_header("retry-after", RETRY_DATE)) + .up_to_n_times(1) + .with_priority(1) + .mount(server) + .await; + Mock::given(method(verb)) + .and(path(route)) + .respond_with(then) + .with_priority(2) + .mount(server) + .await; + } + + /// Every vendor-service retry loop (the reference POST, the archive + /// GET, the capped artifact GET, and a deferred artifact GET resumed) + /// waits out an HTTP-date `Retry-After` through `api::retry`: 2 s, not + /// the policy's millisecond backoff. Before the shared parser, each + /// read delta-seconds only and backed off 1-5 ms here. + #[tokio::test] + async fn every_vendor_retry_loop_honors_an_http_date_retry_after() { + let server = MockServer::start().await; + mount_once_then(&server, "POST", POST_PATH, 429, granted(&server, UUID_A)).await; + let bytes = || ResponseTemplate::new(200).set_body_bytes(BYTES.to_vec()); + for route in ["/archive", "/capped", "/deferred"] { + mount_once_then(&server, "GET", route, 503, bytes()).await; + } + let url = |route: &str| format!("{}{route}", server.uri()); + + let (c, slept) = hooked(&server.uri(), roomy(), 7); + c.request_vendor_references(&[UUID_A.to_string()], false, None) + .await + .expect("the retry after the 429 succeeds"); + assert_eq!( + *slept.lock().unwrap(), + [Duration::from_secs(2)], + "reference POST" + ); + + let (c, slept) = hooked(&server.uri(), roomy(), 7); + let (outcome, _) = c.download_vendor_archive_retrying(&url("/archive")).await; + assert!(matches!(outcome, ServeDownload::Ok(ref b) if b == BYTES)); + assert_eq!( + *slept.lock().unwrap(), + [Duration::from_secs(2)], + "archive GET" + ); + + let (c, slept) = hooked(&server.uri(), roomy(), 7); + let got = c.download_artifact_capped(&url("/capped"), 1 << 20).await; + assert_eq!(got.expect("the retry succeeds"), BYTES); + assert_eq!( + *slept.lock().unwrap(), + [Duration::from_secs(2)], + "capped artifact GET" + ); + + let (c, slept) = hooked(&server.uri(), roomy(), 7); + let deferred = c + .download_artifact_first_attempt(&url("/deferred")) + .await + .expect_err("a 503 defers"); + assert!(slept.lock().unwrap().is_empty(), "a deferral doesn't wait"); + let got = c + .download_artifact_resuming(&url("/deferred"), deferred) + .await; + assert_eq!(got.expect("the resumed retry succeeds"), BYTES); + assert_eq!( + *slept.lock().unwrap(), + [Duration::from_secs(2)], + "resumed artifact GET" + ); + } + + /// An HTTP-date `Retry-After` is still capped at the policy's + /// `max_delay`, as delta-seconds always were. + #[tokio::test] + async fn an_http_date_retry_after_is_capped_at_max_delay() { + let server = MockServer::start().await; + mount_once_then(&server, "POST", POST_PATH, 503, granted(&server, UUID_A)).await; + let (c, slept) = hooked(&server.uri(), fast(), 7); + c.request_vendor_references(&[UUID_A.to_string()], false, None) + .await + .expect("the retry succeeds"); + assert_eq!(*slept.lock().unwrap(), [fast().max_delay]); + } + + /// Without a `Retry-After`, the backoff's ±25% jitter comes from the + /// hooks' seed: one seed replays the same waits, and every wait stays + /// inside the policy's spread. + #[tokio::test] + async fn vendor_backoff_jitter_replays_from_the_hook_seed() { + let server = MockServer::start().await; + Mock::given(method("GET")) + .and(path("/down")) + .respond_with(ResponseTemplate::new(503)) + .mount(&server) + .await; + let policy = VendorRetryPolicy { + attempts: 4, + base: Duration::from_millis(400), + max_delay: Duration::from_secs(10), + ..VendorRetryPolicy::default() + }; + let url = format!("{}/down", server.uri()); + let waits = |seed: u64| { + let (uri, url) = (server.uri(), url.clone()); + async move { + let (c, slept) = hooked(&uri, policy, seed); + c.download_artifact(&url).await.expect_err("always 503"); + let waits = slept.lock().unwrap().clone(); + waits + } + }; + let a = waits(1).await; + assert_eq!(a, waits(1).await, "one seed, one sequence"); + assert_ne!(a, waits(2).await, "another seed, other waits"); + assert_eq!(a.len(), 3); + for (i, w) in a.iter().enumerate() { + let nominal = policy.base * (1u32 << i); + assert!( + *w >= nominal.mul_f64(0.75) && *w < nominal.mul_f64(1.25), + "retry {}: {w:?}", + i + 1 + ); + } + } } #[cfg(test)] diff --git a/crates/socket-patch-core/src/patch/jvm_jar.rs b/crates/socket-patch-core/src/patch/jvm_jar.rs index 82d679406..d4f940a22 100644 --- a/crates/socket-patch-core/src/patch/jvm_jar.rs +++ b/crates/socket-patch-core/src/patch/jvm_jar.rs @@ -25,8 +25,6 @@ use std::collections::HashMap; use std::path::{Path, PathBuf}; -use sha1::Digest as _; - use crate::crawlers::gradle_cache; use crate::hash::git_sha256::compute_git_sha256_from_bytes; use crate::manifest::schema::PatchFileInfo; @@ -36,6 +34,7 @@ use crate::patch::apply::{ }; use crate::patch::rollback::{RollbackResult, VerifyRollbackResult, VerifyRollbackStatus}; use crate::patch::sidecars::{self, maven as maven_sidecars}; +use crate::utils::digest::{sha1_hex_of, sha256_hex_of}; use crate::utils::purl::{parse_maven_purl, purl_qualifier}; use crate::vendor::VendorServiceConfig; @@ -270,7 +269,7 @@ pub fn derived_copies_in( // reader fails fast on a FIFO or device instead of wedging in open(2). let current = crate::utils::fs::read_regular_to_bytes_sync(&jar) .ok() - .map(|b| sha1_hex(&b)); + .map(|b| sha1_hex_of(&b)); let written = std::fs::metadata(&jar).and_then(|m| m.modified()).ok(); let unverified = found .unknown @@ -279,7 +278,7 @@ pub fn derived_copies_in( let Ok(copy) = crate::utils::fs::read_regular_to_bytes_sync(p) else { return true; }; - if Some(sha1_hex(©)) == current { + if Some(sha1_hex_of(©)) == current { return false; } let made = std::fs::metadata(p).and_then(|m| m.modified()).ok(); @@ -352,20 +351,11 @@ fn unpatched_members( Ok(members) } -fn sha256_hex(bytes: &[u8]) -> String { - use sha2::Digest as _; - hex::encode(sha2::Sha256::digest(bytes)) -} - -fn sha1_hex(bytes: &[u8]) -> String { - hex::encode(sha1::Sha1::digest(bytes)) -} - /// `/jvm-originals/.jar`. pub fn backup_path(socket_dir: &Path, original: &[u8]) -> PathBuf { socket_dir .join(ORIGINALS_DIR) - .join(format!("{}.jar", sha256_hex(original))) + .join(format!("{}.jar", sha256_hex_of(original))) } /// Keep `original` under [`ORIGINALS_DIR`] (content-addressed: an existing @@ -625,7 +615,7 @@ async fn find_backup(restore: &JarRestore<'_>, dir: &Path, current: &[u8]) -> Op continue; }; if let Some(hash) = &gradle_hash { - if !gradle_cache::hash_eq(hash, &sha1_hex(&bytes)) { + if !gradle_cache::hash_eq(hash, &sha1_hex_of(&bytes)) { continue; } } @@ -663,7 +653,7 @@ async fn upstream_for_gradle_copy(purl: &str, jar_leaf: &str, dir: &Path) -> Opt ) .await .ok()?; - gradle_cache::hash_eq(hash, &sha1_hex(&bytes)).then_some(bytes) + gradle_cache::hash_eq(hash, &sha1_hex_of(&bytes)).then_some(bytes) } fn rollback_result(purl: &str, dir: &Path) -> RollbackResult { @@ -1066,7 +1056,7 @@ mod tests { std::fs::create_dir_all(©).unwrap(); let original = pristine_jar(); std::fs::write(copy.join("lib-1.0.jar"), &original).unwrap(); - let sha1_text = format!("{}\n", sha1_hex(&original)); + let sha1_text = format!("{}\n", sha1_hex_of(&original)); std::fs::write(copy.join("lib-1.0.jar.sha1"), &sha1_text).unwrap(); let files = record(); @@ -1082,7 +1072,7 @@ mod tests { ); assert_eq!( std::fs::read_to_string(copy.join("lib-1.0.jar.sha1")).unwrap(), - format!("{}\n", sha1_hex(&service)) + format!("{}\n", sha1_hex_of(&service)) ); let restore = JarRestore { @@ -1123,7 +1113,7 @@ mod tests { let version = d .path() .join(".gradle/caches/modules-2/files-2.1/com.example/lib/1.0"); - let hash_dir = version.join(sha1_hex(&original)); + let hash_dir = version.join(sha1_hex_of(&original)); std::fs::create_dir_all(&hash_dir).unwrap(); std::fs::write(hash_dir.join("lib-1.0.jar"), patched_jar()).unwrap(); diff --git a/crates/socket-patch-core/src/patch/sidecars/maven.rs b/crates/socket-patch-core/src/patch/sidecars/maven.rs index f2f5a2466..8798bfce6 100644 --- a/crates/socket-patch-core/src/patch/sidecars/maven.rs +++ b/crates/socket-patch-core/src/patch/sidecars/maven.rs @@ -17,8 +17,6 @@ use std::path::{Path, PathBuf}; -use sha1::Digest as _; - use super::{ SidecarAdvisory, SidecarAdvisoryCode, SidecarError, SidecarFile, SidecarFileAction, SidecarPayload, SidecarSeverity, @@ -44,7 +42,7 @@ impl Algo { fn digest(self, bytes: &[u8]) -> String { match self { - Algo::Sha1 => hex::encode(sha1::Sha1::digest(bytes)), + Algo::Sha1 => crate::utils::digest::sha1_hex_of(bytes), Algo::Md5 => hex::encode(md5(bytes)), } } diff --git a/crates/socket-patch-core/src/utils/digest.rs b/crates/socket-patch-core/src/utils/digest.rs index 630adefa1..2b13bd787 100644 --- a/crates/socket-patch-core/src/utils/digest.rs +++ b/crates/socket-patch-core/src/utils/digest.rs @@ -135,6 +135,7 @@ mod tests { /// when you move it onto the helpers above; the test fails on a stale /// entry as well as on a new inline copy. const PENDING_INLINE_DIGESTS: &[&str] = &[ + "crawlers/gradle_cache.rs", "utils/group_commit.rs", "vendor/jvm/mod.rs", "vendor/maven_repo.rs",