deps_core/cache.rs
1//! HTTP response cache shared by every ecosystem's registry client.
2//!
3//! Wraps outbound registry requests with RFC 7232 conditional-request
4//! validation (`ETag`/`If-None-Match`, `Last-Modified`/`If-Modified-Since`) so
5//! that unchanged registry data is served from a bounded in-memory cache
6//! instead of re-fetched. Entry count and total retained bytes are both
7//! capped to keep memory use predictable under long-running LSP sessions.
8
9use crate::cache_policy::CACHE_EVICTION_PERCENTAGE;
10use crate::error::{DepsError, RateLimitEvidence, Result};
11use crate::net_policy::{RegistryAccessPolicy, WorkspaceRegistryAccess};
12use crate::redact::RedactedUrl;
13use bytes::{Bytes, BytesMut};
14use dashmap::DashMap;
15use reqwest::{Client, Response, StatusCode, Url, header};
16use serde::Serialize;
17use std::borrow::Cow;
18use std::hash::{Hash, Hasher};
19use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
20use std::sync::{Arc, RwLock};
21use std::time::Instant;
22
23/// Maximum number of cached entries to prevent unbounded memory growth.
24const MAX_CACHE_ENTRIES: usize = 1000;
25
26/// Maximum total bytes retained across all cached response bodies.
27///
28/// `MAX_CACHE_ENTRIES` alone bounds entry *count*, not size: since a single
29/// response body may be as large as [`MAX_RESPONSE_BYTES`] (32 MiB), a cache
30/// full of near-cap entries could retain tens of gigabytes even though real
31/// registry payloads are typically well under 1 MB. This budget is a
32/// defense-in-depth cap (CWE-400) against that worst case, evicted
33/// alongside the count-based limit in [`HttpCache::evict_entries`]. 64 MiB
34/// comfortably holds thousands of typical registry responses while still
35/// bounding the pathological case.
36///
37/// This is a best-effort bound, not a hard guarantee: it is checked once
38/// per request in [`HttpCache::get_cached_with_headers_via`], so multiple
39/// requests already in flight when the budget is crossed can each finish
40/// inserting before the next check fires. [`MAX_CACHEABLE_ENTRY_BYTES`]
41/// keeps that per-request overshoot small (at most one admission-cap-sized
42/// insert per concurrent in-flight request) rather than bounding it exactly.
43const MAX_CACHE_BYTES: usize = 64 * 1024 * 1024;
44
45/// Maximum size of a single response body that will be retained in the
46/// cache; larger bodies are still returned to the caller, just never
47/// stored.
48///
49/// Without this cap, [`MAX_CACHE_BYTES`] alone lets a handful of
50/// large-but-legitimate responses (up to [`MAX_RESPONSE_BYTES`], 32 MiB
51/// each) evict the *entire* rest of the cache: at a 2x ratio between the
52/// two constants, just two max-size entries would saturate the whole
53/// budget. Set to an eighth of [`MAX_CACHE_BYTES`] (8 MiB) so no single
54/// entry can claim more than 1/8 of the budget — a handful of large
55/// responses degrade to "not cached" instead of "evicts the small-payload
56/// working set".
57const MAX_CACHEABLE_ENTRY_BYTES: usize = MAX_CACHE_BYTES / 8;
58
59/// HTTP request timeout in seconds.
60const HTTP_TIMEOUT_SECS: u64 = 30;
61
62/// Whether [`HttpCache`] may issue outbound network requests (issue #483).
63///
64/// Passed to [`HttpCache::set_offline`]. `Offline` also overrides
65/// [`CacheMode`] to behave as [`CacheMode::Enabled`] — see that method's docs.
66#[derive(Debug, Clone, Copy, PartialEq, Eq)]
67pub enum NetworkMode {
68 /// Outbound registry requests are attempted normally.
69 Online,
70 /// All outbound requests are blocked; only cached/warm entries are served.
71 Offline,
72}
73
74impl NetworkMode {
75 /// Builds a `NetworkMode` from a `network.offline` config/CLI flag (`true` means
76 /// [`Self::Offline`]).
77 ///
78 /// The single, explicitly named conversion point from that boundary's `bool`
79 /// representation (issue #1436 S1) — deliberately not a `From<bool>` impl, which would
80 /// let any unrelated `bool` (e.g. a `cache.enabled` flag transposed at the call site)
81 /// silently convert too, defeating the point of typing this API in the first place.
82 #[must_use]
83 pub fn from_offline_flag(offline: bool) -> Self {
84 if offline { Self::Offline } else { Self::Online }
85 }
86
87 /// Whether outbound network requests may be attempted (`self` is [`Self::Online`]).
88 ///
89 /// Centralizes the `network == NetworkMode::Online` check every call site otherwise
90 /// reimplements inline (issue #1557 code-review finding).
91 ///
92 /// # Examples
93 ///
94 /// ```
95 /// use deps_core::NetworkMode;
96 ///
97 /// assert!(NetworkMode::Online.is_online());
98 /// assert!(!NetworkMode::Offline.is_online());
99 /// ```
100 #[must_use]
101 pub const fn is_online(self) -> bool {
102 matches!(self, Self::Online)
103 }
104
105 /// Whether outbound network requests are blocked (`self` is [`Self::Offline`]).
106 ///
107 /// # Examples
108 ///
109 /// ```
110 /// use deps_core::NetworkMode;
111 ///
112 /// assert!(NetworkMode::Offline.is_offline());
113 /// assert!(!NetworkMode::Online.is_offline());
114 /// ```
115 #[must_use]
116 pub const fn is_offline(self) -> bool {
117 matches!(self, Self::Offline)
118 }
119}
120
121/// Whether [`HttpCache`] uses its entry-map cache to serve warm entries (issue #482).
122///
123/// Passed to [`HttpCache::set_cache_enabled`]. See that method's docs for the override
124/// [`NetworkMode::Offline`] has on this value.
125#[derive(Debug, Clone, Copy, PartialEq, Eq)]
126pub enum CacheMode {
127 /// The entry map is consulted and populated as usual.
128 Enabled,
129 /// The entry map is bypassed entirely, except while [`NetworkMode::Offline`] overrides
130 /// this to behave as `Enabled`.
131 Disabled,
132}
133
134impl CacheMode {
135 /// Builds a `CacheMode` from a `cache.enabled` config/CLI flag (`true` means
136 /// [`Self::Enabled`]).
137 ///
138 /// The single, explicitly named conversion point from that boundary's `bool`
139 /// representation — see [`NetworkMode::from_offline_flag`]'s doc for why this is a named
140 /// constructor rather than a `From<bool>` impl.
141 #[must_use]
142 pub fn from_enabled_flag(enabled: bool) -> Self {
143 if enabled {
144 Self::Enabled
145 } else {
146 Self::Disabled
147 }
148 }
149}
150
151/// Maximum decompressed response body size accepted from a single request.
152///
153/// `reqwest`'s `gzip` feature strips `Content-Length`/`Content-Encoding` after
154/// decoding a response, so a header-based pre-check cannot bound body size
155/// (`response.content_length()` is `None` for every decoded response). This
156/// cap is instead enforced by counting bytes as the body streams in, aborting
157/// as soon as the running total would exceed the limit.
158const MAX_RESPONSE_BYTES: usize = 32 * 1024 * 1024;
159
160/// Ceiling every [`BodyLimit`] is clamped to at construction, so no caller can weaken
161/// the size guard [`read_body_capped`] enforces past this value.
162///
163/// 128 MiB comfortably covers the largest known caller (`deps-pypi`'s PyPI
164/// Simple API full index, ~43 MB decompressed today, capped at 96 MiB for organic
165/// growth) while still bounding the pathological case.
166const ABSOLUTE_MAX_RESPONSE_BYTES: usize = 128 * 1024 * 1024;
167
168/// Upper bound on a single response body, clamped at construction so no caller can
169/// weaken the guard `read_body_capped` enforces past `ABSOLUTE_MAX_RESPONSE_BYTES`.
170///
171/// Every cache method that previously read `MAX_RESPONSE_BYTES` directly now takes
172/// this newtype instead (defaulting to it via [`Self::DEFAULT`]), so a caller that
173/// legitimately needs a larger cap — e.g. a full-index fetch that bypasses the entry
174/// cache entirely, like [`HttpCache::get_transport_only_with_headers_limited`] — can
175/// request one without touching the shared constant every other registry client
176/// relies on.
177///
178/// # Examples
179///
180/// ```
181/// use deps_core::cache::BodyLimit;
182///
183/// let default_limit = BodyLimit::DEFAULT;
184/// let clamped = BodyLimit::new(usize::MAX);
185/// assert_ne!(clamped, default_limit);
186/// ```
187#[derive(Debug, Clone, Copy, PartialEq, Eq)]
188pub struct BodyLimit(usize);
189
190impl BodyLimit {
191 /// The default limit (`MAX_RESPONSE_BYTES`), used by every cache method that
192 /// does not take an explicit [`BodyLimit`].
193 pub const DEFAULT: Self = Self(MAX_RESPONSE_BYTES);
194
195 /// Creates a limit of `bytes`, clamped down to `ABSOLUTE_MAX_RESPONSE_BYTES` if
196 /// `bytes` exceeds it.
197 #[must_use]
198 pub const fn new(bytes: usize) -> Self {
199 if bytes > ABSOLUTE_MAX_RESPONSE_BYTES {
200 Self(ABSOLUTE_MAX_RESPONSE_BYTES)
201 } else {
202 Self(bytes)
203 }
204 }
205
206 /// The clamped byte value this limit enforces.
207 #[must_use]
208 pub const fn bytes(self) -> usize {
209 self.0
210 }
211}
212
213/// Whether `url`'s host is loopback (`127.0.0.1`, `localhost`, or `::1`), with an
214/// `http`/`https` scheme and an optional port — the shape every `mockito::Server` binds to.
215///
216/// Only compiled into test builds (see [`ensure_https`]): a non-loopback host must never
217/// be allowed to bypass the HTTPS requirement, even under `cfg(test)`/`test-util`. See
218/// [`crate::net_policy::validate_index_url`]'s own private loopback check (`is_loopback_url`)
219/// for the counterpart this was modeled on — `is_loopback_url` now matches [`Url::host`]'s
220/// structured `Host` enum (`Ipv6Addr::LOCALHOST` comparison, #1568) rather than comparing the
221/// bracketed string form, closing the IPv6-loopback gap this doc used to describe as
222/// pre-existing/out of scope. The two still differ in accepted scheme: this
223/// function matches both `http` and `https` (symmetric with [`ensure_https`]'s own scheme
224/// check), while `is_loopback_url` only ever matches `http` — its caller already treats
225/// `https` as satisfying the requirement outright, so an `https` loopback host has no need for
226/// the carve-out.
227///
228/// Parses with [`Url::parse`] and compares [`Url::host_str`] rather than splitting the raw
229/// string on `:` — a naive split misreads userinfo as the host boundary (e.g.
230/// `http://localhost:80@evil.com/x`'s actual host is `evil.com`, but splitting on the first
231/// `:` yields `localhost`), which let a public HTTP host bypass the HTTPS requirement by
232/// prefixing a loopback-looking userinfo (#1562).
233#[cfg(any(test, feature = "test-util"))]
234fn is_loopback_host(url: &str) -> bool {
235 let Ok(parsed) = Url::parse(url) else {
236 return false;
237 };
238 if !matches!(parsed.scheme(), "http" | "https") {
239 return false;
240 }
241 // `Url::host_str` keeps the brackets on an IPv6 literal (`"[::1]"`, not `"::1"`).
242 let host = parsed
243 .host_str()
244 .map(|h| h.trim_start_matches('[').trim_end_matches(']'));
245 matches!(host, Some("127.0.0.1" | "localhost" | "::1"))
246}
247
248/// Validates that a URL uses HTTPS protocol.
249///
250/// Returns an error if the URL doesn't start with "https://".
251/// This ensures all network requests are encrypted.
252///
253/// A loopback HTTP URL (`127.0.0.1`/`localhost`/`::1`, the shape every `mockito::Server`
254/// binds to) is allowed in `deps-core`'s own test builds (`cfg(test)`) and in other
255/// workspace crates' test builds via the `test-util` feature — `cfg(test)` alone does not
256/// apply there, since those crates depend on `deps-core` as a normal, non-dev dependency.
257/// Any other HTTP host is still rejected even under those cfgs: `test-util` is a public,
258/// independently-enablable crates.io feature, so this must not become "any host, any
259/// environment" just because the feature is on.
260#[inline]
261fn ensure_https(url: &str) -> Result<()> {
262 if url.starts_with("https://") {
263 return Ok(());
264 }
265 #[cfg(any(test, feature = "test-util"))]
266 if is_loopback_host(url) {
267 return Ok(());
268 }
269 Err(DepsError::CacheError(format!(
270 "URL must use HTTPS: {}",
271 RedactedUrl::new(url)
272 )))
273}
274
275/// Whether a non-2xx response carries explicit evidence of genuine rate-limit exhaustion —
276/// a 403/429 with `X-RateLimit-Remaining: 0`, or a `Retry-After` header on either (#1295) —
277/// as opposed to a 403 for some other reason (abuse-detection false positive, an
278/// access-restricted resource).
279///
280/// `Retry-After` covers GitHub's *secondary* rate limit, which arrives as 403 or 429 with a
281/// non-zero (or absent) `X-RateLimit-Remaining` — a `remaining == 0` check alone misses it
282/// (critic S4). `429` is included alongside `403` since GitHub returns 429 for some rate-limit
283/// responses too, not only 403.
284///
285/// **Residual false positive** (critic N4): `Retry-After` alone on a 403, with no
286/// `X-RateLimit-Remaining` at all, is not *unambiguous* evidence — a WAF/Cloudflare
287/// bot-challenge 403 can also carry `Retry-After`, and this predicate cannot distinguish that
288/// from a genuine secondary rate limit. This is a narrower false-positive surface than the
289/// bug #1295 fixes (which treated *every* untokened 403 as a rate limit with zero
290/// corroborating evidence), and considered an acceptable trade-off rather than a bug to
291/// eliminate here — see the issue for the full evidence-strength discussion.
292///
293/// Checked generically here (any header-carrying non-2xx response, not pinned to a GitHub
294/// host) rather than in `crate::github`: the response's headers are only available at this
295/// live-fetch chokepoint — by the time a caller like `GithubTagsClient` sees the error, it has
296/// already collapsed to [`DepsError::HttpStatus`] with no header data left (the exact gap
297/// `crate::test_util::unwrap_or_skip_github_rate_limit`'s doc used to describe). Not pinning
298/// to a GitHub-shaped `url` also means a `mockito`-backed test can exercise this without a
299/// real `api.github.com` request — including from non-GitHub ecosystem crates (`deps-gitlab-ci`)
300/// whose registries can send the same evidence shape.
301#[inline]
302fn confirmed_rate_limit_exhaustion(status: StatusCode, headers: &header::HeaderMap) -> bool {
303 if !matches!(
304 status,
305 StatusCode::FORBIDDEN | StatusCode::TOO_MANY_REQUESTS
306 ) {
307 return false;
308 }
309 let remaining_exhausted = headers
310 .get("x-ratelimit-remaining")
311 .and_then(|v| v.to_str().ok())
312 .and_then(|v| v.parse::<u64>().ok())
313 == Some(0);
314 remaining_exhausted || headers.contains_key("retry-after")
315}
316
317/// Fixed, registry-neutral message for a confirmed rate-limit-exhaustion classification
318/// (#1295, critic S1). Deliberately generic — [`http_status_error`] runs at every
319/// [`HttpCache`] live-fetch site, shared by all 14 ecosystems, not only GitHub, so it must not
320/// assume a GitHub-specific remedy (`GITHUB_TOKEN`) applies to whichever registry actually
321/// sent the confirming evidence. A GitHub-aware caller
322/// (`crate::github::classify_tags_fetch_error`) swaps in the GitHub-specific hint on top of
323/// this while keeping `verified: RateLimitEvidence::Confirmed`; `deps-gitlab-ci` does the
324/// equivalent for its own gate.
325const CONFIRMED_RATE_LIMIT_MESSAGE: &str =
326 "registry rate limit exceeded (confirmed by the response)";
327
328/// Builds the `Err` for a non-2xx `response`: a *verified* [`DepsError::RateLimited`] when
329/// [`confirmed_rate_limit_exhaustion`] holds, else the usual [`DepsError::HttpStatus`]. Shared
330/// by every live-fetch call site in this module so the check can't be forgotten at a new one
331/// (#1295).
332#[inline]
333fn http_status_error(url: &str, status: StatusCode, headers: &header::HeaderMap) -> DepsError {
334 if confirmed_rate_limit_exhaustion(status, headers) {
335 return DepsError::RateLimited {
336 message: CONFIRMED_RATE_LIMIT_MESSAGE.to_string(),
337 verified: RateLimitEvidence::Confirmed,
338 source_status: Some(status.as_u16()),
339 };
340 }
341 DepsError::HttpStatus {
342 url: RedactedUrl::new(url),
343 status: status.as_u16(),
344 }
345}
346
347/// True when a redirect hop moves from an `https` origin to a plain `http` one.
348///
349/// A redirect to any scheme other than `http`/`https` is already rejected by reqwest
350/// itself once a hop is followed, so the downgrade case is the only one this needs to
351/// catch here.
352fn is_https_downgrade(previous: &Url, next: &Url) -> bool {
353 previous.scheme() == "https" && next.scheme() == "http"
354}
355
356/// Whether a redirect hop's target host is one [`crate::net_policy::HostClass::never_a_registry`]
357/// blocks, exempting `Loopback` in test builds — the identical carve-out [`ensure_https`]
358/// already uses, without which every mockito redirect chain in this workspace's tests would
359/// break.
360fn hop_targets_blocked_host(url: &Url) -> bool {
361 let class = crate::net_policy::classify_host(url);
362 #[cfg(any(test, feature = "test-util"))]
363 let class_blocked = class.never_a_registry() && class != crate::net_policy::HostClass::Loopback;
364 #[cfg(not(any(test, feature = "test-util")))]
365 let class_blocked = class.never_a_registry();
366 class_blocked
367}
368
369/// Which cache-key namespace and [`AddrGuard`] tier a [`Transport`] enforces.
370///
371/// `Baseline` is every non-workspace request (all 11 ecosystems' registry/redirect/API
372/// traffic); `WorkspaceDeclared` is Cargo's workspace-declared-registry traffic, carrying the
373/// same [`WorkspaceRegistryAccess`] **value snapshot** an [`AddrGuard::WorkspaceDeclared`]
374/// carries (see that variant's docs) — [`HttpCache::cache_key`] reads the digit from this
375/// snapshot, never from a live `Arc<RegistryAccessPolicy>` read, so a request's cache key and
376/// its guard always agree on which policy era they were constructed under.
377#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
378enum CacheTier {
379 Baseline,
380 WorkspaceDeclared(WorkspaceRegistryAccess),
381 /// An origin-pinned, connect-address-guarded tier (issue #561/#562) — see
382 /// [`Transport::origin_pinned_guarded`]. `digest` identifies the `(trusted_origin,
383 /// policy_snapshot)` pair this transport was built for (**never** a credential identity —
384 /// see [`HttpCache::get_cached_pinned_with_headers`]'s separate `auth_id` argument, folded
385 /// into the cache key only). `authenticated` distinguishes an authenticated fetch (#561)
386 /// from #562's unauthenticated workspace-declared one for cache-eviction purposes
387 /// ([`CacheTier::is_authenticated`]) without affecting pooling — the shipped
388 /// [`Transport::origin_pinned`] (public path) stays on [`CacheTier::Baseline`], unaffected.
389 Pinned {
390 digest: u64,
391 authenticated: bool,
392 },
393}
394
395impl CacheTier {
396 /// Whether a 401/403 revalidation response on an entry under this tier should evict the
397 /// entry rather than serve the default stale-while-revalidate fallback (FR-015/NFR-004).
398 fn is_authenticated(self) -> bool {
399 matches!(
400 self,
401 Self::Pinned {
402 authenticated: true,
403 ..
404 }
405 )
406 }
407}
408
409/// The tier a [`Transport`]'s redirect policy and DNS resolver both enforce.
410///
411/// `WorkspaceDeclared` holds a [`WorkspaceRegistryAccess`] **value snapshot**, taken once at
412/// [`Transport::workspace`] construction time — not a live `Arc<RegistryAccessPolicy>` read on
413/// every [`Self::tier_allows`] call. [`Transport::workspace`] takes this same snapshot for its
414/// paired [`CacheTier::WorkspaceDeclared`], and [`HttpCache::set_registry_policy`] rebuilds the
415/// whole `Transport` (both snapshots included) on every actual policy transition, so a
416/// request's cache key and its guard always come from one consistent construction-time value,
417/// with no read-skew window between them.
418#[derive(Debug, Clone, Copy, PartialEq, Eq)]
419enum AddrGuard {
420 Baseline,
421 WorkspaceDeclared(WorkspaceRegistryAccess),
422}
423
424impl AddrGuard {
425 /// Whether a host classified as `class` may be reached under this guard's tier.
426 fn tier_allows(self, class: crate::net_policy::HostClass) -> bool {
427 match self {
428 Self::Baseline => true,
429 Self::WorkspaceDeclared(policy) => policy.allows(class),
430 }
431 }
432}
433
434/// Redirect policy for a [`Transport`]'s client, parameterized by the [`AddrGuard`] tier that
435/// client enforces.
436///
437/// [`ensure_https`] only validates the *initial* request URL; a `3xx` response can still
438/// redirect the actual connection anywhere, including down to plain HTTP — or, per spec
439/// `.local/specs/023-cargo-custom-registries/spec.md` NFR-003/plan-1b §1.1, straight to a
440/// cloud metadata endpoint or other host no legitimate registry redirect ever targets. This
441/// policy stops the redirect chain (rather than erroring) the moment a hop would do either,
442/// so the caller sees the last successful `3xx` response and handles it exactly like any
443/// other non-2xx status (`DepsError::HttpStatus`) instead of needing a distinct
444/// "redirect blocked" error variant.
445///
446/// The blocked-host check (`hop_targets_blocked_host`) is unconditional and
447/// policy-independent — it does not consult `guard` at all, since
448/// [`HostClass::never_a_registry`](crate::net_policy::HostClass::never_a_registry)
449/// is deliberately narrower than any workspace-registry policy setting: it blocks only the
450/// classes (loopback, link-local, cloud metadata, unspecified, reserved) that are never a
451/// legitimate registry redirect target for *any* ecosystem, benefiting every one of the eleven
452/// crates sharing the baseline client, not only Cargo's workspace-declared indexes.
453///
454/// `guard`'s own [`AddrGuard::tier_allows`] term additionally rejects a hop whose target class
455/// the *tier* does not allow — under [`AddrGuard::Baseline`] this term is constant `false`, so
456/// every non-Cargo ecosystem (and Cargo's own `$CARGO_HOME`-provenance traffic) is
457/// bit-for-bit unaffected; under [`AddrGuard::WorkspaceDeclared`] it closes the redirect-hop
458/// half of issue #455 (an IP-literal hop to an RFC1918/CGNAT address, which `hyper-util`
459/// parses directly and never routes through a resolver).
460///
461/// This only classifies the redirect target's URL string, not its DNS-resolved address — that
462/// residual gap (issue #449, "D1" in PR #447's plan) is closed for the name-hop case by
463/// [`BlockedAddrResolver`], which [`build_guarded_client`] wires into every client this module
464/// builds: a redirect hop reuses the same `Client`, so its target's resolved address is
465/// validated too, for free (FR-007).
466///
467/// Every other redirect — including cross-host ones, which are out of scope for this
468/// policy — falls through to reqwest's default (`Policy::limited(10)`), preserving the
469/// existing hop-count limit and mockito's plain-`http://` loopback chains
470/// used throughout this module's tests.
471fn redirect_policy(guard: AddrGuard) -> reqwest::redirect::Policy {
472 reqwest::redirect::Policy::custom(move |attempt| {
473 let downgraded = attempt
474 .previous()
475 .last()
476 .is_some_and(|previous| is_https_downgrade(previous, attempt.url()));
477 if downgraded
478 || hop_targets_blocked_host(attempt.url())
479 || !guard.tier_allows(crate::net_policy::classify_host(attempt.url()))
480 {
481 attempt.stop()
482 } else {
483 reqwest::redirect::Policy::default().redirect(attempt)
484 }
485 })
486}
487
488/// Redirect policy for a [`HttpCache::transport_for_origin`]-scoped client.
489///
490/// Stops any hop unless [`crate::net_policy::is_trusted_prefix`] accepts it against
491/// `trusted_origin`, parsed once at [`Transport`] construction — origin equality (scheme,
492/// host, and port must match exactly) **and** the hop's path lying at or under
493/// `trusted_origin`'s own path, at a proper path-segment boundary.
494///
495/// Origin equality is the issue #795 fix: the pre-#795 behavior
496/// (`attempt.url().as_str().starts_with(&trusted_origin)`, a raw string-prefix test) is
497/// satisfied by `https://gitlab.mycorp.dev.evil.com/...`, `https://gitlab.mycorp.dev-evil.com/...`,
498/// and `https://gitlab.mycorp.dev@evil.com/...` alike, even though only the last of those
499/// three actually shares a host with `evil.com` — `Url::origin()` ignores userinfo and
500/// matches scheme+host+port exactly, closing all three shapes at once. This also still
501/// covers a downgrade to plain `http://`: an `http://` hop's origin can never equal an
502/// `https://`-scheme trusted origin, which every current caller passes — a separate scheme
503/// check (as [`redirect_policy`] has, for its no-trusted-origin case) would be dead code
504/// here.
505///
506/// The path-segment-boundary check is **not** subsumed by origin equality, and deliberately
507/// kept: a caller (NuGet's registration-hive/flat-container paging, or `deps-cargo`'s sparse
508/// index, whose `RegistryIndex::as_str()` carries no trailing-slash guarantee either way —
509/// issue #795 S1) pins to a specific *path* on a registry host that also serves other,
510/// less-trusted paths, not merely to the host itself. A plain `str::starts_with` on the raw
511/// path (this function's pre-S1-fix shape) is itself vulnerable to the same class of bug one
512/// level down: a trusted path of `/cargo/index` would wrongly accept the same-origin sibling
513/// `/cargo/index-public/steal` or `/cargo/indexEVIL`, since both start with the trusted
514/// string textually. [`crate::net_policy::is_trusted_prefix`] requires a hop's path to equal
515/// the trusted path or continue immediately after a `/` following it, closing that
516/// regardless of whether `trusted_origin`'s own path happens to end in `/`.
517///
518/// A `trusted_origin` that fails to parse matches no hop — every redirect is stopped
519/// (fail-closed) rather than treated as "no restriction" — logged once at construction time
520/// via `tracing::warn!` rather than silently, since a caller-side bug producing an
521/// unparseable `trusted_origin` would otherwise present only as every redirect on that
522/// transport mysteriously failing. This is reachable today: `deps-nuget`'s Public tier builds
523/// `trusted_prefix` from a service-index `@id` string that is never `Url`-validated (that
524/// validation only runs for [`crate::net_policy::validate_index_url`]'s
525/// `NuGetRegistryTier::WorkspaceDeclared` path), so a malformed `@id` reaches this parse.
526fn trusted_origin_redirect_policy(trusted_origin: &str) -> reqwest::redirect::Policy {
527 let trusted = match Url::parse(trusted_origin) {
528 Ok(url) => Some(url),
529 Err(error) => {
530 tracing::warn!(
531 trusted_origin = crate::redact::url_for_tracing(trusted_origin),
532 %error,
533 "trusted_origin failed to parse; every redirect hop on this transport will be rejected"
534 );
535 None
536 }
537 };
538 reqwest::redirect::Policy::custom(move |attempt| {
539 if is_trusted_origin(attempt.url(), trusted.as_ref()) {
540 reqwest::redirect::Policy::default().redirect(attempt)
541 } else {
542 attempt.stop()
543 }
544 })
545}
546
547/// The origin-and-path decision [`trusted_origin_redirect_policy`]'s closure makes on every
548/// redirect hop, extracted as a pure function so the bypass shapes from issue #795 can be
549/// regression-tested directly against it — `reqwest::redirect::Attempt`'s fields are private
550/// outside the `reqwest` crate, so the closure itself cannot be unit-tested without going
551/// through a real HTTP round trip. Delegates to [`crate::net_policy::is_trusted_prefix`],
552/// shared with `deps-nuget`'s registration-hive page `@id` pre-check (issue #795 S2).
553fn is_trusted_origin(hop_url: &Url, trusted: Option<&Url>) -> bool {
554 trusted.is_some_and(|t| crate::net_policy::is_trusted_prefix(hop_url, t))
555}
556
557/// Error returned by [`BlockedAddrResolver`] when a DNS resolution cannot be trusted for
558/// connection use — either it produced no address, or at least one resolved address falls into
559/// a blocked [`crate::net_policy::HostClass`].
560///
561/// Kept distinct from [`DepsError`] since this crosses into `reqwest::dns::Resolve`'s own
562/// `BoxError` (`Box<dyn std::error::Error + Send + Sync>`) return type, not this crate's own
563/// error type.
564#[derive(Debug, thiserror::Error)]
565enum ResolveGuardError {
566 /// The resolver returned zero addresses for `host` — fail-closed (NFR-004) rather than
567 /// silently treating "nothing resolved" as "nothing to block".
568 #[error("DNS resolution for {host} returned no addresses")]
569 NoAddresses { host: String },
570 /// `addr`, resolved for `host`, falls into `class`, one of the
571 /// [`crate::net_policy::HostClass::never_a_registry`] classes no legitimate registry index
572 /// (or a redirect from one) could ever target.
573 #[error("resolved address {addr} for host {host} is {class}, blocked by net_policy")]
574 Blocked {
575 host: String,
576 addr: std::net::IpAddr,
577 class: crate::net_policy::HostClass,
578 },
579}
580
581/// Validates every address `tokio::net::lookup_host` returned for `host`, rejecting the whole
582/// resolution if any is blocked — an attacker's public A record alongside a blocked one must not
583/// keep the probe alive (FR-003). `guard`'s tier additionally rejects a resolved address whose
584/// class the *tier* does not allow (issue #455: an RFC1918/CGNAT-range name rebound at
585/// connect time), on top of the policy-independent [`HostClass::never_a_registry`](crate::net_policy::HostClass::never_a_registry)
586/// check every tier enforces.
587fn validate_resolved_addrs(
588 host: &str,
589 addrs: Vec<std::net::SocketAddr>,
590 guard: AddrGuard,
591) -> std::result::Result<Vec<std::net::SocketAddr>, ResolveGuardError> {
592 if addrs.is_empty() {
593 tracing::warn!(host, "DNS resolution returned no addresses");
594 return Err(ResolveGuardError::NoAddresses {
595 host: host.to_string(),
596 });
597 }
598 for addr in &addrs {
599 let class = crate::net_policy::classify_addr(addr.ip());
600 if class.never_a_registry() || !guard.tier_allows(class) {
601 tracing::warn!(host, addr = %addr.ip(), %class, "blocking DNS-resolved address");
602 return Err(ResolveGuardError::Blocked {
603 host: host.to_string(),
604 addr: addr.ip(),
605 class,
606 });
607 }
608 }
609 Ok(addrs)
610}
611
612/// The synthetic-lookup function signature [`TestLookup`] wraps.
613#[cfg(test)]
614type SyntheticLookupFn = dyn Fn(&str) -> Vec<std::net::SocketAddr> + Send + Sync;
615
616/// Test-only override for [`BlockedAddrResolver::resolve`], replacing `tokio::net::lookup_host`
617/// with a synthetic lookup — lets a test exercise the resolver-guard wiring against an address
618/// chosen by the test (e.g. an RFC1918 literal) without depending on real DNS. A newtype rather
619/// than a hand-written `Debug` directly on [`BlockedAddrResolver`], so that struct's own
620/// `#[derive(Debug)]` stays valid under both `cfg(test)` and not.
621#[cfg(test)]
622#[derive(Clone)]
623struct TestLookup(Arc<SyntheticLookupFn>);
624
625#[cfg(test)]
626impl std::fmt::Debug for TestLookup {
627 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
628 f.write_str("TestLookup(..)")
629 }
630}
631
632/// Connect-time DNS resolver that closes the rebinding TOCTOU gap left by [`ensure_https`]/
633/// [`hop_targets_blocked_host`]'s URL-string-only classification (issue #449): those check the
634/// declared hostname, but `reqwest`'s connector resolves DNS independently, later, and an
635/// attacker who controls the hostname's DNS can rebind it to a blocked address in between.
636///
637/// Wired into every client [`build_guarded_client`] returns, so all 11 ecosystem crates sharing
638/// the baseline client pool inherit it with zero per-crate plumbing (FR-006/NFR-002).
639///
640/// # Scope
641///
642/// This resolver's guard tier decides how much a resolved address is scrutinized: under
643/// [`AddrGuard::Baseline`] this enforces only the policy-independent
644/// [`crate::net_policy::HostClass::never_a_registry`] tier (loopback, link-local,
645/// cloud-metadata, unspecified, reserved), the same tier [`hop_targets_blocked_host`] already
646/// applies — closing issue #449's filed exploit (cloud-metadata rebinding) but not full `PublicOnly`
647/// semantics. Under [`AddrGuard::WorkspaceDeclared`], [`validate_resolved_addrs`] additionally
648/// rejects any resolved address outside the snapshotted [`crate::net_policy::WorkspaceRegistryAccess`]
649/// policy's allowed classes — closing issue #455 (a workspace-declared name that legitimately
650/// resolves to `HostClass::Global` at parse time, then rebinds to an RFC1918/CGNAT address at
651/// connect time).
652///
653/// # Fail-closed (NFR-004)
654///
655/// Returns `Err` — never `Ok`, never a fallback resolver — on a `lookup_host` error, zero
656/// addresses, or any resolved address [`validate_resolved_addrs`] rejects for `self.guard`'s
657/// tier.
658///
659/// # Known limitations
660///
661/// - [`ClientBuilder::resolve`](reqwest::ClientBuilder::resolve)/
662/// [`resolve_to_addrs`](reqwest::ClientBuilder::resolve_to_addrs) overrides wrap *outside* the
663/// configured resolver (`reqwest`'s `DnsResolverWithOverrides`) and would bypass this guard
664/// entirely if ever called — this workspace does not call them today.
665/// - A configured system proxy (`HTTPS_PROXY`) resolves the target hostname itself; this
666/// resolver then only ever sees the proxy's own address. Operator configuration, not
667/// attacker-controlled, so NOT claimed as a defended case.
668/// - Never applies to an IP-literal host (`https://169.254.169.254/`): `hyper-util`'s connector
669/// parses those directly and never calls the configured resolver, so
670/// [`classify_host`](crate::net_policy::classify_host) (via [`ensure_https`]/
671/// [`hop_targets_blocked_host`]/[`redirect_policy`]'s tier term) remains the sole guard for
672/// literals — a disjoint domain from this resolver's name-based one, not a gap. This is also
673/// why the cache layer gives [`HttpCache::get_cached_workspace`] zero protection against an
674/// *initial* request URL that is itself an IP literal — see that method's docs.
675/// - Unlike [`ensure_https`]/[`hop_targets_blocked_host`], this resolver has **no** `test-util`
676/// carve-out for `Loopback`: it blocks a `localhost`/`127.0.0.1` *name* unconditionally, in
677/// every build. A downstream `test-util` consumer that mocks by binding an IP literal (as
678/// this workspace's own `mockito` usage does) is unaffected — literals never reach this
679/// resolver at all — but one that mocks via a `localhost` *name* would be newly blocked.
680#[derive(Debug, Clone)]
681struct BlockedAddrResolver {
682 guard: AddrGuard,
683 #[cfg(test)]
684 lookup: Option<TestLookup>,
685}
686
687impl BlockedAddrResolver {
688 fn new(guard: AddrGuard) -> Self {
689 Self {
690 guard,
691 #[cfg(test)]
692 lookup: None,
693 }
694 }
695
696 #[cfg(test)]
697 fn with_lookup(guard: AddrGuard, lookup: TestLookup) -> Self {
698 Self {
699 guard,
700 lookup: Some(lookup),
701 }
702 }
703}
704
705impl reqwest::dns::Resolve for BlockedAddrResolver {
706 fn resolve(&self, name: reqwest::dns::Name) -> reqwest::dns::Resolving {
707 let host = name.as_str().to_string();
708 let guard = self.guard;
709 #[cfg(test)]
710 let lookup = self.lookup.clone();
711 Box::pin(async move {
712 #[cfg(test)]
713 let addrs: Vec<std::net::SocketAddr> = match &lookup {
714 Some(lookup) => (lookup.0)(&host),
715 None => tokio::net::lookup_host((host.as_str(), 0)).await?.collect(),
716 };
717 #[cfg(not(test))]
718 let addrs: Vec<std::net::SocketAddr> =
719 tokio::net::lookup_host((host.as_str(), 0)).await?.collect();
720
721 let addrs = validate_resolved_addrs(&host, addrs, guard)?;
722 Ok(Box::new(addrs.into_iter()) as reqwest::dns::Addrs)
723 })
724 }
725}
726
727/// Builds a client with `HttpCache`'s shared configuration (user agent, timeout), varying the
728/// redirect policy and resolver — kept in one place so a future client-wide setting (proxy,
729/// connection pool sizing, etc.) can't silently miss any [`Transport`] this module builds. This
730/// is also the workspace's only `Client::builder()` call site.
731fn build_client_inner(
732 redirect: reqwest::redirect::Policy,
733 resolver: BlockedAddrResolver,
734) -> Client {
735 #[expect(
736 clippy::expect_used,
737 reason = "fixed, hardcoded client configuration — no attacker-influenced input; can \
738 only fail on a genuinely broken TLS backend, which is unrecoverable anyway"
739 )]
740 Client::builder()
741 .user_agent(format!("deps-lsp/{}", env!("CARGO_PKG_VERSION")))
742 .timeout(std::time::Duration::from_secs(HTTP_TIMEOUT_SECS))
743 .redirect(redirect)
744 .dns_resolver(resolver)
745 .build()
746 .expect("failed to create HTTP client")
747}
748
749/// Pairs one [`AddrGuard`] value with both halves it governs — its redirect policy and its
750/// resolver — so a client whose redirect policy and resolver enforce different tiers cannot be
751/// built by this function.
752fn build_guarded_client(guard: AddrGuard) -> Client {
753 build_client_inner(redirect_policy(guard), BlockedAddrResolver::new(guard))
754}
755
756/// Test-only variant of [`build_guarded_client`] that substitutes a synthetic DNS lookup for
757/// `tokio::net::lookup_host` — shares [`build_client_inner`] with the production constructor, so
758/// deleting the `.dns_resolver(...)` wiring from that shared function fails any test built on
759/// this too, not just the production path.
760#[cfg(test)]
761fn build_guarded_client_with_lookup(guard: AddrGuard, lookup: TestLookup) -> Client {
762 build_client_inner(
763 redirect_policy(guard),
764 BlockedAddrResolver::with_lookup(guard, lookup),
765 )
766}
767
768/// A `Client` welded to the [`CacheTier`] its guard enforces.
769///
770/// [`Self::baseline`], [`Self::workspace`] and [`Self::origin_pinned`] are the only sanctioned
771/// way to build one: each derives its redirect policy, its resolver and its tier from a single
772/// [`AddrGuard`] value, so a mismatched pairing (e.g. a baseline-guarded client keyed under the
773/// workspace cache namespace) never arises through normal construction — though `cache.rs` is
774/// one module, so a hand-written `Transport { .. }` literal elsewhere in this file could still
775/// mismatch them; the three constructors are what make that a deliberate act, not an accident
776/// reachable by passing the wrong argument to an existing function.
777#[derive(Clone)]
778struct Transport {
779 client: Client,
780 tier: CacheTier,
781}
782
783impl Transport {
784 /// The shared, unauthenticated transport used by every non-workspace request.
785 fn baseline() -> Self {
786 Self {
787 client: build_guarded_client(AddrGuard::Baseline),
788 tier: CacheTier::Baseline,
789 }
790 }
791
792 /// The transport for Cargo's workspace-declared-registry requests, snapshotting `policy`'s
793 /// current value once and sharing that single snapshot between the guard and the cache-key
794 /// tier — see [`AddrGuard::WorkspaceDeclared`]'s docs for why this is a value snapshot, not
795 /// a live `Arc` read, and why the guard and the tier must never snapshot independently.
796 fn workspace(policy: &Arc<RegistryAccessPolicy>) -> Self {
797 let snapshot = policy.get();
798 Self {
799 client: build_guarded_client(AddrGuard::WorkspaceDeclared(snapshot)),
800 tier: CacheTier::WorkspaceDeclared(snapshot),
801 }
802 }
803
804 /// The transport for one [`HttpCache::transport_for_origin`]-pinned origin: a plain
805 /// [`AddrGuard::Baseline`] resolver, paired with [`trusted_origin_redirect_policy`] instead
806 /// of [`redirect_policy`] — that policy pins by URL prefix, which already subsumes the
807 /// blocked-host hop check, so this is the one documented caller of [`build_client_inner`]
808 /// directly rather than [`build_guarded_client`].
809 fn origin_pinned(trusted_origin: &str) -> Self {
810 Self {
811 client: build_client_inner(
812 trusted_origin_redirect_policy(trusted_origin),
813 BlockedAddrResolver::new(AddrGuard::Baseline),
814 ),
815 tier: CacheTier::Baseline,
816 }
817 }
818
819 /// The transport for one origin-pinned **workspace-declared** host (issue #561/#562),
820 /// optionally carrying a credential. Pairs [`trusted_origin_redirect_policy`] (send-scope
821 /// confinement — no redirect hop may leave `trusted_origin`) with
822 /// [`AddrGuard::WorkspaceDeclared`] (the connect-address policy guard, #455-class
823 /// protection) and the namespaced [`CacheTier::Pinned`] tier. One constructor serves both
824 /// #562's unauthenticated workspace-declared fetches (`authenticated: false`) and #561's
825 /// authenticated ones (`authenticated: true`) — `authenticated` only affects
826 /// [`HttpCache`]'s revalidation-eviction rule (FR-015), never `AddrGuard`/redirect
827 /// confinement. The shipped [`Self::origin_pinned`] (public `api.nuget.org` path) is
828 /// unaffected — this is a distinct constructor, not a modification of that one.
829 fn origin_pinned_guarded(
830 trusted_origin: &str,
831 policy: &Arc<RegistryAccessPolicy>,
832 authenticated: bool,
833 ) -> Self {
834 let snapshot = policy.get();
835 Self {
836 client: build_client_inner(
837 trusted_origin_redirect_policy(trusted_origin),
838 BlockedAddrResolver::new(AddrGuard::WorkspaceDeclared(snapshot)),
839 ),
840 tier: CacheTier::Pinned {
841 digest: pinned_digest(trusted_origin, snapshot),
842 authenticated,
843 },
844 }
845 }
846}
847
848/// Identifies a `(trusted_origin, policy_snapshot)` pair for [`CacheTier::Pinned`] — **never**
849/// a credential identity (see that variant's docs). Not cryptographically salted: unlike a
850/// caller's own credential-header digest (e.g. `deps_nuget`'s salted `auth_id`), this value is
851/// not attacker-observable secret material, only a pool/cache-key discriminant.
852fn pinned_digest(trusted_origin: &str, snapshot: WorkspaceRegistryAccess) -> u64 {
853 let mut hasher = std::collections::hash_map::DefaultHasher::new();
854 trusted_origin.hash(&mut hasher);
855 snapshot.hash(&mut hasher);
856 hasher.finish()
857}
858
859/// Reads a response body incrementally, aborting once it exceeds `limit`.
860///
861/// Chunked reading (via [`Response::chunk`]) is required because the
862/// decompressed body size is not known upfront: `gzip` decoding strips
863/// `Content-Length`, so the only reliable guard against an oversized or
864/// maliciously amplified (decompression-bomb) response is counting bytes
865/// as they arrive and bailing before the whole body is buffered.
866async fn read_body_capped(url: &str, mut response: Response, limit: BodyLimit) -> Result<Bytes> {
867 let mut body = BytesMut::new();
868 let limit = limit.bytes();
869
870 while let Some(chunk) = response
871 .chunk()
872 .await
873 .map_err(|e| DepsError::RegistryError {
874 package: RedactedUrl::new(url),
875 source: e.into(),
876 })?
877 {
878 if body.len() + chunk.len() > limit {
879 return Err(DepsError::ResponseTooLarge {
880 url: RedactedUrl::new(url),
881 limit,
882 });
883 }
884 body.extend_from_slice(&chunk);
885 }
886
887 Ok(body.freeze())
888}
889
890/// Cached HTTP response with validation headers.
891///
892/// Stores response body and cache validation headers (ETag, Last-Modified)
893/// for efficient conditional requests. The body uses `Bytes` which is an
894/// Arc-like type optimized for network data, enabling zero-cost cloning
895/// across multiple consumers without copying.
896///
897/// # Examples
898///
899/// ```
900/// use deps_core::cache::CachedResponse;
901/// use bytes::Bytes;
902/// use std::time::Instant;
903///
904/// let response = CachedResponse::new(Bytes::from("response data")).with_etag("\"abc123\"");
905///
906/// // Clone is cheap - only increments reference count
907/// let cloned = response.clone();
908/// ```
909#[non_exhaustive]
910#[derive(Debug, Clone)]
911pub struct CachedResponse {
912 /// Raw response body, shareable across consumers without copying.
913 pub body: Bytes,
914 /// `ETag` header from the response, used for `If-None-Match` revalidation.
915 pub etag: Option<String>,
916 /// `Last-Modified` header from the response, used for `If-Modified-Since` revalidation.
917 pub last_modified: Option<String>,
918 /// Local time the response was fetched, used for TTL expiry checks.
919 pub fetched_at: Instant,
920}
921
922impl CachedResponse {
923 /// Constructs a `CachedResponse` for `body`, fetched now, with [`Self::etag`] and
924 /// [`Self::last_modified`] left `None` — chain [`Self::with_etag`] and/or
925 /// [`Self::with_last_modified`] to attach them.
926 ///
927 /// Needed because [`Self`] is `#[non_exhaustive]`: a struct literal only works inside
928 /// this crate, so every other crate must go through this constructor instead.
929 ///
930 /// # Examples
931 ///
932 /// ```
933 /// use deps_core::cache::CachedResponse;
934 /// use bytes::Bytes;
935 ///
936 /// let response = CachedResponse::new(Bytes::from("response data")).with_etag("\"abc123\"");
937 /// assert_eq!(response.etag.as_deref(), Some("\"abc123\""));
938 /// ```
939 #[must_use]
940 pub fn new(body: Bytes) -> Self {
941 Self {
942 body,
943 etag: None,
944 last_modified: None,
945 fetched_at: Instant::now(),
946 }
947 }
948
949 /// Attaches the response's `ETag` header. See [`Self::etag`].
950 #[must_use]
951 pub fn with_etag(mut self, etag: impl Into<String>) -> Self {
952 self.etag = Some(etag.into());
953 self
954 }
955
956 /// Attaches the response's `Last-Modified` header. See [`Self::last_modified`].
957 #[must_use]
958 pub fn with_last_modified(mut self, last_modified: impl Into<String>) -> Self {
959 self.last_modified = Some(last_modified.into());
960 self
961 }
962
963 /// Overrides when the response was fetched — otherwise [`Self::new`] stamps
964 /// [`Instant::now`]. Lets external code (tests, benches) construct a backdated entry to
965 /// exercise TTL-expiry behavior, the one thing a direct struct literal used to allow.
966 #[must_use]
967 pub const fn with_fetched_at(mut self, fetched_at: Instant) -> Self {
968 self.fetched_at = fetched_at;
969 self
970 }
971}
972
973/// HTTP cache with ETag and Last-Modified validation.
974///
975/// Implements RFC 7232 conditional requests to minimize network traffic.
976/// All responses are cached with their validation headers, and subsequent
977/// requests use `If-None-Match` (ETag) or `If-Modified-Since` headers
978/// to check for updates.
979///
980/// The cache uses `Bytes` for response bodies, enabling efficient sharing
981/// of cached data across multiple consumers without copying. `Bytes` is
982/// an Arc-like type optimized for network I/O.
983///
984/// # Examples
985///
986/// ```no_run
987/// use deps_core::cache::HttpCache;
988///
989/// # async fn example() -> deps_core::error::Result<()> {
990/// let cache = HttpCache::new();
991///
992/// // First request - fetches from network
993/// let data1 = cache.get_cached("https://index.crates.io/se/rd/serde").await?;
994///
995/// // Second request - uses conditional GET (304 Not Modified if unchanged)
996/// let data2 = cache.get_cached("https://index.crates.io/se/rd/serde").await?;
997/// # Ok(())
998/// # }
999/// ```
1000///
1001/// # Cache key
1002///
1003/// Entries are keyed by URL alone (see `Self::cache_key`, private) — `extra_headers` (see
1004/// [`HttpCache::get_cached_with_headers`]) play no part in the cache key.
1005/// This is safe only as long as "same URL" implies "same representation":
1006/// a content-negotiating header (e.g. a per-request `Accept`) that can vary
1007/// the response body for an otherwise-identical URL requires giving each
1008/// distinct representation its own URL (e.g. a query parameter or distinct
1009/// path), not just a distinct header value, or callers requesting different
1010/// representations of the same URL will silently share one cache entry.
1011///
1012/// Likewise, the key doesn't encode *which* client (and so which redirect policy)
1013/// produced an entry — [`HttpCache::get_cached`] and [`HttpCache::get_cached_trusted_origin`]
1014/// share one entry map. No caller today requests the same URL through both, but one that did
1015/// could observe the other's cached (and differently redirect-validated) body.
1016///
1017/// [`Self::get_cached_workspace`] is the one exception: it is namespaced under a distinct,
1018/// policy-scoped key prefix (see `Self::cache_key`, private) so a body fetched under a looser
1019/// [`crate::net_policy::WorkspaceRegistryAccess`] can never be served back once the policy
1020/// tightens.
1021pub struct HttpCache {
1022 entries: DashMap<String, CachedResponse>,
1023 /// Running total of `body.len()` across all `entries`, kept in sync by
1024 /// [`HttpCache::store_entry`], [`HttpCache::evict_entries`], and
1025 /// [`HttpCache::clear`] via relative `fetch_add`/`fetch_sub` only —
1026 /// never an absolute `store` after the initial `0`, since that would
1027 /// silently discard any concurrent relative update racing with it. Used
1028 /// to trigger byte-bounded eviction without summing every entry on each
1029 /// check; advisory (see [`MAX_CACHE_BYTES`]), not an exact live count
1030 /// under concurrent access.
1031 total_bytes: AtomicUsize,
1032 /// The shared, unauthenticated transport used by every non-workspace request.
1033 baseline: Transport,
1034 /// Per-`(trusted_origin, tier)` transport pool backing [`Self::get_cached_trusted_origin`]
1035 /// and [`Self::get_cached_pinned`] alike (issue #561/#562, FR-017), keyed by the exact
1036 /// `trusted_origin` prefix string passed to that call, paired with the [`CacheTier`] it was
1037 /// built for. reqwest's redirect policy is fixed per-`Client`, so a distinct client is
1038 /// unavoidable per distinct origin; pooled here so repeated calls against the same
1039 /// `(origin, tier)` reuse one transport (and its connection pool) instead of rebuilding on
1040 /// every call. Deliberately **uncapped** — see [`Self::set_registry_policy`]'s docs for why
1041 /// a capacity cap was considered and dropped for this pool.
1042 trusted_clients: DashMap<(String, CacheTier), Transport>,
1043 /// Live-updatable Cargo workspace-registry policy the `workspace` transport field below
1044 /// and cache-key namespace are derived from. Kept alongside that field (not just read once
1045 /// at construction) so [`Self::cache_key`] and [`Self::set_registry_policy`] both read the
1046 /// same discriminant.
1047 policy: Arc<RegistryAccessPolicy>,
1048 /// The transport for [`Self::get_cached_workspace`], rebuilt in place by
1049 /// [`Self::set_registry_policy`] on every actual policy transition.
1050 workspace: RwLock<Transport>,
1051 /// Test-only counter of how many times [`Self::set_registry_policy`] has actually rebuilt
1052 /// the `workspace` transport field above (as opposed to no-op'ing on an unchanged value) —
1053 /// asserts C4's rebuild-only-on-change behavior directly.
1054 #[cfg(test)]
1055 workspace_rebuilds: AtomicUsize,
1056 /// Live-updatable "no outbound requests" flag (issue #483). Enforced by
1057 /// [`Self::ensure_online`] at all 4 send sites. See [`Self::set_offline`]'s docs for
1058 /// the override this has on `cache_enabled` below.
1059 offline: AtomicBool,
1060 /// Live-updatable entry-map toggle (issue #482): `false` bypasses the entry map
1061 /// entirely (see `get_cached_with_headers_via`). Overridden to effectively `true`
1062 /// whenever `offline` is set — see [`Self::set_offline`]'s docs.
1063 cache_enabled: AtomicBool,
1064}
1065
1066impl HttpCache {
1067 /// Creates a new HTTP cache with default configuration and the default
1068 /// [`crate::net_policy::WorkspaceRegistryAccess`] policy (`PublicOnly`).
1069 ///
1070 /// The cache uses a configurable timeout for all requests and identifies
1071 /// itself with an auto-versioned user agent.
1072 pub fn new() -> Self {
1073 Self::with_policy(Arc::new(RegistryAccessPolicy::default()))
1074 }
1075
1076 /// Creates a new HTTP cache whose [`Self::get_cached_workspace`] requests are governed by
1077 /// `policy`'s live value.
1078 ///
1079 /// A later [`Self::set_registry_policy`] call rebuilds the workspace transport (and its
1080 /// cache-key namespace) in place, so every caller holding this `HttpCache` sees the new
1081 /// policy take effect immediately, with no need to reconstruct the cache.
1082 ///
1083 /// # Examples
1084 ///
1085 /// ```
1086 /// use deps_core::HttpCache;
1087 /// use deps_core::net_policy::RegistryAccessPolicy;
1088 /// use std::sync::Arc;
1089 ///
1090 /// let policy = Arc::new(RegistryAccessPolicy::default());
1091 /// let cache = HttpCache::with_policy(Arc::clone(&policy));
1092 /// assert!(cache.is_empty());
1093 /// ```
1094 pub fn with_policy(policy: Arc<RegistryAccessPolicy>) -> Self {
1095 let workspace = Transport::workspace(&policy);
1096 Self {
1097 entries: DashMap::new(),
1098 total_bytes: AtomicUsize::new(0),
1099 baseline: Transport::baseline(),
1100 trusted_clients: DashMap::new(),
1101 policy,
1102 workspace: RwLock::new(workspace),
1103 #[cfg(test)]
1104 workspace_rebuilds: AtomicUsize::new(0),
1105 offline: AtomicBool::new(false),
1106 cache_enabled: AtomicBool::new(true),
1107 }
1108 }
1109
1110 /// Sets whether outbound network requests are permitted (issue #483).
1111 ///
1112 /// Enforced by `Self::ensure_online` (private) at every one of this module's 4 send sites —
1113 /// effective for every call after this returns. While `mode` is [`NetworkMode::Offline`],
1114 /// this also overrides `cache_enabled` (see [`Self::set_cache_enabled`]) to behave as
1115 /// [`CacheMode::Enabled`] on both the read and write path in `get_cached_with_headers_via`:
1116 /// without this, a warm entry fetched before going offline could never have been stored in
1117 /// the first place if caching was disabled, leaving the offline warm-cache path with
1118 /// nothing to serve — the exact combination `cache.enabled: false` + `network.offline: true`
1119 /// is meant to survive.
1120 pub fn set_offline(&self, mode: NetworkMode) {
1121 self.offline
1122 .store(mode == NetworkMode::Offline, Ordering::Relaxed);
1123 }
1124
1125 /// Returns whether outbound network requests are currently blocked.
1126 #[must_use]
1127 pub fn is_offline(&self) -> bool {
1128 self.offline.load(Ordering::Relaxed)
1129 }
1130
1131 /// Sets whether the entry-map cache is used (issue #482). See [`Self::set_offline`]'s
1132 /// docs for the override `offline` has on this flag while set.
1133 pub fn set_cache_enabled(&self, mode: CacheMode) {
1134 self.cache_enabled
1135 .store(mode == CacheMode::Enabled, Ordering::Relaxed);
1136 }
1137
1138 /// Returns `Err(DepsError::Offline)` when `network.offline` is set, without making any
1139 /// request — the last check before a socket opens at each of this module's 4 send
1140 /// sites, placed beside the existing [`ensure_https`] call at each.
1141 fn ensure_online(&self, url: &str) -> Result<()> {
1142 if self.is_offline() {
1143 return Err(DepsError::Offline {
1144 url: RedactedUrl::new(url),
1145 });
1146 }
1147 Ok(())
1148 }
1149
1150 /// Returns the transport scoped to `trusted_origin`, building and pooling one on first use.
1151 ///
1152 /// The `get` fast path (a shared read lock) serves the common case — a `trusted_origin`
1153 /// already pooled — without ever taking `trusted_clients`' write-capable `entry` lock;
1154 /// `entry().or_insert_with()` only runs on a miss, so two callers racing on the same new
1155 /// origin still only ever build and store one [`Transport`] for it, not one each.
1156 fn transport_for_origin(&self, trusted_origin: &str) -> Transport {
1157 let key = (trusted_origin.to_string(), CacheTier::Baseline);
1158 if let Some(existing) = self.trusted_clients.get(&key) {
1159 return existing.clone();
1160 }
1161
1162 self.trusted_clients
1163 .entry(key)
1164 .or_insert_with(|| Transport::origin_pinned(trusted_origin))
1165 .clone()
1166 }
1167
1168 /// Like [`Self::transport_for_origin`], but for an origin-pinned, connect-address-guarded
1169 /// [`CacheTier::Pinned`] transport (issue #561/#562) — building and pooling one on first
1170 /// use, keyed by `(trusted_origin, CacheTier::Pinned { .. })` so an authenticated and
1171 /// unauthenticated transport for the same origin are pooled separately.
1172 fn transport_for_pinned(&self, trusted_origin: &str, authenticated: bool) -> Transport {
1173 let digest = pinned_digest(trusted_origin, self.policy.get());
1174 let key = (
1175 trusted_origin.to_string(),
1176 CacheTier::Pinned {
1177 digest,
1178 authenticated,
1179 },
1180 );
1181 if let Some(existing) = self.trusted_clients.get(&key) {
1182 return existing.clone();
1183 }
1184
1185 self.trusted_clients
1186 .entry(key)
1187 .or_insert_with(|| {
1188 Transport::origin_pinned_guarded(trusted_origin, &self.policy, authenticated)
1189 })
1190 .clone()
1191 }
1192
1193 /// The prefix marking a workspace-tier cache key, chosen as a control character that can
1194 /// never appear at the start of a URL string this module writes: production code paths
1195 /// only ever write keys derived from this function, which are either the bare URL (starting
1196 /// `https://`, or — test cfgs only — `http://` on loopback) or this prefix followed by a
1197 /// policy digit. No in-process caller other than [`Self::insert_for_bench`] (a
1198 /// `#[doc(hidden)]` test/bench helper that accepts a caller-chosen key) can write an
1199 /// arbitrary key, so a `Baseline`-tier and `WorkspaceDeclared`-tier entry can never collide
1200 /// in production use.
1201 const WS_KEY_PREFIX: char = '\u{1}';
1202
1203 /// The prefix marking a [`CacheTier::Pinned`]-tier cache key (issue #561/#562) — distinct
1204 /// from [`Self::WS_KEY_PREFIX`] so the two namespaces can never collide, chosen as another
1205 /// control character no URL string this module writes can start with.
1206 const PINNED_KEY_PREFIX: char = '\u{2}';
1207
1208 /// Computes the cache-map key for `url` under `tier` — `Cow::Borrowed(url)` for
1209 /// [`CacheTier::Baseline`] (allocation-free, and identical to every entry this cache wrote
1210 /// before this policy-tier split existed), or a policy-digit-prefixed owned key for
1211 /// [`CacheTier::WorkspaceDeclared`] so a policy tightening can never serve a body fetched
1212 /// under a looser policy (C5): the digit comes from `tier`'s own snapshot — the exact same
1213 /// value the paired [`Transport`]'s [`AddrGuard`] enforced for this request, taken together
1214 /// at [`Transport::workspace`] construction time — never a separate live `self.policy` read,
1215 /// which would open a read-skew window between the guard that let a fetch through and the
1216 /// key that fetch's body gets stored under. `self.policy` remains this cache's live handle
1217 /// for [`Self::set_registry_policy`]'s change detection and the write-through source that
1218 /// method snapshots from when rebuilding the workspace transport — never consulted here.
1219 ///
1220 /// Callers must compute this once per request and thread the result through, never
1221 /// recompute mid-request — a policy flip between two recomputations would read and write
1222 /// under different keys for what should be one atomic operation.
1223 ///
1224 /// `auth_id` (FR-014) is folded in only for [`CacheTier::Pinned`] — a separate credential
1225 /// identity from `digest` (see that variant's docs), so a rotated or distinct credential
1226 /// against the same origin never reads back a body fetched under a different one.
1227 /// `authenticated` (the tier's own field) is folded in too, so an authenticated and an
1228 /// unauthenticated fetch against the same origin/`auth_id` can never share an entry either
1229 /// — relevant only in principle today (every current caller correlates `authenticated` and
1230 /// `auth_id.is_some()`), but [`HttpCache::get_cached_pinned_with_headers`] is a public API
1231 /// with no enforced invariant tying the two together, so the key format does not assume one.
1232 ///
1233 /// Both `auth_id` and `authenticated` are encoded as one-char tags immediately preceding a
1234 /// fixed-width field (`0`/`1` before a 16-hex-digit `auth_id` field, `U`/`A` right before
1235 /// `url`), rather than collapsing `auth_id: None` to a `0` sentinel value:
1236 /// [`auth_digest`](crate::secret::auth_digest) is a non-cryptographic hash, so `Some(0)` is
1237 /// a legitimate (if astronomically unlikely) digest a real credential can produce, and
1238 /// `unwrap_or(0)` would make that fetch's cache key identical to an unauthenticated one —
1239 /// letting a cached authenticated body be served to an unauthenticated request, or vice
1240 /// versa (issue #1025). Every hex field is written at a fixed width (`digest`, `auth_id`),
1241 /// so no tag or field can be confused with part of `digest`, the `auth_id` field, or `url`
1242 /// — a variable-width field would let two distinct `(digest, auth_id)` pairs produce the
1243 /// same concatenated string.
1244 fn cache_key<'a>(&self, url: &'a str, tier: CacheTier, auth_id: Option<u64>) -> Cow<'a, str> {
1245 match tier {
1246 CacheTier::Baseline => Cow::Borrowed(url),
1247 CacheTier::WorkspaceDeclared(snapshot) => {
1248 Cow::Owned(format!("{}{}{url}", Self::WS_KEY_PREFIX, snapshot.to_u8()))
1249 }
1250 CacheTier::Pinned {
1251 digest,
1252 authenticated,
1253 } => {
1254 let (auth_tag, id) = match auth_id {
1255 Some(id) => ('1', id),
1256 None => ('0', 0),
1257 };
1258 let authenticated_tag = if authenticated { 'A' } else { 'U' };
1259 Cow::Owned(format!(
1260 "{}{digest:016x}{auth_tag}{id:016x}{authenticated_tag}{url}",
1261 Self::PINNED_KEY_PREFIX,
1262 ))
1263 }
1264 }
1265 }
1266
1267 /// Retrieves data from URL with intelligent caching.
1268 ///
1269 /// On first request, fetches data from the network and caches it.
1270 /// On subsequent requests, performs a conditional GET request using
1271 /// cached ETag or Last-Modified headers. If the server responds with
1272 /// 304 Not Modified, returns the cached data. Otherwise, fetches and
1273 /// caches the new data.
1274 ///
1275 /// If the conditional request fails due to network errors, falls back
1276 /// to the cached data (stale-while-revalidate pattern).
1277 ///
1278 /// # Returns
1279 ///
1280 /// Returns `Bytes` containing the response body. Multiple calls for the
1281 /// same URL return cheap clones (reference counting) without copying data.
1282 ///
1283 /// # Errors
1284 ///
1285 /// Returns `DepsError::RegistryError` if the initial fetch fails and no
1286 /// cached data exists, `DepsError::HttpStatus` if the server returns a
1287 /// non-2xx status on that initial fetch, or `DepsError::ResponseTooLarge`
1288 /// if the response body exceeds the configured size cap.
1289 ///
1290 /// # Examples
1291 ///
1292 /// ```no_run
1293 /// # use deps_core::cache::HttpCache;
1294 /// # async fn example() -> deps_core::error::Result<()> {
1295 /// let cache = HttpCache::new();
1296 /// let data = cache.get_cached("https://example.com/api/data").await?;
1297 /// println!("Fetched {} bytes", data.len());
1298 /// # Ok(())
1299 /// # }
1300 /// ```
1301 pub async fn get_cached(&self, url: &str) -> Result<Bytes> {
1302 self.get_cached_with_headers(url, &[]).await
1303 }
1304
1305 /// Returns the cached body for `url` without making any network request.
1306 ///
1307 /// Unlike `get_cached`'s own stale-while-revalidate fallback (the `Err` arm of
1308 /// `conditional_request_with_headers`'s match in `get_cached_with_headers_via`), this
1309 /// is reachable even when a caller wraps `get_cached` in a short outer timeout: a
1310 /// hung conditional request that never resolves within that timeout gets its whole
1311 /// future cancelled, so `get_cached`'s internal fallback logic never runs and the
1312 /// caller sees a timeout instead of stale data. A caller in that position can call
1313 /// this instead — a synchronous map lookup, no I/O — to serve the last known-good
1314 /// body itself. Returns `None` if `url` has never been successfully cached.
1315 ///
1316 /// The returned body carries no age bound: this bypasses `get_cached`'s own
1317 /// freshness/revalidation logic entirely, so a caller that surfaces this body to
1318 /// the user (e.g. inserting it into a manifest edit) should treat it as
1319 /// arbitrarily stale, not just-expired.
1320 ///
1321 /// Reads the baseline (unprefixed) cache-key namespace only (see `Self::cache_key`, private) — a
1322 /// body fetched via [`Self::get_cached_workspace`] is never visible through this method.
1323 #[must_use]
1324 pub fn peek_cached(&self, url: &str) -> Option<Bytes> {
1325 self.entries.get(url).map(|r| r.body.clone())
1326 }
1327
1328 /// Fetches a URL with additional request headers, using the cache.
1329 ///
1330 /// Works the same as `get_cached` but injects extra headers (e.g., Authorization)
1331 /// into every request. Useful for APIs that require authentication tokens.
1332 ///
1333 /// # Errors
1334 ///
1335 /// Returns `DepsError::RegistryError` if the initial fetch fails and no
1336 /// cached data exists, `DepsError::HttpStatus` if the server returns a
1337 /// non-2xx status on that initial fetch, or `DepsError::ResponseTooLarge`
1338 /// if the response body exceeds the configured size cap.
1339 pub async fn get_cached_with_headers(
1340 &self,
1341 url: &str,
1342 extra_headers: &[(header::HeaderName, &str)],
1343 ) -> Result<Bytes> {
1344 self.get_cached_with_headers_via(url, extra_headers, &self.baseline, None)
1345 .await
1346 }
1347
1348 /// Like [`Self::get_cached`], but additionally stops any redirect hop that does not match
1349 /// `trusted_origin` via [`crate::net_policy::is_trusted_prefix`] (origin equality plus a
1350 /// path-segment-boundary prefix check, e.g. against `https://api.nuget.org/v3/registration5-gz/`).
1351 ///
1352 /// For a caller that already validated the *initial* request URL against a trusted
1353 /// prefix (NuGet's registration-hive paging validates `page.id` this way) and needs
1354 /// that guarantee to hold through any redirect too, not just the first hop —
1355 /// [`Self::get_cached`]'s own policy deliberately does not enforce this, since
1356 /// cross-host redirects are legitimate for the other registry clients sharing this
1357 /// cache; the stricter check is opt-in per call rather than global.
1358 ///
1359 /// The block only surfaces as an error on a cold cache: like [`Self::get_cached`]'s own
1360 /// stale-while-revalidate fallback, a warm entry for `url` still returns the last
1361 /// known-good body (itself already fetched and origin-validated on a prior call)
1362 /// instead of propagating a blocked-redirect `HttpStatus` from a revalidation attempt.
1363 ///
1364 /// # Errors
1365 ///
1366 /// Same as [`Self::get_cached`].
1367 pub async fn get_cached_trusted_origin(
1368 &self,
1369 url: &str,
1370 trusted_origin: &str,
1371 ) -> Result<Bytes> {
1372 let transport = self.transport_for_origin(trusted_origin);
1373 self.get_cached_with_headers_via(url, &[], &transport, None)
1374 .await
1375 }
1376
1377 /// Like [`Self::get_cached_trusted_origin`], but additionally injects `extra_headers`
1378 /// (e.g. an `Authorization` bearer token) into every request — the authenticated
1379 /// counterpart to [`Self::get_cached_with_headers`], composed with the same
1380 /// origin-pinned redirect policy [`Self::get_cached_trusted_origin`] uses.
1381 ///
1382 /// This exists specifically so a header carrying a credential can never survive a
1383 /// cross-origin redirect hop: [`Self::get_cached_with_headers`] attaches
1384 /// `extra_headers` to the *initial* request only and follows reqwest's default
1385 /// (same-scheme, any-host) redirect policy for every hop after that, which is
1386 /// exactly the shape a hostile or misconfigured redirect on the resolved index
1387 /// itself could exploit to exfiltrate a bearer token to an attacker-controlled
1388 /// host. Composing `Self::transport_for_origin`'s (private) pinned-origin transport with header
1389 /// injection closes that by construction — no empirical redirect test is needed
1390 /// to prove the header cannot leak, since the client stops following before a
1391 /// cross-origin hop would ever be sent.
1392 ///
1393 /// # Errors
1394 ///
1395 /// Same as [`Self::get_cached_trusted_origin`].
1396 ///
1397 /// # Examples
1398 ///
1399 /// ```no_run
1400 /// use deps_core::cache::HttpCache;
1401 /// use reqwest::header;
1402 ///
1403 /// # async fn example() -> deps_core::error::Result<()> {
1404 /// let cache = HttpCache::new();
1405 /// let data = cache
1406 /// .get_cached_trusted_origin_with_headers(
1407 /// "https://index.mycorp.dev/se/rd/serde",
1408 /// "https://index.mycorp.dev/",
1409 /// &[(header::AUTHORIZATION, "Bearer secret-token")],
1410 /// )
1411 /// .await?;
1412 /// println!("Fetched {} bytes", data.len());
1413 /// # Ok(())
1414 /// # }
1415 /// ```
1416 pub async fn get_cached_trusted_origin_with_headers(
1417 &self,
1418 url: &str,
1419 trusted_origin: &str,
1420 extra_headers: &[(header::HeaderName, &str)],
1421 ) -> Result<Bytes> {
1422 let transport = self.transport_for_origin(trusted_origin);
1423 self.get_cached_with_headers_via(url, extra_headers, &transport, None)
1424 .await
1425 }
1426
1427 /// Like [`Self::get_cached_trusted_origin_with_headers`], but for an origin-pinned,
1428 /// connect-address-guarded `CacheTier::Pinned` transport (issue #561/#562) instead of the
1429 /// baseline-guarded [`Self::get_cached_trusted_origin`] one — the only sanctioned way to
1430 /// send a credential to a workspace-declared host. Delegates to
1431 /// [`Self::get_cached_pinned_with_headers`] with no extra headers.
1432 ///
1433 /// # Errors
1434 ///
1435 /// Same as [`Self::get_cached`].
1436 pub async fn get_cached_pinned(
1437 &self,
1438 url: &str,
1439 trusted_origin: &str,
1440 authenticated: bool,
1441 auth_id: Option<u64>,
1442 ) -> Result<Bytes> {
1443 self.get_cached_pinned_with_headers(url, trusted_origin, authenticated, auth_id, &[])
1444 .await
1445 }
1446
1447 /// Like [`Self::get_cached_pinned`], but additionally injects `extra_headers` (e.g. an
1448 /// `Authorization` header carrying a credential) into every request — composed with the
1449 /// same origin-pinned, connect-address-guarded transport [`Self::get_cached_pinned`] uses,
1450 /// so a credential header can never survive a cross-origin redirect hop, exactly like
1451 /// [`Self::get_cached_trusted_origin_with_headers`]'s identical closure argument for the
1452 /// baseline-guarded tier.
1453 ///
1454 /// `auth_id` (FR-014) — a caller-computed, salted digest of the credential actually being
1455 /// attached (`None` for an unauthenticated #562 fetch) — is folded into the cache key only,
1456 /// never into `CacheTier`/the transport-pool key, so a rotated or distinct credential
1457 /// against the same origin never reads back a body fetched under a different one.
1458 ///
1459 /// # Errors
1460 ///
1461 /// Same as [`Self::get_cached`].
1462 pub async fn get_cached_pinned_with_headers(
1463 &self,
1464 url: &str,
1465 trusted_origin: &str,
1466 authenticated: bool,
1467 auth_id: Option<u64>,
1468 extra_headers: &[(header::HeaderName, &str)],
1469 ) -> Result<Bytes> {
1470 let transport = self.transport_for_pinned(trusted_origin, authenticated);
1471 self.get_cached_with_headers_via(url, extra_headers, &transport, auth_id)
1472 .await
1473 }
1474
1475 /// Like [`Self::get_cached`], but for Cargo's workspace-declared-registry requests: routes
1476 /// through the workspace transport field, whose guard enforces the live
1477 /// [`crate::net_policy::WorkspaceRegistryAccess`] policy on both the resolved connect-time
1478 /// address (issue #455) and any redirect hop, and keys the entry under a
1479 /// policy-scoped namespace (see `Self::cache_key`, private) distinct from every other method on
1480 /// this cache.
1481 ///
1482 /// This gives the *resolved address* the same policy scrutiny `deps_cargo::config::RegistryIndex::new`
1483 /// already gives the declared URL string at parse time — it does **not**
1484 /// re-check the initial request URL itself: a caller passing an IP-literal `url` whose
1485 /// class the policy would reject connects anyway, since `hyper-util`'s connector parses an
1486 /// IP literal directly and never calls the configured resolver (see
1487 /// `BlockedAddrResolver`'s docs, private). `RegistryIndex::new` is the sole, by-design gate for
1488 /// that residual — every caller of this method already went through it.
1489 ///
1490 /// If an entry is already cached, a revalidation failure — including a guard rejection
1491 /// from a since-rebound or since-tightened-policy address — falls back to serving the
1492 /// cached body, logging only a `tracing::warn!` (pre-existing behavior, unrelated to
1493 /// this method). This is not a bypass — no new connection to the blocked address is
1494 /// made — but it means such a block is invisible to the caller whenever an entry already
1495 /// exists for that URL.
1496 ///
1497 /// # Errors
1498 ///
1499 /// Same as [`Self::get_cached`].
1500 ///
1501 /// # Examples
1502 ///
1503 /// ```no_run
1504 /// use deps_core::HttpCache;
1505 /// use deps_core::net_policy::RegistryAccessPolicy;
1506 /// use std::sync::Arc;
1507 ///
1508 /// # async fn example() -> deps_core::error::Result<()> {
1509 /// let policy = Arc::new(RegistryAccessPolicy::default());
1510 /// let cache = HttpCache::with_policy(policy);
1511 /// let data = cache
1512 /// .get_cached_workspace("https://index.mycorp.dev/se/rd/serde")
1513 /// .await?;
1514 /// println!("Fetched {} bytes", data.len());
1515 /// # Ok(())
1516 /// # }
1517 /// ```
1518 pub async fn get_cached_workspace(&self, url: &str) -> Result<Bytes> {
1519 self.get_cached_workspace_with_headers(url, &[]).await
1520 }
1521
1522 /// Like [`Self::get_cached_workspace`], but additionally forwards `extra_headers` to the
1523 /// underlying request — the headered form needed by a registry client whose
1524 /// workspace-declared fetch requires a non-default header (e.g. `deps-npm`'s abbreviated-
1525 /// packument `Accept` header for an alternate npm registry).
1526 ///
1527 /// # Security
1528 ///
1529 /// `extra_headers` are attached to the **initial** request only. The workspace transport
1530 /// pins by [`crate::net_policy::HostClass`], not origin — unlike
1531 /// [`Self::get_cached_trusted_origin_with_headers`], which exists precisely to close this
1532 /// gap for a caller that needs it — so a cross-origin redirect hop to any other
1533 /// policy-permitted host is followed with `extra_headers` re-sent by reqwest's default
1534 /// redirect policy. **This method must never carry a credential.** Harmless for its
1535 /// current sole caller (a fixed `Accept` header), but directly load-bearing for any
1536 /// future auth-wiring work: reach for [`Self::get_cached_trusted_origin_with_headers`]
1537 /// instead if a header ever needs to stay pinned to one origin.
1538 ///
1539 /// # Errors
1540 ///
1541 /// Same as [`Self::get_cached`].
1542 pub async fn get_cached_workspace_with_headers(
1543 &self,
1544 url: &str,
1545 extra_headers: &[(header::HeaderName, &str)],
1546 ) -> Result<Bytes> {
1547 #[expect(
1548 clippy::expect_used,
1549 reason = "poisoned only if another thread already panicked holding the lock — \
1550 propagate via panic, matching RwLock's poisoning contract"
1551 )]
1552 let transport = self
1553 .workspace
1554 .read()
1555 .expect("workspace transport lock poisoned")
1556 .clone();
1557 self.get_cached_with_headers_via(url, extra_headers, &transport, None)
1558 .await
1559 }
1560
1561 /// Updates the policy governing [`Self::get_cached_workspace`], rebuilding the workspace
1562 /// transport field (and so its cache-key namespace and guard together) when
1563 /// `value` actually differs from the current setting — a no-op call does not rebuild, so a
1564 /// caller that re-applies an unchanged configuration does not pay for a fresh `Client` and
1565 /// its connection pool.
1566 ///
1567 /// Effective for every [`Self::get_cached_workspace`] call after this returns. Note this
1568 /// only gates *future* fetches: an `All -> PublicOnly`/`Off` tightening does not purge
1569 /// already-registered `deps-cargo` alternate-registry clients resolved under the looser
1570 /// policy (pre-existing, documented on [`crate::net_policy::RegistryAccessPolicy::set`]).
1571 ///
1572 /// Unlike that pre-existing gap, every `CacheTier::Pinned` cache entry (issue #561/#562)
1573 /// **is** purged on every actual policy transition, along with every pinned-tier pooled
1574 /// `Transport` — substantially narrowing the `All -> PublicOnly -> All` round-trip hole for
1575 /// credential-carrying entries specifically (NFR-004): re-namespacing alone (as the
1576 /// pre-existing workspace-tier digit prefix does) would leave an old-era authenticated
1577 /// body reachable once the policy round-trips back to a value whose digest happens to
1578 /// collide again. Not an absolute close: a fetch already in flight when the transition
1579 /// happens can still land its response and re-insert an old-era key after the purge —
1580 /// harmless (readable only under the era it was legitimately fetched in), just not
1581 /// prevented by this purge alone.
1582 pub fn set_registry_policy(&self, value: WorkspaceRegistryAccess) {
1583 if self.policy.get() == value {
1584 return;
1585 }
1586 self.policy.set(value);
1587 let rebuilt = Transport::workspace(&self.policy);
1588 #[expect(
1589 clippy::expect_used,
1590 reason = "poisoned only if another thread panicked while holding the lock — see \
1591 the matching justification on get_cached_workspace_with_headers above"
1592 )]
1593 {
1594 *self
1595 .workspace
1596 .write()
1597 .expect("workspace transport lock poisoned") = rebuilt;
1598 }
1599 #[cfg(test)]
1600 self.workspace_rebuilds.fetch_add(1, Ordering::Relaxed);
1601
1602 let mut freed_bytes = 0usize;
1603 self.entries.retain(|k, v| {
1604 let keep = !k.starts_with(Self::PINNED_KEY_PREFIX);
1605 if !keep {
1606 freed_bytes += v.body.len();
1607 }
1608 keep
1609 });
1610 self.total_bytes.fetch_sub(freed_bytes, Ordering::Relaxed);
1611 self.trusted_clients
1612 .retain(|(_, tier), _| !matches!(tier, CacheTier::Pinned { .. }));
1613 }
1614
1615 /// `auth_id` (FR-014) is meaningful only under [`CacheTier::Pinned`] — every other tier
1616 /// ignores it (see [`Self::cache_key`]'s docs).
1617 #[tracing::instrument(
1618 skip(self, extra_headers, transport, auth_id),
1619 fields(url = %RedactedUrl::new(url), cache = tracing::field::Empty)
1620 )]
1621 async fn get_cached_with_headers_via(
1622 &self,
1623 url: &str,
1624 extra_headers: &[(header::HeaderName, &str)],
1625 transport: &Transport,
1626 auth_id: Option<u64>,
1627 ) -> Result<Bytes> {
1628 if self.entries.len() >= MAX_CACHE_ENTRIES
1629 || self.total_bytes.load(Ordering::Relaxed) >= MAX_CACHE_BYTES
1630 {
1631 self.evict_entries();
1632 }
1633
1634 let offline = self.is_offline();
1635 // `offline` forces `cache_enabled` true (S1 fix): otherwise `cache.enabled: false`
1636 // would leave nothing to serve offline for a URL fetched while caching was disabled.
1637 // See `Self::set_offline`'s docs.
1638 let cache_enabled = self.cache_enabled.load(Ordering::Relaxed) || offline;
1639
1640 // Computed once and threaded through every call — recomputing mid-request could read
1641 // and write under different keys (see `Self::cache_key`'s docs).
1642 let cache_key = self.cache_key(url, transport.tier, auth_id);
1643
1644 if !cache_enabled {
1645 // Explicit, not `Empty`: an empty `cache` field would be indistinguishable from
1646 // broken instrumentation (#756 S3) on a path that's neither a hit nor a miss.
1647 tracing::Span::current().record("cache", "disabled");
1648 return self
1649 .transport_only_via(url, extra_headers, BodyLimit::DEFAULT, &transport.client)
1650 .await;
1651 }
1652
1653 // Clone+drop the Ref immediately: holding it across `.await` can deadlock a concurrent
1654 // task needing write access to the same shard.
1655 if let Some(cached) = self.entries.get(cache_key.as_ref()).map(|r| r.clone()) {
1656 if offline {
1657 // Skips the conditional-request attempt: `ensure_online` would block it and
1658 // fall back to this same body anyway, just with a spurious warn + wasted
1659 // allocation. The only branch below that is a genuine zero-network "hit"
1660 // (#756 S3); every other branch still issued a request.
1661 tracing::Span::current().record("cache", "hit");
1662 return Ok(cached.body);
1663 }
1664 match self
1665 .conditional_request_with_headers(
1666 url,
1667 &cached,
1668 extra_headers,
1669 &transport.client,
1670 &cache_key,
1671 )
1672 .await
1673 {
1674 // 304: a round trip happened but no body was re-transferred — distinct from
1675 // `hit` (no request) and `refreshed` (full re-fetch).
1676 Ok(None) => {
1677 tracing::Span::current().record("cache", "revalidated");
1678 return Ok(cached.body);
1679 }
1680 // Stale entry cost a full re-fetch, same as a miss — must not report as a hit.
1681 Ok(Some(new_body)) => {
1682 tracing::Span::current().record("cache", "refreshed");
1683 return Ok(new_body);
1684 }
1685 Err(e) => {
1686 debug_assert!(
1687 !matches!(
1688 &e,
1689 DepsError::RateLimited {
1690 source_status: None,
1691 ..
1692 }
1693 ),
1694 "a RateLimited reaching this eviction guard must always carry \
1695 source_status — only http_status_error's confirmed-evidence branch \
1696 produces RateLimited here; None would silently bypass FR-015/NFR-004 \
1697 eviction on a genuine 401/403 credential-revocation signal"
1698 );
1699 // FR-015/NFR-004: a 401/403 revalidation against an *authenticated*
1700 // pinned-tier entry must evict rather than serve the possibly-revoked
1701 // credential's last-known-good body — every other tier keeps today's
1702 // stale-while-revalidate fallback unchanged.
1703 //
1704 // `RateLimited { source_status: Some(401 | 403), .. }` is included here
1705 // (#1295 critic C1): `e` is always this match's own
1706 // `conditional_request_with_headers`'s `Err`, whose only source of that
1707 // variant is `http_status_error`'s confirmed-evidence branch — without this
1708 // arm, a confirmed-evidence 403 would silently bypass eviction and keep
1709 // serving the possibly-revoked credential's stale body. Deliberately
1710 // narrowed to `source_status` 401/403 only, not a bare `RateLimited { .. }`
1711 // (critic N1 regression fix): `confirmed_rate_limit_exhaustion` also
1712 // classifies a 429 this way, but a 429 is mere throttling, not a
1713 // credential-revocation signal — NFR-004's scope is "401/403", and evicting
1714 // on 429 would drop a still-good cached body the client can no longer
1715 // re-fetch until the throttle clears, purely because of a signal unrelated
1716 // to the credential's validity.
1717 if transport.tier.is_authenticated()
1718 && matches!(
1719 &e,
1720 DepsError::HttpStatus {
1721 status: 401 | 403,
1722 ..
1723 } | DepsError::RateLimited {
1724 source_status: Some(401 | 403),
1725 ..
1726 }
1727 )
1728 {
1729 if let Some((_, old)) = self.entries.remove(cache_key.as_ref()) {
1730 self.total_bytes
1731 .fetch_sub(old.body.len(), Ordering::Relaxed);
1732 }
1733 tracing::Span::current().record("cache", "evicted");
1734 // #756 round 2 S1: never interpolate `e` — its `Display`/`Debug` embed
1735 // the raw unredacted `url`, defeating this span's `RedactedUrl` field.
1736 // `safe_tracing_summary` extracts only the safe status+cause instead.
1737 let (status, cause) = e.safe_tracing_summary();
1738 tracing::warn!(
1739 status = ?status,
1740 cause,
1741 "evicting authenticated cache entry after revalidation failure"
1742 );
1743 return Err(e);
1744 }
1745 tracing::Span::current().record("cache", "stale-fallback");
1746 // Same rationale as above: no `e` interpolation.
1747 let (status, cause) = e.safe_tracing_summary();
1748 tracing::warn!(
1749 status = ?status,
1750 cause,
1751 "conditional request failed, using cache"
1752 );
1753 return Ok(cached.body);
1754 }
1755 }
1756 }
1757
1758 tracing::Span::current().record("cache", "miss");
1759 self.fetch_and_store_with_headers(url, extra_headers, &transport.client, &cache_key)
1760 .await
1761 }
1762
1763 /// Performs conditional HTTP request using cached validation headers.
1764 ///
1765 /// Sends `If-None-Match` (ETag) and/or `If-Modified-Since` headers
1766 /// to check if the cached content is still valid.
1767 ///
1768 /// # Returns
1769 ///
1770 /// - `Ok(Some(Bytes))` - Server returned 200 OK with new content
1771 /// - `Ok(None)` - Server returned 304 Not Modified (cache is valid)
1772 /// - `Err(_)` - Network or HTTP error occurred
1773 async fn conditional_request_with_headers(
1774 &self,
1775 url: &str,
1776 cached: &CachedResponse,
1777 extra_headers: &[(header::HeaderName, &str)],
1778 client: &Client,
1779 cache_key: &str,
1780 ) -> Result<Option<Bytes>> {
1781 self.ensure_online(url)?;
1782 ensure_https(url)?;
1783 let mut request = client.get(url);
1784
1785 for (name, value) in extra_headers {
1786 request = request.header(name, *value);
1787 }
1788 if let Some(etag) = &cached.etag {
1789 request = request.header(header::IF_NONE_MATCH, etag);
1790 }
1791 if let Some(last_modified) = &cached.last_modified {
1792 request = request.header(header::IF_MODIFIED_SINCE, last_modified);
1793 }
1794
1795 let response = request.send().await.map_err(|e| DepsError::RegistryError {
1796 package: RedactedUrl::new(url),
1797 source: e.into(),
1798 })?;
1799
1800 if response.status() == StatusCode::NOT_MODIFIED {
1801 return Ok(None);
1802 }
1803
1804 if !response.status().is_success() {
1805 return Err(http_status_error(
1806 url,
1807 response.status(),
1808 response.headers(),
1809 ));
1810 }
1811
1812 let etag = response
1813 .headers()
1814 .get(header::ETAG)
1815 .and_then(|v| v.to_str().ok())
1816 .map(String::from);
1817 let last_modified = response
1818 .headers()
1819 .get(header::LAST_MODIFIED)
1820 .and_then(|v| v.to_str().ok())
1821 .map(String::from);
1822 let body = read_body_capped(url, response, BodyLimit::DEFAULT).await?;
1823
1824 self.store_entry(
1825 cache_key.to_string(),
1826 CachedResponse {
1827 body: body.clone(),
1828 etag,
1829 last_modified,
1830 fetched_at: Instant::now(),
1831 },
1832 );
1833
1834 Ok(Some(body))
1835 }
1836
1837 /// Fetches a fresh response from the network and stores it in the cache.
1838 ///
1839 /// This method bypasses the cache and always makes a network request.
1840 /// The response is stored with its ETag and Last-Modified headers for
1841 /// future conditional requests.
1842 ///
1843 /// # Errors
1844 ///
1845 /// Returns `DepsError::HttpStatus` if the server returns a non-2xx status code
1846 /// (or `DepsError::RateLimited` when that status is a 403 with confirmed
1847 /// `X-RateLimit-Remaining: 0` evidence — see [`http_status_error`], #1295),
1848 /// `DepsError::RegistryError` if the network request fails, or
1849 /// `DepsError::ResponseTooLarge` if the response body exceeds the
1850 /// configured size cap.
1851 async fn fetch_and_store_with_headers(
1852 &self,
1853 url: &str,
1854 extra_headers: &[(header::HeaderName, &str)],
1855 client: &Client,
1856 cache_key: &str,
1857 ) -> Result<Bytes> {
1858 self.ensure_online(url)?;
1859 ensure_https(url)?;
1860 // #756 security follow-up (S-A): `RedactedUrl`, not the raw `url` — this is the
1861 // direct callee of the now-hardened `get_cached_with_headers_via`, and at `debug`
1862 // level, which is this project's own continuous-improvement convention
1863 // (`RUST_LOG=debug`).
1864 tracing::debug!(
1865 extra_headers = extra_headers.len(),
1866 "fetching fresh: {}",
1867 RedactedUrl::new(url)
1868 );
1869
1870 let mut request = client.get(url);
1871 for (name, value) in extra_headers {
1872 request = request.header(name, *value);
1873 }
1874
1875 let response = request.send().await.map_err(|e| DepsError::RegistryError {
1876 package: RedactedUrl::new(url),
1877 source: e.into(),
1878 })?;
1879
1880 if !response.status().is_success() {
1881 return Err(http_status_error(
1882 url,
1883 response.status(),
1884 response.headers(),
1885 ));
1886 }
1887
1888 let etag = response
1889 .headers()
1890 .get(header::ETAG)
1891 .and_then(|v| v.to_str().ok())
1892 .map(String::from);
1893 let last_modified = response
1894 .headers()
1895 .get(header::LAST_MODIFIED)
1896 .and_then(|v| v.to_str().ok())
1897 .map(String::from);
1898 let body = read_body_capped(url, response, BodyLimit::DEFAULT).await?;
1899
1900 self.store_entry(
1901 cache_key.to_string(),
1902 CachedResponse {
1903 body: body.clone(),
1904 etag,
1905 last_modified,
1906 fetched_at: Instant::now(),
1907 },
1908 );
1909
1910 Ok(body)
1911 }
1912
1913 /// POSTs `body` as JSON and returns the response body.
1914 ///
1915 /// Deliberately does not cache: the OSV batch endpoint is a POST with a
1916 /// request-body-dependent response and sends no `ETag`/`Last-Modified`
1917 /// validators, so entry-map caching would be meaningless here — every
1918 /// call reuses the client, HTTPS enforcement, size cap, and timeout
1919 /// (via `read_body_capped`) without touching the entry map or
1920 /// [`Self::total_bytes`].
1921 ///
1922 /// # Errors
1923 ///
1924 /// Returns `DepsError::HttpStatus` if the server returns a non-2xx
1925 /// status, `DepsError::RegistryError` if the request fails, or
1926 /// `DepsError::ResponseTooLarge` if the response body exceeds the
1927 /// configured size cap.
1928 pub async fn post_json<T: Serialize + Sync + ?Sized>(
1929 &self,
1930 url: &str,
1931 body: &T,
1932 ) -> Result<Bytes> {
1933 self.post_json_via(url, body, BodyLimit::DEFAULT, &self.baseline.client)
1934 .await
1935 }
1936
1937 /// Same as [`Self::post_json`], but additionally takes an explicit [`BodyLimit`]
1938 /// (rather than the hardcoded [`BodyLimit::DEFAULT`]) and confines every redirect hop to
1939 /// `trusted_origin` via [`crate::net_policy::is_trusted_prefix`] — the same guarantee
1940 /// [`Self::get_transport_only_with_headers_limited_trusted_origin`] gives the GET path.
1941 ///
1942 /// [`Self::post_json`] itself has neither of these: it sends through
1943 /// `Self::baseline`'s client (the generic, non-origin-pinned redirect policy) and
1944 /// always caps the response at [`BodyLimit::DEFAULT`] (32 MiB) regardless of how much
1945 /// smaller a caller's own payloads actually are — a gap for a caller (deps.dev's GOSSIP
1946 /// batch endpoint) whose response is expected to be a few KB and whose request body
1947 /// carries every declared dependency's name, where an unconfined redirect is a bigger
1948 /// concern than for `post_json`'s existing callers.
1949 ///
1950 /// # Errors
1951 ///
1952 /// Same as [`Self::post_json`].
1953 pub async fn post_json_limited_trusted_origin<T: Serialize + Sync + ?Sized>(
1954 &self,
1955 url: &str,
1956 body: &T,
1957 limit: BodyLimit,
1958 trusted_origin: &str,
1959 ) -> Result<Bytes> {
1960 let transport = self.transport_for_origin(trusted_origin);
1961 self.post_json_via(url, body, limit, &transport.client)
1962 .await
1963 }
1964
1965 /// Shared POST body for [`Self::post_json`] and [`Self::post_json_limited_trusted_origin`]
1966 /// — mirrors [`Self::transport_only_via`]'s role on the GET path exactly: both public POST
1967 /// methods differ only in which `client` (origin-pinned or not) and [`BodyLimit`] they pass
1968 /// in, so this is the single place that sends the request, checks the status, and reads the
1969 /// capped body.
1970 ///
1971 /// # Errors
1972 ///
1973 /// Same as [`Self::post_json`].
1974 #[tracing::instrument(
1975 skip(self, body, client),
1976 fields(url = %RedactedUrl::new(url))
1977 )]
1978 async fn post_json_via<T: Serialize + Sync + ?Sized>(
1979 &self,
1980 url: &str,
1981 body: &T,
1982 limit: BodyLimit,
1983 client: &Client,
1984 ) -> Result<Bytes> {
1985 self.ensure_online(url)?;
1986 ensure_https(url)?;
1987
1988 let response =
1989 client
1990 .post(url)
1991 .json(body)
1992 .send()
1993 .await
1994 .map_err(|e| DepsError::RegistryError {
1995 package: RedactedUrl::new(url),
1996 source: e.into(),
1997 })?;
1998
1999 if !response.status().is_success() {
2000 return Err(http_status_error(
2001 url,
2002 response.status(),
2003 response.headers(),
2004 ));
2005 }
2006
2007 read_body_capped(url, response, limit).await
2008 }
2009
2010 /// GETs `url` and returns the response body, bypassing the entry-map
2011 /// cache entirely — reuses the client, HTTPS enforcement, size cap, and
2012 /// timeout, exactly like [`Self::post_json`], but for a plain GET.
2013 ///
2014 /// For a caller whose own values are already cached elsewhere (e.g.
2015 /// `OsvClient`'s record cache, validated by a `modified` timestamp
2016 /// rather than `ETag`/`Last-Modified`): reusing [`Self::get_cached`]
2017 /// there would double-cache every fetched record in *this* cache's byte
2018 /// budget too, competing with registry responses for it even though
2019 /// nothing here ever reads that cached copy back.
2020 ///
2021 /// # Errors
2022 ///
2023 /// Returns `DepsError::HttpStatus` if the server returns a non-2xx
2024 /// status, `DepsError::RegistryError` if the request fails, or
2025 /// `DepsError::ResponseTooLarge` if the response body exceeds the
2026 /// configured size cap.
2027 pub async fn get_transport_only(&self, url: &str) -> Result<Bytes> {
2028 self.get_transport_only_with_headers(url, &[]).await
2029 }
2030
2031 /// Same as [`Self::get_transport_only`], but injects extra request headers (e.g. a
2032 /// content-negotiating `Accept`) — mirrors how [`Self::get_cached_with_headers`] relates
2033 /// to [`Self::get_cached`].
2034 ///
2035 /// # Errors
2036 ///
2037 /// Returns `DepsError::HttpStatus` if the server returns a non-2xx
2038 /// status, `DepsError::RegistryError` if the request fails, or
2039 /// `DepsError::ResponseTooLarge` if the response body exceeds the
2040 /// configured size cap.
2041 pub async fn get_transport_only_with_headers(
2042 &self,
2043 url: &str,
2044 extra_headers: &[(header::HeaderName, &str)],
2045 ) -> Result<Bytes> {
2046 self.get_transport_only_with_headers_limited(url, extra_headers, BodyLimit::DEFAULT)
2047 .await
2048 }
2049
2050 /// Same as [`Self::get_transport_only_with_headers`], but takes an explicit
2051 /// [`BodyLimit`] instead of the [`BodyLimit::DEFAULT`] (`MAX_RESPONSE_BYTES`) cap.
2052 ///
2053 /// For a caller whose response is legitimately larger than every other registry
2054 /// payload — e.g. `deps-pypi`'s full Simple API project index — without weakening
2055 /// the cap every other caller of this cache relies on. `BodyLimit` clamps at
2056 /// construction, so this can never be widened past `ABSOLUTE_MAX_RESPONSE_BYTES`
2057 /// regardless of what the caller passes in.
2058 ///
2059 /// # Errors
2060 ///
2061 /// Same as [`Self::get_transport_only_with_headers`].
2062 pub async fn get_transport_only_with_headers_limited(
2063 &self,
2064 url: &str,
2065 extra_headers: &[(header::HeaderName, &str)],
2066 limit: BodyLimit,
2067 ) -> Result<Bytes> {
2068 self.transport_only_via(url, extra_headers, limit, &self.baseline.client)
2069 .await
2070 }
2071
2072 /// Same as [`Self::get_transport_only_with_headers_limited`], but additionally
2073 /// stops any redirect hop that does not match `trusted_origin` via
2074 /// [`crate::net_policy::is_trusted_prefix`]
2075 /// (see [`Self::get_cached_trusted_origin`], which applies the identical policy
2076 /// to the entry-cached path). For a caller carrying a materially larger
2077 /// [`BodyLimit`] than [`BodyLimit::DEFAULT`] — the bigger the budget, the more
2078 /// worth pinning the origin an arbitrary cross-host redirect could point it at.
2079 ///
2080 /// # Errors
2081 ///
2082 /// Same as [`Self::get_transport_only_with_headers_limited`].
2083 pub async fn get_transport_only_with_headers_limited_trusted_origin(
2084 &self,
2085 url: &str,
2086 extra_headers: &[(header::HeaderName, &str)],
2087 limit: BodyLimit,
2088 trusted_origin: &str,
2089 ) -> Result<Bytes> {
2090 let transport = self.transport_for_origin(trusted_origin);
2091 self.transport_only_via(url, extra_headers, limit, &transport.client)
2092 .await
2093 }
2094
2095 /// Note: "transport" here means "bypasses the entry-map cache" (see this method's
2096 /// callers' docs) — a different axis from the [`Transport`] type, which pairs a `Client`
2097 /// with its [`CacheTier`]. The name predates that type and is kept as-is to avoid
2098 /// churning every `get_transport_only*` call site for a naming collision that causes no
2099 /// actual ambiguity at the call sites themselves.
2100 #[tracing::instrument(
2101 skip(self, extra_headers, limit, client),
2102 fields(url = %RedactedUrl::new(url))
2103 )]
2104 async fn transport_only_via(
2105 &self,
2106 url: &str,
2107 extra_headers: &[(header::HeaderName, &str)],
2108 limit: BodyLimit,
2109 client: &Client,
2110 ) -> Result<Bytes> {
2111 self.ensure_online(url)?;
2112 ensure_https(url)?;
2113
2114 let mut request = client.get(url);
2115 for (name, value) in extra_headers {
2116 request = request.header(name, *value);
2117 }
2118
2119 let response = request.send().await.map_err(|e| DepsError::RegistryError {
2120 package: RedactedUrl::new(url),
2121 source: e.into(),
2122 })?;
2123
2124 if !response.status().is_success() {
2125 return Err(http_status_error(
2126 url,
2127 response.status(),
2128 response.headers(),
2129 ));
2130 }
2131
2132 read_body_capped(url, response, limit).await
2133 }
2134
2135 /// Inserts (or replaces) a cache entry, keeping [`Self::total_bytes`] in sync.
2136 ///
2137 /// `DashMap::insert` returns the replaced value, if any, so the byte
2138 /// delta is computed from a single insert rather than a separate
2139 /// lookup-then-insert (which would race with concurrent writers).
2140 ///
2141 /// A body larger than [`MAX_CACHEABLE_ENTRY_BYTES`] is not inserted at
2142 /// all (the caller already has it from the network response; only
2143 /// caching is skipped), and any stale entry previously cached for this
2144 /// URL is dropped rather than left to serve increasingly outdated data.
2145 fn store_entry(&self, url: String, response: CachedResponse) {
2146 let new_len = response.body.len();
2147
2148 if new_len > MAX_CACHEABLE_ENTRY_BYTES {
2149 if let Some((_, old)) = self.entries.remove(&url) {
2150 self.total_bytes
2151 .fetch_sub(old.body.len(), Ordering::Relaxed);
2152 }
2153 return;
2154 }
2155
2156 let old_len = self
2157 .entries
2158 .insert(url, response)
2159 .map_or(0, |old| old.body.len());
2160 self.total_bytes.fetch_add(new_len, Ordering::Relaxed);
2161 self.total_bytes.fetch_sub(old_len, Ordering::Relaxed);
2162 }
2163
2164 /// Clears all cached entries.
2165 ///
2166 /// This removes all cached responses, forcing the next request for
2167 /// any URL to fetch fresh data from the network.
2168 pub fn clear(&self) {
2169 self.entries.clear();
2170 self.total_bytes.store(0, Ordering::Relaxed);
2171 }
2172
2173 /// Returns the number of cached entries.
2174 pub fn len(&self) -> usize {
2175 self.entries.len()
2176 }
2177
2178 /// Returns `true` if the cache contains no entries.
2179 pub fn is_empty(&self) -> bool {
2180 self.entries.is_empty()
2181 }
2182
2183 /// Returns the total bytes retained across all cached response bodies.
2184 pub fn total_bytes(&self) -> usize {
2185 self.total_bytes.load(Ordering::Relaxed)
2186 }
2187
2188 /// Evicts the oldest cache entries when either capacity limit is reached.
2189 ///
2190 /// When the entry count is at or over `MAX_CACHE_ENTRIES`, evicts at
2191 /// least `CACHE_EVICTION_PERCENTAGE`% of entries (by count). Note this
2192 /// is a *fix*, not a preserved behavior: the original count-only
2193 /// eviction built its bounded min-heap with an inverted comparison
2194 /// (`peek()` returns the oldest entry, but the old code treated it as
2195 /// the newest-of-the-oldest-so-far and only replaced it when a *newer*
2196 /// candidate came along that was still older than it — backwards), so
2197 /// it evicted roughly the first `target_removals` entries in DashMap
2198 /// hash-iteration order, not the oldest ones. This version evicts
2199 /// genuinely oldest-first.
2200 ///
2201 /// Independently, if the tracked byte total is over [`MAX_CACHE_BYTES`]
2202 /// — which can happen with far fewer than `MAX_CACHE_ENTRIES` entries if
2203 /// a few responses are large — eviction keeps removing the next-oldest
2204 /// entries until the byte budget is satisfied too. A cache that is over
2205 /// the byte budget but well under the entry-count threshold only evicts
2206 /// as many entries as the byte budget requires, not a fixed count-based
2207 /// batch.
2208 ///
2209 /// Builds a min-heap over all entry keys by `fetched_at` (O(N)), then
2210 /// pops the oldest one at a time (O(R log N) for R removals) — unlike a
2211 /// heap bounded to a fixed top-K, the removal count isn't known upfront
2212 /// since it depends on the byte budget as well as the count target.
2213 ///
2214 /// Every byte-count adjustment here is a relative `fetch_sub` applied to
2215 /// exactly the entry [`DashMap::remove`] actually returned — never a
2216 /// snapshot-then-absolute-`store` of a locally computed total. The
2217 /// latter would silently discard any [`Self::store_entry`] delta that
2218 /// lands between this method's start and its end (lost-update race), and
2219 /// under adversarial timing could even underflow `total_bytes` to
2220 /// `usize::MAX`, permanently wedging every future request into
2221 /// evicting the entire cache. Reading `total_bytes` fresh on every loop
2222 /// iteration (rather than maintaining a local mirror) keeps this
2223 /// correct under concurrent `evict_entries`/`store_entry` calls: two
2224 /// callers can race to remove the same key — the second `remove` simply
2225 /// returns `None` and is a no-op, not a double-subtraction.
2226 fn evict_entries(&self) {
2227 use std::cmp::Reverse;
2228 use std::collections::BinaryHeap;
2229
2230 let count_target_removals = if self.entries.len() >= MAX_CACHE_ENTRIES {
2231 (MAX_CACHE_ENTRIES / CACHE_EVICTION_PERCENTAGE).max(1)
2232 } else {
2233 0
2234 };
2235
2236 let mut oldest: BinaryHeap<Reverse<(Instant, String)>> = self
2237 .entries
2238 .iter()
2239 .map(|entry| Reverse((entry.value().fetched_at, entry.key().clone())))
2240 .collect();
2241
2242 let mut removed = 0usize;
2243
2244 while removed < count_target_removals
2245 || self.total_bytes.load(Ordering::Relaxed) > MAX_CACHE_BYTES
2246 {
2247 let Some(Reverse((_, url))) = oldest.pop() else {
2248 break;
2249 };
2250 if let Some((_, old)) = self.entries.remove(&url) {
2251 self.total_bytes
2252 .fetch_sub(old.body.len(), Ordering::Relaxed);
2253 }
2254 removed += 1;
2255 }
2256
2257 tracing::debug!(
2258 "evicted {removed} cache entries ({} bytes remaining)",
2259 self.total_bytes.load(Ordering::Relaxed)
2260 );
2261 }
2262
2263 /// Benchmark-only helper: Direct cache lookup without network requests.
2264 #[cfg(feature = "test-util")]
2265 #[doc(hidden)]
2266 pub fn get_for_bench(&self, url: &str) -> Option<Bytes> {
2267 self.entries.get(url).map(|entry| entry.body.clone())
2268 }
2269
2270 /// Benchmark-only helper: Direct cache insertion.
2271 #[cfg(feature = "test-util")]
2272 #[doc(hidden)]
2273 pub fn insert_for_bench(&self, url: String, response: CachedResponse) {
2274 self.store_entry(url, response);
2275 }
2276}
2277
2278impl Default for HttpCache {
2279 fn default() -> Self {
2280 Self::new()
2281 }
2282}
2283
2284#[cfg(test)]
2285mod tests {
2286 use super::*;
2287
2288 use std::assert_matches;
2289
2290 // Guards the non-loopback path of `ensure_https`: every other test in this
2291 // module reaches it only through loopback `mockito` URLs, so without this
2292 // test the "reject any other HTTP host" branch would never run.
2293 #[test]
2294 fn test_ensure_https_rejects_non_loopback_http() {
2295 assert!(ensure_https("http://example.com").is_err());
2296 }
2297
2298 /// #767 M2: `ensure_https`'s `CacheError` used to bake the raw rejected URL in
2299 /// verbatim, the same message-level leak class fixed in `deps-cargo::sparse`'s
2300 /// fail-closed log on the same `window/showMessage` path.
2301 #[test]
2302 fn test_ensure_https_rejection_message_redacts_query_string() {
2303 let err = ensure_https("http://example.com/pkg?token=super-secret-value").unwrap_err();
2304 assert!(
2305 !err.to_string().contains("super-secret-value"),
2306 "err: {err}"
2307 );
2308 }
2309
2310 // `http://example.com` alone would still pass under a regressed, substring-based
2311 // `is_loopback_host` (e.g. `url.contains("localhost")`) — these hosts embed a
2312 // loopback token without actually being loopback, and must still be rejected.
2313 #[test]
2314 fn test_ensure_https_rejects_hosts_resembling_loopback() {
2315 assert!(ensure_https("http://localhost.evil.com/").is_err());
2316 assert!(ensure_https("http://127.0.0.1.evil.com/").is_err());
2317 assert!(ensure_https("http://evil.com/?cb=127.0.0.1").is_err());
2318 }
2319
2320 // The sole exerciser of `is_loopback_host`'s bracketed-IPv6 branch
2321 // (`strip_prefix('[')`/`split(']')`) — every other test/mockito URL in this
2322 // module uses `127.0.0.1`.
2323 #[test]
2324 fn test_ensure_https_accepts_bracketed_ipv6_loopback() {
2325 assert!(ensure_https("http://[::1]:1234/x").is_ok());
2326 }
2327
2328 // #1562: a naive `:`-split misread userinfo as the host boundary, so a loopback-looking
2329 // userinfo let a plain-HTTP request to a public host slip past the HTTPS requirement.
2330 #[test]
2331 fn test_ensure_https_rejects_userinfo_spoofed_loopback() {
2332 assert!(ensure_https("http://localhost:80@evil.com/x").is_err());
2333 assert!(ensure_https("http://localhost:@evil.com/").is_err());
2334 }
2335
2336 // reqwest's `Attempt` has no public constructor, so the redirect closure can't be
2337 // unit-tested directly; this exercises the pure detection logic it delegates to (mockito is
2338 // http-only, so an actual https->http redirect isn't testable end-to-end here).
2339 #[test]
2340 fn test_is_https_downgrade() {
2341 let https = Url::parse("https://example.com/a").unwrap();
2342 let http = Url::parse("http://example.com/a").unwrap();
2343
2344 assert!(is_https_downgrade(&https, &http));
2345 assert!(!is_https_downgrade(&http, &https));
2346 assert!(!is_https_downgrade(&https, &https));
2347 assert!(!is_https_downgrade(&http, &http));
2348 }
2349
2350 #[test]
2351 fn test_hop_targets_blocked_host_blocks_cloud_metadata() {
2352 let url = Url::parse("https://169.254.169.254/latest/meta-data/").unwrap();
2353 assert!(hop_targets_blocked_host(&url));
2354 }
2355
2356 #[test]
2357 fn test_hop_targets_blocked_host_allows_global() {
2358 let url = Url::parse("https://index.crates.io/").unwrap();
2359 assert!(!hop_targets_blocked_host(&url));
2360 }
2361
2362 // The loopback carve-out (identical to `ensure_https`'s) must still exempt
2363 // loopback hops under `cfg(test)`, or every mockito redirect chain in this
2364 // module's own tests would start failing.
2365 #[test]
2366 fn test_hop_targets_blocked_host_exempts_loopback_under_test_cfg() {
2367 let url = Url::parse("http://127.0.0.1:1234/api/target").unwrap();
2368 assert!(!hop_targets_blocked_host(&url));
2369 }
2370
2371 // Issue #449: the connect-time resolver guard's pure classification core, unit-tested
2372 // directly rather than through `tokio::net::lookup_host` — no live DNS/network needed.
2373 #[test]
2374 fn test_validate_resolved_addrs_blocks_cloud_metadata() {
2375 let addrs = vec!["169.254.169.254:0".parse().unwrap()];
2376 assert_matches!(
2377 validate_resolved_addrs("evil.example", addrs, AddrGuard::Baseline),
2378 Err(ResolveGuardError::Blocked { .. })
2379 );
2380 }
2381
2382 // FR-003: an attacker's public A record alongside a blocked one must not keep the probe
2383 // alive — the whole resolution is rejected, not filtered down to the public address.
2384 #[test]
2385 fn test_validate_resolved_addrs_blocks_when_any_address_is_blocked() {
2386 let addrs = vec![
2387 "1.1.1.1:0".parse().unwrap(),
2388 "169.254.169.254:0".parse().unwrap(),
2389 ];
2390 assert_matches!(
2391 validate_resolved_addrs("evil.example", addrs, AddrGuard::Baseline),
2392 Err(ResolveGuardError::Blocked { .. })
2393 );
2394 }
2395
2396 #[test]
2397 fn test_validate_resolved_addrs_allows_global() {
2398 let addrs = vec!["1.1.1.1:0".parse().unwrap()];
2399 assert_eq!(
2400 validate_resolved_addrs("index.crates.io", addrs.clone(), AddrGuard::Baseline).unwrap(),
2401 addrs
2402 );
2403 }
2404
2405 // NFR-004: fail-closed on an empty resolution rather than silently treating "nothing
2406 // resolved" as "nothing to block".
2407 #[test]
2408 fn test_validate_resolved_addrs_fails_closed_on_empty() {
2409 assert_matches!(
2410 validate_resolved_addrs("evil.example", vec![], AddrGuard::Baseline),
2411 Err(ResolveGuardError::NoAddresses { .. })
2412 );
2413 }
2414
2415 #[test]
2416 fn test_validate_resolved_addrs_unwraps_mapped_v4() {
2417 let addrs = vec!["[::ffff:169.254.169.254]:0".parse().unwrap()];
2418 assert_matches!(
2419 validate_resolved_addrs("evil.example", addrs, AddrGuard::Baseline),
2420 Err(ResolveGuardError::Blocked { .. })
2421 );
2422 }
2423
2424 // Issue #455, test-plan item 1: under `Baseline`, an RFC1918/CGNAT/unique-local address is
2425 // allowed (today's pre-#455 behavior) — only `never_a_registry` classes are blocked.
2426 #[test]
2427 fn test_validate_resolved_addrs_baseline_allows_private_ranges() {
2428 for addr_str in ["10.0.0.1:0", "100.64.0.1:0", "[fc00::1]:0"] {
2429 let addrs = vec![addr_str.parse().unwrap()];
2430 assert!(
2431 validate_resolved_addrs("corp.example", addrs, AddrGuard::Baseline).is_ok(),
2432 "{addr_str} must be allowed under Baseline"
2433 );
2434 }
2435 }
2436
2437 // Issue #455, test-plan item 1: under `WorkspaceDeclared(PublicOnly)`, every RFC1918/CGNAT/
2438 // unique-local address is blocked, while a `Global` address is still allowed.
2439 #[test]
2440 fn test_validate_resolved_addrs_workspace_public_only_blocks_private_ranges() {
2441 let guard = AddrGuard::WorkspaceDeclared(WorkspaceRegistryAccess::PublicOnly);
2442 for addr_str in ["10.0.0.1:0", "100.64.0.1:0", "[fc00::1]:0"] {
2443 let addrs = vec![addr_str.parse().unwrap()];
2444 assert!(
2445 matches!(
2446 validate_resolved_addrs("evil.example", addrs, guard),
2447 Err(ResolveGuardError::Blocked { .. })
2448 ),
2449 "{addr_str} must be blocked under WorkspaceDeclared(PublicOnly)"
2450 );
2451 }
2452
2453 let global = vec!["1.1.1.1:0".parse().unwrap()];
2454 assert!(validate_resolved_addrs("index.crates.io", global, guard).is_ok());
2455 }
2456
2457 // Test-plan item 2: `WorkspaceDeclared(All)` admits a private-range address.
2458 #[test]
2459 fn test_validate_resolved_addrs_workspace_all_allows_private_ranges() {
2460 let guard = AddrGuard::WorkspaceDeclared(WorkspaceRegistryAccess::All);
2461 let addrs = vec!["10.0.0.1:0".parse().unwrap()];
2462 assert!(validate_resolved_addrs("corp.example", addrs, guard).is_ok());
2463 }
2464
2465 // Test-plan item 2: `WorkspaceDeclared(Off)` rejects even a `Global` address.
2466 #[test]
2467 fn test_validate_resolved_addrs_workspace_off_rejects_global() {
2468 let guard = AddrGuard::WorkspaceDeclared(WorkspaceRegistryAccess::Off);
2469 let addrs = vec!["1.1.1.1:0".parse().unwrap()];
2470 assert_matches!(
2471 validate_resolved_addrs("index.crates.io", addrs, guard),
2472 Err(ResolveGuardError::Blocked { .. })
2473 );
2474 }
2475
2476 // Test-plan item 3: `build_guarded_client_with_lookup` shares `build_client_inner` with the
2477 // production `build_guarded_client`, so deleting the `.dns_resolver(...)` wiring from that
2478 // shared function fails this test too, not just the production-path wiring test above. The
2479 // synthetic lookup returns an RFC1918 address, resolved (not connected — the resolver
2480 // guard rejects before any TCP attempt) purely through the built `Client`.
2481 #[tokio::test]
2482 async fn test_build_guarded_client_with_lookup_blocks_private_range_under_workspace_tier() {
2483 let lookup = TestLookup(Arc::new(|_host: &str| vec!["10.0.0.1:0".parse().unwrap()]));
2484 let guard = AddrGuard::WorkspaceDeclared(WorkspaceRegistryAccess::PublicOnly);
2485 let client = build_guarded_client_with_lookup(guard, lookup);
2486 let result = client.get("https://corp.example/").send().await;
2487 let err = result.expect_err(
2488 "a private-range synthetic lookup must be rejected under WorkspaceDeclared(PublicOnly)",
2489 );
2490 assert!(
2491 format!("{err:?}").contains("Blocked"),
2492 "expected rejection at the resolver-guard step, got: {err:?}"
2493 );
2494 }
2495
2496 // Test-plan item 3, `Baseline` contrast: the same synthetic private-range lookup is not
2497 // blocked at the resolver-guard step under `Baseline` — asserted directly against the
2498 // resolver (not a full `Client`) to avoid depending on any real network behavior of
2499 // actually connecting to the synthetic address.
2500 #[tokio::test]
2501 async fn test_blocked_addr_resolver_allows_private_range_under_baseline() {
2502 use reqwest::dns::Resolve;
2503
2504 let lookup = TestLookup(Arc::new(|_host: &str| vec!["10.0.0.1:0".parse().unwrap()]));
2505 let resolver = BlockedAddrResolver::with_lookup(AddrGuard::Baseline, lookup);
2506 let result = resolver.resolve("corp.example".parse().unwrap()).await;
2507 assert!(
2508 result.is_ok(),
2509 "Baseline must allow a private-range address through"
2510 );
2511 }
2512
2513 // Direct unit coverage of `BlockedAddrResolver::resolve` on a *name* (not an IP literal —
2514 // that path never reaches any resolver in production, see the struct's `# Known
2515 // limitations` doc). `localhost` resolves via the OS's own hosts file, no network needed.
2516 // This alone does not prove the resolver is wired into `build_guarded_client` — see the
2517 // sibling test below (critic S1) for that.
2518 #[tokio::test]
2519 async fn test_blocked_addr_resolver_rejects_loopback_name_directly() {
2520 use reqwest::dns::Resolve;
2521
2522 let addrs = BlockedAddrResolver::new(AddrGuard::Baseline)
2523 .resolve("localhost".parse().unwrap())
2524 .await;
2525 assert!(addrs.is_err());
2526 }
2527
2528 // Issue #449 critic S1: the prior version called `BlockedAddrResolver::resolve` directly and
2529 // never went through `build_guarded_client`, so deleting `.dns_resolver(...)` from
2530 // `build_client_inner` left it green. This proves actual wiring: a real mockito listener
2531 // answers on `server.socket_address()`'s port, reached via the `localhost` *name* so the
2532 // request actually reaches the configured resolver (unlike an IP literal — see
2533 // `BlockedAddrResolver`'s `# Known limitations` doc).
2534 #[tokio::test]
2535 async fn test_build_client_wires_in_blocked_addr_resolver() {
2536 let mut server = mockito::Server::new_async().await;
2537 let _mock = server
2538 .mock("GET", "/")
2539 .with_status(200)
2540 .create_async()
2541 .await;
2542 let port = server.socket_address().port();
2543
2544 let client = build_guarded_client(AddrGuard::Baseline);
2545 let result = client.get(format!("http://localhost:{port}/")).send().await;
2546
2547 let err = result.expect_err(
2548 "expected the wired-in resolver guard to reject a loopback-resolving name even \
2549 though a real listener answers at this port",
2550 );
2551 // `Debug` (unlike `Display`) surfaces the boxed source chain, confirming
2552 // `ResolveGuardError::Blocked` produced the error rather than an unrelated failure.
2553 let debug = format!("{err:?}");
2554 assert!(
2555 debug.contains("Blocked") && debug.contains("Loopback"),
2556 "expected the failure to originate from ResolveGuardError::Blocked with class \
2557 Loopback, got: {debug}"
2558 );
2559 }
2560
2561 // S5 (plan-1b §1.1/§4): a 302 to the cloud-metadata IP must be stopped by the
2562 // *unconditional* redirect policy, not just the trusted-origin one — this is the
2563 // empirical proof that #443's default unauthenticated client also closes the
2564 // redirect-hop bypass, not only `get_cached_trusted_origin`.
2565 #[tokio::test]
2566 async fn test_get_cached_stops_redirect_to_cloud_metadata() {
2567 let mut server = mockito::Server::new_async().await;
2568
2569 let _redirect = server
2570 .mock("GET", "/api/source")
2571 .with_status(302)
2572 .with_header("location", "https://169.254.169.254/latest/meta-data/")
2573 .create_async()
2574 .await;
2575
2576 let cache = HttpCache::new();
2577 let source_url = format!("{}/api/source", server.url());
2578 let result: Result<Bytes> = cache.get_cached(&source_url).await;
2579
2580 assert!(
2581 matches!(result, Err(DepsError::HttpStatus { status: 302, .. })),
2582 "expected the redirect to be stopped and surfaced as HttpStatus(302)"
2583 );
2584 }
2585
2586 #[tokio::test]
2587 async fn test_get_cached_follows_same_scheme_redirect() {
2588 let mut server = mockito::Server::new_async().await;
2589 let target_url = format!("{}/api/target", server.url());
2590
2591 let _redirect = server
2592 .mock("GET", "/api/source")
2593 .with_status(302)
2594 .with_header("location", &target_url)
2595 .create_async()
2596 .await;
2597 let _target = server
2598 .mock("GET", "/api/target")
2599 .with_status(200)
2600 .with_body("redirected data")
2601 .create_async()
2602 .await;
2603
2604 let cache = HttpCache::new();
2605 let source_url = format!("{}/api/source", server.url());
2606 let result: Bytes = cache.get_cached(&source_url).await.unwrap();
2607
2608 assert_eq!(result.as_ref(), b"redirected data");
2609 }
2610
2611 // Issue #455, test-plan item 4(a): loopback -> loopback. `Baseline` follows the hop
2612 // (test-cfg carve-out for `Loopback`); `WorkspaceDeclared(PublicOnly)` stops it since
2613 // `PublicOnly.allows(Loopback) == false` — the contrast proving the tier split is real.
2614 #[tokio::test]
2615 async fn test_workspace_transport_stops_loopback_redirect_baseline_follows() {
2616 let mut server_a = mockito::Server::new_async().await;
2617 let mut server_b = mockito::Server::new_async().await;
2618 let target_url = format!("{}/api/target", server_b.url());
2619
2620 let _redirect = server_a
2621 .mock("GET", "/api/source")
2622 .with_status(302)
2623 .with_header("location", &target_url)
2624 .create_async()
2625 .await;
2626 let _target = server_b
2627 .mock("GET", "/api/target")
2628 .with_status(200)
2629 .with_body("redirected data")
2630 .create_async()
2631 .await;
2632
2633 let source_url = format!("{}/api/source", server_a.url());
2634
2635 let cache = HttpCache::new();
2636 let result: Bytes = cache
2637 .get_cached_with_headers_via(&source_url, &[], &Transport::baseline(), None)
2638 .await
2639 .unwrap();
2640 assert_eq!(result.as_ref(), b"redirected data");
2641
2642 let policy = Arc::new(RegistryAccessPolicy::new(
2643 WorkspaceRegistryAccess::PublicOnly,
2644 ));
2645 let workspace_cache = HttpCache::with_policy(Arc::clone(&policy));
2646 let result: Result<Bytes> = workspace_cache
2647 .get_cached_with_headers_via(&source_url, &[], &Transport::workspace(&policy), None)
2648 .await;
2649 assert!(
2650 matches!(result, Err(DepsError::HttpStatus { status: 302, .. })),
2651 "expected the workspace transport to stop the loopback hop, got {result:?}"
2652 );
2653 }
2654
2655 // Issue #455, test-plan item 4(b): a workspace-blocked literal. The redirect-policy tier
2656 // term rejects an RFC1918-literal target from its URL string alone (no resolver involved),
2657 // so the caller sees `HttpStatus{302}` with no `HTTP_TIMEOUT_SECS` stall.
2658 #[tokio::test]
2659 async fn test_workspace_transport_stops_redirect_to_private_literal() {
2660 let mut server = mockito::Server::new_async().await;
2661
2662 let _redirect = server
2663 .mock("GET", "/api/source")
2664 .with_status(302)
2665 .with_header("location", "https://10.0.0.1/x")
2666 .create_async()
2667 .await;
2668
2669 let policy = Arc::new(RegistryAccessPolicy::new(
2670 WorkspaceRegistryAccess::PublicOnly,
2671 ));
2672 let cache = HttpCache::with_policy(Arc::clone(&policy));
2673 let source_url = format!("{}/api/source", server.url());
2674 let result: Result<Bytes> = cache
2675 .get_cached_with_headers_via(&source_url, &[], &Transport::workspace(&policy), None)
2676 .await;
2677
2678 assert!(
2679 matches!(result, Err(DepsError::HttpStatus { status: 302, .. })),
2680 "expected the redirect to 10.0.0.1 to be stopped, got {result:?}"
2681 );
2682 }
2683
2684 // Unlike the https->http downgrade case, cross-origin redirect blocking IS reachable
2685 // through mockito: two separate `mockito::Server` instances bind to distinct ports,
2686 // and a distinct port is a distinct origin (scheme+host+port), so a 302 from one to
2687 // the other is a genuine cross-origin redirect the trusted-origin policy must stop.
2688 #[tokio::test]
2689 async fn test_get_cached_trusted_origin_stops_cross_origin_redirect() {
2690 let mut trusted_server = mockito::Server::new_async().await;
2691 let mut other_server = mockito::Server::new_async().await;
2692
2693 let trusted_origin = format!("{}/", trusted_server.url());
2694 let escape_target = format!("{}/api/stolen", other_server.url());
2695
2696 let _redirect = trusted_server
2697 .mock("GET", "/api/source")
2698 .with_status(302)
2699 .with_header("location", &escape_target)
2700 .create_async()
2701 .await;
2702 let escape = other_server
2703 .mock("GET", "/api/stolen")
2704 .with_status(200)
2705 .with_body("must not be returned")
2706 .expect(0)
2707 .create_async()
2708 .await;
2709
2710 let cache = HttpCache::new();
2711 let source_url = format!("{}/api/source", trusted_server.url());
2712 let result: Result<Bytes> = cache
2713 .get_cached_trusted_origin(&source_url, &trusted_origin)
2714 .await;
2715
2716 // The stopped redirect surfaces as the 302 response, like any other non-2xx status.
2717 // `matches!` rather than debug-formatting `result`: the `Ok` arm holds the raw body.
2718 assert!(
2719 matches!(result, Err(DepsError::HttpStatus { status: 302, .. })),
2720 "expected HttpStatus(302)"
2721 );
2722
2723 // Proves the escape origin was never contacted, not just that the result is a 302.
2724 escape.assert_async().await;
2725 }
2726
2727 #[tokio::test]
2728 async fn test_get_cached_trusted_origin_follows_same_origin_redirect() {
2729 let mut server = mockito::Server::new_async().await;
2730 let trusted_origin = format!("{}/", server.url());
2731 let target_url = format!("{}/api/target", server.url());
2732
2733 let _redirect = server
2734 .mock("GET", "/api/source")
2735 .with_status(302)
2736 .with_header("location", &target_url)
2737 .create_async()
2738 .await;
2739 let _target = server
2740 .mock("GET", "/api/target")
2741 .with_status(200)
2742 .with_body("trusted data")
2743 .create_async()
2744 .await;
2745
2746 let cache = HttpCache::new();
2747 let source_url = format!("{}/api/source", server.url());
2748 let result: Bytes = cache
2749 .get_cached_trusted_origin(&source_url, &trusted_origin)
2750 .await
2751 .unwrap();
2752
2753 assert_eq!(result.as_ref(), b"trusted data");
2754 }
2755
2756 /// Issue #795: a raw `str::starts_with` prefix test (the pre-fix behavior) is satisfied
2757 /// by a subdomain-suffix bypass — `gitlab.mycorp.dev.evil.com` starts with
2758 /// `https://gitlab.mycorp.dev` as a string, even though its actual host is
2759 /// `gitlab.mycorp.dev.evil.com`, entirely under attacker control. Parsed-origin equality
2760 /// must reject it.
2761 #[test]
2762 fn test_is_trusted_origin_rejects_subdomain_suffix_bypass() {
2763 let trusted = Url::parse("https://gitlab.mycorp.dev").unwrap();
2764 let hop = Url::parse("https://gitlab.mycorp.dev.evil.com/steal").unwrap();
2765 assert!(!is_trusted_origin(&hop, Some(&trusted)));
2766 }
2767
2768 /// Issue #795: a userinfo bypass — `https://gitlab.mycorp.dev@evil.com/...` also starts
2769 /// with the trusted origin as a string, but its host is `evil.com`; `gitlab.mycorp.dev`
2770 /// is merely a (discarded) username. `Url::origin()` ignores userinfo entirely, so this
2771 /// must be rejected.
2772 #[test]
2773 fn test_is_trusted_origin_rejects_userinfo_bypass() {
2774 let trusted = Url::parse("https://gitlab.mycorp.dev").unwrap();
2775 let hop = Url::parse("https://gitlab.mycorp.dev@evil.com/steal").unwrap();
2776 assert!(!is_trusted_origin(&hop, Some(&trusted)));
2777 }
2778
2779 /// Issue #795: a hyphen-suffix bypass — `gitlab.mycorp.dev-evil.com` again starts with
2780 /// the trusted origin as a string while being an entirely distinct, attacker-controlled
2781 /// host.
2782 #[test]
2783 fn test_is_trusted_origin_rejects_hyphen_suffix_bypass() {
2784 let trusted = Url::parse("https://gitlab.mycorp.dev").unwrap();
2785 let hop = Url::parse("https://gitlab.mycorp.dev-evil.com/steal").unwrap();
2786 assert!(!is_trusted_origin(&hop, Some(&trusted)));
2787 }
2788
2789 /// Companion to the three bypass-rejection tests above: the legitimate same-origin case
2790 /// (a different path, same scheme/host/port) must still be accepted.
2791 #[test]
2792 fn test_is_trusted_origin_accepts_exact_origin_match() {
2793 let trusted = Url::parse("https://gitlab.mycorp.dev").unwrap();
2794 let hop = Url::parse("https://gitlab.mycorp.dev/api/v4/x").unwrap();
2795 assert!(is_trusted_origin(&hop, Some(&trusted)));
2796 }
2797
2798 /// A `trusted_origin` that fails to parse must fail closed — every hop is rejected,
2799 /// never treated as "no restriction".
2800 #[test]
2801 fn test_is_trusted_origin_rejects_when_trusted_origin_unparseable() {
2802 let hop = Url::parse("https://gitlab.mycorp.dev/api/v4/x").unwrap();
2803 assert!(!is_trusted_origin(&hop, None));
2804 }
2805
2806 /// Path-prefix scoping (NuGet's registration-hive/flat-container pinning) must survive
2807 /// the #795 origin-equality fix: same origin, but a hop outside the trusted path, is
2808 /// still rejected — this is what `test_get_cached_trusted_origin_rejects_sibling_path_prefix`
2809 /// exercises end-to-end; this is the same property pinned at the unit level.
2810 #[test]
2811 fn test_is_trusted_origin_rejects_same_origin_sibling_path() {
2812 let trusted = Url::parse("https://registry.example/v3/registration5-gz/").unwrap();
2813 let hop = Url::parse("https://registry.example/v3/registration5-gzX/evil").unwrap();
2814 assert!(!is_trusted_origin(&hop, Some(&trusted)));
2815 }
2816
2817 /// Companion: same origin, hop path under the trusted path prefix, is still accepted.
2818 #[test]
2819 fn test_is_trusted_origin_accepts_same_origin_nested_path() {
2820 let trusted = Url::parse("https://registry.example/v3/registration5-gz/").unwrap();
2821 let hop =
2822 Url::parse("https://registry.example/v3/registration5-gz/serde/page1.json").unwrap();
2823 assert!(is_trusted_origin(&hop, Some(&trusted)));
2824 }
2825
2826 /// Issue #795 S1: unlike the two trailing-slash tests above (which, with a trailing `/`
2827 /// already present in the trusted path, would have passed even under the pre-S1-fix raw
2828 /// `str::starts_with` check — they do not actually exercise the segment-boundary fix),
2829 /// `RegistryIndex::as_str()` (`deps-cargo`'s sparse-index trusted origin) carries **no**
2830 /// trailing-slash guarantee. This reproduces that exact shape and the critic's repro: a
2831 /// same-origin sibling whose path merely shares a textual prefix must still be rejected.
2832 #[test]
2833 fn test_is_trusted_origin_rejects_same_origin_sibling_path_no_trailing_slash() {
2834 let trusted = Url::parse("https://artifacts.corp/cargo/index").unwrap();
2835 for sibling in [
2836 "https://artifacts.corp/cargo/index-public/steal",
2837 "https://artifacts.corp/cargo/indexEVIL",
2838 "https://artifacts.corp/cargo/index.evil/x",
2839 ] {
2840 let hop = Url::parse(sibling).unwrap();
2841 assert!(
2842 !is_trusted_origin(&hop, Some(&trusted)),
2843 "expected {sibling} to be rejected"
2844 );
2845 }
2846 }
2847
2848 /// Companion: the trusted path itself, and a proper child path, are still accepted when
2849 /// the trusted path carries no trailing slash — the real `deps-cargo` request shape
2850 /// (`sparse_index_url` appends `/{crate_path}` to the trimmed base).
2851 #[test]
2852 fn test_is_trusted_origin_accepts_self_and_child_no_trailing_slash() {
2853 let trusted = Url::parse("https://artifacts.corp/cargo/index").unwrap();
2854 let itself = Url::parse("https://artifacts.corp/cargo/index").unwrap();
2855 let child = Url::parse("https://artifacts.corp/cargo/index/se/rd/serde").unwrap();
2856 assert!(is_trusted_origin(&itself, Some(&trusted)));
2857 assert!(is_trusted_origin(&child, Some(&trusted)));
2858 }
2859
2860 // Proves every hop is re-checked, not just the first: a same-origin hop is followed,
2861 // then a second, cross-origin hop from that (already-followed) intermediate is stopped.
2862 #[tokio::test]
2863 async fn test_get_cached_trusted_origin_stops_second_hop_of_multi_hop_chain() {
2864 let mut trusted_server = mockito::Server::new_async().await;
2865 let mut other_server = mockito::Server::new_async().await;
2866
2867 let trusted_origin = format!("{}/", trusted_server.url());
2868 let intermediate_url = format!("{}/api/intermediate", trusted_server.url());
2869 let escape_target = format!("{}/api/stolen", other_server.url());
2870
2871 let _first_hop = trusted_server
2872 .mock("GET", "/api/source")
2873 .with_status(302)
2874 .with_header("location", &intermediate_url)
2875 .create_async()
2876 .await;
2877 let _second_hop = trusted_server
2878 .mock("GET", "/api/intermediate")
2879 .with_status(302)
2880 .with_header("location", &escape_target)
2881 .create_async()
2882 .await;
2883 let escape = other_server
2884 .mock("GET", "/api/stolen")
2885 .with_status(200)
2886 .with_body("must not be returned")
2887 .expect(0)
2888 .create_async()
2889 .await;
2890
2891 let cache = HttpCache::new();
2892 let source_url = format!("{}/api/source", trusted_server.url());
2893 let result: Result<Bytes> = cache
2894 .get_cached_trusted_origin(&source_url, &trusted_origin)
2895 .await;
2896
2897 // matches! rather than debug-formatting result: the Ok arm holds the raw body.
2898 assert!(
2899 matches!(result, Err(DepsError::HttpStatus { status: 302, .. })),
2900 "expected HttpStatus(302)"
2901 );
2902 escape.assert_async().await;
2903 }
2904
2905 // Sibling path-prefix rejection: `.../api/` must not accept `.../apiX/...`. The other
2906 // trusted-origin tests use a bare-host prefix, which never exercises this boundary.
2907 #[tokio::test]
2908 async fn test_get_cached_trusted_origin_rejects_sibling_path_prefix() {
2909 let mut server = mockito::Server::new_async().await;
2910 let trusted_origin = format!("{}/api/", server.url());
2911 let escape_target = format!("{}/apiX/evil", server.url());
2912
2913 let _redirect = server
2914 .mock("GET", "/api/source")
2915 .with_status(302)
2916 .with_header("location", &escape_target)
2917 .create_async()
2918 .await;
2919 let escape = server
2920 .mock("GET", "/apiX/evil")
2921 .with_status(200)
2922 .with_body("must not be returned")
2923 .expect(0)
2924 .create_async()
2925 .await;
2926
2927 let cache = HttpCache::new();
2928 let source_url = format!("{}/api/source", server.url());
2929 let result: Result<Bytes> = cache
2930 .get_cached_trusted_origin(&source_url, &trusted_origin)
2931 .await;
2932
2933 // matches! rather than debug-formatting result: the Ok arm holds the raw body.
2934 assert!(
2935 matches!(result, Err(DepsError::HttpStatus { status: 302, .. })),
2936 "expected HttpStatus(302)"
2937 );
2938 escape.assert_async().await;
2939 }
2940
2941 #[tokio::test]
2942 async fn test_get_cached_trusted_origin_with_headers_sends_extra_header() {
2943 let mut server = mockito::Server::new_async().await;
2944 let trusted_origin = format!("{}/", server.url());
2945
2946 let _m = server
2947 .mock("GET", "/api/data")
2948 .match_header("authorization", "Bearer secret-token")
2949 .with_status(200)
2950 .with_body("authenticated data")
2951 .create_async()
2952 .await;
2953
2954 let cache = HttpCache::new();
2955 let url = format!("{}/api/data", server.url());
2956 let result: Bytes = cache
2957 .get_cached_trusted_origin_with_headers(
2958 &url,
2959 &trusted_origin,
2960 &[(header::AUTHORIZATION, "Bearer secret-token")],
2961 )
2962 .await
2963 .unwrap();
2964
2965 assert_eq!(result.as_ref(), b"authenticated data");
2966 }
2967
2968 // A credential header must never survive a cross-origin redirect hop — proven by the
2969 // escape origin never being contacted, not just the header being absent on a landed request.
2970 #[tokio::test]
2971 async fn test_get_cached_trusted_origin_with_headers_stops_cross_origin_redirect() {
2972 let mut trusted_server = mockito::Server::new_async().await;
2973 let mut other_server = mockito::Server::new_async().await;
2974
2975 let trusted_origin = format!("{}/", trusted_server.url());
2976 let escape_target = format!("{}/api/stolen", other_server.url());
2977
2978 let _redirect = trusted_server
2979 .mock("GET", "/api/source")
2980 .with_status(302)
2981 .with_header("location", &escape_target)
2982 .create_async()
2983 .await;
2984 let escape = other_server
2985 .mock("GET", "/api/stolen")
2986 .with_status(200)
2987 .with_body("must not be returned")
2988 .expect(0)
2989 .create_async()
2990 .await;
2991
2992 let cache = HttpCache::new();
2993 let source_url = format!("{}/api/source", trusted_server.url());
2994 let result: Result<Bytes> = cache
2995 .get_cached_trusted_origin_with_headers(
2996 &source_url,
2997 &trusted_origin,
2998 &[(header::AUTHORIZATION, "Bearer secret-token")],
2999 )
3000 .await;
3001
3002 assert!(
3003 matches!(result, Err(DepsError::HttpStatus { status: 302, .. })),
3004 "expected HttpStatus(302)"
3005 );
3006 escape.assert_async().await;
3007 }
3008
3009 // Proves `redirect_policy`'s delegation is actually wired in and live: without it, this
3010 // chain would keep following past 10 hops instead of erroring — a single-hop test can't
3011 // distinguish "delegation is live" from "no policy at all".
3012 #[tokio::test]
3013 async fn test_get_cached_default_client_enforces_ten_hop_redirect_limit() {
3014 let mut server = mockito::Server::new_async().await;
3015 let base = server.url();
3016
3017 // reqwest errors once `previous.len() > 10` (the 11th hop), so 11 redirecting steps
3018 // (step/0..step/10) are needed to trigger it; step/11 must never be requested.
3019 let mut hop_mocks = Vec::new();
3020 for i in 0..11u32 {
3021 let path = format!("/step/{i}");
3022 let next = format!("{base}/step/{}", i + 1);
3023 hop_mocks.push(
3024 server
3025 .mock("GET", path.as_str())
3026 .with_status(302)
3027 .with_header("location", &next)
3028 .create_async()
3029 .await,
3030 );
3031 }
3032 let final_step = server
3033 .mock("GET", "/step/11")
3034 .with_status(200)
3035 .with_body("unreachable")
3036 .expect(0)
3037 .create_async()
3038 .await;
3039
3040 // Kept alive until here: each `Mock` deregisters on drop, turning hops 404 otherwise.
3041 assert_eq!(hop_mocks.len(), 11);
3042
3043 let cache = HttpCache::new();
3044 let start_url = format!("{base}/step/0");
3045 let result: Result<Bytes> = cache.get_cached(&start_url).await;
3046
3047 assert!(
3048 matches!(result, Err(DepsError::RegistryError { .. })),
3049 "expected a too-many-redirects network error, got {result:?}"
3050 );
3051 final_step.assert_async().await;
3052 }
3053
3054 #[test]
3055 fn test_cache_creation() {
3056 let cache = HttpCache::new();
3057 assert_eq!(cache.len(), 0);
3058 assert!(cache.is_empty());
3059 }
3060
3061 #[test]
3062 fn test_cache_clear() {
3063 let cache = HttpCache::new();
3064 cache.entries.insert(
3065 "test".into(),
3066 CachedResponse {
3067 body: Bytes::from_static(&[1, 2, 3]),
3068 etag: None,
3069 last_modified: None,
3070 fetched_at: Instant::now(),
3071 },
3072 );
3073 assert_eq!(cache.len(), 1);
3074 cache.clear();
3075 assert_eq!(cache.len(), 0);
3076 }
3077
3078 #[test]
3079 fn test_cached_response_clone() {
3080 let response = CachedResponse {
3081 body: Bytes::from_static(&[1, 2, 3]),
3082 etag: Some("test".into()),
3083 last_modified: Some("date".into()),
3084 fetched_at: Instant::now(),
3085 };
3086 let cloned = response.clone();
3087 // Bytes clone is cheap (reference counting)
3088 assert_eq!(response.body, cloned.body);
3089 assert_eq!(response.etag, cloned.etag);
3090 }
3091
3092 #[test]
3093 fn test_cache_len() {
3094 let cache = HttpCache::new();
3095 assert_eq!(cache.len(), 0);
3096
3097 cache.entries.insert(
3098 "url1".into(),
3099 CachedResponse {
3100 body: Bytes::new(),
3101 etag: None,
3102 last_modified: None,
3103 fetched_at: Instant::now(),
3104 },
3105 );
3106
3107 assert_eq!(cache.len(), 1);
3108 }
3109
3110 #[tokio::test]
3111 async fn test_get_cached_fresh_fetch() {
3112 let mut server = mockito::Server::new_async().await;
3113
3114 let _m = server
3115 .mock("GET", "/api/data")
3116 .with_status(200)
3117 .with_header("etag", "\"abc123\"")
3118 .with_body("test data")
3119 .create_async()
3120 .await;
3121
3122 let cache = HttpCache::new();
3123 let url = format!("{}/api/data", server.url());
3124 let result: Bytes = cache.get_cached(&url).await.unwrap();
3125
3126 assert_eq!(result.as_ref(), b"test data");
3127 assert_eq!(cache.len(), 1);
3128 }
3129
3130 #[tokio::test]
3131 async fn test_get_cached_cache_hit() {
3132 let mut server = mockito::Server::new_async().await;
3133 let url = format!("{}/api/data", server.url());
3134
3135 let cache = HttpCache::new();
3136
3137 let _m1 = server
3138 .mock("GET", "/api/data")
3139 .with_status(200)
3140 .with_header("etag", "\"abc123\"")
3141 .with_body("original data")
3142 .create_async()
3143 .await;
3144
3145 let result1: Bytes = cache.get_cached(&url).await.unwrap();
3146 assert_eq!(result1.as_ref(), b"original data");
3147 assert_eq!(cache.len(), 1);
3148
3149 drop(_m1);
3150
3151 let _m2 = server
3152 .mock("GET", "/api/data")
3153 .match_header("if-none-match", "\"abc123\"")
3154 .with_status(304)
3155 .create_async()
3156 .await;
3157
3158 let result2: Bytes = cache.get_cached(&url).await.unwrap();
3159 assert_eq!(result2.as_ref(), b"original data");
3160 }
3161
3162 #[tokio::test]
3163 async fn test_get_cached_304_not_modified() {
3164 let mut server = mockito::Server::new_async().await;
3165 let url = format!("{}/api/data", server.url());
3166
3167 let cache = HttpCache::new();
3168
3169 let _m1 = server
3170 .mock("GET", "/api/data")
3171 .with_status(200)
3172 .with_header("etag", "\"abc123\"")
3173 .with_body("original data")
3174 .create_async()
3175 .await;
3176
3177 let result1: Bytes = cache.get_cached(&url).await.unwrap();
3178 assert_eq!(result1.as_ref(), b"original data");
3179
3180 drop(_m1);
3181
3182 let _m2 = server
3183 .mock("GET", "/api/data")
3184 .match_header("if-none-match", "\"abc123\"")
3185 .with_status(304)
3186 .create_async()
3187 .await;
3188
3189 let result2: Bytes = cache.get_cached(&url).await.unwrap();
3190 assert_eq!(result2.as_ref(), b"original data");
3191 }
3192
3193 #[tokio::test]
3194 async fn test_get_cached_etag_validation() {
3195 let mut server = mockito::Server::new_async().await;
3196 let url = format!("{}/api/data", server.url());
3197
3198 let cache = HttpCache::new();
3199
3200 cache.entries.insert(
3201 url.clone(),
3202 CachedResponse {
3203 body: Bytes::from_static(b"cached"),
3204 etag: Some("\"tag123\"".into()),
3205 last_modified: None,
3206 fetched_at: Instant::now(),
3207 },
3208 );
3209
3210 let _m = server
3211 .mock("GET", "/api/data")
3212 .match_header("if-none-match", "\"tag123\"")
3213 .with_status(304)
3214 .create_async()
3215 .await;
3216
3217 let result: Bytes = cache.get_cached(&url).await.unwrap();
3218 assert_eq!(result.as_ref(), b"cached");
3219 }
3220
3221 #[tokio::test]
3222 async fn test_get_cached_last_modified_validation() {
3223 let mut server = mockito::Server::new_async().await;
3224 let url = format!("{}/api/data", server.url());
3225
3226 let cache = HttpCache::new();
3227
3228 cache.entries.insert(
3229 url.clone(),
3230 CachedResponse {
3231 body: Bytes::from_static(b"cached"),
3232 etag: None,
3233 last_modified: Some("Wed, 21 Oct 2024 07:28:00 GMT".into()),
3234 fetched_at: Instant::now(),
3235 },
3236 );
3237
3238 let _m = server
3239 .mock("GET", "/api/data")
3240 .match_header("if-modified-since", "Wed, 21 Oct 2024 07:28:00 GMT")
3241 .with_status(304)
3242 .create_async()
3243 .await;
3244
3245 let result: Bytes = cache.get_cached(&url).await.unwrap();
3246 assert_eq!(result.as_ref(), b"cached");
3247 }
3248
3249 #[tokio::test]
3250 async fn test_get_cached_network_error_fallback() {
3251 let cache = HttpCache::new();
3252 // https:// (not http://) so this exercises DNS-resolution failure, not the
3253 // HTTPS-only policy enforced by `ensure_https`.
3254 let url = "https://invalid.localhost.test/data";
3255
3256 cache.entries.insert(
3257 url.to_string(),
3258 CachedResponse {
3259 body: Bytes::from_static(b"stale data"),
3260 etag: Some("\"old\"".into()),
3261 last_modified: None,
3262 fetched_at: Instant::now(),
3263 },
3264 );
3265
3266 let result: Bytes = cache.get_cached(url).await.unwrap();
3267 assert_eq!(result.as_ref(), b"stale data");
3268 }
3269
3270 #[tokio::test]
3271 async fn test_fetch_and_store_http_error() {
3272 let mut server = mockito::Server::new_async().await;
3273
3274 let _m = server
3275 .mock("GET", "/api/missing")
3276 .with_status(404)
3277 .with_body("Not Found")
3278 .create_async()
3279 .await;
3280
3281 let cache = HttpCache::new();
3282 let url = format!("{}/api/missing", server.url());
3283 let result: Result<Bytes> = cache
3284 .fetch_and_store_with_headers(&url, &[], &cache.baseline.client, &url)
3285 .await;
3286
3287 assert!(result.is_err());
3288 match result {
3289 Err(DepsError::HttpStatus { status, .. }) => {
3290 assert_eq!(status, 404);
3291 }
3292 _ => panic!("Expected HttpStatus"),
3293 }
3294 }
3295
3296 /// #1295 (a): a 403 carrying confirmed `X-RateLimit-Remaining: 0` evidence classifies as
3297 /// a *verified* rate limit, not a bare `HttpStatus`.
3298 #[tokio::test]
3299 async fn test_fetch_and_store_403_with_confirmed_evidence_is_verified_rate_limited() {
3300 let mut server = mockito::Server::new_async().await;
3301
3302 let _m = server
3303 .mock("GET", "/repos/owner/repo/tags")
3304 .with_status(403)
3305 .with_header("x-ratelimit-remaining", "0")
3306 .with_body(r#"{"message":"API rate limit exceeded"}"#)
3307 .create_async()
3308 .await;
3309
3310 let cache = HttpCache::new();
3311 let url = format!("{}/repos/owner/repo/tags", server.url());
3312 let result: Result<Bytes> = cache
3313 .fetch_and_store_with_headers(&url, &[], &cache.baseline.client, &url)
3314 .await;
3315
3316 match result {
3317 Err(DepsError::RateLimited { verified, .. }) => {
3318 assert_eq!(verified, RateLimitEvidence::Confirmed);
3319 }
3320 other => panic!("expected RateLimited, got {other:?}"),
3321 }
3322 }
3323
3324 /// #1295 (b): a plain 403 with no `X-RateLimit-Remaining` header at all — the shape a
3325 /// non-rate-limit 403 cause (abuse-detection false positive, secondary rate limit, an
3326 /// access-restricted repo) would have — stays a bare `HttpStatus`, distinguishable from
3327 /// the confirmed case above.
3328 #[tokio::test]
3329 async fn test_fetch_and_store_403_without_evidence_stays_http_status() {
3330 let mut server = mockito::Server::new_async().await;
3331
3332 let _m = server
3333 .mock("GET", "/repos/owner/repo/tags")
3334 .with_status(403)
3335 .with_body(r#"{"message":"Resource not accessible by integration"}"#)
3336 .create_async()
3337 .await;
3338
3339 let cache = HttpCache::new();
3340 let url = format!("{}/repos/owner/repo/tags", server.url());
3341 let result: Result<Bytes> = cache
3342 .fetch_and_store_with_headers(&url, &[], &cache.baseline.client, &url)
3343 .await;
3344
3345 match result {
3346 Err(DepsError::HttpStatus { status, .. }) => assert_eq!(status, 403),
3347 other => panic!("expected HttpStatus, got {other:?}"),
3348 }
3349 }
3350
3351 /// #1295: a non-zero `X-RateLimit-Remaining` on a 403 is not confirming evidence either —
3352 /// the request was rejected for some other reason while quota remains.
3353 #[tokio::test]
3354 async fn test_fetch_and_store_403_with_nonzero_remaining_stays_http_status() {
3355 let mut server = mockito::Server::new_async().await;
3356
3357 let _m = server
3358 .mock("GET", "/repos/owner/repo/tags")
3359 .with_status(403)
3360 .with_header("x-ratelimit-remaining", "42")
3361 .create_async()
3362 .await;
3363
3364 let cache = HttpCache::new();
3365 let url = format!("{}/repos/owner/repo/tags", server.url());
3366 let result: Result<Bytes> = cache
3367 .fetch_and_store_with_headers(&url, &[], &cache.baseline.client, &url)
3368 .await;
3369
3370 match result {
3371 Err(DepsError::HttpStatus { status, .. }) => assert_eq!(status, 403),
3372 other => panic!("expected HttpStatus, got {other:?}"),
3373 }
3374 }
3375
3376 /// #1295: a `X-RateLimit-Remaining: 0` header on a status other than 403/429 is not
3377 /// rate-limit evidence — [`confirmed_rate_limit_exhaustion`] must stay status-gated.
3378 #[test]
3379 fn test_confirmed_rate_limit_exhaustion_requires_403_or_429_status() {
3380 let mut headers = header::HeaderMap::new();
3381 headers.insert(
3382 "x-ratelimit-remaining",
3383 header::HeaderValue::from_static("0"),
3384 );
3385 assert!(!confirmed_rate_limit_exhaustion(
3386 StatusCode::SERVICE_UNAVAILABLE,
3387 &headers
3388 ));
3389 assert!(confirmed_rate_limit_exhaustion(
3390 StatusCode::FORBIDDEN,
3391 &headers
3392 ));
3393 }
3394
3395 /// #1295 critic S4: GitHub's primary rate limit can also arrive as 429 (not only 403),
3396 /// with `X-RateLimit-Remaining: 0`.
3397 #[test]
3398 fn test_confirmed_rate_limit_exhaustion_accepts_429_with_zero_remaining() {
3399 let mut headers = header::HeaderMap::new();
3400 headers.insert(
3401 "x-ratelimit-remaining",
3402 header::HeaderValue::from_static("0"),
3403 );
3404 assert!(confirmed_rate_limit_exhaustion(
3405 StatusCode::TOO_MANY_REQUESTS,
3406 &headers
3407 ));
3408 }
3409
3410 /// #1295 critic S4: a secondary rate limit arrives as 403/429 with `Retry-After` and a
3411 /// *non-zero* (or absent) `X-RateLimit-Remaining` — `Retry-After` alone must count as
3412 /// evidence, since `remaining == 0` alone would miss this shape entirely.
3413 #[test]
3414 fn test_confirmed_rate_limit_exhaustion_accepts_retry_after_without_remaining() {
3415 let mut headers = header::HeaderMap::new();
3416 headers.insert("retry-after", header::HeaderValue::from_static("60"));
3417 assert!(confirmed_rate_limit_exhaustion(
3418 StatusCode::FORBIDDEN,
3419 &headers
3420 ));
3421 assert!(confirmed_rate_limit_exhaustion(
3422 StatusCode::TOO_MANY_REQUESTS,
3423 &headers
3424 ));
3425 }
3426
3427 /// #1295 critic S4: `Retry-After` evidence still requires a 403/429 status — an unrelated
3428 /// 503 with a `Retry-After` header (ordinary server-maintenance semantics) is not a
3429 /// rate-limit confirmation.
3430 #[test]
3431 fn test_confirmed_rate_limit_exhaustion_retry_after_still_status_gated() {
3432 let mut headers = header::HeaderMap::new();
3433 headers.insert("retry-after", header::HeaderValue::from_static("60"));
3434 assert!(!confirmed_rate_limit_exhaustion(
3435 StatusCode::SERVICE_UNAVAILABLE,
3436 &headers
3437 ));
3438 }
3439
3440 /// #1295 critic M4: a malformed/padded `X-RateLimit-Remaining` value (not a bare `"0"`)
3441 /// must not be treated as confirming evidence — `confirmed_rate_limit_exhaustion` parses
3442 /// the value rather than comparing it as an exact string, so this is evidence-neutral
3443 /// (falls through to `HttpStatus`) rather than a false positive or a panic.
3444 #[test]
3445 fn test_confirmed_rate_limit_exhaustion_rejects_malformed_remaining() {
3446 let mut headers = header::HeaderMap::new();
3447 headers.insert(
3448 "x-ratelimit-remaining",
3449 header::HeaderValue::from_static("not-a-number"),
3450 );
3451 assert!(!confirmed_rate_limit_exhaustion(
3452 StatusCode::FORBIDDEN,
3453 &headers
3454 ));
3455 }
3456
3457 /// #1295 critic M4: a padded numeric value (e.g. `"00"`) parses to the same integer `0`
3458 /// and must still count as evidence — the point of switching from string equality to
3459 /// `parse::<u64>()`.
3460 #[test]
3461 fn test_confirmed_rate_limit_exhaustion_accepts_padded_zero() {
3462 let mut headers = header::HeaderMap::new();
3463 headers.insert(
3464 "x-ratelimit-remaining",
3465 header::HeaderValue::from_static("00"),
3466 );
3467 assert!(confirmed_rate_limit_exhaustion(
3468 StatusCode::FORBIDDEN,
3469 &headers
3470 ));
3471 }
3472
3473 #[tokio::test]
3474 async fn test_fetch_and_store_stores_headers() {
3475 let mut server = mockito::Server::new_async().await;
3476
3477 let _m = server
3478 .mock("GET", "/api/data")
3479 .with_status(200)
3480 .with_header("etag", "\"abc123\"")
3481 .with_header("last-modified", "Wed, 21 Oct 2024 07:28:00 GMT")
3482 .with_body("test")
3483 .create_async()
3484 .await;
3485
3486 let cache = HttpCache::new();
3487 let url = format!("{}/api/data", server.url());
3488 let _: Bytes = cache
3489 .fetch_and_store_with_headers(&url, &[], &cache.baseline.client, &url)
3490 .await
3491 .unwrap();
3492
3493 let cached = cache.entries.get(&url).unwrap();
3494 assert_eq!(cached.etag, Some("\"abc123\"".into()));
3495 assert_eq!(
3496 cached.last_modified,
3497 Some("Wed, 21 Oct 2024 07:28:00 GMT".into())
3498 );
3499 }
3500
3501 /// #756 security follow-up S-A: the "fetching fresh: {url}" debug log in
3502 /// `fetch_and_store_with_headers` — the direct callee `get_cached_with_headers_via`'s
3503 /// `miss` branch delegates to — must never carry a token embedded in the URL's query
3504 /// string. This is `debug`-level, the level this project's own continuous-improvement
3505 /// convention runs at (`RUST_LOG=debug`), so it is not merely a theoretical exposure.
3506 #[cfg(feature = "test-util")]
3507 #[tokio::test]
3508 async fn test_fetch_and_store_fetching_fresh_log_redacts_query_string_token() {
3509 let mut server = mockito::Server::new_async().await;
3510 let url = format!("{}/pkg?token=super-secret-value", server.url());
3511
3512 let _m = server
3513 .mock("GET", "/pkg")
3514 .match_query(mockito::Matcher::UrlEncoded(
3515 "token".into(),
3516 "super-secret-value".into(),
3517 ))
3518 .with_status(200)
3519 .with_body("ok")
3520 .create_async()
3521 .await;
3522
3523 let cache = HttpCache::new();
3524 let output =
3525 crate::test_util::capture_tracing_output_async_at(tracing::Level::DEBUG, async {
3526 let result: Bytes = cache
3527 .fetch_and_store_with_headers(&url, &[], &cache.baseline.client, &url)
3528 .await
3529 .unwrap();
3530 assert_eq!(result.as_ref(), b"ok");
3531 })
3532 .await;
3533
3534 assert!(
3535 !output.contains("super-secret-value"),
3536 "leaked token via 'fetching fresh' debug log: {output:?}"
3537 );
3538 }
3539
3540 /// #767 S3 / #789: proves `SanitizedRegistryError`'s `From<reqwest::Error>` actually
3541 /// strips the URL from the wrapped `reqwest::Error`'s own `Display`, not just that
3542 /// `RegistryError::package` is redacted — a genuine transport-level error (connection
3543 /// refused on a closed loopback port, so `.url()` is populated the way a builder-only
3544 /// error like `Client::get("not a url").build().unwrap_err()` never is) is required to
3545 /// exercise this: reverting any of the 5 `source: e.into()` call sites in this file back
3546 /// to a bare `reqwest::Error` must fail to compile, and reverting
3547 /// `SanitizedRegistryError::from`'s `.without_url()` call must fail this test.
3548 #[tokio::test]
3549 async fn test_registry_error_source_redacts_url_on_real_transport_error() {
3550 // Bind then immediately drop a loopback listener: nothing accepts connections on
3551 // this port afterward, so a request to it fails fast with connection-refused
3552 // instead of hanging or needing a real unreachable host.
3553 let port = {
3554 let listener = std::net::TcpListener::bind("127.0.0.1:0").unwrap();
3555 listener.local_addr().unwrap().port()
3556 };
3557 let url = format!("http://127.0.0.1:{port}/pkg?token=super-secret-value");
3558
3559 let cache = HttpCache::new();
3560 let err = cache.get_cached(&url).await.unwrap_err();
3561
3562 assert_matches!(err, DepsError::RegistryError { .. });
3563 assert!(
3564 !err.to_string().contains("super-secret-value"),
3565 "err: {err}"
3566 );
3567 }
3568
3569 #[tokio::test]
3570 async fn test_get_cached_with_headers_sends_extra_headers() {
3571 let mut server = mockito::Server::new_async().await;
3572 let url = format!("{}/api/data", server.url());
3573
3574 let _m = server
3575 .mock("GET", "/api/data")
3576 .match_header("authorization", "Bearer token123")
3577 .with_status(200)
3578 .with_header("etag", "\"abc123\"")
3579 .with_body("authed data")
3580 .create_async()
3581 .await;
3582
3583 let cache = HttpCache::new();
3584 let headers = [(header::AUTHORIZATION, "Bearer token123")];
3585 let result: Bytes = cache.get_cached_with_headers(&url, &headers).await.unwrap();
3586
3587 assert_eq!(result.as_ref(), b"authed data");
3588 }
3589
3590 /// The headered form of `get_cached_workspace` used by `deps-npm`'s alternate-registry
3591 /// client (A1): forwards `extra_headers` while still going through the workspace-tier
3592 /// transport (mirrors `test_get_cached_and_get_cached_workspace_do_not_share_an_entry`'s
3593 /// use of the unheadered `get_cached_workspace` against a loopback mockito server under
3594 /// the default policy — an IP-literal host like mockito's has no DNS resolution step for
3595 /// the connect-time `AddrGuard` to intercept, so no policy elevation is needed here
3596 /// either; that guard's actual job is catching a *hostname* that resolves differently at
3597 /// connect time than its parse-time classification, see `validate_resolved_addrs`).
3598 #[tokio::test]
3599 async fn test_get_cached_workspace_with_headers_sends_extra_headers() {
3600 let mut server = mockito::Server::new_async().await;
3601 let url = format!("{}/api/data", server.url());
3602
3603 let _m = server
3604 .mock("GET", "/api/data")
3605 .match_header("accept", "application/vnd.npm.install-v1+json")
3606 .with_status(200)
3607 .with_body("abbreviated packument")
3608 .create_async()
3609 .await;
3610
3611 let cache = HttpCache::new();
3612 let headers = [(header::ACCEPT, "application/vnd.npm.install-v1+json")];
3613 let result: Bytes = cache
3614 .get_cached_workspace_with_headers(&url, &headers)
3615 .await
3616 .unwrap();
3617
3618 assert_eq!(result.as_ref(), b"abbreviated packument");
3619 }
3620
3621 #[tokio::test]
3622 async fn test_fetch_and_store_rejects_oversized_response() {
3623 let mut server = mockito::Server::new_async().await;
3624 let oversized_body = vec![0u8; MAX_RESPONSE_BYTES + 1];
3625
3626 let _m = server
3627 .mock("GET", "/api/huge")
3628 .with_status(200)
3629 .with_body(oversized_body)
3630 .create_async()
3631 .await;
3632
3633 let cache = HttpCache::new();
3634 let url = format!("{}/api/huge", server.url());
3635 let result: Result<Bytes> = cache
3636 .fetch_and_store_with_headers(&url, &[], &cache.baseline.client, &url)
3637 .await;
3638
3639 match result {
3640 Err(DepsError::ResponseTooLarge { limit, .. }) => {
3641 assert_eq!(limit, MAX_RESPONSE_BYTES);
3642 }
3643 other => panic!("expected ResponseTooLarge, got {other:?}"),
3644 }
3645
3646 assert!(cache.entries.get(&url).is_none());
3647 }
3648
3649 #[tokio::test]
3650 async fn test_fetch_and_store_accepts_response_at_exact_cap() {
3651 // MAX_RESPONSE_BYTES is well over MAX_CACHEABLE_ENTRY_BYTES, so the network-layer cap
3652 // and the cache admission cap are independent (see test_store_entry_skips_caching_oversized_entry).
3653 let mut server = mockito::Server::new_async().await;
3654 let exact_cap_body = vec![0u8; MAX_RESPONSE_BYTES];
3655
3656 let _m = server
3657 .mock("GET", "/api/exact")
3658 .with_status(200)
3659 .with_body(exact_cap_body)
3660 .create_async()
3661 .await;
3662
3663 let cache = HttpCache::new();
3664 let url = format!("{}/api/exact", server.url());
3665 let result: Bytes = cache
3666 .fetch_and_store_with_headers(&url, &[], &cache.baseline.client, &url)
3667 .await
3668 .unwrap();
3669
3670 assert_eq!(result.len(), MAX_RESPONSE_BYTES);
3671 }
3672
3673 #[tokio::test]
3674 async fn test_get_cached_non_2xx_on_refresh_preserves_stale_cache() {
3675 let mut server = mockito::Server::new_async().await;
3676 let url = format!("{}/api/data", server.url());
3677
3678 let cache = HttpCache::new();
3679 cache.entries.insert(
3680 url.clone(),
3681 CachedResponse {
3682 body: Bytes::from_static(b"stale but good"),
3683 etag: Some("\"stale-etag\"".into()),
3684 last_modified: None,
3685 fetched_at: Instant::now(),
3686 },
3687 );
3688
3689 // Registry down for maintenance: a non-2xx, non-304 response instead of "unchanged"
3690 // or "here's the new body".
3691 let _m = server
3692 .mock("GET", "/api/data")
3693 .match_header("if-none-match", "\"stale-etag\"")
3694 .with_status(503)
3695 .with_body("<html>maintenance</html>")
3696 .create_async()
3697 .await;
3698
3699 let result: Bytes = cache.get_cached(&url).await.unwrap();
3700
3701 // Stale-while-revalidate: last-known-good body returned, entry untouched.
3702 assert_eq!(result.as_ref(), b"stale but good");
3703 let cached = cache.entries.get(&url).unwrap();
3704 assert_eq!(cached.etag, Some("\"stale-etag\"".into()));
3705 }
3706
3707 /// #756 round 2 S1 regression: the "conditional request failed, using cache" warn (fired
3708 /// on exactly this stale-while-revalidate path) must never interpolate the `DepsError`
3709 /// itself — `DepsError::HttpStatus`'s `Display` embeds the full, unredacted URL (including
3710 /// the query string), which would defeat `RedactedUrl`'s redaction on this same span's
3711 /// `url` field two lines above it. Reuses the mock/seeding shape of
3712 /// `test_get_cached_non_2xx_on_refresh_preserves_stale_cache` with a token-bearing query
3713 /// string, wrapped in a real tracing capture (an actual `warn!` event fires here, unlike
3714 /// the offline-hit case covered by `test_get_cached_span_url_field_redacts_query_string_token`).
3715 #[cfg(feature = "test-util")]
3716 #[tokio::test]
3717 async fn test_get_cached_conditional_request_failure_does_not_leak_token_via_warn() {
3718 let mut server = mockito::Server::new_async().await;
3719 let url = format!("{}/pkg?token=super-secret-value", server.url());
3720
3721 let cache = HttpCache::new();
3722 cache.entries.insert(
3723 url.clone(),
3724 CachedResponse {
3725 body: Bytes::from_static(b"stale but good"),
3726 etag: Some("\"stale-etag\"".into()),
3727 last_modified: None,
3728 fetched_at: Instant::now(),
3729 },
3730 );
3731
3732 let _m = server
3733 .mock("GET", "/pkg")
3734 .match_query(mockito::Matcher::UrlEncoded(
3735 "token".into(),
3736 "super-secret-value".into(),
3737 ))
3738 .match_header("if-none-match", "\"stale-etag\"")
3739 .with_status(503)
3740 .with_body("<html>maintenance</html>")
3741 .create_async()
3742 .await;
3743
3744 let output = crate::test_util::capture_tracing_output_async(async {
3745 let result: Bytes = cache.get_cached(&url).await.unwrap();
3746 assert_eq!(result.as_ref(), b"stale but good");
3747 })
3748 .await;
3749
3750 assert!(
3751 !output.contains("super-secret-value"),
3752 "leaked token via warn! output: {output:?}"
3753 );
3754 }
3755
3756 /// #756 C1 regression: `get_cached_with_headers_via`'s `url` span field must never carry
3757 /// a token embedded in the URL's query string — the same shape as an `.npmrc`-style
3758 /// `${VAR}`-expanded `registry=` URL (see `deps-npm`'s `NpmRegistryIndex` security model).
3759 /// The critic's exact repro was the span's own `url={...}` prefix on a captured log line,
3760 /// so this enables `FmtSpan::NEW` to capture that prefix directly, at span-creation time
3761 /// (before the function body runs), rather than relying on some other event firing inside
3762 /// the span. Deliberately exercises the offline-hit branch — no `DepsError` is ever
3763 /// constructed on that path — so this is isolated from that type's own (separate,
3764 /// pre-existing) URL-embedding `Display` impl and tests the span field in isolation.
3765 #[cfg(feature = "test-util")]
3766 #[tokio::test]
3767 async fn test_get_cached_span_url_field_redacts_query_string_token() {
3768 let url = "https://npm.internal/pkg?token=super-secret-value".to_string();
3769 let cache = HttpCache::new();
3770 cache.set_offline(NetworkMode::Offline);
3771 cache.entries.insert(
3772 url.clone(),
3773 CachedResponse {
3774 body: Bytes::from_static(b"cached"),
3775 etag: None,
3776 last_modified: None,
3777 fetched_at: Instant::now(),
3778 },
3779 );
3780
3781 let output = crate::test_util::capture_tracing_span_fields_async(async {
3782 let result: Bytes = cache.get_cached(&url).await.unwrap();
3783 assert_eq!(result.as_ref(), b"cached");
3784 })
3785 .await;
3786
3787 assert!(
3788 !output.contains("super-secret-value"),
3789 "leaked token into span output: {output:?}"
3790 );
3791 assert!(
3792 output.contains("https://npm.internal/pkg"),
3793 "expected the redacted host+path in span output: {output:?}"
3794 );
3795 }
3796
3797 #[tokio::test]
3798 async fn test_post_json_success_returns_body_and_does_not_cache() {
3799 let mut server = mockito::Server::new_async().await;
3800 let url = format!("{}/v1/querybatch", server.url());
3801
3802 let _m = server
3803 .mock("POST", "/v1/querybatch")
3804 .match_header("content-type", "application/json")
3805 .with_status(200)
3806 .with_body(r#"{"results":[{}]}"#)
3807 .create_async()
3808 .await;
3809
3810 let cache = HttpCache::new();
3811 let body = serde_json::json!({ "queries": [] });
3812 let result: Bytes = cache.post_json(&url, &body).await.unwrap();
3813
3814 assert_eq!(result.as_ref(), br#"{"results":[{}]}"#);
3815 assert!(
3816 cache.is_empty(),
3817 "post_json must not populate the entry-map cache"
3818 );
3819 }
3820
3821 #[tokio::test]
3822 async fn test_post_json_non_2xx_returns_http_status_error() {
3823 let mut server = mockito::Server::new_async().await;
3824 let url = format!("{}/v1/querybatch", server.url());
3825
3826 let _m = server
3827 .mock("POST", "/v1/querybatch")
3828 .with_status(400)
3829 .create_async()
3830 .await;
3831
3832 let cache = HttpCache::new();
3833 let body = serde_json::json!({ "queries": [] });
3834 let result: Result<Bytes> = cache.post_json(&url, &body).await;
3835
3836 match result {
3837 Err(DepsError::HttpStatus { status, .. }) => assert_eq!(status, 400),
3838 other => panic!("expected HttpStatus, got {other:?}"),
3839 }
3840 }
3841
3842 #[tokio::test]
3843 async fn test_post_json_limited_trusted_origin_success_returns_body_and_does_not_cache() {
3844 let mut server = mockito::Server::new_async().await;
3845 let trusted_origin = format!("{}/", server.url());
3846 let url = format!("{}/v3alpha/findings:batchGet", server.url());
3847
3848 let _m = server
3849 .mock("POST", "/v3alpha/findings:batchGet")
3850 .match_header("content-type", "application/json")
3851 .with_status(200)
3852 .with_body(r#"{"findings":[]}"#)
3853 .create_async()
3854 .await;
3855
3856 let cache = HttpCache::new();
3857 let body = serde_json::json!({ "queries": [] });
3858 let result: Bytes = cache
3859 .post_json_limited_trusted_origin(&url, &body, BodyLimit::new(1024), &trusted_origin)
3860 .await
3861 .unwrap();
3862
3863 assert_eq!(result.as_ref(), br#"{"findings":[]}"#);
3864 assert!(
3865 cache.is_empty(),
3866 "post_json_limited_trusted_origin must not populate the entry-map cache"
3867 );
3868 }
3869
3870 #[tokio::test]
3871 async fn test_post_json_limited_trusted_origin_enforces_body_limit() {
3872 let mut server = mockito::Server::new_async().await;
3873 let trusted_origin = format!("{}/", server.url());
3874 let url = format!("{}/v3alpha/findings:batchGet", server.url());
3875
3876 let _m = server
3877 .mock("POST", "/v3alpha/findings:batchGet")
3878 .with_status(200)
3879 .with_body("x".repeat(64))
3880 .create_async()
3881 .await;
3882
3883 let cache = HttpCache::new();
3884 let body = serde_json::json!({ "queries": [] });
3885 let result: Result<Bytes> = cache
3886 .post_json_limited_trusted_origin(&url, &body, BodyLimit::new(8), &trusted_origin)
3887 .await;
3888
3889 // Assert via `matches!` rather than debug-formatting `result` in a panic message
3890 // (#409's established pattern): on the `Ok` arm that value is the raw response
3891 // body, which would otherwise be written to the test log by the panic machinery.
3892 assert!(
3893 matches!(result, Err(DepsError::ResponseTooLarge { .. })),
3894 "expected ResponseTooLarge"
3895 );
3896 }
3897
3898 #[tokio::test]
3899 async fn test_post_json_limited_trusted_origin_rejects_untrusted_redirect() {
3900 let mut server = mockito::Server::new_async().await;
3901 let mut evil = mockito::Server::new_async().await;
3902 let trusted_origin = format!("{}/", server.url());
3903 let url = format!("{}/v3alpha/findings:batchGet", server.url());
3904
3905 let _redirect = server
3906 .mock("POST", "/v3alpha/findings:batchGet")
3907 .with_status(302)
3908 .with_header("location", &format!("{}/steal", evil.url()))
3909 .create_async()
3910 .await;
3911 let evil_call = evil.mock("GET", "/steal").expect(0).create_async().await;
3912
3913 let cache = HttpCache::new();
3914 let body = serde_json::json!({ "queries": [] });
3915 let result: Result<Bytes> = cache
3916 .post_json_limited_trusted_origin(&url, &body, BodyLimit::DEFAULT, &trusted_origin)
3917 .await;
3918
3919 assert!(
3920 result.is_err(),
3921 "an untrusted redirect hop must not be followed"
3922 );
3923 evil_call.assert_async().await;
3924 }
3925
3926 #[tokio::test]
3927 async fn test_get_transport_only_success_returns_body_and_does_not_cache() {
3928 let mut server = mockito::Server::new_async().await;
3929 let url = format!("{}/v1/vulns/RUSTSEC-2020-0071", server.url());
3930
3931 let _m = server
3932 .mock("GET", "/v1/vulns/RUSTSEC-2020-0071")
3933 .with_status(200)
3934 .with_body(r#"{"id":"RUSTSEC-2020-0071"}"#)
3935 .create_async()
3936 .await;
3937
3938 let cache = HttpCache::new();
3939 let result: Bytes = cache.get_transport_only(&url).await.unwrap();
3940
3941 assert_eq!(result.as_ref(), br#"{"id":"RUSTSEC-2020-0071"}"#);
3942 assert!(
3943 cache.is_empty(),
3944 "get_transport_only must not populate the entry-map cache"
3945 );
3946 }
3947
3948 #[tokio::test]
3949 async fn test_get_transport_only_with_headers_sends_extra_headers() {
3950 let mut server = mockito::Server::new_async().await;
3951 let url = format!("{}/v1/vulns/RUSTSEC-2020-0071", server.url());
3952
3953 let _m = server
3954 .mock("GET", "/v1/vulns/RUSTSEC-2020-0071")
3955 .match_header("accept", "application/json")
3956 .with_status(200)
3957 .with_body(r#"{"id":"RUSTSEC-2020-0071"}"#)
3958 .create_async()
3959 .await;
3960
3961 let cache = HttpCache::new();
3962 let headers = [(header::ACCEPT, "application/json")];
3963 let result: Bytes = cache
3964 .get_transport_only_with_headers(&url, &headers)
3965 .await
3966 .unwrap();
3967
3968 assert_eq!(result.as_ref(), br#"{"id":"RUSTSEC-2020-0071"}"#);
3969 assert!(
3970 cache.is_empty(),
3971 "get_transport_only_with_headers must not populate the entry-map cache"
3972 );
3973 }
3974
3975 #[tokio::test]
3976 async fn test_get_transport_only_non_2xx_returns_http_status_error() {
3977 let mut server = mockito::Server::new_async().await;
3978 let url = format!("{}/v1/vulns/missing", server.url());
3979
3980 let _m = server
3981 .mock("GET", "/v1/vulns/missing")
3982 .with_status(404)
3983 .create_async()
3984 .await;
3985
3986 let cache = HttpCache::new();
3987 let result: Result<Bytes> = cache.get_transport_only(&url).await;
3988
3989 match result {
3990 Err(DepsError::HttpStatus { status, .. }) => assert_eq!(status, 404),
3991 other => panic!("expected HttpStatus, got {other:?}"),
3992 }
3993 }
3994
3995 fn dummy_response(size: usize) -> CachedResponse {
3996 CachedResponse {
3997 body: Bytes::from(vec![0u8; size]),
3998 etag: None,
3999 last_modified: None,
4000 fetched_at: Instant::now(),
4001 }
4002 }
4003
4004 #[test]
4005 fn test_total_bytes_tracks_inserts_and_replacement() {
4006 let cache = HttpCache::new();
4007 cache.store_entry("url1".into(), dummy_response(100));
4008 assert_eq!(cache.total_bytes(), 100);
4009
4010 // Replacing the same key must account for the delta, not just add.
4011 cache.store_entry("url1".into(), dummy_response(40));
4012 assert_eq!(cache.total_bytes(), 40);
4013
4014 cache.store_entry("url2".into(), dummy_response(60));
4015 assert_eq!(cache.total_bytes(), 100);
4016 }
4017
4018 #[test]
4019 fn test_clear_resets_total_bytes() {
4020 let cache = HttpCache::new();
4021 cache.store_entry("url1".into(), dummy_response(1000));
4022 assert_eq!(cache.total_bytes(), 1000);
4023
4024 cache.clear();
4025 assert_eq!(cache.total_bytes(), 0);
4026 }
4027
4028 #[test]
4029 fn test_small_payloads_do_not_trigger_eviction() {
4030 let cache = HttpCache::new();
4031 for i in 0..50 {
4032 cache.store_entry(format!("url{i}"), dummy_response(1024));
4033 }
4034
4035 assert_eq!(cache.len(), 50);
4036 assert_eq!(cache.total_bytes(), 50 * 1024);
4037 }
4038
4039 #[test]
4040 fn test_evict_entries_triggers_on_byte_budget_with_few_entries() {
4041 let cache = HttpCache::new();
4042
4043 // 9 entries at the per-entry admission cap: far below MAX_CACHE_ENTRIES by count, but
4044 // their combined size (72 MiB) overshoots MAX_CACHE_BYTES (64 MiB).
4045 for i in 0..9 {
4046 cache.store_entry(format!("url{i}"), dummy_response(MAX_CACHEABLE_ENTRY_BYTES));
4047 }
4048 assert_eq!(cache.len(), 9);
4049 assert!(cache.total_bytes() > MAX_CACHE_BYTES);
4050
4051 cache.evict_entries();
4052
4053 // Only as many oldest entries as needed to clear the byte budget are removed, not a
4054 // fixed count-based batch.
4055 assert!(cache.total_bytes() <= MAX_CACHE_BYTES);
4056 assert_eq!(cache.len(), 8);
4057 }
4058
4059 #[test]
4060 fn test_evict_entries_removes_oldest_first_for_bytes() {
4061 let cache = HttpCache::new();
4062
4063 // Evicting just the single oldest entry should restore the cache to within budget,
4064 // proving eviction picks the genuinely oldest entry, not hash-iteration order.
4065 cache.store_entry("oldest".into(), dummy_response(MAX_CACHEABLE_ENTRY_BYTES));
4066 std::thread::sleep(std::time::Duration::from_millis(5));
4067 for i in 0..8 {
4068 cache.store_entry(
4069 format!("newer{i}"),
4070 dummy_response(MAX_CACHEABLE_ENTRY_BYTES),
4071 );
4072 }
4073 assert_eq!(cache.len(), 9);
4074
4075 cache.evict_entries();
4076
4077 assert_eq!(cache.len(), 8);
4078 assert!(cache.entries.get("oldest").is_none());
4079 for i in 0..8 {
4080 assert!(cache.entries.get(&format!("newer{i}")).is_some());
4081 }
4082 }
4083
4084 #[tokio::test]
4085 async fn test_get_cached_with_headers_evicts_on_byte_budget() {
4086 let mut server = mockito::Server::new_async().await;
4087 let url = format!("{}/api/data", server.url());
4088
4089 let cache = HttpCache::new();
4090
4091 // Pre-fill past the byte budget, staying under MAX_CACHE_ENTRIES by count.
4092 for i in 0..9 {
4093 cache.store_entry(
4094 format!("stale{i}"),
4095 dummy_response(MAX_CACHEABLE_ENTRY_BYTES),
4096 );
4097 }
4098 assert!(cache.total_bytes() > MAX_CACHE_BYTES);
4099
4100 let _m = server
4101 .mock("GET", "/api/data")
4102 .with_status(200)
4103 .with_body("fresh")
4104 .create_async()
4105 .await;
4106
4107 let result: Bytes = cache.get_cached(&url).await.unwrap();
4108 assert_eq!(result.as_ref(), b"fresh");
4109
4110 // The pre-request byte-budget check evicted stale entries before fetching.
4111 assert!(cache.total_bytes() <= MAX_CACHE_BYTES + result.len());
4112 }
4113
4114 #[test]
4115 fn test_store_entry_skips_caching_oversized_entry() {
4116 let cache = HttpCache::new();
4117
4118 // A body over the per-entry admission cap is not retained, even though the caller
4119 // still gets it back (store_entry's caller already holds `body` independently).
4120 cache.store_entry("big".into(), dummy_response(MAX_CACHEABLE_ENTRY_BYTES + 1));
4121 assert!(cache.entries.get("big").is_none());
4122 assert_eq!(cache.total_bytes(), 0);
4123
4124 // Replacing an existing small entry with an oversized one drops the stale entry too,
4125 // rather than leaving it to keep serving increasingly outdated data.
4126 cache.store_entry("small".into(), dummy_response(100));
4127 assert_eq!(cache.total_bytes(), 100);
4128
4129 cache.store_entry(
4130 "small".into(),
4131 dummy_response(MAX_CACHEABLE_ENTRY_BYTES + 1),
4132 );
4133 assert!(cache.entries.get("small").is_none());
4134 assert_eq!(cache.total_bytes(), 0);
4135 }
4136
4137 #[test]
4138 fn test_concurrent_store_and_evict_keeps_total_bytes_consistent() {
4139 use std::sync::Arc;
4140 use std::thread;
4141
4142 let cache = Arc::new(HttpCache::new());
4143 let handles: Vec<_> = (0..8)
4144 .map(|t| {
4145 let cache = Arc::clone(&cache);
4146 thread::spawn(move || {
4147 for i in 0..150 {
4148 cache.store_entry(format!("t{t}-{i}"), dummy_response(4096));
4149 if i % 10 == 0 {
4150 cache.evict_entries();
4151 }
4152 }
4153 })
4154 })
4155 .collect();
4156
4157 for handle in handles {
4158 handle.join().unwrap();
4159 }
4160 cache.evict_entries();
4161
4162 // Regression guard: evict_entries used to snapshot total_bytes once and overwrite it
4163 // with an absolute store, silently discarding any concurrent store_entry delta.
4164 let actual: usize = cache
4165 .entries
4166 .iter()
4167 .map(|entry| entry.value().body.len())
4168 .sum();
4169 assert_eq!(cache.total_bytes(), actual);
4170 }
4171
4172 // Issue #455, test-plan item 5: `get_cached(url)` then `get_cached_workspace(url)` do not
4173 // share an entry — each is keyed under a distinct namespace (see `HttpCache::cache_key`),
4174 // so the mockito mock is hit twice and the cache ends up with two entries for one URL.
4175 #[tokio::test]
4176 async fn test_get_cached_and_get_cached_workspace_do_not_share_an_entry() {
4177 let mut server = mockito::Server::new_async().await;
4178 let url = format!("{}/api/data", server.url());
4179
4180 let mock = server
4181 .mock("GET", "/api/data")
4182 .with_status(200)
4183 .with_body("shared url, distinct tiers")
4184 .expect(2)
4185 .create_async()
4186 .await;
4187
4188 let cache = HttpCache::new();
4189 let baseline_result: Bytes = cache.get_cached(&url).await.unwrap();
4190 let workspace_result: Bytes = cache.get_cached_workspace(&url).await.unwrap();
4191
4192 assert_eq!(baseline_result.as_ref(), b"shared url, distinct tiers");
4193 assert_eq!(workspace_result.as_ref(), b"shared url, distinct tiers");
4194 assert_eq!(cache.len(), 2);
4195 mock.assert_async().await;
4196 }
4197
4198 // Issue #455, test-plan item 6 (C5): fetch under `All`, tighten to `PublicOnly`, re-fetch
4199 // the same URL — the `All`-era body must not be served, since the policy-scoped key
4200 // namespace changes with the policy.
4201 #[tokio::test]
4202 async fn test_set_registry_policy_change_does_not_serve_stale_era_body() {
4203 let mut server = mockito::Server::new_async().await;
4204 let url = format!("{}/api/data", server.url());
4205
4206 let _first = server
4207 .mock("GET", "/api/data")
4208 .with_status(200)
4209 .with_body("all-era body")
4210 .create_async()
4211 .await;
4212
4213 let policy = Arc::new(RegistryAccessPolicy::new(WorkspaceRegistryAccess::All));
4214 let cache = HttpCache::with_policy(Arc::clone(&policy));
4215 let first: Bytes = cache.get_cached_workspace(&url).await.unwrap();
4216 assert_eq!(first.as_ref(), b"all-era body");
4217
4218 cache.set_registry_policy(WorkspaceRegistryAccess::PublicOnly);
4219
4220 let _second = server
4221 .mock("GET", "/api/data")
4222 .with_status(200)
4223 .with_body("public-only-era body")
4224 .create_async()
4225 .await;
4226
4227 let second: Bytes = cache.get_cached_workspace(&url).await.unwrap();
4228 assert_eq!(
4229 second.as_ref(),
4230 b"public-only-era body",
4231 "the All-era cached body must not be served after tightening to PublicOnly"
4232 );
4233 // mockito's second mock answers regardless of which entry was hit, so the body
4234 // assertion alone wouldn't prove a policy-blind key — this length check does.
4235 assert_eq!(cache.len(), 2);
4236 }
4237
4238 // Issue #455, test-plan item 7 (C4): `set_registry_policy` rebuilds the workspace transport
4239 // only on an actual change, not on a no-op re-application of the same value.
4240 #[test]
4241 fn test_set_registry_policy_rebuilds_only_on_change() {
4242 let policy = Arc::new(RegistryAccessPolicy::new(
4243 WorkspaceRegistryAccess::PublicOnly,
4244 ));
4245 let cache = HttpCache::with_policy(policy);
4246 assert_eq!(cache.workspace_rebuilds.load(Ordering::Relaxed), 0);
4247
4248 cache.set_registry_policy(WorkspaceRegistryAccess::PublicOnly);
4249 assert_eq!(
4250 cache.workspace_rebuilds.load(Ordering::Relaxed),
4251 0,
4252 "re-applying the unchanged policy must not rebuild the workspace transport"
4253 );
4254
4255 cache.set_registry_policy(WorkspaceRegistryAccess::All);
4256 assert_eq!(cache.workspace_rebuilds.load(Ordering::Relaxed), 1);
4257
4258 cache.set_registry_policy(WorkspaceRegistryAccess::All);
4259 assert_eq!(
4260 cache.workspace_rebuilds.load(Ordering::Relaxed),
4261 1,
4262 "re-applying the unchanged (new) policy must not rebuild again"
4263 );
4264
4265 cache.set_registry_policy(WorkspaceRegistryAccess::Off);
4266 assert_eq!(cache.workspace_rebuilds.load(Ordering::Relaxed), 2);
4267 }
4268
4269 // Issue #483: `.expect(0)` proves nothing reached *this mock*, not that zero sockets
4270 // ever opened — adequate only because the loopback carve-out lets mockito stand in here.
4271
4272 #[tokio::test]
4273 async fn test_offline_cold_get_cached_errors_without_network() {
4274 let mut server = mockito::Server::new_async().await;
4275 let url = format!("{}/api/data", server.url());
4276 let mock = server
4277 .mock("GET", "/api/data")
4278 .with_status(200)
4279 .with_body("must not be fetched")
4280 .expect(0)
4281 .create_async()
4282 .await;
4283
4284 let cache = HttpCache::new();
4285 cache.set_offline(NetworkMode::Offline);
4286
4287 let result: Result<Bytes> = cache.get_cached(&url).await;
4288 match result {
4289 Err(DepsError::Offline { url: blocked }) => {
4290 assert_eq!(blocked, RedactedUrl::new(&url));
4291 }
4292 other => panic!("expected Offline, got {other:?}"),
4293 }
4294 mock.assert_async().await;
4295 }
4296
4297 #[tokio::test]
4298 async fn test_offline_cold_post_json_errors_without_network() {
4299 let mut server = mockito::Server::new_async().await;
4300 let url = format!("{}/v1/querybatch", server.url());
4301 let mock = server
4302 .mock("POST", "/v1/querybatch")
4303 .with_status(200)
4304 .expect(0)
4305 .create_async()
4306 .await;
4307
4308 let cache = HttpCache::new();
4309 cache.set_offline(NetworkMode::Offline);
4310
4311 let body = serde_json::json!({ "queries": [] });
4312 let result: Result<Bytes> = cache.post_json(&url, &body).await;
4313 assert_matches!(result, Err(DepsError::Offline { .. }));
4314 mock.assert_async().await;
4315 }
4316
4317 #[tokio::test]
4318 async fn test_offline_cold_get_transport_only_errors_without_network() {
4319 let mut server = mockito::Server::new_async().await;
4320 let url = format!("{}/v1/vulns/RUSTSEC-2020-0071", server.url());
4321 let mock = server
4322 .mock("GET", "/v1/vulns/RUSTSEC-2020-0071")
4323 .with_status(200)
4324 .expect(0)
4325 .create_async()
4326 .await;
4327
4328 let cache = HttpCache::new();
4329 cache.set_offline(NetworkMode::Offline);
4330
4331 let result: Result<Bytes> = cache.get_transport_only(&url).await;
4332 assert_matches!(result, Err(DepsError::Offline { .. }));
4333 mock.assert_async().await;
4334 }
4335
4336 #[tokio::test]
4337 async fn test_offline_warm_serves_cached_body_without_network() {
4338 let mut server = mockito::Server::new_async().await;
4339 let url = format!("{}/api/data", server.url());
4340 let mock = server
4341 .mock("GET", "/api/data")
4342 .with_status(200)
4343 .with_body("must not be fetched")
4344 .expect(0)
4345 .create_async()
4346 .await;
4347
4348 let cache = HttpCache::new();
4349 cache.entries.insert(
4350 url.clone(),
4351 CachedResponse {
4352 body: Bytes::from_static(b"warm cached body"),
4353 etag: Some("\"tag123\"".into()),
4354 last_modified: None,
4355 fetched_at: Instant::now(),
4356 },
4357 );
4358 cache.set_offline(NetworkMode::Offline);
4359
4360 let result: Bytes = cache.get_cached(&url).await.unwrap();
4361 assert_eq!(result.as_ref(), b"warm cached body");
4362 mock.assert_async().await;
4363 }
4364
4365 #[tokio::test]
4366 async fn test_cache_disabled_two_calls_each_hit_the_server() {
4367 let mut server = mockito::Server::new_async().await;
4368 let url = format!("{}/api/data", server.url());
4369 let mock = server
4370 .mock("GET", "/api/data")
4371 .with_status(200)
4372 .with_body("fresh every time")
4373 .expect(2)
4374 .create_async()
4375 .await;
4376
4377 let cache = HttpCache::new();
4378 cache.set_cache_enabled(CacheMode::Disabled);
4379
4380 let first: Bytes = cache.get_cached(&url).await.unwrap();
4381 let second: Bytes = cache.get_cached(&url).await.unwrap();
4382 assert_eq!(first.as_ref(), b"fresh every time");
4383 assert_eq!(second.as_ref(), b"fresh every time");
4384 assert!(
4385 cache.is_empty(),
4386 "cache.enabled: false must never populate the entry map"
4387 );
4388 mock.assert_async().await;
4389 }
4390
4391 // S1 fix (critic-corrected): the naive design's `!cache_enabled` bypass ran *before*
4392 // any offline check, so `cache.enabled: false` + `network.offline: true` cold always
4393 // took the network-only bypass path — which `ensure_online` then blocked — even though
4394 // this combination is meant to still surface a clean, immediate signal rather than
4395 // hang or silently return empty data forever.
4396 #[tokio::test]
4397 async fn test_offline_and_cache_disabled_cold_start_errors_cleanly() {
4398 let mut server = mockito::Server::new_async().await;
4399 let url = format!("{}/api/data", server.url());
4400 let mock = server
4401 .mock("GET", "/api/data")
4402 .with_status(200)
4403 .expect(0)
4404 .create_async()
4405 .await;
4406
4407 let cache = HttpCache::new();
4408 cache.set_cache_enabled(CacheMode::Disabled);
4409 cache.set_offline(NetworkMode::Offline);
4410
4411 let result: Result<Bytes> = cache.get_cached(&url).await;
4412 assert_matches!(result, Err(DepsError::Offline { .. }));
4413 mock.assert_async().await;
4414 }
4415
4416 // The scenario S1 actually exists to fix: an entry stored while caching was enabled
4417 // must still be servable once `cache.enabled` is later turned off *and* the cache goes
4418 // offline in the same breath — proving `offline` truly overrides `cache_enabled` on the
4419 // read path, not just when the two flags never change together.
4420 #[tokio::test]
4421 async fn test_offline_overrides_disabled_cache_to_serve_warm_entry() {
4422 let mut server = mockito::Server::new_async().await;
4423 let url = format!("{}/api/data", server.url());
4424 let mock = server
4425 .mock("GET", "/api/data")
4426 .with_status(200)
4427 .with_body("fetched while online")
4428 .expect(1)
4429 .create_async()
4430 .await;
4431
4432 let cache = HttpCache::new();
4433 let first: Bytes = cache.get_cached(&url).await.unwrap();
4434 assert_eq!(first.as_ref(), b"fetched while online");
4435
4436 cache.set_cache_enabled(CacheMode::Disabled);
4437 cache.set_offline(NetworkMode::Offline);
4438
4439 let second: Bytes = cache.get_cached(&url).await.unwrap();
4440 assert_eq!(
4441 second.as_ref(),
4442 b"fetched while online",
4443 "offline must override cache_enabled:false and still serve the warm entry"
4444 );
4445 mock.assert_async().await;
4446 }
4447
4448 // The primary UX case (critic M6a): a full online -> offline transition on an
4449 // otherwise-default cache (cache.enabled stays true throughout) must keep serving what
4450 // was already fetched.
4451 #[tokio::test]
4452 async fn test_online_to_offline_transition_serves_previously_fetched_entry() {
4453 let mut server = mockito::Server::new_async().await;
4454 let url = format!("{}/api/data", server.url());
4455 let mock = server
4456 .mock("GET", "/api/data")
4457 .with_status(200)
4458 .with_body("fetched while online")
4459 .expect(1)
4460 .create_async()
4461 .await;
4462
4463 let cache = HttpCache::new();
4464 let online: Bytes = cache.get_cached(&url).await.unwrap();
4465 assert_eq!(online.as_ref(), b"fetched while online");
4466
4467 cache.set_offline(NetworkMode::Offline);
4468
4469 let offline: Bytes = cache.get_cached(&url).await.unwrap();
4470 assert_eq!(offline.as_ref(), b"fetched while online");
4471 mock.assert_async().await;
4472 }
4473
4474 // Offline -> online restores live fetches (the flag's other half of critic M6a),
4475 // exercised here through `get_cached`'s conditional-revalidation path directly (the
4476 // live `did_change_configuration` toggle is covered by `deps-lsp`'s own test).
4477 #[tokio::test]
4478 async fn test_offline_to_online_transition_resumes_fetching() {
4479 let mut server = mockito::Server::new_async().await;
4480 let url = format!("{}/api/data", server.url());
4481 let mock = server
4482 .mock("GET", "/api/data")
4483 .with_status(200)
4484 .with_header("etag", "\"abc123\"")
4485 .with_body("fetched while online")
4486 .expect(1)
4487 .create_async()
4488 .await;
4489
4490 let cache = HttpCache::new();
4491 let online: Bytes = cache.get_cached(&url).await.unwrap();
4492 assert_eq!(online.as_ref(), b"fetched while online");
4493 mock.assert_async().await;
4494
4495 cache.set_offline(NetworkMode::Offline);
4496 let offline: Bytes = cache.get_cached(&url).await.unwrap();
4497 assert_eq!(offline.as_ref(), b"fetched while online");
4498
4499 cache.set_offline(NetworkMode::Online);
4500 drop(mock);
4501 let revalidate = server
4502 .mock("GET", "/api/data")
4503 .match_header("if-none-match", "\"abc123\"")
4504 .with_status(304)
4505 .expect(1)
4506 .create_async()
4507 .await;
4508 let resumed: Bytes = cache.get_cached(&url).await.unwrap();
4509 assert_eq!(
4510 resumed.as_ref(),
4511 b"fetched while online",
4512 "returning online must resume live requests, not stay pinned to the cached body"
4513 );
4514 revalidate.assert_async().await;
4515 }
4516
4517 // Critic M6b: `get_cached_workspace` and `get_cached_trusted_origin` key entries under
4518 // distinct namespaces from `get_cached`'s baseline tier — the offline warm-cache path
4519 // needs its own proof it holds for each.
4520 #[tokio::test]
4521 async fn test_offline_warm_serves_workspace_tier_without_network() {
4522 let mut server = mockito::Server::new_async().await;
4523 let url = format!("{}/api/data", server.url());
4524 let mock = server
4525 .mock("GET", "/api/data")
4526 .with_status(200)
4527 .with_body("workspace fetch")
4528 .expect(1)
4529 .create_async()
4530 .await;
4531
4532 let cache = HttpCache::new();
4533 let online: Bytes = cache.get_cached_workspace(&url).await.unwrap();
4534 assert_eq!(online.as_ref(), b"workspace fetch");
4535
4536 cache.set_offline(NetworkMode::Offline);
4537 let offline: Bytes = cache.get_cached_workspace(&url).await.unwrap();
4538 assert_eq!(offline.as_ref(), b"workspace fetch");
4539 mock.assert_async().await;
4540 }
4541
4542 #[tokio::test]
4543 async fn test_offline_warm_serves_trusted_origin_tier_without_network() {
4544 let mut server = mockito::Server::new_async().await;
4545 let trusted_origin = format!("{}/", server.url());
4546 let url = format!("{}/api/data", server.url());
4547 let mock = server
4548 .mock("GET", "/api/data")
4549 .with_status(200)
4550 .with_body("trusted-origin fetch")
4551 .expect(1)
4552 .create_async()
4553 .await;
4554
4555 let cache = HttpCache::new();
4556 let online: Bytes = cache
4557 .get_cached_trusted_origin(&url, &trusted_origin)
4558 .await
4559 .unwrap();
4560 assert_eq!(online.as_ref(), b"trusted-origin fetch");
4561
4562 cache.set_offline(NetworkMode::Offline);
4563 let offline: Bytes = cache
4564 .get_cached_trusted_origin(&url, &trusted_origin)
4565 .await
4566 .unwrap();
4567 assert_eq!(offline.as_ref(), b"trusted-origin fetch");
4568 mock.assert_async().await;
4569 }
4570
4571 // --- issue #561/#562: CacheTier::Pinned, get_cached_pinned{,_with_headers} ---
4572
4573 #[tokio::test]
4574 async fn test_get_cached_pinned_attaches_auth_header() {
4575 let mut server = mockito::Server::new_async().await;
4576 let trusted_origin = format!("{}/", server.url());
4577 let url = format!("{}/api/data", server.url());
4578
4579 let _m = server
4580 .mock("GET", "/api/data")
4581 .match_header("authorization", "Basic dXNlcjpwYXQ=")
4582 .with_status(200)
4583 .with_body("authenticated data")
4584 .create_async()
4585 .await;
4586
4587 let cache = HttpCache::new();
4588 let result = cache
4589 .get_cached_pinned_with_headers(
4590 &url,
4591 &trusted_origin,
4592 true,
4593 Some(42),
4594 &[(header::AUTHORIZATION, "Basic dXNlcjpwYXQ=")],
4595 )
4596 .await
4597 .unwrap();
4598
4599 assert_eq!(result.as_ref(), b"authenticated data");
4600 }
4601
4602 /// FR-014: distinct `auth_id` values against the same `(url, trusted_origin)` never share
4603 /// a cache entry — a rotated or distinct credential never reads back a body fetched under a
4604 /// different one.
4605 #[tokio::test]
4606 async fn test_get_cached_pinned_distinct_auth_id_never_shares_cache_entry() {
4607 let mut server = mockito::Server::new_async().await;
4608 let trusted_origin = format!("{}/", server.url());
4609 let url = format!("{}/api/data", server.url());
4610
4611 let _m1 = server
4612 .mock("GET", "/api/data")
4613 .with_status(200)
4614 .with_body("body-for-credential-a")
4615 .create_async()
4616 .await;
4617
4618 let cache = HttpCache::new();
4619 let a = cache
4620 .get_cached_pinned(&url, &trusted_origin, true, Some(1))
4621 .await
4622 .unwrap();
4623 assert_eq!(a.as_ref(), b"body-for-credential-a");
4624 drop(_m1);
4625
4626 let _m2 = server
4627 .mock("GET", "/api/data")
4628 .with_status(200)
4629 .with_body("body-for-credential-b")
4630 .create_async()
4631 .await;
4632
4633 let b = cache
4634 .get_cached_pinned(&url, &trusted_origin, true, Some(2))
4635 .await
4636 .unwrap();
4637 assert_eq!(
4638 b.as_ref(),
4639 b"body-for-credential-b",
4640 "a distinct auth_id must not read back credential A's cached body"
4641 );
4642 }
4643
4644 /// #1025: `auth_id: None` (unauthenticated) and `auth_id: Some(0)` (a credential whose
4645 /// salted digest happens to hash to exactly `0`) must produce different `Pinned`-tier
4646 /// cache keys — collapsing `None` to the same `0` sentinel `Some(0)` digests to would let
4647 /// a cached authenticated body be served to a subsequent unauthenticated request, or vice
4648 /// versa.
4649 #[test]
4650 fn test_cache_key_pinned_none_and_some_zero_auth_id_never_collide() {
4651 let cache = HttpCache::new();
4652 let tier = CacheTier::Pinned {
4653 digest: 42,
4654 authenticated: true,
4655 };
4656
4657 let unauthenticated = cache.cache_key("https://example.com/pkg", tier, None);
4658 let zero_digest_credential = cache.cache_key("https://example.com/pkg", tier, Some(0));
4659
4660 assert_ne!(unauthenticated, zero_digest_credential);
4661 }
4662
4663 /// #1025 regression guard: non-zero `auth_id` values keep producing the pre-existing,
4664 /// distinct-per-id cache key behavior (only the `None` vs `Some(0)` collision was fixed).
4665 #[test]
4666 fn test_cache_key_pinned_nonzero_auth_id_still_distinct() {
4667 let cache = HttpCache::new();
4668 let tier = CacheTier::Pinned {
4669 digest: 42,
4670 authenticated: true,
4671 };
4672
4673 let a = cache.cache_key("https://example.com/pkg", tier, Some(1));
4674 let b = cache.cache_key("https://example.com/pkg", tier, Some(2));
4675 let unauthenticated = cache.cache_key("https://example.com/pkg", tier, None);
4676
4677 assert_ne!(a, b);
4678 assert_ne!(a, unauthenticated);
4679 assert_ne!(b, unauthenticated);
4680 }
4681
4682 /// #1025 M2: proves the fixed-width invariant itself, not just a couple of hand-picked
4683 /// values — a regression that dropped zero-padding (e.g. `{digest:x}` instead of
4684 /// `{digest:016x}`) would still pass the two tests above (`None != Some(0)`,
4685 /// `Some(1) != Some(2)`) but would let differently-shaped `(digest, auth_id)` pairs collide,
4686 /// e.g. `digest=0x1, auth_id=Some(0x11)` vs. `digest=0x11, auth_id=Some(0x1)` concatenating
4687 /// to the same string once padding is gone. Every `(digest, auth_id)` combination drawn
4688 /// from these boundary-value sets must produce a pairwise-distinct key.
4689 #[test]
4690 fn test_cache_key_pinned_digest_and_auth_id_matrix_never_collide() {
4691 let cache = HttpCache::new();
4692 let digests = [0u64, 1, 0x11, u64::MAX];
4693 let auth_ids = [None, Some(0u64), Some(1), Some(0x11), Some(u64::MAX)];
4694
4695 let mut keys = Vec::new();
4696 for &digest in &digests {
4697 for &auth_id in &auth_ids {
4698 let tier = CacheTier::Pinned {
4699 digest,
4700 authenticated: true,
4701 };
4702 keys.push((
4703 (digest, auth_id),
4704 cache
4705 .cache_key("https://example.com/pkg", tier, auth_id)
4706 .into_owned(),
4707 ));
4708 }
4709 }
4710
4711 for i in 0..keys.len() {
4712 for j in (i + 1)..keys.len() {
4713 assert_ne!(
4714 keys[i].1, keys[j].1,
4715 "inputs {:?} and {:?} produced the same cache key",
4716 keys[i].0, keys[j].0
4717 );
4718 }
4719 }
4720 }
4721
4722 /// #1025 M1: `authenticated` is folded into the `Pinned`-tier cache key too, not just
4723 /// `auth_id` — otherwise `(authenticated: true, auth_id: None)` and
4724 /// `(authenticated: false, auth_id: None)` against the same origin would collide. Not
4725 /// reachable through any current caller (every ecosystem correlates the two), but
4726 /// `get_cached_pinned_with_headers` is public API with no enforced invariant tying them
4727 /// together, so the key format must not assume one.
4728 #[test]
4729 fn test_cache_key_pinned_authenticated_flag_never_collides_with_unauthenticated() {
4730 let cache = HttpCache::new();
4731 let authenticated_tier = CacheTier::Pinned {
4732 digest: 42,
4733 authenticated: true,
4734 };
4735 let unauthenticated_tier = CacheTier::Pinned {
4736 digest: 42,
4737 authenticated: false,
4738 };
4739
4740 let a = cache.cache_key("https://example.com/pkg", authenticated_tier, None);
4741 let b = cache.cache_key("https://example.com/pkg", unauthenticated_tier, None);
4742
4743 assert_ne!(a, b);
4744 }
4745
4746 /// FR-015/NFR-004: a 401 revalidation response against an authenticated `Pinned`-tier
4747 /// entry evicts the entry and returns the error — never the default
4748 /// stale-while-revalidate fallback that would serve the possibly-revoked credential's
4749 /// last-known-good body.
4750 #[tokio::test]
4751 async fn test_pinned_authenticated_401_revalidation_evicts_instead_of_stale_serve() {
4752 let mut server = mockito::Server::new_async().await;
4753 let trusted_origin = format!("{}/", server.url());
4754 let url = format!("{}/api/data", server.url());
4755
4756 let _m1 = server
4757 .mock("GET", "/api/data")
4758 .with_status(200)
4759 .with_header("etag", "\"abc123\"")
4760 .with_body("private data")
4761 .create_async()
4762 .await;
4763
4764 let cache = HttpCache::new();
4765 let first = cache
4766 .get_cached_pinned(&url, &trusted_origin, true, Some(7))
4767 .await
4768 .unwrap();
4769 assert_eq!(first.as_ref(), b"private data");
4770 assert_eq!(cache.len(), 1);
4771 drop(_m1);
4772
4773 let _m2 = server
4774 .mock("GET", "/api/data")
4775 .match_header("if-none-match", "\"abc123\"")
4776 .with_status(401)
4777 .create_async()
4778 .await;
4779
4780 let result = cache
4781 .get_cached_pinned(&url, &trusted_origin, true, Some(7))
4782 .await;
4783
4784 assert!(
4785 matches!(result, Err(DepsError::HttpStatus { status: 401, .. })),
4786 "expected the 401 to surface as an error, not a stale-served body: {result:?}"
4787 );
4788 assert_eq!(
4789 cache.len(),
4790 0,
4791 "the revoked-credential entry must be evicted, not left cached"
4792 );
4793 }
4794
4795 /// #1295 critic C1 regression test: a 403 revalidation carrying confirmed
4796 /// `X-RateLimit-Remaining: 0` evidence now classifies as `DepsError::RateLimited` rather
4797 /// than `HttpStatus` (see `http_status_error`) — without the eviction guard's
4798 /// `RateLimited` arm, this would silently bypass FR-015/NFR-004 and keep serving the
4799 /// possibly-revoked credential's stale body.
4800 #[tokio::test]
4801 async fn test_pinned_authenticated_403_with_confirmed_evidence_still_evicts() {
4802 let mut server = mockito::Server::new_async().await;
4803 let trusted_origin = format!("{}/", server.url());
4804 let url = format!("{}/api/data", server.url());
4805
4806 let _m1 = server
4807 .mock("GET", "/api/data")
4808 .with_status(200)
4809 .with_header("etag", "\"abc123\"")
4810 .with_body("private data")
4811 .create_async()
4812 .await;
4813
4814 let cache = HttpCache::new();
4815 let first = cache
4816 .get_cached_pinned(&url, &trusted_origin, true, Some(7))
4817 .await
4818 .unwrap();
4819 assert_eq!(first.as_ref(), b"private data");
4820 assert_eq!(cache.len(), 1);
4821 drop(_m1);
4822
4823 let _m2 = server
4824 .mock("GET", "/api/data")
4825 .match_header("if-none-match", "\"abc123\"")
4826 .with_status(403)
4827 .with_header("x-ratelimit-remaining", "0")
4828 .create_async()
4829 .await;
4830
4831 let result = cache
4832 .get_cached_pinned(&url, &trusted_origin, true, Some(7))
4833 .await;
4834
4835 assert!(
4836 matches!(
4837 result,
4838 Err(DepsError::RateLimited {
4839 verified: RateLimitEvidence::Confirmed,
4840 ..
4841 })
4842 ),
4843 "expected the confirmed-evidence 403 to surface as verified RateLimited: {result:?}"
4844 );
4845 assert_eq!(
4846 cache.len(),
4847 0,
4848 "the revoked-credential entry must still be evicted on a confirmed-evidence 403, \
4849 not left cached"
4850 );
4851 }
4852
4853 /// #1295 critic N1 regression test: a confirmed-evidence 429 (mere throttling, not a
4854 /// credential-revocation signal) must NOT evict an authenticated pinned-tier entry —
4855 /// NFR-004's scope is 401/403 only. Counterpart to the 403 test above: same setup, only
4856 /// the revalidation status differs, and the assertion flips (kept, not evicted).
4857 #[tokio::test]
4858 async fn test_pinned_authenticated_429_with_confirmed_evidence_does_not_evict() {
4859 let mut server = mockito::Server::new_async().await;
4860 let trusted_origin = format!("{}/", server.url());
4861 let url = format!("{}/api/data", server.url());
4862
4863 let _m1 = server
4864 .mock("GET", "/api/data")
4865 .with_status(200)
4866 .with_header("etag", "\"abc123\"")
4867 .with_body("private data")
4868 .create_async()
4869 .await;
4870
4871 let cache = HttpCache::new();
4872 let first = cache
4873 .get_cached_pinned(&url, &trusted_origin, true, Some(7))
4874 .await
4875 .unwrap();
4876 assert_eq!(first.as_ref(), b"private data");
4877 assert_eq!(cache.len(), 1);
4878 drop(_m1);
4879
4880 let _m2 = server
4881 .mock("GET", "/api/data")
4882 .match_header("if-none-match", "\"abc123\"")
4883 .with_status(429)
4884 .with_header("retry-after", "60")
4885 .create_async()
4886 .await;
4887
4888 let result = cache
4889 .get_cached_pinned(&url, &trusted_origin, true, Some(7))
4890 .await;
4891
4892 assert_eq!(
4893 result.unwrap().as_ref(),
4894 b"private data",
4895 "a confirmed-evidence 429 must still fall back to stale-while-revalidate, not \
4896 propagate an error"
4897 );
4898 assert_eq!(
4899 cache.len(),
4900 1,
4901 "a 429 (throttling, not a credential-revocation signal) must not evict the \
4902 authenticated pinned-tier entry"
4903 );
4904 }
4905
4906 /// Every other tier keeps today's stale-while-revalidate fallback unchanged — only an
4907 /// *authenticated* `Pinned` entry evicts on 401/403 (FR-015's scope is deliberately
4908 /// narrow).
4909 #[tokio::test]
4910 async fn test_unauthenticated_pinned_401_revalidation_still_serves_stale() {
4911 let mut server = mockito::Server::new_async().await;
4912 let trusted_origin = format!("{}/", server.url());
4913 let url = format!("{}/api/data", server.url());
4914
4915 let _m1 = server
4916 .mock("GET", "/api/data")
4917 .with_status(200)
4918 .with_header("etag", "\"abc123\"")
4919 .with_body("workspace data")
4920 .create_async()
4921 .await;
4922
4923 let cache = HttpCache::new();
4924 let first = cache
4925 .get_cached_pinned(&url, &trusted_origin, false, None)
4926 .await
4927 .unwrap();
4928 assert_eq!(first.as_ref(), b"workspace data");
4929 drop(_m1);
4930
4931 let _m2 = server
4932 .mock("GET", "/api/data")
4933 .match_header("if-none-match", "\"abc123\"")
4934 .with_status(401)
4935 .create_async()
4936 .await;
4937
4938 let second = cache
4939 .get_cached_pinned(&url, &trusted_origin, false, None)
4940 .await
4941 .unwrap();
4942 assert_eq!(
4943 second.as_ref(),
4944 b"workspace data",
4945 "an unauthenticated Pinned entry must keep the default stale-while-revalidate fallback"
4946 );
4947 assert_eq!(cache.len(), 1);
4948 }
4949
4950 /// `set_registry_policy` purges every `Pinned`-tier cache entry (and pooled transport) on
4951 /// an actual policy transition — closing the round-trip hole for credential-carrying
4952 /// entries (NFR-004), unlike the pre-existing `WorkspaceDeclared` non-purge behavior.
4953 #[tokio::test]
4954 async fn test_set_registry_policy_purges_pinned_tier_entries() {
4955 let mut server = mockito::Server::new_async().await;
4956 let trusted_origin = format!("{}/", server.url());
4957 let url = format!("{}/api/data", server.url());
4958
4959 let _m = server
4960 .mock("GET", "/api/data")
4961 .with_status(200)
4962 .with_body("private data")
4963 .create_async()
4964 .await;
4965
4966 let policy = Arc::new(RegistryAccessPolicy::new(WorkspaceRegistryAccess::All));
4967 let cache = HttpCache::with_policy(Arc::clone(&policy));
4968 cache
4969 .get_cached_pinned(&url, &trusted_origin, true, Some(1))
4970 .await
4971 .unwrap();
4972 assert_eq!(cache.len(), 1);
4973
4974 cache.set_registry_policy(WorkspaceRegistryAccess::PublicOnly);
4975
4976 assert_eq!(
4977 cache.len(),
4978 0,
4979 "a Pinned-tier entry must be purged on any actual policy transition"
4980 );
4981 }
4982}