Skip to main content

apify_client/
http_client.rs

1//! The HTTP layer of the client.
2//!
3//! The [`HttpBackend`] trait defines the minimal contract for sending a single HTTP
4//! request and receiving a response. It is the *replaceable component* of the client:
5//! the default implementation [`ReqwestBackend`] uses [`reqwest`], but a custom
6//! backend (e.g. for testing or for a different runtime) can be plugged in via
7//! [`ApifyClientBuilder::http_backend`](crate::ApifyClientBuilder::http_backend).
8//!
9//! [`HttpClient`] wraps a backend and adds the cross-cutting concerns shared by every
10//! endpoint: authentication, the `User-Agent` header, query-parameter serialization,
11//! timeouts and retries with exponential backoff (mirroring the JavaScript and Python
12//! reference clients).
13
14use std::collections::HashMap;
15use std::sync::Arc;
16use std::time::Duration;
17
18use async_trait::async_trait;
19
20use crate::error::{ApiError, ApiErrorBody, ApifyClientError, ApifyClientResult};
21
22/// HTTP status code returned by the API when the per-resource rate limit is exceeded.
23const RATE_LIMIT_EXCEEDED_STATUS_CODE: u16 = 429;
24/// Statuses `>= 500` are considered retryable internal server errors.
25const MIN_SERVER_ERROR_STATUS_CODE: u16 = 500;
26/// Responses with status `< 300` are treated as success.
27const MAX_SUCCESS_STATUS_CODE: u16 = 300;
28/// Multiplier applied to the inter-retry delay after each attempt (exponential backoff).
29/// Matches the reference client's `async-retry` default factor of 2.
30const BACKOFF_FACTOR: u32 = 2;
31
32/// Request bodies at least this many bytes are compressed before sending. Smaller bodies are
33/// left uncompressed because the CPU cost outweighs the transfer savings. Matches the reference
34/// client's `MIN_COMPRESS_BYTES` threshold.
35const MIN_COMPRESS_BYTES: usize = 1024;
36/// `Content-Encoding` value used for brotli-compressed request bodies.
37const CONTENT_ENCODING_BROTLI: &str = "br";
38/// `Content-Encoding` value used for gzip-compressed request bodies.
39const CONTENT_ENCODING_GZIP: &str = "gzip";
40/// Brotli quality level (0–11). The reference client compresses request bodies at quality 6,
41/// which balances ratio against CPU cost; we mirror that.
42const BROTLI_QUALITY: u32 = 6;
43/// Brotli sliding-window size (log2), 22 is the library default (a 4 MiB window).
44const BROTLI_WINDOW_SIZE: u32 = 22;
45/// Internal buffer size for the brotli encoder.
46const BROTLI_BUFFER_SIZE: usize = 4096;
47/// Gzip compression level (0–9). Matches the reference client, which gzips using `node:zlib`'s
48/// default level (6).
49const GZIP_COMPRESSION_LEVEL: u32 = 6;
50
51/// Media-type prefixes whose payloads already carry their own compression, so running them
52/// through brotli/gzip burns CPU and memory for a result that is usually no smaller — and the
53/// request also keeps the intended `Content-Type`, rather than becoming an unreadable
54/// `Content-Encoding: br` blob the destination may not expect for these formats.
55/// Matches the reference client's `ALREADY_COMPRESSED_MEDIA_TYPE_PREFIXES`.
56const ALREADY_COMPRESSED_MEDIA_TYPE_PREFIXES: [&str; 3] = ["audio/", "image/", "video/"];
57
58/// Exact media types (outside the prefixes above) whose payloads already carry their own
59/// compression. Matches the reference client's `ALREADY_COMPRESSED_MEDIA_TYPES`.
60const ALREADY_COMPRESSED_MEDIA_TYPES: [&str; 19] = [
61    "application/epub+zip",
62    "application/gzip",
63    "application/java-archive",
64    "application/vnd.android.package-archive",
65    "application/vnd.openxmlformats-officedocument.presentationml.presentation",
66    "application/vnd.openxmlformats-officedocument.spreadsheetml.sheet",
67    "application/vnd.openxmlformats-officedocument.wordprocessingml.document",
68    "application/vnd.rar",
69    "application/x-7z-compressed",
70    "application/x-bzip",
71    "application/x-bzip2",
72    "application/x-gzip",
73    "application/x-rar-compressed",
74    "application/x-xz",
75    "application/x-zip-compressed",
76    "application/zip",
77    "application/zstd",
78    "font/woff",
79    "font/woff2",
80];
81
82/// Uncompressed media types that sit under an already-compressed prefix above, so compressing
83/// them still pays off. Matches the reference client's `COMPRESSIBLE_MEDIA_TYPES`.
84const COMPRESSIBLE_MEDIA_TYPES: [&str; 16] = [
85    "audio/aiff",
86    "audio/basic",
87    "audio/l16",
88    "audio/l24",
89    "audio/midi",
90    "audio/vnd.wave",
91    "audio/wav",
92    "audio/wave",
93    "audio/x-aiff",
94    "audio/x-wav",
95    "image/bmp",
96    "image/tiff",
97    "image/vnd.adobe.photoshop",
98    "image/vnd.microsoft.icon",
99    "image/x-icon",
100    "image/x-ms-bmp",
101];
102
103/// Structured-syntax suffixes that mark a media type as text even under an already-compressed
104/// prefix (e.g. `image/svg+xml`). Matches the reference client's
105/// `COMPRESSIBLE_MEDIA_TYPE_SUFFIXES`.
106const COMPRESSIBLE_MEDIA_TYPE_SUFFIXES: [&str; 2] = ["+json", "+xml"];
107
108/// Algorithm used to compress large request bodies before they are sent.
109///
110/// The Apify API accepts both brotli (`br`) and gzip (`gzip`) request bodies. The reference JS
111/// client picks between them automatically (brotli when available, gzip otherwise); this client
112/// exposes the choice explicitly via
113/// [`ApifyClientBuilder::request_compression`](crate::ApifyClientBuilder::request_compression),
114/// because Rust's brotli support is always compiled in and would leave no runtime path to gzip.
115/// [`Brotli`](RequestCompression::Brotli) is the default (best ratio).
116#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
117#[non_exhaustive]
118pub enum RequestCompression {
119    /// Brotli (`Content-Encoding: br`). The default: best compression ratio, and the encoding the
120    /// reference client prefers.
121    #[default]
122    Brotli,
123    /// Gzip (`Content-Encoding: gzip`). Choose this for environments or intermediaries that do not
124    /// handle brotli.
125    Gzip,
126}
127
128/// HTTP method of a request.
129#[derive(Debug, Clone, Copy, PartialEq, Eq)]
130pub enum HttpMethod {
131    /// HTTP `GET`.
132    Get,
133    /// HTTP `POST`.
134    Post,
135    /// HTTP `PUT`.
136    Put,
137    /// HTTP `DELETE`.
138    Delete,
139    /// HTTP `HEAD`.
140    Head,
141}
142
143impl HttpMethod {
144    /// Returns the uppercase string representation, e.g. `"GET"`.
145    pub fn as_str(&self) -> &'static str {
146        match self {
147            HttpMethod::Get => "GET",
148            HttpMethod::Post => "POST",
149            HttpMethod::Put => "PUT",
150            HttpMethod::Delete => "DELETE",
151            HttpMethod::Head => "HEAD",
152        }
153    }
154}
155
156/// A fully-resolved HTTP request, ready to be sent by an [`HttpBackend`].
157///
158/// All cross-cutting concerns (auth header, user agent, retry policy) are applied by
159/// [`HttpClient`] before the request reaches the backend.
160#[derive(Debug, Clone)]
161pub struct HttpRequest {
162    /// The HTTP method.
163    pub method: HttpMethod,
164    /// The fully-qualified request URL (including query string).
165    pub url: String,
166    /// Request headers.
167    pub headers: HashMap<String, String>,
168    /// Raw request body bytes (already serialized).
169    pub body: Option<Vec<u8>>,
170    /// Per-request timeout. The backend should abort the request after this duration.
171    pub timeout: Duration,
172}
173
174/// An HTTP response returned by an [`HttpBackend`].
175#[derive(Debug, Clone)]
176pub struct HttpResponse {
177    /// HTTP status code.
178    pub status: u16,
179    /// Response headers.
180    pub headers: HashMap<String, String>,
181    /// Raw response body bytes.
182    pub body: Vec<u8>,
183}
184
185impl HttpResponse {
186    /// Returns the value of a response header (case-insensitive lookup).
187    pub fn header(&self, name: &str) -> Option<&str> {
188        let lower = name.to_ascii_lowercase();
189        self.headers
190            .iter()
191            .find(|(k, _)| k.to_ascii_lowercase() == lower)
192            .map(|(_, v)| v.as_str())
193    }
194}
195
196/// The replaceable transport contract.
197///
198/// Implementors are responsible only for sending a single request and returning the
199/// raw response. Retries, authentication and serialization are handled by
200/// [`HttpClient`], so a backend only needs to perform one network round-trip.
201#[async_trait]
202pub trait HttpBackend: Send + Sync + std::fmt::Debug {
203    /// Sends a single HTTP request and returns the response.
204    ///
205    /// Network-level failures (connection refused, DNS, timeout) should be returned as
206    /// [`ApifyClientError::Http`] or [`ApifyClientError::Timeout`]. A non-2xx HTTP
207    /// status is *not* an error at this layer — return it as a normal [`HttpResponse`].
208    async fn send(&self, request: HttpRequest) -> ApifyClientResult<HttpResponse>;
209}
210
211/// The default [`HttpBackend`] implementation, backed by [`reqwest`].
212#[derive(Debug, Clone)]
213pub struct ReqwestBackend {
214    client: reqwest::Client,
215}
216
217impl ReqwestBackend {
218    /// Creates a new backend with a default `reqwest::Client`.
219    pub fn new() -> Self {
220        Self {
221            client: reqwest::Client::new(),
222        }
223    }
224
225    /// Creates a backend wrapping a caller-provided `reqwest::Client`.
226    ///
227    /// Useful for sharing a connection pool or customizing proxy/TLS settings.
228    pub fn with_client(client: reqwest::Client) -> Self {
229        Self { client }
230    }
231}
232
233impl Default for ReqwestBackend {
234    fn default() -> Self {
235        Self::new()
236    }
237}
238
239#[async_trait]
240impl HttpBackend for ReqwestBackend {
241    async fn send(&self, request: HttpRequest) -> ApifyClientResult<HttpResponse> {
242        let method = match request.method {
243            HttpMethod::Get => reqwest::Method::GET,
244            HttpMethod::Post => reqwest::Method::POST,
245            HttpMethod::Put => reqwest::Method::PUT,
246            HttpMethod::Delete => reqwest::Method::DELETE,
247            HttpMethod::Head => reqwest::Method::HEAD,
248        };
249
250        let mut builder = self
251            .client
252            .request(method, &request.url)
253            .timeout(request.timeout);
254
255        for (key, value) in &request.headers {
256            builder = builder.header(key, value);
257        }
258        if let Some(body) = request.body {
259            builder = builder.body(body);
260        }
261
262        let response = builder.send().await?;
263        let status = response.status().as_u16();
264
265        let mut headers = HashMap::new();
266        for (name, value) in response.headers().iter() {
267            if let Ok(v) = value.to_str() {
268                headers.insert(name.as_str().to_string(), v.to_string());
269            }
270        }
271
272        let body = response.bytes().await?.to_vec();
273        Ok(HttpResponse {
274            status,
275            headers,
276            body,
277        })
278    }
279}
280
281/// Configuration for the retry/timeout behaviour of the [`HttpClient`].
282#[derive(Debug, Clone)]
283pub struct RetryConfig {
284    /// Maximum number of *retries* (i.e. the request is attempted up to `max_retries + 1` times).
285    pub max_retries: u32,
286    /// Minimum delay between retries; doubled on each subsequent retry (exponential backoff).
287    pub min_delay_between_retries: Duration,
288    /// Overall per-request timeout budget. Each attempt's timeout grows but is capped here.
289    pub timeout: Duration,
290}
291
292/// The orchestrating HTTP client shared by every resource client.
293///
294/// It owns the [`HttpBackend`], the optional API token, and the retry/timeout policy.
295/// It is cheap to clone (everything is reference-counted) so each resource client can
296/// hold its own handle.
297#[derive(Debug, Clone)]
298pub struct HttpClient {
299    backend: Arc<dyn HttpBackend>,
300    token: Option<String>,
301    user_agent: String,
302    retry: RetryConfig,
303    compression: RequestCompression,
304}
305
306impl HttpClient {
307    pub(crate) fn new(
308        backend: Arc<dyn HttpBackend>,
309        token: Option<String>,
310        user_agent: String,
311        retry: RetryConfig,
312        compression: RequestCompression,
313    ) -> Self {
314        Self {
315            backend,
316            token,
317            user_agent,
318            retry,
319            compression,
320        }
321    }
322
323    /// Sends `request` with authentication, the user-agent header, and the retry policy
324    /// applied. Returns the first successful response, or the final error.
325    pub async fn call(&self, mut request: HttpRequest) -> ApifyClientResult<HttpResponse> {
326        // Inject auth + user-agent headers shared by every endpoint.
327        request
328            .headers
329            .insert("User-Agent".to_string(), self.user_agent.clone());
330        if let Some(token) = &self.token {
331            request
332                .headers
333                .insert("Authorization".to_string(), format!("Bearer {token}"));
334        }
335
336        // Compress the request body once (not per attempt) when it is large enough, mirroring the
337        // reference client. The API accepts both brotli- and gzip-encoded request bodies.
338        maybe_compress_request(&mut request, self.compression);
339
340        let method_str = request.method.as_str().to_string();
341        let path = extract_path(&request.url);
342
343        // The caller-supplied `request.timeout` is the per-endpoint base; it grows with each
344        // attempt up to the client's overall timeout budget.
345        let base_timeout = request.timeout;
346        let mut delay = self.retry.min_delay_between_retries;
347        // `saturating_add` so an extreme `max_retries` can't overflow the attempt count.
348        let max_attempts = self.retry.max_retries.saturating_add(1);
349
350        let mut attempt = 1;
351        loop {
352            // Grow per-attempt timeout with each attempt, capped at the overall budget.
353            let mut attempt_request = request.clone();
354            attempt_request.timeout = self.attempt_timeout(base_timeout, attempt);
355
356            let outcome = match self.backend.send(attempt_request).await {
357                Ok(response) => {
358                    if response.status < MAX_SUCCESS_STATUS_CODE {
359                        return Ok(response);
360                    }
361                    let api_error = build_api_error(&response, attempt, &method_str, &path);
362                    let retryable = is_status_retryable(response.status);
363                    (ApifyClientError::from(api_error), retryable)
364                }
365                Err(err) => {
366                    let retryable = is_error_retryable(&err);
367                    (err, retryable)
368                }
369            };
370
371            let (error, retryable) = outcome;
372            // Give up immediately on non-retryable errors or after the last attempt.
373            if !retryable || attempt == max_attempts {
374                return Err(error);
375            }
376
377            // Sleep with randomized exponential backoff before the next attempt. The backoff
378            // doubles each retry (matching the reference client, which uses `async-retry` with
379            // a factor of 2) and is capped at the overall request timeout so a single backoff
380            // can never exceed the budget the whole request is allowed.
381            sleep(randomized_delay(delay)).await;
382            // `saturating_mul` mirrors the saturating arithmetic in `attempt_timeout`.
383            delay = delay.saturating_mul(BACKOFF_FACTOR).min(self.retry.timeout);
384            attempt += 1;
385        }
386    }
387
388    /// Per-attempt timeout: `min(overall_timeout, base * 2^(attempt-1))`.
389    ///
390    /// The first attempt uses the per-endpoint `base` timeout; each retry doubles it so a
391    /// slow-but-progressing connection gets more time, while never exceeding the client's
392    /// overall timeout budget. Mirrors the reference clients.
393    fn attempt_timeout(&self, base: Duration, attempt: u32) -> Duration {
394        let scaled = base.saturating_mul(2u32.saturating_pow(attempt.saturating_sub(1)));
395        scaled.min(self.retry.timeout)
396    }
397
398    pub(crate) fn user_agent(&self) -> &str {
399        &self.user_agent
400    }
401
402    /// Returns the token and user-agent, for endpoints (like log streaming) that must
403    /// open a raw connection outside the buffered backend.
404    pub(crate) fn stream_credentials(&self) -> (Option<String>, String) {
405        (self.token.clone(), self.user_agent.clone())
406    }
407}
408
409/// Compresses `request.body` in place when it is present, at least [`MIN_COMPRESS_BYTES`] long,
410/// and no `Content-Encoding` is already set, adding the matching `Content-Encoding` header.
411///
412/// The algorithm is chosen by `compression` (defaulting to brotli). The size threshold and the
413/// "compress once, before retries" behaviour mirror the reference client.
414fn maybe_compress_request(request: &mut HttpRequest, compression: RequestCompression) {
415    let Some(body) = request.body.as_ref() else {
416        return;
417    };
418    if body.len() < MIN_COMPRESS_BYTES {
419        return;
420    }
421    // Respect a caller-provided `Content-Encoding` (case-insensitive): the body is then assumed to
422    // already be encoded, so re-compressing it would corrupt it.
423    let already_encoded = request
424        .headers
425        .keys()
426        .any(|k| k.eq_ignore_ascii_case("Content-Encoding"));
427    if already_encoded {
428        return;
429    }
430    // Skip media types that already carry their own compression (images, audio, video,
431    // archives): compressing them again burns CPU for little to no size reduction. This is
432    // mainly relevant to a raw-bytes Actor input (e.g. `ActorClient::start_raw` with
433    // `content_type: Some("application/zip")`); JSON request bodies always pass this check.
434    let content_type = request
435        .headers
436        .iter()
437        .find(|(k, _)| k.eq_ignore_ascii_case("Content-Type"))
438        .map(|(_, v)| v.as_str());
439    if !is_compressible_content_type(content_type) {
440        return;
441    }
442
443    let (encoding, compressed) = match compression {
444        RequestCompression::Brotli => (CONTENT_ENCODING_BROTLI, brotli_compress(body)),
445        RequestCompression::Gzip => (CONTENT_ENCODING_GZIP, gzip_compress(body)),
446    };
447    request
448        .headers
449        .insert("Content-Encoding".to_string(), encoding.to_string());
450    request.body = Some(compressed);
451}
452
453/// Decides whether a request body with the given content type is worth compressing.
454///
455/// Images, audio, video and archives already carry their own compression: running them through
456/// brotli or gzip burns CPU, holds a second full copy of the body in memory, and usually produces
457/// output no smaller than the input (sometimes larger). Formats that are raw despite such a media
458/// type, e.g. `image/bmp` or `audio/wav`, are still compressed. A body with no content type is
459/// assumed to be compressible. Matches the reference client's `isCompressibleContentType`.
460fn is_compressible_content_type(content_type: Option<&str>) -> bool {
461    let Some(content_type) = content_type else {
462        return true;
463    };
464    // `Content-Type` is case-insensitive and may carry parameters, e.g. `text/plain; charset=utf-8`.
465    let media_type = content_type
466        .split(';')
467        .next()
468        .unwrap_or(content_type)
469        .trim()
470        .to_ascii_lowercase();
471
472    if COMPRESSIBLE_MEDIA_TYPES.contains(&media_type.as_str()) {
473        return true;
474    }
475    if COMPRESSIBLE_MEDIA_TYPE_SUFFIXES
476        .iter()
477        .any(|suffix| media_type.ends_with(suffix))
478    {
479        return true;
480    }
481    if ALREADY_COMPRESSED_MEDIA_TYPES.contains(&media_type.as_str()) {
482        return false;
483    }
484    !ALREADY_COMPRESSED_MEDIA_TYPE_PREFIXES
485        .iter()
486        .any(|prefix| media_type.starts_with(prefix))
487}
488
489/// Brotli-compresses `data`. Writing to an in-memory `Vec` is infallible, so this cannot fail.
490fn brotli_compress(data: &[u8]) -> Vec<u8> {
491    use std::io::Write;
492
493    let mut writer = brotli::CompressorWriter::new(
494        Vec::new(),
495        BROTLI_BUFFER_SIZE,
496        BROTLI_QUALITY,
497        BROTLI_WINDOW_SIZE,
498    );
499    writer
500        .write_all(data)
501        .expect("writing to an in-memory Vec never fails");
502    writer.into_inner()
503}
504
505/// Gzip-compresses `data`. Writing to and finishing an in-memory `Vec` is infallible, so this
506/// cannot fail.
507fn gzip_compress(data: &[u8]) -> Vec<u8> {
508    use flate2::{write::GzEncoder, Compression};
509    use std::io::Write;
510
511    let mut encoder = GzEncoder::new(Vec::new(), Compression::new(GZIP_COMPRESSION_LEVEL));
512    encoder
513        .write_all(data)
514        .expect("writing to an in-memory Vec never fails");
515    encoder
516        .finish()
517        .expect("finishing an in-memory Vec never fails")
518}
519
520/// Returns the path + query portion of a URL, for error reporting.
521pub(crate) fn extract_path(url: &str) -> Option<String> {
522    // Find the start of the path after the scheme+host.
523    let after_scheme = url.split_once("://").map(|(_, rest)| rest).unwrap_or(url);
524    after_scheme
525        .find('/')
526        .map(|idx| after_scheme[idx..].to_string())
527}
528
529/// We retry `429` (rate limit) and `5xx` (internal server errors), matching the
530/// reference client policy. Other `4xx` statuses are caller errors and are not retried.
531fn is_status_retryable(status: u16) -> bool {
532    status == RATE_LIMIT_EXCEEDED_STATUS_CODE || status >= MIN_SERVER_ERROR_STATUS_CODE
533}
534
535/// Only transport-level failures are retryable. Programming errors (serde, invalid
536/// argument) and already-classified API errors are handled elsewhere, so they are not
537/// retried here — matching the reference clients, which retry only network/timeout errors.
538fn is_error_retryable(err: &ApifyClientError) -> bool {
539    matches!(err, ApifyClientError::Http(_) | ApifyClientError::Timeout)
540}
541
542/// Parses the API error body (if present) into an [`ApiError`].
543pub(crate) fn build_api_error(
544    response: &HttpResponse,
545    attempt: u32,
546    method: &str,
547    path: &Option<String>,
548) -> ApiError {
549    let parsed: Option<ApiErrorBody> = serde_json::from_slice(&response.body).ok();
550    let (error_type, message, data) = match parsed {
551        Some(body) => (
552            body.error.error_type,
553            body.error
554                .message
555                .unwrap_or_else(|| format!("Unexpected error with status {}", response.status)),
556            body.error.data,
557        ),
558        None => {
559            let raw = String::from_utf8_lossy(&response.body);
560            let message = if raw.trim().is_empty() {
561                format!("Unexpected error with status {}", response.status)
562            } else {
563                format!("Unexpected error: {raw}")
564            };
565            (None, message, None)
566        }
567    };
568
569    ApiError {
570        status_code: response.status,
571        error_type,
572        message,
573        attempt,
574        http_method: Some(method.to_string()),
575        path: path.clone(),
576        data,
577    }
578}
579
580/// Returns a delay chosen randomly from the interval `[delay, 2*delay)`, matching the
581/// exponential-backoff-with-jitter algorithm described in the API docs.
582///
583/// `pub(crate)` so other retrying call sites (e.g. `RequestQueueClient::batch_add_requests`'s
584/// unprocessed-request retries) can reuse the same jitter source instead of duplicating it.
585pub(crate) fn randomized_delay(delay: Duration) -> Duration {
586    let base = delay.as_millis() as u64;
587    if base == 0 {
588        return delay;
589    }
590    let extra = next_jitter() % base;
591    Duration::from_millis(base + extra)
592}
593
594/// A process-wide pseudo-random source for backoff jitter.
595///
596/// Backoff jitter does not need cryptographic quality, but it must be well-distributed and
597/// uncorrelated across concurrent retries (otherwise many clients retry in lockstep). A
598/// shared atomically-advanced SplitMix64 generator, seeded once from the clock, gives each
599/// caller a distinct value without pulling in a heavyweight RNG dependency.
600fn next_jitter() -> u64 {
601    use std::sync::atomic::{AtomicU64, Ordering};
602    static STATE: AtomicU64 = AtomicU64::new(0);
603
604    const GOLDEN_GAMMA: u64 = 0x9E3779B97F4A7C15;
605
606    // Lazily seed from the clock on first use. A racing double-seed is harmless: both
607    // candidate seeds are valid SplitMix64 stream starting points.
608    if STATE.load(Ordering::Relaxed) == 0 {
609        let seed = std::time::SystemTime::now()
610            .duration_since(std::time::UNIX_EPOCH)
611            .map(|d| d.as_nanos() as u64)
612            .unwrap_or(GOLDEN_GAMMA)
613            | 1;
614        let _ = STATE.compare_exchange(0, seed, Ordering::Relaxed, Ordering::Relaxed);
615    }
616
617    // SplitMix64: advance the shared state by the golden-ratio increment in a single atomic
618    // read-modify-write (`fetch_add`) so concurrent callers each observe a distinct value —
619    // a plain load-then-store could hand two racing retries the same number. Then scramble.
620    let mut z = STATE
621        .fetch_add(GOLDEN_GAMMA, Ordering::Relaxed)
622        .wrapping_add(GOLDEN_GAMMA);
623    z = (z ^ (z >> 30)).wrapping_mul(0xBF58476D1CE4E5B9);
624    z = (z ^ (z >> 27)).wrapping_mul(0x94D049BB133111EB);
625    z ^ (z >> 31)
626}
627
628/// Sleeps for the given duration (public crate-internal helper for poll loops).
629pub(crate) async fn sleep_public(duration: Duration) {
630    sleep(duration).await;
631}
632
633/// Sleeps for the given duration using the Tokio timer (the runtime `reqwest` requires).
634async fn sleep(duration: Duration) {
635    tokio::time::sleep(duration).await;
636}