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}