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