Skip to main content

feather_reader/
net.rs

1//! Hardened outbound HTTP for **untrusted, user-supplied feed URLs**.
2//!
3//! A feed reader fetches arbitrary URLs on behalf of its users: the add-feed
4//! flow, the background poller, and OPML import all hand a *user-controlled*
5//! host to `reqwest`. Left unguarded that is a textbook **SSRF** primitive — a
6//! subscribed feed can `302` to `http://169.254.169.254/` (cloud metadata) or
7//! `http://127.0.0.1:<port>/` (an internal service), and because the body is
8//! reflected back into the reader UI the exfiltration is *non-blind*.
9//!
10//! This module centralises the defence so every fetch path shares one guard:
11//!
12//! 1. **Scheme allow-list** — only `http` / `https`. No `file:`, `gopher:`, …
13//! 2. **IP allow-list** — the target host is resolved to IP(s) and rejected if
14//!    *any* resolved address is loopback, link-local (`169.254.0.0/16`,
15//!    `fe80::/10`), private (`10/8`, `172.16/12`, `192.168/16`), ULA
16//!    (`fc00::/7`), multicast, unspecified, or broadcast.
17//! 3. **Per-hop re-validation** — auto-redirect is disabled and redirects are
18//!    followed manually, re-running (1) and (2) on **every** hop, so a benign
19//!    first host cannot bounce us onto an internal one.
20//! 4. **Capped streaming body** — the response body is streamed and aborted the
21//!    moment it exceeds [`MAX_BODY_BYTES`], so a gzip decompression bomb cannot
22//!    materialise gigabytes before a post-hoc size check (a Content-Length
23//!    guard is useless once gzip strips the header).
24//!
25//! Resolution happens immediately before each request. The vetted IP is then
26//! **pinned** onto the connection (reqwest `.resolve(host, addr)`), so `connect`
27//! reuses the exact address that passed [`is_forbidden_ip`] rather than doing an
28//! independent second DNS lookup. That closes the DNS-rebinding TOCTOU window: an
29//! attacker-controlled resolver cannot answer "public IP" for the check and
30//! "127.0.0.1" for the connect, because there is no second resolution.
31
32use std::collections::HashMap;
33use std::net::{IpAddr, Ipv4Addr, Ipv6Addr, SocketAddr};
34use std::sync::atomic::{AtomicUsize, Ordering};
35use std::sync::{LazyLock, Mutex};
36use std::time::{Duration, Instant};
37
38use anyhow::{bail, Context, Result};
39use reqwest::header::{
40    HeaderName, HeaderValue, AUTHORIZATION, CONTENT_TYPE, COOKIE, PROXY_AUTHORIZATION,
41    WWW_AUTHENTICATE,
42};
43use reqwest::{Client, Response};
44use url::{Host, Url};
45
46/// Cap on how many bytes we will read from any body, streamed. 8 MiB is
47/// comfortably above any sane feed; a body that exceeds it is aborted mid-stream
48/// (never fully buffered), which is what defeats a gzip decompression bomb.
49pub const MAX_BODY_BYTES: usize = 8 * 1024 * 1024;
50
51/// Total per-request timeout for a guarded fetch. Matches
52/// [`crate::feed::build_client`]'s `FETCH_TIMEOUT` so the poller's per-hop
53/// pinned client is bounded the same way the feed client is — an unattended
54/// poll can't hang forever on a slow/silent upstream.
55///
56/// `pub(crate)` so a caller that wraps a *multi-request* walk in its own
57/// deadline can size that deadline against this per-request bound rather than
58/// hardcoding a second copy of the number — see [`crate::network::RelayClient`].
59pub(crate) const FETCH_TIMEOUT: Duration = Duration::from_secs(30);
60
61/// Per-read idle timeout: cap the wait for the *next* body chunk, so a server
62/// that trickles bytes forever (slowloris) can't tie up a fetch under the total
63/// timeout. Matches [`crate::feed::build_client`]'s `READ_TIMEOUT`.
64const READ_TIMEOUT: Duration = Duration::from_secs(15);
65
66/// Maximum number of redirect hops we will follow (each re-validated).
67///
68/// `pub(crate)` alongside [`FETCH_TIMEOUT`] because the two together give the
69/// real worst-case cost of ONE guarded request: the redirect loop runs
70/// `0..=MAX_REDIRECTS`, and every hop builds a fresh [`pinned_client`] carrying
71/// its own full [`FETCH_TIMEOUT`]. A caller that wraps a guarded request in an
72/// outer deadline must budget `(MAX_REDIRECTS + 1) * FETCH_TIMEOUT`, not one
73/// `FETCH_TIMEOUT` — getting that wrong silently pre-empts the inner logic.
74pub(crate) const MAX_REDIRECTS: usize = 5;
75
76/// Worst-case wall-clock cost of a single guarded request, redirects included.
77/// The number an outer deadline has to respect; see [`MAX_REDIRECTS`].
78pub(crate) const WORST_CASE_REQUEST: Duration =
79    Duration::from_secs(FETCH_TIMEOUT.as_secs() * (MAX_REDIRECTS as u64 + 1));
80
81/// Whether an already-resolved IP address is one we must never connect to on
82/// behalf of an untrusted URL (SSRF sinks): loopback, link-local, private,
83/// ULA, multicast, unspecified, or broadcast.
84pub fn is_forbidden_ip(ip: &IpAddr) -> bool {
85    match ip {
86        IpAddr::V4(v4) => is_forbidden_v4(v4),
87        IpAddr::V6(v6) => is_forbidden_v6(v6),
88    }
89}
90
91fn is_forbidden_v4(ip: &Ipv4Addr) -> bool {
92    ip.is_loopback()            // 127.0.0.0/8
93        || ip.is_private()      // 10/8, 172.16/12, 192.168/16
94        || ip.is_link_local()   // 169.254.0.0/16 (cloud metadata)
95        || ip.is_unspecified()  // 0.0.0.0
96        || ip.is_broadcast()    // 255.255.255.255
97        || ip.is_multicast()    // 224.0.0.0/4
98        // Carrier-grade NAT / "this-host" / benchmarking ranges — not routable
99        // to a legitimate public feed, but reachable internally.
100        || matches!(ip.octets(), [0, ..])
101        || matches!(ip.octets(), [100, b, ..] if (64..=127).contains(&b)) // 100.64/10 CGNAT (also used by overlay VPNs)
102        || matches!(ip.octets(), [192, 0, 0, _])
103        || matches!(ip.octets(), [198, 18..=19, _, _])
104}
105
106fn is_forbidden_v6(ip: &Ipv6Addr) -> bool {
107    if ip.is_loopback() || ip.is_unspecified() || ip.is_multicast() {
108        return true;
109    }
110    // Unwrap IPv4-mapped / -compatible addresses and re-check against the v4
111    // rules, so `::ffff:127.0.0.1` and friends can't slip past.
112    if let Some(v4) = ip.to_ipv4() {
113        return is_forbidden_v4(&v4);
114    }
115    // **And every OTHER way an IPv6 address carries an IPv4 one.** `to_ipv4()`
116    // stops at the mapped and compatible forms; five more families embed an
117    // address this function would refuse on sight, and all five were getting
118    // through. See [`embedded_v4`].
119    //
120    // This arm only ever returns `true`, so an address whose embedded IPv4 is
121    // public still falls through to the link-local and ULA checks below — which
122    // is what keeps `fe80::5efe:8.8.8.8` refused for being link-local.
123    if embedded_v4(ip).iter().any(is_forbidden_v4) {
124        return true;
125    }
126    let seg = ip.segments();
127    // fe80::/10 link-local (incl. RFC-4291 metadata equivalents).
128    let link_local = (seg[0] & 0xffc0) == 0xfe80;
129    // fc00::/7 unique-local addresses.
130    let ula = (seg[0] & 0xfe00) == 0xfc00;
131    link_local || ula
132}
133
134/// Every IPv4 address `ip` embeds under a translation scheme, for re-checking
135/// against the v4 rules.
136///
137/// **`to_ipv4()` is not the whole story, and the gap was a live SSRF hole.** It
138/// handles `::ffff:a.b.c.d` and `::a.b.c.d`. These it does not:
139///
140/// * **NAT64** — `64:ff9b::/32`. RFC 6052's well-known prefix is defined as a
141///   `/96`, so inside it the IPv4 is the last 32 bits and is decoded. The rest
142///   of the `/32`, including RFC 8215's local-use `64:ff9b:1::/48`, is refused
143///   outright, because RFC 6052 §2.2 allows six embedding lengths and which one
144///   a local deployment used is not something this code can know.
145/// * **6to4** — `2002::/16` (RFC 3056), IPv4 in the next two groups.
146/// * **IPv4-translated** — `::ffff:0:0:0/96` (RFC 2765), one group away from the
147///   mapped form.
148/// * **Teredo** — `2001::/32` (RFC 4380): the relay's IPv4 in groups 2-3 and the
149///   client's in groups 6-7, the latter obfuscated by XOR with all-ones. Both are
150///   returned; either one reaching an internal address is enough to refuse.
151/// * **ISATAP** — RFC 5214, and the odd one out: **no prefix to anchor on.** The
152///   IPv4 is the low 32 bits behind the IANA-reserved `00-00-5E-FE` OUI, under
153///   ANY /64, so an ordinary-looking global address can carry one. A link-local
154///   ISATAP address was already refused for being `fe80::/10`; one under a
155///   global prefix was not refused at all.
156///
157/// Decoded rather than blanket-refused for 6to4, IPv4-translated and Teredo,
158/// because those prefixes carry public addresses too and a blocklist would take
159/// out ordinary traffic. `allows_ipv6_that_embeds_a_public_ipv4` holds that line.
160///
161/// Found while bumping a JavaScript dependency whose advisory was this class:
162/// "no classifier recognizes the NAT64 local-use range". Ours did not either.
163fn embedded_v4(ip: &Ipv6Addr) -> Vec<Ipv4Addr> {
164    let seg = ip.segments();
165    let v4 = |hi: u16, lo: u16| {
166        Ipv4Addr::new(
167            (hi >> 8) as u8,
168            (hi & 0xff) as u8,
169            (lo >> 8) as u8,
170            (lo & 0xff) as u8,
171        )
172    };
173    // NAT64, split by prefix length because only one of the two is unambiguous.
174    //
175    // RFC 6052's well-known prefix is DEFINED as `64:ff9b::/96`, so inside it
176    // the IPv4 is unambiguously the last 32 bits and decodes like the others.
177    // That matters for availability, not just tidiness: a DNS64 resolver
178    // (RFC 6147) synthesises a well-known-prefix AAAA for every IPv4-only host,
179    // and `first_vetted` rejects a whole DNS answer set if ANY address in it is
180    // forbidden — so refusing the prefix outright makes every IPv4-only feed
181    // publisher unfetchable on an IPv6-only network. Which is the network this
182    // guard was written for.
183    //
184    // Nothing is given up by decoding: `64:ff9b::a9fe:a9fe` still refuses,
185    // because 169.254.169.254 refuses on its own merits. RFC 6052 §3.1 also
186    // forbids the well-known prefix from carrying a non-global IPv4 at all, so
187    // such an address is malformed as well as hostile.
188    //
189    // The REST of `64:ff9b::/32` — notably RFC 8215's local-use
190    // `64:ff9b:1::/48` — stays refused outright, and that is deliberate rather
191    // than lazy. RFC 6052 §2.2 defines six embedding lengths, and which one a
192    // local-use deployment chose is a property of that deployment. Guessing
193    // wrong reads the wrong bits, which could turn an internal target into a
194    // public-looking one — so for anything but the /96 the conservative answer
195    // is the only safe one. `LOCALHOST` there is a stand-in for "forbidden",
196    // not a claim about where the address points.
197    if seg[0] == 0x0064 && seg[1] == 0xff9b {
198        if seg[2..6] == [0, 0, 0, 0] {
199            return vec![v4(seg[6], seg[7])];
200        }
201        return vec![Ipv4Addr::LOCALHOST];
202    }
203    // **From here the arms ACCUMULATE instead of returning.** An early return
204    // was a bypass: 6to4 delegates `2002:<site-v4>::/48` to whoever owns that
205    // IPv4 and the site assigns identifiers inside it, so a site running ISATAP
206    // in its own 6to4 space produces `2002:<site-v4>:0:0:5efe:<internal-v4>` —
207    // two readings of DISJOINT bits, both true at once. Returning the 6to4 site
208    // address alone reported a public address and skipped the tunnel endpoint.
209    //
210    // This is the opposite of the Teredo case below, and the difference is
211    // which bits each family claims, not which is more important.
212    let mut out = Vec::new();
213    // 6to4: the site's IPv4 in groups 1-2, disjoint from the identifier.
214    if seg[0] == 0x2002 {
215        out.push(v4(seg[1], seg[2]));
216    }
217    // IPv4-translated: `::ffff:0:a.b.c.d`.
218    //
219    // This one accumulates for uniformity rather than necessity, and the
220    // distinction is worth recording: it requires `seg[5] == 0` where ISATAP
221    // requires `0x5efe`, and `seg[..4] == 0` excludes 6to4 and Teredo too, so
222    // it can only ever be the sole match. Returning here instead is an
223    // EQUIVALENT mutant — no test can tell the difference, and one written to
224    // try would be asserting on nothing. It pushes so that an arm added below
225    // it later is not silently skipped, which is the mistake the 6to4 arm
226    // above made.
227    if seg[..4] == [0, 0, 0, 0] && seg[4] == 0xffff && seg[5] == 0 {
228        out.push(v4(seg[6], seg[7]));
229    }
230    // Teredo: relay, then the client with the RFC 4380 obfuscation undone.
231    //
232    // **The one arm that still returns, because it CLAIMS the identifier's
233    // bits.** Teredo's client address lives in groups 6-7 complemented — the
234    // same bits ISATAP reads uncomplemented — so the two readings are of one
235    // field and contradict each other. Accumulating both would refuse a
236    // legitimate Teredo address whenever the inverse of its client address
237    // happens to be internal. Returning here resolves that in favour of the
238    // prefix, which `a_teredo_address_is_read_as_teredo_not_as_isatap` pins.
239    if seg[0] == 0x2001 && seg[1] == 0 {
240        out.push(v4(seg[2], seg[3]));
241        out.push(v4(seg[6] ^ 0xffff, seg[7] ^ 0xffff));
242        return out;
243    }
244    // ISATAP, and it is last so that Teredo above can suppress it.
245    //
246    // The others are prefix-anchored; this one is not — the IPv4 sits in the low
247    // 32 bits behind the IANA `00-00-5E-FE` OUI under ANY /64, so the test is on
248    // the interface identifier and matches whatever the prefix. Being reachable
249    // under another family's prefix is the point, not an edge case: it is why
250    // the arms above accumulate rather than return.
251    //
252    // **`seg[4]` is deliberately not constrained.** RFC 5214 spells the
253    // identifier as `000000ug 00000000 0x5E 0xFE` + the IPv4, so a spec-exact
254    // test would require all of `seg[4]` except the `u` and `g` bits to be
255    // zero. Two drafts of this arm tried to be that precise and the first was
256    // wrong: it enumerated `0x0000` and `0x0200`, missed the two values with
257    // `g` set, and so was bypassable by flipping one bit while reading as
258    // complete.
259    //
260    // The asymmetry decides it. Reading the marker loosely costs a false
261    // positive only when a non-ISATAP address happens to carry `0x5efe` in
262    // group 5 AND its low 32 bits decode to a forbidden IPv4 — and `00-00-5E`
263    // is IANA's own OUI, reserved for this, so a real interface identifier does
264    // not land there. Reading it strictly costs a total bypass if any tunnel
265    // driver is more lenient than the RFC about the reserved bits. A guard
266    // should be conservative about what it accepts as safe, which here means
267    // the simpler condition, not the more exact one.
268    if seg[5] == 0x5efe {
269        out.push(v4(seg[6], seg[7]));
270    }
271    out
272}
273
274/// Validate a URL's scheme (http/https only). Returns the host as a string.
275fn check_scheme(url: &Url) -> Result<()> {
276    match url.scheme() {
277        "http" | "https" => Ok(()),
278        other => bail!("refusing non-http(s) URL scheme {other:?}"),
279    }
280}
281
282/// Resolve a URL's host to socket addresses, reject if *any* resolved IP is a
283/// forbidden (SSRF) target, and return the **vetted** `SocketAddr` to pin the
284/// connection to.
285///
286/// An IP literal host is checked directly (no DNS); a named host is resolved via
287/// the async resolver and *every* answer must pass — but the returned address is
288/// the specific one `connect` must use, so no independent second resolution can
289/// slip a rebound IP past the check (DNS-rebinding TOCTOU). Handles both IPv4 and
290/// IPv6 answers.
291async fn resolve_and_check(url: &Url) -> Result<SocketAddr> {
292    let host = url.host().context("URL has no host")?;
293    let port = url
294        .port_or_known_default()
295        .context("URL has no usable port")?;
296
297    match host {
298        Host::Ipv4(ip) => {
299            if is_forbidden_ip(&IpAddr::V4(ip)) {
300                bail!("refusing to fetch forbidden (internal) address {ip}");
301            }
302            Ok(SocketAddr::new(IpAddr::V4(ip), port))
303        }
304        Host::Ipv6(ip) => {
305            if is_forbidden_ip(&IpAddr::V6(ip)) {
306                bail!("refusing to fetch forbidden (internal) address {ip}");
307            }
308            Ok(SocketAddr::new(IpAddr::V6(ip), port))
309        }
310        Host::Domain(name) => {
311            // **Test seam — `#[cfg(test)]`, so it does not exist in a release
312            // build at all.** Not a parameter, not an env var, not a feature
313            // flag: the compiler removes it, so there is no runtime bypass to
314            // reason about. It exists because the guard is otherwise untestable
315            // end-to-end — a local test server lives on loopback, which
316            // `is_forbidden_ip` correctly refuses, so nothing could ever drive a
317            // real redirect through this function. See `test_host_override`.
318            #[cfg(test)]
319            if let Some(addr) = test_override_for(name, port) {
320                return Ok(addr);
321            }
322            let addrs = tokio::net::lookup_host((name, port))
323                .await
324                .with_context(|| format!("resolving host {name:?}"))?;
325            first_vetted(name, addrs)
326        }
327    }
328}
329
330/// Pick the address to pin to from a host's DNS answers, rejecting the whole
331/// set if ANY answer is forbidden.
332///
333/// **Extracted so it can be tested.** Inline in the resolver it was unreachable
334/// without real DNS returning a mixed answer set, and a mutation that checked
335/// only the FIRST answer left the entire suite green — a DNS-rebinding style
336/// attack that publishes `1.2.3.4, 127.0.0.1` would have been accepted on the
337/// strength of the first record.
338///
339/// Rejecting wholesale rather than filtering is deliberate: a host that resolves
340/// to any internal address is not a host we want to talk to, even on the answers
341/// that look fine.
342fn first_vetted(name: &str, addrs: impl Iterator<Item = SocketAddr>) -> Result<SocketAddr> {
343    let mut vetted: Option<SocketAddr> = None;
344    for sa in addrs {
345        let ip = sa.ip();
346        if is_forbidden_ip(&ip) {
347            bail!("refusing to fetch {name:?}: resolves to forbidden address {ip}");
348        }
349        // Keep the FIRST vetted answer as the address to pin the connect to.
350        // Every answer is still checked (the loop continues), so a mixed A/AAAA
351        // set with any forbidden entry is rejected wholesale.
352        if vetted.is_none() {
353            vetted = Some(sa);
354        }
355    }
356    vetted.ok_or_else(|| anyhow::anyhow!("host {name:?} did not resolve to any address"))
357}
358
359/// Test-only host→address overrides, consulted by [`resolve_and_check`] before
360/// real DNS. Keyed by host so tests using distinct hostnames never collide, and
361/// gone entirely from a release build.
362#[cfg(test)]
363static TEST_HOSTS: std::sync::Mutex<Option<std::collections::HashMap<String, SocketAddr>>> =
364    std::sync::Mutex::new(None);
365
366/// The TEST certificate authority, and the leaf it issues for test hostnames.
367///
368/// **Why this exists at all.** An OAuth issuer is required to be `https`
369/// (`discovery::validate_issuer_form`), so a plain-HTTP loopback server cannot
370/// stand in for an authorization server — which meant the real `login::complete`
371/// could never be driven end to end, and the wiring between its tested core and
372/// the network had no coverage. A review proved that gap was live: the
373/// authorization-server mix-up defence could be disabled in that wiring with the
374/// whole suite green.
375///
376/// **Why a CA rather than relaxing the rule.** The alternative was a test-only
377/// escape from the https requirement. That would *suspend* a security rule; this
378/// *satisfies* it — the server really presents a certificate and the client
379/// really validates the chain. It also keeps the existing tests that assert
380/// `http` issuers are REJECTED meaningful, which a blanket relaxation would not.
381///
382/// Generated once per process. `#[cfg(test)]`, so none of it — not the trust
383/// decision, not the key material — exists in a release build.
384#[cfg(test)]
385pub(crate) struct TestPki {
386    /// PEM of the CA certificate, for `reqwest`'s root store.
387    pub ca_pem: String,
388    /// PEM of the leaf certificate chain, for the server.
389    pub leaf_pem: String,
390    /// PEM of the leaf private key, for the server.
391    pub leaf_key_pem: String,
392}
393
394#[cfg(test)]
395pub(crate) fn test_pki() -> &'static TestPki {
396    static PKI: std::sync::OnceLock<TestPki> = std::sync::OnceLock::new();
397    PKI.get_or_init(|| {
398        use rcgen::{
399            BasicConstraints, CertificateParams, DnType, IsCa, KeyPair, KeyUsagePurpose, SanType,
400        };
401
402        let mut ca_params = CertificateParams::default();
403        ca_params
404            .distinguished_name
405            .push(DnType::CommonName, "featherreader test CA");
406        ca_params.is_ca = IsCa::Ca(BasicConstraints::Constrained(0));
407        ca_params.key_usages = vec![
408            KeyUsagePurpose::KeyCertSign,
409            KeyUsagePurpose::CrlSign,
410            KeyUsagePurpose::DigitalSignature,
411        ];
412        let ca_key = KeyPair::generate().expect("test CA key");
413        let ca_cert = ca_params
414            .clone()
415            .self_signed(&ca_key)
416            .expect("test CA cert");
417        let issuer = rcgen::Issuer::new(ca_params, ca_key);
418
419        // SANs for the hostnames the tests register with `test_host_override`.
420        // A wildcard would not cover the multi-label names, so they are listed.
421        let mut leaf_params = CertificateParams::default();
422        leaf_params
423            .distinguished_name
424            .push(DnType::CommonName, "featherreader test leaf");
425        leaf_params.subject_alt_names = TEST_TLS_HOSTS
426            .iter()
427            .map(|h| SanType::DnsName((*h).try_into().expect("test SAN")))
428            .collect();
429        let leaf_key = KeyPair::generate().expect("test leaf key");
430        let leaf_cert = leaf_params
431            .signed_by(&leaf_key, &issuer)
432            .expect("test leaf cert");
433
434        TestPki {
435            ca_pem: ca_cert.pem(),
436            leaf_pem: leaf_cert.pem(),
437            leaf_key_pem: leaf_key.serialize_pem(),
438        }
439    })
440}
441
442#[cfg(test)]
443/// A loopback HTTPS server presenting the test CA's leaf, routing by path.
444///
445/// The point of the TLS is not TLS: it is that an OAuth issuer must be
446/// `https`, so nothing could drive the real `login::complete` against a
447/// local server. The client validates this chain for real — no invalid-cert
448/// acceptance anywhere.
449///
450/// `routes` maps a path to a canned `(status, body)`. Unknown paths 404.
451/// Every request line is recorded.
452pub(crate) async fn spawn_tls<F>(
453    build_routes: F,
454) -> (SocketAddr, std::sync::Arc<std::sync::Mutex<Vec<String>>>)
455where
456    F: FnOnce(SocketAddr) -> std::collections::HashMap<String, Vec<TestResponse>>,
457{
458    use tokio::io::{AsyncReadExt, AsyncWriteExt};
459    use tokio_rustls::rustls::pki_types::{CertificateDer, PrivateKeyDer};
460
461    // Both `ring` and `aws-lc-rs` are reachable in this tree, so rustls refuses
462    // to guess a process-level provider for the SERVER side here. Install ring.
463    //
464    // **Two earlier versions of this comment were wrong in opposite directions;
465    // this is what reqwest 0.13 actually does** (`async_impl/client.rs`):
466    //
467    //     let provider = rustls::crypto::CryptoProvider::get_default()
468    //         .map(|arc| arc.clone())
469    //         .unwrap_or_else(default_rustls_crypto_provider);
470    //
471    // So it READS the process default and falls back to aws-lc-rs. Installing
472    // ring here therefore DOES affect reqwest clients built afterwards in the
473    // same test binary — which makes the client's provider depend on whether any
474    // test called `spawn_tls` first. Benign (ring and aws-lc-rs interoperate),
475    // and absent from release builds, where nothing installs a default and
476    // production is genuinely aws-lc-rs. Recorded precisely because two previous
477    // attempts at this comment stated a checkable fact without checking it.
478    //
479    // `install_default` errors if something got there first, which is fine.
480    static PROVIDER: std::sync::Once = std::sync::Once::new();
481    PROVIDER.call_once(|| {
482        let _ = tokio_rustls::rustls::crypto::ring::default_provider().install_default();
483    });
484
485    let pki = test_pki();
486    let certs: Vec<CertificateDer<'static>> = rustls_pemfile_certs(pki.leaf_pem.as_bytes());
487    let key: PrivateKeyDer<'static> = rustls_pemfile_key(pki.leaf_key_pem.as_bytes());
488
489    let config = tokio_rustls::rustls::ServerConfig::builder()
490        .with_no_client_auth()
491        .with_single_cert(certs, key)
492        .expect("test server TLS config");
493    let acceptor = tokio_rustls::TlsAcceptor::from(std::sync::Arc::new(config));
494
495    let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
496    let addr = listener.local_addr().unwrap();
497    // Routes are built from the bound address: the documents have to name their
498    // own port, and the port is not known until the listener exists.
499    let routes = build_routes(addr);
500    let log = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
501    let sink = std::sync::Arc::clone(&log);
502    let hits: std::sync::Arc<std::sync::Mutex<std::collections::HashMap<String, usize>>> =
503        Default::default();
504
505    tokio::spawn(async move {
506        loop {
507            let Ok((sock, _)) = listener.accept().await else {
508                break;
509            };
510            let acceptor = acceptor.clone();
511            let routes = routes.clone();
512            let sink = std::sync::Arc::clone(&sink);
513            let hits = std::sync::Arc::clone(&hits);
514            tokio::spawn(async move {
515                let Ok(mut tls) = acceptor.accept(sock).await else {
516                    return;
517                };
518                // **The head, and the body when `content-length` says how
519                // much — before replying.**
520                //
521                // A single `read` is what the capturing sidecar on the #138
522                // branch did, and the review of that branch found the trap: if
523                // the head and the body land in separate segments, the capture
524                // holds only the head, and every `contains` assertion over it
525                // then passes for the wrong reason.
526                //
527                // Nothing here asserts on a body today — the assertions are the
528                // request line and the DPoP header, and both panic loudly when
529                // absent rather than passing — so that false green is not live
530                // in this harness. It is the NEXT body assertion that would
531                // inherit one, which is the whole reason the same shape was
532                // worth fixing there.
533                //
534                // **A chunked body is NOT drained.** With no `content-length`
535                // there is nothing to wait for, so this stops after the head —
536                // exactly what the single read did. Every request this harness
537                // sees is a GET or a reqwest-buffered form and carries a length,
538                // but nothing here enforces that, so a future chunked request
539                // would be captured short and quietly. Said plainly rather than
540                // left inside a claim to have read "the whole request".
541                //
542                // Draining also stops the reply being written while the client
543                // is still sending, which would make a split request a broken
544                // pipe rather than a response.
545                let mut raw: Vec<u8> = Vec::new();
546                let mut chunk = [0u8; 4096];
547                loop {
548                    let Ok(n) = tls.read(&mut chunk).await else {
549                        return;
550                    };
551                    if n == 0 {
552                        break;
553                    }
554                    raw.extend_from_slice(&chunk[..n]);
555                    let Some(split) = raw.windows(4).position(|w| w == b"\r\n\r\n") else {
556                        continue;
557                    };
558                    let (head, body) = raw.split_at(split + 4);
559                    let want = String::from_utf8_lossy(head).lines().find_map(|l| {
560                        let (k, v) = l.split_once(':')?;
561                        k.eq_ignore_ascii_case("content-length")
562                            .then(|| v.trim().parse::<usize>().ok())?
563                    });
564                    if want.is_none_or(|want| body.len() >= want) {
565                        break;
566                    }
567                }
568                let req = String::from_utf8_lossy(&raw).to_string();
569                let path = req
570                    .lines()
571                    .next()
572                    .and_then(|l| l.split_whitespace().nth(1))
573                    .unwrap_or("/")
574                    .to_string();
575                sink.lock().unwrap().push(req);
576                // Nth hit on this path picks the Nth canned reply; the last one
577                // repeats. That is what lets a route answer a nonce challenge
578                // once and something else afterwards — the only way to observe
579                // whether a request was RETRIED.
580                let n = {
581                    let mut c = hits.lock().unwrap();
582                    let e = c.entry(path.clone()).or_insert(0usize);
583                    let n = *e;
584                    *e += 1;
585                    n
586                };
587                let reply = routes
588                    .get(&path)
589                    .and_then(|v| v.get(n.min(v.len().saturating_sub(1))))
590                    .cloned()
591                    .unwrap_or_else(|| TestResponse::json(404, "not found"));
592                let extra: String = reply
593                    .headers
594                    .iter()
595                    .map(|(k, v)| format!("{k}: {v}\r\n"))
596                    .collect();
597                let (status, body) = (reply.status, reply.body);
598                let resp = format!(
599                    "HTTP/1.1 {status} X\r\nContent-Type: application/json\r\n{extra}\
600                     Content-Length: {}\r\nConnection: close\r\n\r\n{body}",
601                    body.len()
602                );
603                let _ = tls.write_all(resp.as_bytes()).await;
604                let _ = tls.shutdown().await;
605            });
606        }
607    });
608    (addr, log)
609}
610
611#[cfg(test)]
612fn rustls_pemfile_certs(
613    pem: &[u8],
614) -> Vec<tokio_rustls::rustls::pki_types::CertificateDer<'static>> {
615    // Minimal PEM splitter — avoids another dependency for two blocks.
616    decode_pem_blocks(pem, "CERTIFICATE")
617        .into_iter()
618        .map(Into::into)
619        .collect()
620}
621
622#[cfg(test)]
623fn rustls_pemfile_key(pem: &[u8]) -> tokio_rustls::rustls::pki_types::PrivateKeyDer<'static> {
624    let der = decode_pem_blocks(pem, "PRIVATE KEY")
625        .into_iter()
626        .next()
627        .expect("a private key block");
628    tokio_rustls::rustls::pki_types::PrivatePkcs8KeyDer::from(der).into()
629}
630
631#[cfg(test)]
632fn decode_pem_blocks(pem: &[u8], label: &str) -> Vec<Vec<u8>> {
633    use base64::Engine as _;
634    let text = String::from_utf8_lossy(pem);
635    let begin = format!("-----BEGIN {label}-----");
636    let end = format!("-----END {label}-----");
637    let mut out = Vec::new();
638    let mut rest = text.as_ref();
639    while let Some(i) = rest.find(&begin) {
640        let after = &rest[i + begin.len()..];
641        let Some(j) = after.find(&end) else { break };
642        let b64: String = after[..j].chars().filter(|c| !c.is_whitespace()).collect();
643        out.push(
644            base64::engine::general_purpose::STANDARD
645                .decode(b64)
646                .expect("valid base64 in test PEM"),
647        );
648        rest = &after[j + end.len()..];
649    }
650    out
651}
652
653/// One canned reply from the TLS test server.
654#[cfg(test)]
655#[derive(Clone)]
656pub(crate) struct TestResponse {
657    pub status: u16,
658    pub body: String,
659    pub headers: Vec<(String, String)>,
660}
661
662#[cfg(test)]
663impl TestResponse {
664    pub fn json(status: u16, body: impl Into<String>) -> Self {
665        Self {
666            status,
667            body: body.into(),
668            headers: Vec::new(),
669        }
670    }
671
672    pub fn with_header(mut self, k: &str, v: &str) -> Self {
673        self.headers.push((k.to_string(), v.to_string()));
674        self
675    }
676}
677
678/// Hostnames the test leaf is valid for. Adding a new `.test` host to a test
679/// means adding it here, which is deliberate friction: the certificate is
680/// supposed to be narrow.
681#[cfg(test)]
682pub(crate) const TEST_TLS_HOSTS: &[&str] = &[
683    "pds-e2e.test",
684    "as-e2e.test",
685    "feed-tls.test",
686    "hop-tls.test",
687    "as-evil.test",
688];
689
690/// Point `host` at `addr` for the rest of the process, bypassing DNS **and** the
691/// forbidden-IP check for that host only.
692///
693/// **Registrations are process-global and last-write-wins.** Several tests
694/// register the SAME hostnames to different servers concurrently. What keeps
695/// them apart is not the host key — an earlier comment claimed it was — but that
696/// reqwest's `.resolve()` ignores the port, so every registration collapses to
697/// `127.0.0.1` and each test's URL port routes it back to its own listener. That
698/// is incidental, and would break the moment a test server bound anything other
699/// than loopback.
700///
701/// Bypassing the IP check is the entire point: the test server is on loopback,
702/// which the guard is right to refuse. Only the registered host is exempt —
703/// anything else in the same test, including every redirect target, still goes
704/// through the real check. That is what makes a redirect test meaningful.
705#[cfg(test)]
706pub(crate) fn test_host_override(host: &str, addr: SocketAddr) {
707    TEST_HOSTS
708        .lock()
709        .unwrap()
710        .get_or_insert_with(Default::default)
711        .insert(host.to_string(), addr);
712}
713
714#[cfg(test)]
715fn test_override_for(name: &str, port: u16) -> Option<SocketAddr> {
716    let guard = TEST_HOSTS.lock().unwrap();
717    let map = guard.as_ref()?;
718    map.get(name)
719        .copied()
720        .or_else(|| map.get(&format!("{name}:{port}")).copied())
721}
722
723/// How long an idle pinned client may be kept before it is rebuilt.
724///
725/// Not a security boundary — the address is re-resolved and re-checked on every
726/// single request, and a changed address misses the cache by construction. This
727/// only bounds how long a pooled connection to a once-vetted address may live,
728/// and keeps the map from holding entries for hosts nobody fetches any more.
729const PINNED_CLIENT_TTL: Duration = Duration::from_secs(300);
730
731/// Most distinct (host, address) pairs kept. A bound, not a target: the reader
732/// talks to one PDS, while the poller talks to as many hosts as there are feeds.
733const MAX_PINNED_CLIENTS: usize = 256;
734
735/// How long a pinned client may hold an IDLE socket open.
736///
737/// Deliberately shorter than [`PINNED_CLIENT_TTL`] so a client releases its
738/// sockets before the cache releases the client — otherwise the last minute of
739/// an entry's life is pure socket rent. See [`build_pinned_client`] for why the
740/// pool needs bounding at all.
741const POOL_IDLE_TIMEOUT: Duration = Duration::from_secs(60);
742
743/// Pinned clients, keyed by the **vetted address** they are pinned to.
744///
745/// ## Why this is safe to reuse
746///
747/// Building a fresh client per request meant a fresh connection pool, so every
748/// PDS call paid a full TCP + TLS handshake: measured at 91 ms against this
749/// project's PDS versus 30 ms on a warm connection. That is most of why the
750/// Rust repo backend measured ~3x slower than the Node sidecar, which pools.
751///
752/// Reuse does NOT weaken the DNS-rebinding defence, because the defence does not
753/// live in the client's lifetime:
754///
755/// * every request still resolves the host and runs [`is_forbidden_ip`] over
756///   EVERY answer before this cache is consulted — a host that now resolves to
757///   an internal address is refused before a pooled client could be returned;
758/// * the key includes the vetted [`SocketAddr`], so a host that legitimately
759///   moves to a different address MISSES the cache and gets a client pinned to
760///   the new one. A pooled connection can only ever be reused for an address
761///   that was just re-vetted this request.
762struct PinnedClients {
763    entries: Mutex<HashMap<(String, SocketAddr), (Client, Instant)>>,
764    /// How many clients have actually been constructed. Test-only bookkeeping:
765    /// it is the only way to observe that a hit avoided a rebuild, since
766    /// `reqwest::Client` exposes no identity.
767    builds: AtomicUsize,
768}
769
770impl PinnedClients {
771    fn new() -> Self {
772        Self {
773            entries: Mutex::new(HashMap::new()),
774            builds: AtomicUsize::new(0),
775        }
776    }
777
778    /// A client pinned to `addr` for `host`, reusing a pooled one when the
779    /// address is unchanged and the entry is fresh.
780    fn get(&self, host: &str, addr: SocketAddr, now: Instant) -> Result<Client> {
781        let key = (host.to_string(), addr);
782        // A poisoned lock here is NOT fatal and must not be treated as fatal: the
783        // guard is held across a fallible builder, so one panic inside it would
784        // otherwise make EVERY subsequent outbound request panic, forever, with a
785        // live-looking process and a green /health. Recover the data like the rate
786        // limiter already does — a torn entry is a cache entry, worst case a rebuild.
787        let mut entries = self.entries.lock().unwrap_or_else(|p| p.into_inner());
788
789        if let Some((client, last_used)) = entries.get_mut(&key) {
790            if now.duration_since(*last_used) < PINNED_CLIENT_TTL {
791                *last_used = now;
792                // Cloning a `reqwest::Client` shares its connection pool, which
793                // is the entire point — a clone is a handle, not a new pool.
794                return Ok(client.clone());
795            }
796        }
797
798        let client = build_pinned_client(host, addr)?;
799        self.builds.fetch_add(1, Ordering::Relaxed);
800
801        // Drop anything idle past the TTL before considering the bound, so a
802        // burst of one-off hosts does not evict the PDS client we use constantly.
803        entries.retain(|_, (_, last_used)| now.duration_since(*last_used) < PINNED_CLIENT_TTL);
804        if entries.len() >= MAX_PINNED_CLIENTS {
805            if let Some(oldest) = entries
806                .iter()
807                .min_by_key(|(_, (_, last_used))| *last_used)
808                .map(|(k, _)| k.clone())
809            {
810                entries.remove(&oldest);
811            }
812        }
813        entries.insert(key, (client.clone(), now));
814        Ok(client)
815    }
816}
817
818static PINNED_CLIENTS: LazyLock<PinnedClients> = LazyLock::new(PinnedClients::new);
819
820/// Build a per-hop client that **pins** DNS for `host` to the already-vetted
821/// `addr`, so reqwest's `connect` reuses the exact IP that passed the SSRF check
822/// instead of doing its own second resolution (the DNS-rebinding fix). The pin is
823/// scoped to `host`, keyed to the address family of `addr` (works for both IPv4
824/// and IPv6). Mirrors [`crate::feed::build_client`]'s policy: the same total
825/// [`FETCH_TIMEOUT`] + per-read [`READ_TIMEOUT`] (so the unattended poller keeps
826/// its slowloris / slow-upstream defence even though each hop is a freshly built
827/// client), and auto-redirect off — [`guarded_get`] follows + re-validates each
828/// hop itself.
829fn build_pinned_client(host: &str, addr: SocketAddr) -> Result<Client> {
830    let builder = Client::builder()
831        .user_agent(crate::USER_AGENT)
832        // Bound each hop the same way the feed client is bounded: a total
833        // request timeout plus a per-read idle timeout. Without these the
834        // per-hop client the poller actually connects through had NO timeouts,
835        // leaving the unattended poll with no defence against a slowloris /
836        // never-finishing upstream.
837        .timeout(FETCH_TIMEOUT)
838        .read_timeout(READ_TIMEOUT)
839        // Bound the idle connection pool too.
840        //
841        // These clients are CACHED — up to `MAX_PINNED_CLIENTS` of them, each
842        // holding its own pool — and every entry keeps live keep-alive TLS
843        // connections open until it is evicted. With reqwest's defaults
844        // (unlimited idle per host, no idle timeout) a poller touching many
845        // distinct feed hosts drives the cache toward its bound and each entry
846        // toward an unbounded number of sockets, on a 512 MB box with one shared
847        // core. The cache was given a size bound for the same reason; its pools
848        // were not.
849        //
850        // One idle connection per host is the right number here: reuse across
851        // the ~300 s TTL is what the cache exists for (measured 91 ms cold
852        // versus 30 ms warm), and nothing in this codebase issues concurrent
853        // requests to the SAME host through one client — `guarded_get` walks
854        // redirect hops sequentially, and the poller's concurrency is across
855        // DIFFERENT feeds. The idle timeout is well under the cache TTL so
856        // sockets are released before the client itself is.
857        .pool_max_idle_per_host(1)
858        .pool_idle_timeout(POOL_IDLE_TIMEOUT)
859        // **Ignore ambient proxy configuration.** reqwest defaults
860        // `auto_sys_proxy: true`, so `HTTP_PROXY` / `HTTPS_PROXY` / `ALL_PROXY`
861        // in the process environment silently route every request through a
862        // proxy — and a proxied request is sent in absolute form for the PROXY
863        // to resolve the hostname. That defeats the two mechanisms this whole
864        // module rests on at once: the `.resolve()` pin below never sees the
865        // connection, and `is_forbidden_ip` never sees the address, because we
866        // no longer do the resolving.
867        //
868        // Measured before this line existed, with `HTTP_PROXY` set: the vetted
869        // address received ZERO requests, the proxy received
870        // `GET http://pinned.invalid/feed HTTP/1.1`, and the call returned
871        // `Ok(200)`. It failed OPEN and silently.
872        //
873        // Not remotely triggerable — it needs a proxy variable in the server's
874        // own environment — but that is one `fly secrets set`, one debugging
875        // session, or one base image away, and nothing would have reported the
876        // guard had stopped working.
877        .no_proxy()
878        // Override reqwest's resolver for this host only: connect goes straight
879        // to the vetted socket address — no independent re-resolution.
880        .resolve(host, addr)
881        // No auto-redirect: guarded_get follows + re-validates each hop.
882        .redirect(reqwest::redirect::Policy::none());
883
884    // **The TEST certificate authority — `#[cfg(test)]`, so a release build has
885    // neither this call nor the certificate.**
886    //
887    // This is the one test seam in this file that touches TLS TRUST, so it is
888    // worth being exact about what it does and does not do. It ADDS one root:
889    // the built-in roots stay, nothing is disabled, and `danger_accept_invalid_
890    // certs` is NOT used — a server still has to present a chain that validates,
891    // and a hostname still has to match a SAN. What it buys is that a loopback
892    // test server can hold a certificate the client will accept, which is what
893    // makes it possible to drive the real `login::complete` (and the real
894    // redirect path) against a server at all: the OAuth issuer must be `https`.
895    //
896    // `test_pki()` is itself `#[cfg(test)]`, so removing the attribute here
897    // fails to compile rather than silently trusting an extra root in prod.
898    #[cfg(test)]
899    let builder = builder.add_root_certificate(
900        reqwest::Certificate::from_pem(test_pki().ca_pem.as_bytes())
901            .context("parsing the test CA")?,
902    );
903
904    builder
905        .build()
906        .context("failed to build IP-pinned fetch client")
907}
908
909/// The per-hop client for an already-vetted `(host, addr)`, pooled.
910///
911/// Callers must have run [`resolve_and_check`] for THIS request before calling
912/// this — the cache trusts its key, and the key is only as good as the check
913/// that produced it.
914fn pinned_client(host: &str, addr: SocketAddr) -> Result<Client> {
915    PINNED_CLIENTS.get(host, addr, Instant::now())
916}
917
918/// Fetch a user-supplied URL through the full SSRF guard: scheme + IP checks on
919/// the initial URL and on **every** redirect hop, following redirects manually.
920///
921/// The passed `client` is used only as a policy reference; each hop is actually
922/// sent through a freshly-built `pinned_client` whose DNS for the target host
923/// is pinned to the exact IP that just passed `resolve_and_check` — so the
924/// connect can't be rebound onto an internal address between the check and the
925/// TCP handshake.
926///
927/// `extra_headers` are applied to every hop (e.g. the conditional-GET
928/// `If-None-Match` / `If-Modified-Since` validators) — **except** credential
929/// headers (`Authorization`, `Cookie`, …), which are dropped the moment a
930/// redirect leaves the original origin, mirroring what reqwest's own redirect
931/// policy does for the shared client (see `hop_headers`). Returns the final
932/// `Response` (headers only; the body is read separately via [`read_capped`]).
933/// `Err` on a blocked scheme/address, an exhausted redirect budget, or a
934/// transport error.
935pub async fn guarded_get(
936    client: &Client,
937    url: &str,
938    extra_headers: &[(HeaderName, HeaderValue)],
939) -> Result<Response> {
940    guarded_get_inner(client, url, extra_headers, true, MAX_REDIRECTS).await
941}
942
943/// The SSRF core of [`guarded_get`] **without** the feed-privacy layer: scheme +
944/// IP allow-list, connect-pinning, and per-hop re-validation, but no
945/// `classify_feed_privacy` check.
946///
947/// This is the entry point for **non-feed** fetches of *user-influenced* URLs —
948/// notably atproto identity resolution (a handle's PDS host, a `did:web`
949/// well-known document, and a DID document's `serviceEndpoint`). Those are
950/// legitimate atproto XRPC / DID-doc requests, so the feed-privacy heuristic
951/// (which flags Substack/Patreon-style token URLs) must not apply — but the SSRF
952/// guard absolutely must, since a hostile `did:web` or `serviceEndpoint` can
953/// otherwise point the server at `169.254.169.254`, loopback, or a private host.
954pub async fn guarded_get_no_privacy(
955    client: &Client,
956    url: &str,
957    extra_headers: &[(HeaderName, HeaderValue)],
958) -> Result<Response> {
959    guarded_get_inner(client, url, extra_headers, false, MAX_REDIRECTS).await
960}
961
962/// Like [`guarded_get_no_privacy`] but **refuses redirects outright**.
963///
964/// For the OAuth discovery and DID documents, following a redirect is not a
965/// convenience — it is a hole. The mix-up defence rests on comparing a
966/// document's `issuer` against *the URL it was fetched from*; if a `302` can move
967/// the fetch to another origin, that comparison is against the original URL while
968/// the bytes came from somewhere else, and the check silently stops meaning
969/// anything. The reference client sets `redirect: 'manual'`/`'error'` on every
970/// one of these fetches for the same reason.
971///
972/// Applies to: `/.well-known/oauth-protected-resource`,
973/// `/.well-known/oauth-authorization-server`, `plc.directory/<did>`, `did:web`
974/// `did.json`, and the client-metadata self-fetch. It deliberately does NOT
975/// apply to `/.well-known/atproto-did`, where the handle spec explicitly permits
976/// redirects.
977pub async fn guarded_get_no_redirect(
978    client: &Client,
979    url: &str,
980    extra_headers: &[(HeaderName, HeaderValue)],
981) -> Result<Response> {
982    guarded_get_inner(client, url, extra_headers, false, 0).await
983}
984
985/// Whether a header carries credentials that must never follow a redirect onto a
986/// different origin. Mirrors reqwest's own `redirect::remove_sensitive_headers`
987/// set (`Authorization`, `Cookie`, `Cookie2`, `Proxy-Authorization`,
988/// `WWW-Authenticate`), which the shared client applies automatically — and which
989/// [`guarded_get_inner`] must reimplement because it disables auto-redirect and
990/// re-applies `extra_headers` by hand on every manually-followed hop.
991fn is_sensitive_header(name: &HeaderName) -> bool {
992    name == AUTHORIZATION
993        || name == COOKIE
994        || name == PROXY_AUTHORIZATION
995        || name == WWW_AUTHENTICATE
996        || name.as_str() == "cookie2"
997}
998
999/// Same-origin in the web sense: identical scheme, host, and effective port.
1000fn same_origin(a: &Url, b: &Url) -> bool {
1001    a.scheme() == b.scheme()
1002        && a.host_str() == b.host_str()
1003        && a.port_or_known_default() == b.port_or_known_default()
1004}
1005
1006/// The headers to apply on THIS hop: all of `extra` while we are still on the
1007/// original origin, otherwise only the non-sensitive ones.
1008///
1009/// The comparison base is the **original** URL rather than the previous hop (what
1010/// reqwest does). That is strictly stricter: an `a → b → a` redirect chain never
1011/// re-attaches the credential, at the cost of a small, deliberate divergence from
1012/// the stock client's behaviour.
1013fn hop_headers<'a>(
1014    original: &Url,
1015    current: &Url,
1016    extra: &'a [(HeaderName, HeaderValue)],
1017) -> Vec<&'a (HeaderName, HeaderValue)> {
1018    let cross_origin = !same_origin(original, current);
1019    extra
1020        .iter()
1021        .filter(|(name, _)| !(cross_origin && is_sensitive_header(name)))
1022        .collect()
1023}
1024
1025async fn guarded_get_inner(
1026    client: &Client,
1027    url: &str,
1028    extra_headers: &[(HeaderName, HeaderValue)],
1029    check_privacy: bool,
1030    max_redirects: usize,
1031) -> Result<Response> {
1032    // `client` is retained in the signature for API stability + as the policy
1033    // template; the actual send goes through a per-hop IP-pinned client.
1034    let _ = client;
1035    let mut current = Url::parse(url).with_context(|| format!("not a valid URL {url:?}"))?;
1036    // The origin the caller's credentials belong to; a hop off it drops them.
1037    let original = current.clone();
1038
1039    for _ in 0..=max_redirects {
1040        check_scheme(&current)?;
1041        // Re-validate PRIVACY on EVERY hop: a public URL can `30x` to a
1042        // secret-bearing private feed (Substack/Patreon/tokened podcast). Without
1043        // this, the private target would be fetched — its body streamed and
1044        // reflected into the UI — before storage is refused, violating the
1045        // "never fetched" half of the public-feeds-only guarantee. Classify the
1046        // resolved target BEFORE the request and abort the whole fetch if private.
1047        // (Skipped for non-feed atproto identity fetches — see
1048        // [`guarded_get_no_privacy`].)
1049        if check_privacy {
1050            if let crate::feed::FeedPrivacy::Private(reason) =
1051                crate::feed::classify_feed_privacy(current.as_str())
1052            {
1053                bail!("refusing to fetch private/paid feed URL (redirect target): {reason}");
1054            }
1055        }
1056        // Re-validate on EVERY hop and capture the vetted address to pin to.
1057        let vetted = resolve_and_check(&current).await?;
1058        let host = current
1059            .host_str()
1060            .context("URL lost its host between hops")?
1061            .to_string();
1062        let hop_client = pinned_client(&host, vetted)?;
1063
1064        let mut req = hop_client.get(current.clone());
1065        // Sensitive headers (Authorization / Cookie / …) are applied only while
1066        // the hop is still on the ORIGINAL origin: a hostile upstream must not be
1067        // able to `302` a caller's bearer token onto a host it controls.
1068        for (name, value) in hop_headers(&original, &current, extra_headers) {
1069            req = req.header(name.clone(), value.clone());
1070        }
1071        let resp = req
1072            .send()
1073            .await
1074            .with_context(|| format!("fetching {current}"))?;
1075
1076        // **Only the statuses that actually relocate — NOT all of `3xx`.**
1077        //
1078        // `is_redirection()` is `300..=399`, which swallows `304 Not Modified`.
1079        // A 304 carries no `Location` *by definition*, so it fell into the
1080        // branch below and failed the whole fetch with "redirect response
1081        // without a usable Location header". `feed.rs` sends `If-None-Match` /
1082        // `If-Modified-Since` on every poll and has a correct 304 branch — which
1083        // could therefore never be reached. The effect was that "nothing new"
1084        // became a recorded failure plus exponential backoff, punishing exactly
1085        // the feeds that implement conditional GET properly. Observed in
1086        // production against 9to5mac.com, proton.me and kodi.tv, all live.
1087        //
1088        // **304 is the ONLY status carved out.** Everything else in `3xx`
1089        // relocates in some sense, and stays inside this branch — because
1090        // `guarded_get_no_redirect` documents that it "refuses redirects
1091        // outright", and the OAuth mix-up defence rests on that holding for all
1092        // of them, not just the five we would otherwise follow.
1093        if resp.status() != reqwest::StatusCode::NOT_MODIFIED && resp.status().is_redirection() {
1094            if max_redirects == 0 {
1095                bail!(
1096                    "refusing to follow a {} redirect while fetching {url:?} \u{2014} \
1097                     this document's origin is load-bearing and must not be moved",
1098                    resp.status()
1099                );
1100            }
1101            // Of the relocating statuses, only these five name a single target
1102            // worth following. `300 Multiple Choices` names no one target, and
1103            // `305 Use Proxy` names a PROXY — following it would route the
1104            // request through a host the RESPONSE chose.
1105            //
1106            // **Refused, not returned.** Handing one back would be worse than
1107            // erroring: callers do not uniformly check the status.
1108            // `web::resolve_feed_url` reads the body straight into feed
1109            // autodiscovery, so a `305` whose error page carries a
1110            // `<link rel="alternate">` would become a subscription.
1111            if !matches!(resp.status().as_u16(), 301 | 302 | 303 | 307 | 308) {
1112                bail!(
1113                    "refusing to act on a {} response while fetching {url:?} \u{2014} \
1114                     it names no single target that can be followed safely",
1115                    resp.status()
1116                );
1117            }
1118            let location = resp
1119                .headers()
1120                .get(reqwest::header::LOCATION)
1121                .and_then(|v| v.to_str().ok())
1122                .context("redirect response without a usable Location header")?;
1123            // Resolve the (possibly relative) Location against the current URL,
1124            // then loop to re-validate the new hop before touching it.
1125            current = current
1126                .join(location)
1127                .with_context(|| format!("resolving redirect Location {location:?}"))?;
1128            continue;
1129        }
1130
1131        return Ok(resp);
1132    }
1133
1134    bail!("too many redirects (> {max_redirects}) while fetching {url:?}")
1135}
1136
1137/// POST a JSON body to a **user-influenced** URL through the SSRF guard.
1138///
1139/// The write-side counterpart to [`guarded_get_no_privacy`], and the only way
1140/// [`crate::atproto::PdsClient`] is allowed to reach a PDS host it did not
1141/// choose. It runs the same scheme allow-list, the same IP allow-list, and the
1142/// same connect-pinning (via `pinned_client`), so the DNS-rebinding window
1143/// between "`assert_public_target` said this host is public" and "the TCP
1144/// handshake happens" is closed for writes exactly as it is for reads.
1145///
1146/// `Content-Type: application/json` is set here rather than by the caller, so
1147/// the one header the XRPC wire format requires cannot be forgotten; the caller
1148/// passes only its credential header(s).
1149///
1150/// **Redirects are refused, not followed** — the single deliberate divergence
1151/// from [`guarded_get`]. A `307`/`308` re-sends the method *and the body*
1152/// verbatim, and reqwest's cross-origin header sanitisation only strips
1153/// **headers**: an app password or a record body lives in the JSON payload, so a
1154/// hostile PDS answering `307 Location: https://evil.example/collect` would
1155/// exfiltrate it however carefully the headers were handled. There is no
1156/// legitimate reason for a PDS to redirect an `com.atproto.repo.*` write, so the
1157/// safe behaviour and the correct behaviour coincide: `Err`, loudly.
1158pub async fn guarded_post_json(
1159    client: &Client,
1160    url: &str,
1161    extra_headers: &[(HeaderName, HeaderValue)],
1162    body: Vec<u8>,
1163) -> Result<Response> {
1164    guarded_post(client, url, extra_headers, PostBody::Json(body)).await
1165}
1166
1167/// A request body together with the content type that describes it.
1168///
1169/// The two travel as ONE value deliberately. Passing the content type alongside
1170/// the bytes made it possible to send a JSON body labelled as a form, or the
1171/// reverse — a swap no test could see without a live server, and the SSRF guard
1172/// forbids pointing one of these at loopback. Deriving the header from the same
1173/// value that produces the bytes removes the failure mode instead of watching
1174/// for it.
1175pub(crate) enum PostBody<'a> {
1176    Json(Vec<u8>),
1177    Form(&'a [(&'a str, &'a str)]),
1178}
1179
1180impl PostBody<'_> {
1181    fn content_type(&self) -> HeaderValue {
1182        match self {
1183            PostBody::Json(_) => HeaderValue::from_static("application/json"),
1184            PostBody::Form(_) => HeaderValue::from_static("application/x-www-form-urlencoded"),
1185        }
1186    }
1187
1188    /// Every form value goes through the serializer rather than string
1189    /// interpolation: an OAuth form carries the authorization code, the PKCE
1190    /// verifier and the client assertion, and a raw `&` or `=` in any of them
1191    /// would otherwise splice an extra parameter into the request.
1192    fn into_bytes(self) -> Vec<u8> {
1193        match self {
1194            PostBody::Json(bytes) => bytes,
1195            PostBody::Form(params) => {
1196                let mut ser = url::form_urlencoded::Serializer::new(String::new());
1197                for (k, v) in params {
1198                    ser.append_pair(k, v);
1199                }
1200                ser.finish().into_bytes()
1201            }
1202        }
1203    }
1204}
1205
1206/// POST a form-encoded body to a **user-influenced** URL through the SSRF guard.
1207///
1208/// The OAuth counterpart to [`guarded_post_json`]: PAR, token exchange and
1209/// refresh are all `application/x-www-form-urlencoded`. It matters more here
1210/// than anywhere else that the guard applies — these are the requests that
1211/// carry the client assertion and the authorization code, so an issuer URL
1212/// that resolves to loopback or RFC1918 has to fail closed *before* the
1213/// credential leaves the process.
1214///
1215/// Redirects are refused for the same reason as [`guarded_post_json`], and more
1216/// acutely: a `307` would re-send the assertion and code to the new host.
1217pub async fn guarded_post_form(
1218    client: &Client,
1219    url: &str,
1220    extra_headers: &[(HeaderName, HeaderValue)],
1221    params: &[(&str, &str)],
1222) -> Result<Response> {
1223    guarded_post(client, url, extra_headers, PostBody::Form(params)).await
1224}
1225
1226/// The shared body of [`guarded_post_json`] and [`guarded_post_form`]. Kept as
1227/// one function so the guard cannot drift between the two content types.
1228async fn guarded_post(
1229    client: &Client,
1230    url: &str,
1231    extra_headers: &[(HeaderName, HeaderValue)],
1232    body: PostBody<'_>,
1233) -> Result<Response> {
1234    let content_type = body.content_type();
1235    let body = body.into_bytes();
1236    // As in `guarded_get_inner`: `client` is the policy template; the send goes
1237    // through a freshly built, IP-pinned client.
1238    let _ = client;
1239    let target = Url::parse(url).with_context(|| format!("not a valid URL {url:?}"))?;
1240    check_scheme(&target)?;
1241    let vetted = resolve_and_check(&target).await?;
1242    let host = target.host_str().context("URL has no host")?.to_string();
1243    let hop_client = pinned_client(&host, vetted)?;
1244
1245    let mut req = hop_client
1246        .post(target.clone())
1247        .header(CONTENT_TYPE, content_type)
1248        .body(body);
1249    for (name, value) in extra_headers {
1250        req = req.header(name.clone(), value.clone());
1251    }
1252    let resp = req
1253        .send()
1254        .await
1255        .with_context(|| format!("posting to {target}"))?;
1256
1257    // Unlike the GET path this refuses the WHOLE of `3xx`, `304` included, and
1258    // that is deliberate: nothing here sends `If-None-Match`/`If-Modified-Since`,
1259    // so a 304 to a POST is a server protocol violation rather than a
1260    // conditional-GET success, and there is no sane way to act on it.
1261    //
1262    // It is worded separately all the same. Calling a 304 "a redirect we refused
1263    // to follow" is the same misattribution that made the GET bug take a
1264    // production investigation to find — the message should not send the next
1265    // reader looking for a `Location` that was never supposed to exist.
1266    if resp.status().is_redirection() {
1267        if resp.status() == reqwest::StatusCode::NOT_MODIFIED {
1268            bail!(
1269                "a POST to {url:?} answered 304 Not Modified, which is not a valid \
1270                 response to a request carrying no conditional headers"
1271            );
1272        }
1273        let location = resp
1274            .headers()
1275            .get(reqwest::header::LOCATION)
1276            .and_then(|v| v.to_str().ok())
1277            .unwrap_or("<none>");
1278        bail!(
1279            "refusing to follow a {} redirect on a POST to {url:?} (Location: {location}) — \
1280             a 307/308 would re-send the request body to the new host",
1281            resp.status()
1282        );
1283    }
1284
1285    Ok(resp)
1286}
1287
1288/// Read a response body, streaming chunk-by-chunk and **aborting** the moment
1289/// the accumulated size would exceed [`MAX_BODY_BYTES`]. Never trusts
1290/// `Content-Length` (gzip strips it) and never fully buffers an over-cap body —
1291/// this is the decompression-bomb / OOM guard.
1292pub async fn read_capped(mut resp: Response) -> Result<Vec<u8>> {
1293    let mut buf: Vec<u8> = Vec::with_capacity(16 * 1024);
1294    while let Some(chunk) = resp.chunk().await.context("reading response body chunk")? {
1295        if buf.len() + chunk.len() > MAX_BODY_BYTES {
1296            bail!(
1297                "response body exceeded the {} byte cap; aborting",
1298                MAX_BODY_BYTES
1299            );
1300        }
1301        buf.extend_from_slice(&chunk);
1302    }
1303    Ok(buf)
1304}
1305
1306/// Validate that a URL is safe to use as an outbound target: `http`/`https`
1307/// scheme AND every resolved IP passes the SSRF allow-list. Returns `Ok(())` for
1308/// a public target, `Err` for a forbidden one (loopback / link-local / private /
1309/// ULA / CGNAT / metadata) or a bad scheme.
1310///
1311/// Use this to vet a URL *before* it is stashed and later fetched by a client
1312/// that does not itself route through [`guarded_get`] — notably an atproto PDS
1313/// `serviceEndpoint` resolved out of a (hostile-controllable) DID document, so a
1314/// `serviceEndpoint: "http://169.254.169.254/"` is rejected at resolve time
1315/// rather than reaching a raw XRPC client.
1316pub async fn assert_public_target(url: &str) -> Result<()> {
1317    let parsed = Url::parse(url).with_context(|| format!("not a valid URL {url:?}"))?;
1318    check_scheme(&parsed)?;
1319    resolve_and_check(&parsed).await?;
1320    Ok(())
1321}
1322
1323/// Scheme-allow-list a URL destined to be rendered as an `href` (an entry's
1324/// "View original" link, a feed's site link). Accepts only `http`/`https`;
1325/// anything else (notably `javascript:` / `data:` — stored-XSS vectors that
1326/// survive HTML escaping) yields `None` so the caller drops the link.
1327pub fn safe_link(raw: &str) -> Option<String> {
1328    let trimmed = raw.trim();
1329    if trimmed.is_empty() {
1330        return None;
1331    }
1332    match Url::parse(trimmed) {
1333        Ok(u) if matches!(u.scheme(), "http" | "https") => Some(trimmed.to_string()),
1334        _ => None,
1335    }
1336}
1337
1338#[cfg(test)]
1339pub(crate) mod tests {
1340    use super::*;
1341
1342    #[test]
1343    fn forbids_loopback_and_link_local_and_private_v4() {
1344        for ip in [
1345            "127.0.0.1",
1346            "127.1.2.3",
1347            "169.254.169.254", // cloud metadata
1348            "10.0.0.5",
1349            "172.16.9.9",
1350            "192.168.1.1",
1351            "0.0.0.0",
1352            "255.255.255.255",
1353            "100.64.0.1", // 100.64/10 CGNAT range (also overlay VPNs)
1354        ] {
1355            let ip: IpAddr = ip.parse().unwrap();
1356            assert!(is_forbidden_ip(&ip), "{ip} should be forbidden");
1357        }
1358    }
1359
1360    #[test]
1361    fn allows_public_v4() {
1362        for ip in ["1.1.1.1", "8.8.8.8", "93.184.216.34"] {
1363            let ip: IpAddr = ip.parse().unwrap();
1364            assert!(!is_forbidden_ip(&ip), "{ip} should be allowed");
1365        }
1366    }
1367
1368    #[test]
1369    fn forbids_internal_v6() {
1370        for ip in [
1371            "::1",
1372            "fe80::1",
1373            "fc00::1",
1374            "fd00::1",
1375            "::ffff:127.0.0.1",
1376            "::",
1377        ] {
1378            let ip: IpAddr = ip.parse().unwrap();
1379            assert!(is_forbidden_ip(&ip), "{ip} should be forbidden");
1380        }
1381    }
1382
1383    /// **An IPv6 address that EMBEDS a forbidden IPv4 one is a forbidden address,
1384    /// and four families of them were getting through.**
1385    ///
1386    /// `is_forbidden_v6` unwrapped IPv4-mapped (`::ffff:a.b.c.d`) and
1387    /// IPv4-compatible (`::a.b.c.d`) forms, which is where `to_ipv4()` stops. It
1388    /// did not unwrap:
1389    ///
1390    /// * **NAT64**, `64:ff9b::/32` — the well-known prefix of RFC 6052 and the
1391    ///   local-use prefix of RFC 8215. On a NAT64/DNS64 network,
1392    ///   `64:ff9b::a9fe:a9fe` is the cloud metadata service.
1393    /// * **6to4**, `2002::/16` (RFC 3056) — the IPv4 sits in the next two groups,
1394    ///   so `2002:a9fe:a9fe::` is the same address again.
1395    /// * **IPv4-translated**, `::ffff:0:0/96` (RFC 2765) — one group away from the
1396    ///   mapped form `to_ipv4()` does handle.
1397    /// * **Teredo**, `2001::/32` (RFC 4380) — carries the relay's IPv4 in groups
1398    ///   2-3 and the client's, obfuscated by XOR with all-ones, in groups 6-7.
1399    /// * **ISATAP**, RFC 5214 — the IPv4 in the low 32 bits behind the IANA
1400    ///   `00-00-5E-FE` OUI, under ANY /64. The only one of the five with no
1401    ///   prefix to anchor on, so `2606:4700::5efe:c0a8:1` is an entirely
1402    ///   ordinary-looking global address that names 192.168.0.1.
1403    ///
1404    /// Found while bumping a JavaScript dependency whose advisory was the same
1405    /// class: "no classifier recognizes the NAT64 local-use range". Ours did not
1406    /// either.
1407    ///
1408    /// Whether a given deployment can route these depends on a translator being on
1409    /// the path — but the attacker does not need to know that, only to try it, and
1410    /// an IPv6-only network with DNS64 is now the ordinary case rather than the
1411    /// exotic one. This guard is defence in depth against exactly the address that
1412    /// reaches the host's own network without looking like it.
1413    #[test]
1414    fn forbids_ipv6_that_embeds_a_forbidden_ipv4() {
1415        for (ip, what) in [
1416            ("64:ff9b::7f00:1", "NAT64 well-known -> 127.0.0.1"),
1417            ("64:ff9b::a9fe:a9fe", "NAT64 well-known -> 169.254.169.254"),
1418            ("64:ff9b::c0a8:1", "NAT64 well-known -> 192.168.0.1"),
1419            ("64:ff9b:1::7f00:1", "NAT64 local-use, RFC 8215"),
1420            ("64:ff9b:1:ffff::1", "anywhere in the NAT64 /32"),
1421            ("2002:7f00:1::", "6to4 -> 127.0.0.1"),
1422            ("2002:a9fe:a9fe::", "6to4 -> 169.254.169.254"),
1423            ("::ffff:0:7f00:1", "IPv4-translated -> 127.0.0.1"),
1424            // Teredo, laid out the way the format actually is: server IPv4 in
1425            // groups 2-3, client IPv4 in groups 6-7 XORed with all-ones. Each case
1426            // keeps the OTHER field public, so it fails for the reason its label
1427            // claims rather than because a zero field is forbidden anyway.
1428            ("2001:0:7f00:1:0:0:f7f7:fbfb", "Teredo server -> 127.0.0.1"),
1429            ("2001:0:808:808:0:0:80ff:fffe", "Teredo client -> 127.0.0.1"),
1430            (
1431                "2001:0:808:808:0:0:5601:5601",
1432                "Teredo client -> 169.254.169.254",
1433            ),
1434            // ISATAP (RFC 5214): the IPv4 sits in the low 32 bits behind the
1435            // IANA-reserved `00-00-5E-FE` OUI, under ANY /64 — so unlike the
1436            // four above there is no prefix to anchor on, and a perfectly
1437            // ordinary-looking global address can carry one.
1438            ("2001:db8::5efe:7f00:1", "ISATAP -> 127.0.0.1"),
1439            ("2001:db8::5efe:a9fe:a9fe", "ISATAP -> 169.254.169.254"),
1440            // All four values the IID's first byte can take, kept as named
1441            // regressions. RFC 5214 spells it `000000ug`, so `u` and `g` are
1442            // both free. The first draft of this guard enumerated only the two
1443            // with `g` clear, leaving the other two allowed — a one-bit bypass
1444            // of a guard that read as complete. Reverting the arm to that
1445            // enumeration fails on the `g=1` rows below.
1446            (
1447                "2001:db8::200:5efe:a9fe:a9fe",
1448                "ISATAP u=1 g=0 -> 169.254.169.254",
1449            ),
1450            ("2001:db8::100:5efe:7f00:1", "ISATAP u=0 g=1 -> 127.0.0.1"),
1451            ("2001:db8::300:5efe:7f00:1", "ISATAP u=1 g=1 -> 127.0.0.1"),
1452            (
1453                "2606:4700::5efe:c0a8:1",
1454                "ISATAP under a REAL public prefix -> 192.168.0.1",
1455            ),
1456            // **An ISATAP identifier INSIDE another family's prefix.** 6to4
1457            // delegates `2002:<site-v4>::/48` to whoever owns that IPv4, and
1458            // the site assigns identifiers inside it — so a site running ISATAP
1459            // in its own 6to4 space produces exactly this. The two families
1460            // read DISJOINT bits (6to4 the site address in groups 1-2, ISATAP
1461            // the tunnel endpoint in groups 6-7), so both readings are true at
1462            // once and checking only the first is a bypass.
1463            (
1464                "2002:808:808:0:0:5efe:a9fe:a9fe",
1465                "6to4 site 8.8.8.8 + ISATAP -> 169.254.169.254",
1466            ),
1467            (
1468                "2002:808:808:0:0:5efe:7f00:1",
1469                "6to4 site 8.8.8.8 + ISATAP -> 127.0.0.1",
1470            ),
1471            (
1472                "2002:101:101:0:0:5efe:c0a8:1",
1473                "6to4 site 1.1.1.1 + ISATAP -> 192.168.0.1",
1474            ),
1475        ] {
1476            let parsed: IpAddr = ip.parse().unwrap();
1477            assert!(
1478                is_forbidden_ip(&parsed),
1479                "{ip} reaches {what} and was allowed",
1480            );
1481        }
1482    }
1483
1484    /// The other direction, and it is not decoration: refusing every address that
1485    /// merely *looks* translated would take out ordinary public traffic. A 6to4
1486    /// address wrapping a PUBLIC IPv4, and a Teredo address wrapping one, must both
1487    /// still be allowed — that is what makes this a decode rather than a
1488    /// prefix-blocklist.
1489    #[test]
1490    fn allows_ipv6_that_embeds_a_public_ipv4() {
1491        for (ip, what) in [
1492            ("2002:0808:0808::", "6to4 -> 8.8.8.8"),
1493            (
1494                "2001:0:808:808:0:0:f7f7:fbfb",
1495                "Teredo, server 8.8.8.8 and client 8.8.4.4",
1496            ),
1497            ("::ffff:0:808:808", "IPv4-translated -> 8.8.8.8"),
1498            ("2606:4700::5efe:808:808", "ISATAP -> 8.8.8.8"),
1499            // The DNS64 case, and the reason NAT64 is decoded rather than
1500            // prefix-refused: a resolver doing DNS64 synthesises exactly this
1501            // for an IPv4-only host, so refusing the prefix outright makes
1502            // every IPv4-only feed publisher unfetchable on an IPv6-only
1503            // network — the very network that motivated the guard.
1504            ("64:ff9b::808:808", "NAT64 well-known prefix -> 8.8.8.8"),
1505            // Both readings of one address, both public. The arms accumulate,
1506            // so this is the case that keeps that a decode rather than "any
1507            // 6to4 address carrying a `5efe` identifier is refused".
1508            (
1509                "2002:808:808:0:0:5efe:808:404",
1510                "6to4 site 8.8.8.8 + ISATAP 8.8.4.4",
1511            ),
1512        ] {
1513            let parsed: IpAddr = ip.parse().unwrap();
1514            assert!(!is_forbidden_ip(&parsed), "{ip} is {what} and was refused");
1515        }
1516    }
1517
1518    /// **The well-known prefix decodes; the local-use one does not — and the
1519    /// difference is deliberate, so it needs a test and not just a comment.**
1520    ///
1521    /// `64:ff9b::/96` is a fixed-length prefix by definition (RFC 6052 §3.1), so
1522    /// the embedded IPv4 is unambiguously the last 32 bits. RFC 8215's local-use
1523    /// `64:ff9b:1::/48` is not: RFC 6052 §2.2 allows six embedding lengths and
1524    /// which one a deployment chose is a property of that deployment. Guessing
1525    /// wrong reads the wrong bits and could render an internal target as a
1526    /// public-looking address, so everything outside the /96 is refused whole.
1527    ///
1528    /// The cost is real and accepted: a site translating through its local-use
1529    /// prefix cannot fetch through this reader. The alternative is a decode that
1530    /// is wrong whenever the guess is wrong, in the one direction that matters.
1531    ///
1532    /// Extending the decode to the whole `/32` fails this test.
1533    #[test]
1534    fn a_local_use_nat64_prefix_is_refused_even_wrapping_a_public_address() {
1535        let ip: IpAddr = "64:ff9b:1::808:808".parse().unwrap();
1536        assert!(
1537            is_forbidden_ip(&ip),
1538            "the local-use NAT64 prefix was decoded as if its embedding length \
1539             were known",
1540        );
1541    }
1542
1543    /// **The DNS64 path, end to end through the function that rejects answer
1544    /// sets.** This is the interaction the unit cases cannot see.
1545    ///
1546    /// On an IPv6-only network a DNS64 resolver (RFC 6147) synthesises a
1547    /// well-known-prefix AAAA for every IPv4-only host, and that synthesised
1548    /// address is the ONLY answer — there is no "ordinary address we resolve
1549    /// anyway". Since `first_vetted` rejects a whole set if any member is
1550    /// forbidden, refusing `64:ff9b::/96` outright made every IPv4-only feed
1551    /// publisher unfetchable on exactly the network this guard was written for.
1552    ///
1553    /// Both directions, because the fix must not cost the guard: a synthesised
1554    /// answer for a PUBLIC host resolves, and a synthesised answer for the
1555    /// metadata service still poisons the set.
1556    #[test]
1557    fn a_dns64_answer_set_for_an_ipv4_only_host_is_fetchable() {
1558        let synthesised: SocketAddr = "[64:ff9b::808:808]:80".parse().unwrap();
1559        let public_v4: SocketAddr = "1.2.3.4:80".parse().unwrap();
1560
1561        // IPv6-only: the synthesised address is the whole answer.
1562        assert_eq!(
1563            first_vetted("v4only.example", [synthesised].into_iter()).unwrap(),
1564            synthesised,
1565            "a DNS64-synthesised answer for a public host was refused, which \
1566             makes every IPv4-only publisher unfetchable behind NAT64",
1567        );
1568        // Dual-stack with DNS64: the synthesised answer must not poison the set.
1569        assert!(first_vetted("both.example", [public_v4, synthesised].into_iter()).is_ok());
1570
1571        // And the guard still bites: synthesising the metadata service is
1572        // exactly the attack, and one such answer rejects the whole set.
1573        let hostile: SocketAddr = "[64:ff9b::a9fe:a9fe]:80".parse().unwrap();
1574        assert!(
1575            first_vetted("evil.example", [public_v4, hostile].into_iter()).is_err(),
1576            "a NAT64-synthesised metadata address was accepted",
1577        );
1578        assert!(first_vetted("evil.example", [hostile].into_iter()).is_err());
1579    }
1580
1581    /// **The ISATAP marker is read loosely ON PURPOSE, and this is the test that
1582    /// says so.**
1583    ///
1584    /// RFC 5214 spells the interface identifier `000000ug 00000000 0x5E 0xFE` +
1585    /// the IPv4, so a spec-exact test would also require the six reserved bits
1586    /// of `seg[4]` to be zero and would ALLOW the address below. `embedded_v4`
1587    /// tests only for `0x5efe` in group 5, so it refuses it.
1588    ///
1589    /// That is a deliberate over-refusal, and without this test it was a
1590    /// comment and nothing else: restoring the spec-exact mask
1591    /// (`seg[4] & !0x0300 == 0`) passed all 960 tests. The asymmetry is the
1592    /// argument — reading the marker loosely costs a false positive only if a
1593    /// non-ISATAP interface identifier carries IANA's own `00-00-5E` OUI *and*
1594    /// its low 32 bits decode to an internal address, while reading it strictly
1595    /// costs a total bypass if any tunnel driver is more lenient than the RFC.
1596    ///
1597    /// So if a future change tightens this arm toward the spec, that is a
1598    /// decision to take deliberately, by deleting this test and saying why —
1599    /// not something to discover from a bypass.
1600    #[test]
1601    fn a_reserved_bit_in_the_isatap_identifier_does_not_buy_a_bypass() {
1602        let ip: IpAddr = "2001:db8::400:5efe:7f00:1".parse().unwrap();
1603        assert!(
1604            is_forbidden_ip(&ip),
1605            "an identifier carrying 00-00-5E-FE and 127.0.0.1 was allowed \
1606             because a reserved bit was set",
1607        );
1608    }
1609
1610    /// **The ISATAP test is on the interface identifier, so it matches under any
1611    /// prefix — including prefixes that belong to one of the other four.**
1612    ///
1613    /// A Teredo address with zero flags whose obfuscated port happens to be
1614    /// `0x5efe` matches the ISATAP pattern too, and the two readings disagree:
1615    /// Teredo stores the client address complemented, so the ISATAP reading of
1616    /// the same bits is its bitwise inverse. Here the Teredo reading is server
1617    /// 8.8.8.8 and client 128.255.255.254 — both public, so the address is
1618    /// legitimate — while the ISATAP reading of those low 32 bits is 127.0.0.1.
1619    ///
1620    /// Teredo is the one arm that still RETURNS rather than accumulating, which
1621    /// suppresses the ISATAP reading of bits Teredo has already claimed.
1622    /// `2001:0000::/32` is IANA-assigned Teredo space, a real ISATAP host would
1623    /// not be using it, and refusing this would be a false positive on an
1624    /// address whose traffic goes to a Teredo relay rather than to loopback.
1625    ///
1626    /// Making the Teredo arm accumulate like the others — i.e. letting the
1627    /// ISATAP arm also read groups 6-7 here — fails this test. That is the
1628    /// whole difference between this case and the 6to4 one: there the two
1629    /// families read disjoint bits and both readings hold, here they read the
1630    /// same field and contradict each other.
1631    #[test]
1632    fn a_teredo_address_is_read_as_teredo_not_as_isatap() {
1633        let ip: IpAddr = "2001:0:808:808:0:5efe:7f00:1".parse().unwrap();
1634        assert!(
1635            !is_forbidden_ip(&ip),
1636            "an address in Teredo space was read as ISATAP and wrongly refused",
1637        );
1638    }
1639
1640    /// **The embedded-IPv4 arm may only ADD refusals, never grant permission.**
1641    ///
1642    /// It is checked before the link-local and ULA rules, so if it returned a
1643    /// verdict rather than falling through, an ISATAP address wrapping a PUBLIC
1644    /// IPv4 under an `fe80::/10` prefix would come back allowed — a link-local
1645    /// address let through because the thing it embeds happens to be fine.
1646    ///
1647    /// Changing `if embedded_v4(..).any(..) { return true; }` to return the
1648    /// condition fails this test — and also `forbids_internal_v6` and
1649    /// `every_blocklist_branch_is_load_bearing`, which were already standing
1650    /// guard over the fall-through in general. So this case is a NAMED
1651    /// regression for the ISATAP interaction rather than the only thing holding
1652    /// the property down; it is measured, not assumed, and stated that way
1653    /// because a test whose comment claims more than it catches is the defect
1654    /// this file keeps finding.
1655    #[test]
1656    fn a_link_local_isatap_address_is_still_refused_for_being_link_local() {
1657        let ip: IpAddr = "fe80::5efe:808:808".parse().unwrap();
1658        assert!(
1659            is_forbidden_ip(&ip),
1660            "fe80::/10 wrapping a public IPv4 escaped the link-local rule",
1661        );
1662    }
1663
1664    #[test]
1665    fn allows_public_v6() {
1666        let ip: IpAddr = "2606:4700:4700::1111".parse().unwrap();
1667        assert!(!is_forbidden_ip(&ip));
1668    }
1669
1670    #[tokio::test]
1671    async fn resolve_and_check_rejects_ip_literals() {
1672        for bad in [
1673            "http://127.0.0.1/feed.xml",
1674            "http://169.254.169.254/latest/meta-data/",
1675            "http://[::1]:80/x",
1676            "http://192.168.0.1/",
1677        ] {
1678            let u = Url::parse(bad).unwrap();
1679            assert!(
1680                resolve_and_check(&u).await.is_err(),
1681                "{bad} should be rejected"
1682            );
1683        }
1684    }
1685
1686    #[tokio::test]
1687    async fn resolve_and_check_allows_public_ip_literal() {
1688        let u = Url::parse("http://1.1.1.1/").unwrap();
1689        let addr = resolve_and_check(&u).await.unwrap();
1690        // The vetted address is pinned back verbatim (IP literal, no DNS).
1691        assert_eq!(addr, "1.1.1.1:80".parse::<SocketAddr>().unwrap());
1692    }
1693
1694    #[tokio::test]
1695    async fn resolve_and_check_pins_public_ipv6_literal() {
1696        let u = Url::parse("http://[2606:4700:4700::1111]:443/").unwrap();
1697        let addr = resolve_and_check(&u).await.unwrap();
1698        assert_eq!(
1699            addr,
1700            "[2606:4700:4700::1111]:443".parse::<SocketAddr>().unwrap()
1701        );
1702    }
1703
1704    // ── the pinned-client cache ──────────────────────────────────────────────
1705
1706    const V4: &str = "93.184.216.34:443";
1707    const V4_OTHER: &str = "93.184.216.35:443";
1708
1709    fn at(base: Instant, secs: u64) -> Instant {
1710        base + Duration::from_secs(secs)
1711    }
1712
1713    /// A repeat request to the same vetted address REUSES the client, so the
1714    /// connection pool survives and the TLS handshake is paid once.
1715    ///
1716    /// Measured motivation: a fresh connection to this project's PDS costs 91 ms
1717    /// against 30 ms warm, which was most of the ~3x gap between the Rust repo
1718    /// backend and the Node sidecar.
1719    #[test]
1720    fn the_same_vetted_address_reuses_one_client() {
1721        let cache = PinnedClients::new();
1722        let now = Instant::now();
1723        let addr: SocketAddr = V4.parse().unwrap();
1724
1725        for i in 0..5 {
1726            cache.get("example.com", addr, at(now, i)).unwrap();
1727        }
1728        assert_eq!(
1729            cache.builds.load(Ordering::Relaxed),
1730            1,
1731            "each request rebuilt the client, so every call pays a TLS handshake"
1732        );
1733    }
1734
1735    /// **A CHANGED ADDRESS MUST NOT REUSE THE POOL.**
1736    ///
1737    /// This is the property that makes the cache safe. The DNS-rebinding defence
1738    /// is that we connect only to an address vetted for THIS request; a cache
1739    /// keyed on the host alone would hand back a connection pinned to an address
1740    /// vetted minutes ago, quietly undoing it. The key includes the address, so
1741    /// a move is a miss.
1742    #[test]
1743    fn a_changed_address_does_not_reuse_the_pooled_client() {
1744        let cache = PinnedClients::new();
1745        let now = Instant::now();
1746
1747        cache.get("example.com", V4.parse().unwrap(), now).unwrap();
1748        cache
1749            .get("example.com", V4_OTHER.parse().unwrap(), at(now, 1))
1750            .unwrap();
1751
1752        assert_eq!(
1753            cache.builds.load(Ordering::Relaxed),
1754            2,
1755            "the same host at a DIFFERENT address reused a connection pinned to the old one"
1756        );
1757        assert_eq!(cache.entries.lock().unwrap().len(), 2);
1758    }
1759
1760    /// Two hosts that happen to resolve to the same address still get their own
1761    /// clients — the pin is per host, and SNI/Host differ.
1762    #[test]
1763    fn different_hosts_at_one_address_are_separate_clients() {
1764        let cache = PinnedClients::new();
1765        let now = Instant::now();
1766        let addr: SocketAddr = V4.parse().unwrap();
1767
1768        cache.get("a.example.com", addr, now).unwrap();
1769        cache.get("b.example.com", addr, now).unwrap();
1770        assert_eq!(cache.builds.load(Ordering::Relaxed), 2);
1771    }
1772
1773    /// An entry idle past the TTL is rebuilt, bounding how long a pooled
1774    /// connection to a once-vetted address can live.
1775    #[test]
1776    fn an_idle_entry_is_rebuilt_after_the_ttl() {
1777        let cache = PinnedClients::new();
1778        let now = Instant::now();
1779        let addr: SocketAddr = V4.parse().unwrap();
1780
1781        cache.get("example.com", addr, now).unwrap();
1782        cache
1783            .get(
1784                "example.com",
1785                addr,
1786                now + PINNED_CLIENT_TTL + Duration::from_secs(1),
1787            )
1788            .unwrap();
1789        assert_eq!(cache.builds.load(Ordering::Relaxed), 2);
1790    }
1791
1792    /// Use keeps an entry alive: a client fetched every minute must not be
1793    /// rebuilt just because it was first created more than a TTL ago. The TTL is
1794    /// idle time, not total age — otherwise the busiest client in the process
1795    /// would be the one thrown away on a schedule.
1796    #[test]
1797    fn continued_use_keeps_an_entry_alive() {
1798        let cache = PinnedClients::new();
1799        let now = Instant::now();
1800        let addr: SocketAddr = V4.parse().unwrap();
1801
1802        for minute in 0..20 {
1803            cache
1804                .get("example.com", addr, at(now, minute * 60))
1805                .unwrap();
1806        }
1807        assert_eq!(
1808            cache.builds.load(Ordering::Relaxed),
1809            1,
1810            "a continuously-used client was expired by age rather than idleness"
1811        );
1812    }
1813
1814    /// A client must release its idle sockets BEFORE the cache releases the
1815    /// client. The other way round, every entry spends the tail of its life
1816    /// holding connections nothing will reuse — which is the whole cost the pool
1817    /// bound exists to avoid.
1818    #[test]
1819    fn idle_sockets_are_released_before_their_client_is() {
1820        assert!(
1821            POOL_IDLE_TIMEOUT < PINNED_CLIENT_TTL,
1822            "pool idle timeout {POOL_IDLE_TIMEOUT:?} is not shorter than the \
1823             client TTL {PINNED_CLIENT_TTL:?}"
1824        );
1825    }
1826
1827    /// The map is bounded. The poller talks to as many hosts as there are feeds,
1828    /// so an unbounded map would be a slow leak of connection pools.
1829    #[test]
1830    fn the_cache_is_bounded() {
1831        let cache = PinnedClients::new();
1832        let now = Instant::now();
1833        for i in 0..(MAX_PINNED_CLIENTS + 50) {
1834            let addr: SocketAddr = format!("93.184.216.34:{}", 1024 + i).parse().unwrap();
1835            cache.get(&format!("h{i}.example.com"), addr, now).unwrap();
1836        }
1837        assert!(
1838            cache.entries.lock().unwrap().len() <= MAX_PINNED_CLIENTS,
1839            "the cache grew past its bound"
1840        );
1841    }
1842
1843    #[test]
1844    fn pinned_client_builds_for_both_families() {
1845        // Both address families must produce a usable pinned client.
1846        assert!(pinned_client("example.com", "93.184.216.34:80".parse().unwrap()).is_ok());
1847        assert!(
1848            pinned_client("example.com", "[2606:4700:4700::1111]:443".parse().unwrap()).is_ok()
1849        );
1850    }
1851
1852    #[test]
1853    fn scheme_allowlist_rejects_non_http() {
1854        assert!(check_scheme(&Url::parse("http://example.com/").unwrap()).is_ok());
1855        assert!(check_scheme(&Url::parse("https://example.com/").unwrap()).is_ok());
1856        // url::Url::parse rejects `javascript:` as opaque, but file/ftp parse.
1857        assert!(check_scheme(&Url::parse("file:///etc/passwd").unwrap()).is_err());
1858        assert!(check_scheme(&Url::parse("ftp://example.com/").unwrap()).is_err());
1859    }
1860
1861    /// A raw HTTP server on loopback that answers every request with `body`
1862    /// (fixed Content-Length). Returns its `http://127.0.0.1:port/` base URL.
1863    /// Used to exercise [`read_capped`] against a real reqwest `Response`.
1864    /// A JSON server that records every raw request (head and body, read to
1865    /// `content-length`) and answers each with `reply`. For asserting on the
1866    /// BYTES a client sent — the only assertion that catches a record that is
1867    /// wrong on the way out.
1868    pub(crate) async fn serve_json_capturing(
1869        reply: Vec<u8>,
1870    ) -> (String, std::sync::Arc<std::sync::Mutex<Vec<String>>>) {
1871        use tokio::io::{AsyncReadExt, AsyncWriteExt};
1872        let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
1873        let addr = listener.local_addr().unwrap();
1874        let log = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
1875        let sink = std::sync::Arc::clone(&log);
1876        tokio::spawn(async move {
1877            loop {
1878                let Ok((mut sock, _)) = listener.accept().await else {
1879                    break;
1880                };
1881                let mut raw: Vec<u8> = Vec::new();
1882                let mut chunk = [0u8; 4096];
1883                let text = loop {
1884                    let Ok(n) = sock.read(&mut chunk).await else {
1885                        break String::new();
1886                    };
1887                    if n == 0 {
1888                        break String::from_utf8_lossy(&raw).to_string();
1889                    }
1890                    raw.extend_from_slice(&chunk[..n]);
1891                    let Some(split) = raw.windows(4).position(|w| w == b"\r\n\r\n") else {
1892                        continue;
1893                    };
1894                    let (head, body) = raw.split_at(split + 4);
1895                    let want = String::from_utf8_lossy(head).lines().find_map(|l| {
1896                        let (k, v) = l.split_once(':')?;
1897                        k.eq_ignore_ascii_case("content-length")
1898                            .then(|| v.trim().parse::<usize>().ok())?
1899                    });
1900                    if want.is_none_or(|w| body.len() >= w) {
1901                        break String::from_utf8_lossy(&raw).to_string();
1902                    }
1903                };
1904                sink.lock().unwrap().push(text);
1905                let header = format!(
1906                    "HTTP/1.1 200 OK\r\nContent-Type: application/json\r\nContent-Length: {}\r\nConnection: close\r\n\r\n",
1907                    reply.len()
1908                );
1909                let _ = sock.write_all(header.as_bytes()).await;
1910                let _ = sock.write_all(&reply).await;
1911                let _ = sock.flush().await;
1912            }
1913        });
1914        (format!("http://{addr}"), log)
1915    }
1916
1917    /// [`serve_body`] that also counts requests — for an assertion that a URL
1918    /// was NEVER fetched, which a body alone cannot make.
1919    pub(crate) async fn serve_body_counted(
1920        body: Vec<u8>,
1921    ) -> (String, std::sync::Arc<std::sync::atomic::AtomicUsize>) {
1922        use tokio::io::{AsyncReadExt, AsyncWriteExt};
1923        let hits = std::sync::Arc::new(std::sync::atomic::AtomicUsize::new(0));
1924        let counter = std::sync::Arc::clone(&hits);
1925        let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
1926        let addr = listener.local_addr().unwrap();
1927        tokio::spawn(async move {
1928            loop {
1929                let Ok((mut sock, _)) = listener.accept().await else {
1930                    break;
1931                };
1932                counter.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
1933                let body = body.clone();
1934                tokio::spawn(async move {
1935                    let mut buf = [0u8; 1024];
1936                    let _ = sock.read(&mut buf).await;
1937                    let header = format!(
1938                        "HTTP/1.1 200 OK\r\nContent-Type: application/rss+xml\r\nContent-Length: {}\r\nConnection: close\r\n\r\n",
1939                        body.len()
1940                    );
1941                    let _ = sock.write_all(header.as_bytes()).await;
1942                    let _ = sock.write_all(&body).await;
1943                    let _ = sock.flush().await;
1944                });
1945            }
1946        });
1947        (format!("http://{addr}/"), hits)
1948    }
1949
1950    /// Serve a different body per request, in order, repeating the last.
1951    ///
1952    /// **The fixed-body servers cannot test a walk.** `serve_body` answers every
1953    /// request identically, so a paging walk sees the same cursor twice and its
1954    /// repeat-detection guard stops it at two pages. Anything that only happens
1955    /// across pages — a budget accumulating, a cursor advancing — is therefore
1956    /// unreachable with them, which is how a cap that was per-page rather than
1957    /// per-walk once passed an entire suite.
1958    ///
1959    /// Each request is a fresh connection (`Connection: close`), so accept order
1960    /// is request order for the sequential walks that use this.
1961    pub(crate) async fn serve_bodies_in_sequence(bodies: Vec<Vec<u8>>) -> String {
1962        use tokio::io::{AsyncReadExt, AsyncWriteExt};
1963        assert!(!bodies.is_empty(), "serve_bodies_in_sequence needs a body");
1964        let bodies = std::sync::Arc::new(bodies);
1965        let next = std::sync::Arc::new(std::sync::atomic::AtomicUsize::new(0));
1966        let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
1967        let addr = listener.local_addr().unwrap();
1968        tokio::spawn(async move {
1969            loop {
1970                let Ok((mut sock, _)) = listener.accept().await else {
1971                    break;
1972                };
1973                let i = next.fetch_add(1, std::sync::atomic::Ordering::SeqCst);
1974                let body = bodies[i.min(bodies.len() - 1)].clone();
1975                tokio::spawn(async move {
1976                    // **Drain the whole request head, not one fixed read.** A
1977                    // DPoP-signed XRPC request head measures ~880 bytes, so a
1978                    // single 1024-byte read is within a longer NSID or cursor
1979                    // of leaving bytes unread — and closing with data still in
1980                    // the receive queue makes the kernel send RST instead of
1981                    // FIN, which can discard a response the client has not
1982                    // drained. That surfaces as an intermittent connection
1983                    // reset in a test whose failure would read as a budget bug.
1984                    let mut req = Vec::new();
1985                    let mut buf = [0u8; 1024];
1986                    // Drain the head, then the body it declares. Stopping at
1987                    // the head is not enough: the sidecar's list call is a POST,
1988                    // so on any platform that does not coalesce head and body
1989                    // into one segment the body stays in the receive queue, and
1990                    // closing on unread bytes is the RST-instead-of-FIN case
1991                    // this loop exists to avoid.
1992                    let mut want: Option<usize> = None;
1993                    loop {
1994                        match sock.read(&mut buf).await {
1995                            Ok(0) => break,
1996                            Ok(n) => {
1997                                req.extend_from_slice(&buf[..n]);
1998                                let Some(head_end) = req.windows(4).position(|w| w == b"\r\n\r\n")
1999                                else {
2000                                    continue;
2001                                };
2002                                let head_len = head_end + 4;
2003                                if want.is_none() {
2004                                    let head = String::from_utf8_lossy(&req[..head_len]);
2005                                    want = Some(
2006                                        head.lines()
2007                                            .find_map(|l| {
2008                                                let (k, v) = l.split_once(':')?;
2009                                                k.eq_ignore_ascii_case("content-length")
2010                                                    .then(|| v.trim().parse::<usize>().ok())?
2011                                            })
2012                                            .unwrap_or(0),
2013                                    );
2014                                }
2015                                if req.len() >= head_len + want.unwrap_or(0) {
2016                                    break;
2017                                }
2018                            }
2019                            Err(_) => break,
2020                        }
2021                    }
2022                    let header = format!(
2023                        "HTTP/1.1 200 OK\r\nContent-Type: application/json\r\nContent-Length: {}\r\nConnection: close\r\n\r\n",
2024                        body.len()
2025                    );
2026                    let _ = sock.write_all(header.as_bytes()).await;
2027                    let _ = sock.write_all(&body).await;
2028                    let _ = sock.flush().await;
2029                });
2030            }
2031        });
2032        format!("http://{addr}/")
2033    }
2034
2035    pub(crate) async fn serve_body(body: Vec<u8>) -> String {
2036        use tokio::io::{AsyncReadExt, AsyncWriteExt};
2037        let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
2038        let addr = listener.local_addr().unwrap();
2039        tokio::spawn(async move {
2040            loop {
2041                let (mut sock, _) = match listener.accept().await {
2042                    Ok(p) => p,
2043                    Err(_) => break,
2044                };
2045                let body = body.clone();
2046                tokio::spawn(async move {
2047                    // Drain the request headers (best-effort) then reply.
2048                    let mut buf = [0u8; 1024];
2049                    let _ = sock.read(&mut buf).await;
2050                    let header = format!(
2051                        "HTTP/1.1 200 OK\r\nContent-Type: text/plain\r\nContent-Length: {}\r\nConnection: close\r\n\r\n",
2052                        body.len()
2053                    );
2054                    let _ = sock.write_all(header.as_bytes()).await;
2055                    let _ = sock.write_all(&body).await;
2056                    let _ = sock.flush().await;
2057                });
2058            }
2059        });
2060        format!("http://{addr}/")
2061    }
2062
2063    #[tokio::test]
2064    async fn read_capped_rejects_over_cap_body() {
2065        // A body one byte over the cap must be rejected (and never fully
2066        // buffered past the cap). Fetch directly (bypassing the SSRF guard, which
2067        // rightly forbids loopback) to exercise read_capped on a real Response.
2068        let big = vec![b'x'; MAX_BODY_BYTES + 1];
2069        let base = serve_body(big).await;
2070        let client = reqwest::Client::builder().build().unwrap();
2071        let resp = client.get(&base).send().await.unwrap();
2072        let err = read_capped(resp).await.unwrap_err().to_string();
2073        assert!(err.contains("exceeded"), "unexpected error: {err}");
2074    }
2075
2076    #[tokio::test]
2077    async fn read_capped_accepts_small_body() {
2078        let base = serve_body(b"hello world".to_vec()).await;
2079        let client = reqwest::Client::builder().build().unwrap();
2080        let resp = client.get(&base).send().await.unwrap();
2081        let body = read_capped(resp).await.unwrap();
2082        assert_eq!(body, b"hello world");
2083    }
2084
2085    /// A raw HTTP server on loopback that answers `/final` with `200 arrived`
2086    /// and **everything else** with `302 Location: /final`. Returns its bound
2087    /// address, so a caller can pin a client to it by address rather than name.
2088    ///
2089    /// This is the fixture the two tests below need and that the module did not
2090    /// previously have. Note it returns the `SocketAddr`, not a URL: the whole
2091    /// point is to reach it under a hostname that does not resolve.
2092    pub(crate) async fn serve_redirect_to_final() -> SocketAddr {
2093        use tokio::io::{AsyncReadExt, AsyncWriteExt};
2094        let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
2095        let addr = listener.local_addr().unwrap();
2096        tokio::spawn(async move {
2097            loop {
2098                let (mut sock, _) = match listener.accept().await {
2099                    Ok(p) => p,
2100                    Err(_) => break,
2101                };
2102                tokio::spawn(async move {
2103                    let mut buf = [0u8; 1024];
2104                    let n = sock.read(&mut buf).await.unwrap_or(0);
2105                    let req = String::from_utf8_lossy(&buf[..n]).to_string();
2106                    let resp = if req.starts_with("GET /final") {
2107                        "HTTP/1.1 200 OK\r\nContent-Length: 7\r\nConnection: close\r\n\r\narrived"
2108                    } else {
2109                        "HTTP/1.1 302 Found\r\nLocation: /final\r\nContent-Length: 0\r\n\
2110                         Connection: close\r\n\r\n"
2111                    };
2112                    let _ = sock.write_all(resp.as_bytes()).await;
2113                    let _ = sock.flush().await;
2114                });
2115            }
2116        });
2117        addr
2118    }
2119
2120    /// **The connect goes to the address the guard vetted — enforcement, not
2121    /// decision.**
2122    ///
2123    /// `resolve_and_check` vets an address and `build_pinned_client` then
2124    /// `.resolve()`s the host to exactly that address, so the TCP connect cannot
2125    /// be rebound onto an internal one in the window between the two. That is
2126    /// the DNS-rebinding defence the module doc spends 25 lines on.
2127    ///
2128    /// **Nothing observed it.** Deleting `.resolve(host, addr)` left all 659
2129    /// tests green, because every other test either passes an IP literal — where
2130    /// a second resolution is a no-op — or asserts on `is_forbidden_ip`
2131    /// directly. `is_forbidden_ip` is thoroughly tested; what carries its verdict
2132    /// to the socket was not tested at all.
2133    ///
2134    /// This pins it in the one way that cannot silently stop discriminating: the
2135    /// host **resolves nowhere**. `.invalid` is reserved by RFC 2606 and is
2136    /// guaranteed never to exist, so the only route to the stub is the pin. Drop
2137    /// `.resolve()` and the client falls back to real DNS and cannot connect —
2138    /// which is also why this test needs no network.
2139    #[tokio::test]
2140    async fn the_connect_is_pinned_to_the_vetted_address() {
2141        let addr = serve_redirect_to_final().await;
2142        let host = "pinned-target.invalid";
2143        let client = pinned_client(host, addr).expect("building a pinned client");
2144
2145        let resp = client
2146            .get(format!("http://{host}:{}/final", addr.port()))
2147            .send()
2148            .await
2149            .expect(
2150                "a pinned host must reach the vetted address without consulting DNS — \
2151                 if this failed to connect, the `.resolve()` pin is gone",
2152            );
2153        assert_eq!(resp.status(), 200);
2154        assert_eq!(resp.text().await.unwrap(), "arrived");
2155    }
2156
2157    /// **The per-hop client must not follow redirects on its own.**
2158    ///
2159    /// `guarded_get_inner` follows redirects *manually* so it can re-run the
2160    /// scheme check, the privacy check and `resolve_and_check` on every hop, and
2161    /// so it can strip credential headers when a hop leaves the original origin.
2162    /// All of that is bypassed if reqwest follows the redirect internally: the
2163    /// connect to hop 2 happens inside reqwest, against an address nothing
2164    /// vetted. For `guarded_post` it is worse still — reqwest would re-send a
2165    /// `307`'s BODY (a client assertion, an auth code) to the new origin before
2166    /// `guarded_post`'s own 3xx refusal ever ran.
2167    ///
2168    /// Flipping `Policy::none()` to `Policy::limited(10)` left all 659 tests
2169    /// green. The existing redirect test could not catch it: it builds a 302 stub
2170    /// and then discards the address with `let _ = addr`, because the guard
2171    /// forbids loopback and the stub was therefore unreachable *through* the
2172    /// guard. It asserts on a private URL passed directly in, so no redirect ever
2173    /// occurs in it.
2174    ///
2175    /// Pinning by address sidesteps that — `pinned_client` does not consult the
2176    /// guard, so the stub is reachable — and the assertion is on the status the
2177    /// caller receives: `302`, handed back for the loop to re-validate, not the
2178    /// `200` that reqwest would return after quietly following it.
2179    #[tokio::test]
2180    async fn the_pinned_client_does_not_follow_redirects_itself() {
2181        let addr = serve_redirect_to_final().await;
2182        let host = "redirector.invalid";
2183        let client = pinned_client(host, addr).expect("building a pinned client");
2184
2185        let resp = client
2186            .get(format!("http://{host}:{}/start", addr.port()))
2187            .send()
2188            .await
2189            .expect("the stub must answer the first hop");
2190
2191        assert_eq!(
2192            resp.status(),
2193            302,
2194            "the per-hop client must hand the 30x BACK to guarded_get_inner for \
2195             re-validation; a 200 here means reqwest followed it internally and the \
2196             second hop was connected to without passing resolve_and_check",
2197        );
2198        assert_eq!(
2199            resp.headers()
2200                .get(reqwest::header::LOCATION)
2201                .and_then(|v| v.to_str().ok()),
2202            Some("/final"),
2203            "the Location must reach the caller — it is what the next hop re-validates",
2204        );
2205    }
2206
2207    /// **Every branch of the v4/v6 blocklist is load-bearing.**
2208    ///
2209    /// Four branches were unreachable from the existing tests: `is_multicast()`
2210    /// on both families, `192.0.0.0/24` ("this host on this network", IETF
2211    /// protocol assignments) and `198.18.0.0/15` (benchmarking). Deleting all
2212    /// four at once left the suite green, so a quarter of the blocklist could
2213    /// have been dropped in a refactor without a single failure.
2214    ///
2215    /// These are not decorative: multicast to an internal group and the
2216    /// benchmarking range are both reachable on a real network and neither can
2217    /// host a legitimate public feed.
2218    #[test]
2219    fn every_blocklist_branch_is_load_bearing() {
2220        for ip in [
2221            "224.0.0.1",       // v4 multicast, all-systems group
2222            "239.255.255.250", // v4 multicast, SSDP — a real LAN discovery target
2223            "192.0.0.1",       // 192.0.0.0/24, IETF protocol assignments
2224            "192.0.0.171",     // same /24
2225            "198.18.0.1",      // 198.18/15 benchmarking
2226            "198.19.255.255",  // top of the benchmarking range
2227            "0.0.0.0",         // unspecified
2228            "0.1.2.3",         // rest of 0/8
2229            "255.255.255.255", // broadcast
2230            "100.64.0.1",      // CGNAT floor
2231            "100.127.255.255", // CGNAT ceiling
2232        ] {
2233            let parsed: IpAddr = ip.parse().unwrap();
2234            assert!(is_forbidden_ip(&parsed), "{ip} must be forbidden");
2235        }
2236        for ip in [
2237            "ff02::1",                // v6 multicast, all-nodes
2238            "::1",                    // v6 loopback
2239            "::",                     // v6 unspecified
2240            "fe80::1",                // v6 link-local
2241            "fc00::1",                // v6 ULA
2242            "fd00::1",                // v6 ULA
2243            "::ffff:127.0.0.1",       // v4-mapped loopback
2244            "::ffff:169.254.169.254", // v4-mapped cloud metadata
2245            "::ffff:10.0.0.1",        // v4-mapped RFC1918
2246        ] {
2247            let parsed: IpAddr = ip.parse().unwrap();
2248            assert!(is_forbidden_ip(&parsed), "{ip} must be forbidden");
2249        }
2250    }
2251
2252    /// **The blocklist must not over-block, and only boundaries can show that.**
2253    ///
2254    /// A guard that refuses everything passes every all-negative test, and the
2255    /// existing positive cases (`1.1.1.1`, `8.8.8.8`, `93.184.216.34`) sit
2256    /// nowhere near a blocked range, so none of them would notice. Widening
2257    /// `172.16/12` to all of `172/8` and `100.64/10` to all of `100/8` — which
2258    /// would silently refuse Google and AWS address space — left the suite green.
2259    ///
2260    /// Each address here is the one immediately OUTSIDE a blocked range, so an
2261    /// off-by-one in any CIDR boundary fails this test and nothing else.
2262    #[test]
2263    fn the_blocklist_does_not_over_block_adjacent_public_space() {
2264        for ip in [
2265            "9.255.255.255",   // just below 10/8
2266            "11.0.0.0",        // just above 10/8
2267            "172.15.255.255",  // just below 172.16/12
2268            "172.32.0.0",      // just above 172.16/12 (172.217.x is Google)
2269            "192.167.255.255", // just below 192.168/16
2270            "192.169.0.0",     // just above 192.168/16
2271            "169.253.255.255", // just below 169.254/16
2272            "169.255.0.0",     // just above 169.254/16
2273            "126.255.255.255", // just below 127/8
2274            "128.0.0.0",       // just above 127/8
2275            "100.63.255.255",  // just below 100.64/10 CGNAT
2276            "100.128.0.0",     // just above 100.64/10 (100.20.x is AWS)
2277            "192.0.1.0",       // just above 192.0.0.0/24
2278            "198.17.255.255",  // just below 198.18/15
2279            "198.20.0.0",      // just above 198.18/15
2280            "223.255.255.255", // just below 224/4 multicast
2281            "1.0.0.0",         // just above 0/8
2282        ] {
2283            let parsed: IpAddr = ip.parse().unwrap();
2284            assert!(
2285                !is_forbidden_ip(&parsed),
2286                "{ip} is public and adjacent to a blocked range — refusing it means a \
2287                 CIDR boundary is wrong and real feeds are unreachable",
2288            );
2289        }
2290        for ip in ["2606:4700:4700::1111", "2001:4860:4860::8888"] {
2291            let parsed: IpAddr = ip.parse().unwrap();
2292            assert!(
2293                !is_forbidden_ip(&parsed),
2294                "{ip} is public and must be allowed"
2295            );
2296        }
2297    }
2298
2299    /// **Regression (v0.2.8):** the write-side guard must refuse the same targets
2300    /// the read-side one does — loopback, cloud metadata, RFC1918, ULA — *before*
2301    /// the connect, so a rebound PDS host never receives a request body carrying
2302    /// an app password or a session bearer.
2303    ///
2304    /// (The companion rule — a `307`/`308` is refused rather than followed,
2305    /// because a redirect re-sends the BODY and reqwest only sanitises headers —
2306    /// is not exercised here for the same reason `read_capped`'s stub is fetched
2307    /// unguarded: the guard forbids loopback, so a local stub server can never be
2308    /// reached through it. It is enforced by construction in `guarded_post_json`.)
2309    #[tokio::test]
2310    async fn guarded_post_refuses_internal_targets() {
2311        let client = Client::builder().build().unwrap();
2312        for url in [
2313            "http://127.0.0.1:9/xrpc/com.atproto.server.createSession",
2314            "http://169.254.169.254/latest/meta-data/",
2315            "http://10.0.0.5/xrpc/com.atproto.repo.applyWrites",
2316            "http://[::1]/xrpc/com.atproto.repo.deleteRecord",
2317        ] {
2318            let err = guarded_post_json(&client, url, &[], b"{}".to_vec())
2319                .await
2320                .unwrap_err()
2321                .to_string();
2322            assert!(
2323                err.contains("forbidden") || err.contains("internal"),
2324                "{url}: expected an SSRF refusal, got: {err}"
2325            );
2326        }
2327    }
2328
2329    /// A non-http(s) scheme is refused on the write path too.
2330    #[tokio::test]
2331    async fn guarded_post_refuses_bad_schemes() {
2332        let client = Client::builder().build().unwrap();
2333        for url in ["file:///etc/passwd", "gopher://example.com/1"] {
2334            let err = guarded_post_json(&client, url, &[], b"{}".to_vec())
2335                .await
2336                .unwrap_err()
2337                .to_string();
2338            assert!(err.contains("scheme"), "{url}: got: {err}");
2339        }
2340    }
2341
2342    /// **A public first hop that `30x`es to a private, secret-bearing feed is
2343    /// refused BEFORE the private target is fetched** — the per-hop privacy
2344    /// re-check in [`guarded_get`].
2345    ///
2346    /// The previous version of this test spawned a redirecting server and then
2347    /// threw it away (`let _ = addr;`), asserting on the private URL passed in
2348    /// directly — so it exercised the FIRST-hop check only, and a mutation that
2349    /// skipped privacy on every later hop left the whole suite green. Now a
2350    /// real server really redirects, and the assertion that matters is on the
2351    /// private target's request log: **zero**. "Never fetched" is the half of
2352    /// the public-feeds-only guarantee this check exists for.
2353    #[tokio::test]
2354    async fn a_redirect_to_a_private_feed_is_refused_before_it_is_fetched() {
2355        let (target_addr, target_log) = spawn_http(vec![ok_200()]).await;
2356        test_host_override("private-target.test", target_addr);
2357        let (hop_addr, hop_log) = spawn_http(vec![redirect_to(&format!(
2358            "http://private-target.test:{}/feed/private/deadbeefcafe1234",
2359            target_addr.port()
2360        ))])
2361        .await;
2362        test_host_override("private-hop.test", hop_addr);
2363
2364        let err = guarded_get(
2365            &reqwest::Client::builder().build().unwrap(),
2366            &format!("http://private-hop.test:{}/feed.xml", hop_addr.port()),
2367            &[],
2368        )
2369        .await
2370        .expect_err("a redirect to a private feed was followed");
2371        let rendered = format!("{err:#}");
2372        assert!(
2373            rendered.contains("private/paid feed URL (redirect target)"),
2374            "refused for the wrong reason: {rendered}"
2375        );
2376        assert_eq!(
2377            hop_log.lock().unwrap().len(),
2378            1,
2379            "the public first hop is fetched"
2380        );
2381        assert_eq!(
2382            target_log.lock().unwrap().len(),
2383            0,
2384            "the private target was FETCHED before being refused"
2385        );
2386    }
2387
2388    /// Header names/values for the hop-header tests: one credential, one benign
2389    /// conditional-GET validator.
2390    fn hop_fixture() -> Vec<(HeaderName, HeaderValue)> {
2391        vec![
2392            (AUTHORIZATION, HeaderValue::from_static("Bearer secret")),
2393            (
2394                reqwest::header::IF_NONE_MATCH,
2395                HeaderValue::from_static("\"etag\""),
2396            ),
2397        ]
2398    }
2399
2400    #[test]
2401    fn sensitive_header_set() {
2402        assert!(is_sensitive_header(&AUTHORIZATION));
2403        assert!(is_sensitive_header(&COOKIE));
2404        assert!(is_sensitive_header(&PROXY_AUTHORIZATION));
2405        assert!(is_sensitive_header(&WWW_AUTHENTICATE));
2406        assert!(is_sensitive_header(&HeaderName::from_static("cookie2")));
2407        assert!(!is_sensitive_header(&reqwest::header::IF_NONE_MATCH));
2408        assert!(!is_sensitive_header(&reqwest::header::IF_MODIFIED_SINCE));
2409        assert!(!is_sensitive_header(&reqwest::header::ACCEPT));
2410    }
2411
2412    #[test]
2413    fn hop_headers_keeps_all_on_same_origin() {
2414        let extra = hop_fixture();
2415        let original = Url::parse("https://pds.example.com/xrpc/x").unwrap();
2416        // The very first hop (identical URL) keeps everything…
2417        assert_eq!(hop_headers(&original, &original, &extra).len(), 2);
2418        // …and so does a same-origin path change (a `302 /a → /b` on one host).
2419        let same = Url::parse("https://pds.example.com/other/path?q=1").unwrap();
2420        assert_eq!(hop_headers(&original, &same, &extra).len(), 2);
2421    }
2422
2423    #[test]
2424    fn hop_headers_strips_authorization_cross_host() {
2425        let extra = hop_fixture();
2426        let original = Url::parse("https://pds.example.com/x").unwrap();
2427        let evil = Url::parse("https://evil.example.net/y").unwrap();
2428        let kept = hop_headers(&original, &evil, &extra);
2429        assert_eq!(kept.len(), 1, "the bearer must not follow a cross-host 302");
2430        assert_eq!(kept[0].0, reqwest::header::IF_NONE_MATCH);
2431    }
2432
2433    #[test]
2434    fn hop_headers_strips_on_port_and_scheme_change() {
2435        let extra = hop_fixture();
2436        let original = Url::parse("https://a.example/x").unwrap();
2437        for downgraded in ["http://a.example/x", "https://a.example:8443/x"] {
2438            let current = Url::parse(downgraded).unwrap();
2439            let kept = hop_headers(&original, &current, &extra);
2440            assert_eq!(kept.len(), 1, "{downgraded} must drop the credential");
2441            assert_eq!(kept[0].0, reqwest::header::IF_NONE_MATCH);
2442        }
2443        // The default port spelled explicitly is still the same origin.
2444        let explicit = Url::parse("https://a.example:443/x").unwrap();
2445        assert_eq!(hop_headers(&original, &explicit, &extra).len(), 2);
2446    }
2447
2448    #[test]
2449    fn safe_link_allowlist() {
2450        assert_eq!(
2451            safe_link("https://ok.example/x").as_deref(),
2452            Some("https://ok.example/x")
2453        );
2454        assert_eq!(
2455            safe_link("  http://ok.example/  ").as_deref(),
2456            Some("http://ok.example/")
2457        );
2458        assert_eq!(safe_link("javascript:alert(document.domain)"), None);
2459        assert_eq!(safe_link("data:text/html,<script>alert(1)</script>"), None);
2460        assert_eq!(safe_link(""), None);
2461        assert_eq!(safe_link("   "), None);
2462        // A relative/naked path isn't an absolute http(s) URL → dropped.
2463        assert_eq!(safe_link("/relative/path"), None);
2464    }
2465
2466    /// The OAuth token/PAR calls are form POSTs carrying a client assertion and,
2467    /// on the token call, the authorization code. They must go through the SAME
2468    /// SSRF guard as everything else: a PDS or issuer URL that resolves to
2469    /// loopback/RFC1918 has to fail closed BEFORE the credential is sent.
2470    #[tokio::test]
2471    async fn guarded_post_form_fails_closed_on_a_forbidden_target() {
2472        let client = Client::new();
2473        for url in [
2474            "http://127.0.0.1:2583/oauth/token",
2475            "http://[::1]:2583/oauth/token",
2476            "http://169.254.169.254/latest/meta-data/",
2477            "http://10.0.0.5/oauth/token",
2478        ] {
2479            let err = guarded_post_form(&client, url, &[], &[("grant_type", "authorization_code")])
2480                .await
2481                .expect_err("must refuse {url}");
2482            let msg = err.to_string().to_lowercase();
2483            assert!(
2484                msg.contains("forbidden") || msg.contains("refus") || msg.contains("resolve"),
2485                "unexpected error for {url}: {err:#}"
2486            );
2487        }
2488    }
2489
2490    /// A non-http(s) scheme must be rejected before any DNS work.
2491    #[tokio::test]
2492    async fn guarded_post_form_rejects_non_http_schemes() {
2493        let client = Client::new();
2494        assert!(
2495            guarded_post_form(&client, "file:///etc/passwd", &[], &[("a", "b")])
2496                .await
2497                .is_err()
2498        );
2499    }
2500
2501    /// OAuth metadata and DID documents must be fetched WITHOUT following
2502    /// redirects, and still through the SSRF guard.
2503    /// Asserts on the GUARD's error, not merely `is_err()`. Connecting to
2504    /// `127.0.0.1` fails anyway (refused, or a slow timeout for an unrouted
2505    /// RFC1918 address), so an `is_err()`-only assertion passes with
2506    /// `resolve_and_check` deleted and proves nothing.
2507    #[tokio::test]
2508    async fn guarded_get_no_redirect_still_fails_closed_on_forbidden_targets() {
2509        let client = Client::new();
2510        for url in [
2511            "http://127.0.0.1/.well-known/oauth-authorization-server",
2512            "http://169.254.169.254/latest/meta-data/",
2513            "http://192.168.1.1/.well-known/did.json",
2514            "http://[::1]/.well-known/did.json",
2515        ] {
2516            let err = guarded_get_no_redirect(&client, url, &[])
2517                .await
2518                .expect_err("must refuse");
2519            let rendered = format!("{err:#}");
2520            assert!(
2521                rendered.contains("forbidden (internal) address"),
2522                "{url} failed for the wrong reason: {rendered}"
2523            );
2524        }
2525        // And the scheme check, which is a different branch entirely.
2526        let err = guarded_get_no_redirect(&client, "file:///etc/passwd", &[])
2527            .await
2528            .expect_err("must refuse");
2529        assert!(format!("{err:#}").contains("non-http(s) URL scheme"));
2530    }
2531
2532    /// **`guarded_get_no_redirect` does not follow even one hop.** Its sibling
2533    /// above proves the SSRF guard on this path; nothing proved the ZERO. Every
2534    /// case there is an internal address or a bad scheme, so `max_redirects =
2535    /// 0` — the function's reason to exist — was never exercised, and a
2536    /// mutation passing `MAX_REDIRECTS` instead left the whole suite green.
2537    ///
2538    /// That mutation is the authorization-server mix-up defence collapsing:
2539    /// OAuth metadata, `plc.directory`, `did:web` documents and the client
2540    /// metadata self-fetch would all be read from wherever a `302` pointed,
2541    /// while `issuer` is compared against the URL that was asked for.
2542    #[tokio::test]
2543    async fn guarded_get_no_redirect_refuses_to_follow_even_one_hop() {
2544        let (b_addr, b_log) = spawn_http(vec![ok_200()]).await;
2545        test_host_override("no-redirect-b.test", b_addr);
2546        let (a_addr, a_log) = spawn_http(vec![redirect_to(&format!(
2547            "http://no-redirect-b.test:{}/.well-known/oauth-authorization-server",
2548            b_addr.port()
2549        ))])
2550        .await;
2551        test_host_override("no-redirect-a.test", a_addr);
2552
2553        let err = guarded_get_no_redirect(
2554            &reqwest::Client::builder().build().unwrap(),
2555            &format!(
2556                "http://no-redirect-a.test:{}/.well-known/oauth-authorization-server",
2557                a_addr.port()
2558            ),
2559            &[],
2560        )
2561        .await
2562        .expect_err("a redirect was followed on the no-redirect path");
2563        let rendered = format!("{err:#}");
2564        assert!(
2565            rendered.contains("origin is load-bearing"),
2566            "refused for the wrong reason: {rendered}"
2567        );
2568        assert_eq!(a_log.lock().unwrap().len(), 1);
2569        assert_eq!(
2570            b_log.lock().unwrap().len(),
2571            0,
2572            "the redirect target was fetched — the hop was followed"
2573        );
2574    }
2575
2576    /// **A POST is never redirected.** The code comment on `guarded_post` says
2577    /// this branch is "not exercised here"; now it is. A `307` re-sends the
2578    /// method AND the body — an app password or an authorization code — to the
2579    /// host the response chose.
2580    #[tokio::test]
2581    async fn guarded_post_refuses_a_redirect_rather_than_resending_the_body() {
2582        let (elsewhere_addr, elsewhere_log) = spawn_http(vec![ok_200()]).await;
2583        test_host_override("post-elsewhere.test", elsewhere_addr);
2584        let (addr, log) = spawn_http(vec![format!(
2585            "HTTP/1.1 307 Temporary Redirect\r\nLocation: http://post-elsewhere.test:{}/token\r\n\
2586             Content-Length: 0\r\nConnection: close\r\n\r\n",
2587            elsewhere_addr.port()
2588        )])
2589        .await;
2590        test_host_override("post-redirect.test", addr);
2591
2592        let err = guarded_post_form(
2593            &reqwest::Client::builder().build().unwrap(),
2594            &format!("http://post-redirect.test:{}/token", addr.port()),
2595            &[],
2596            &[
2597                ("grant_type", "authorization_code"),
2598                ("code", "SECRET-CODE"),
2599            ],
2600        )
2601        .await
2602        .expect_err("a POST followed a redirect");
2603        let rendered = format!("{err:#}");
2604        assert!(
2605            rendered.contains("re-send the request body"),
2606            "refused for the wrong reason: {rendered}"
2607        );
2608        assert_eq!(log.lock().unwrap().len(), 1);
2609        assert_eq!(
2610            elsewhere_log.lock().unwrap().len(),
2611            0,
2612            "the body was re-sent to the host the response chose"
2613        );
2614    }
2615
2616    /// **The redirect budget is enforced.** `MAX_REDIRECTS` bounds every
2617    /// outbound fetch, and until now nothing drove a chain long enough to
2618    /// reach it. One host redirects to itself `MAX_REDIRECTS + 2` times; the
2619    /// guard must give up after `MAX_REDIRECTS + 1` requests, not loop on.
2620    #[tokio::test]
2621    async fn too_many_redirects_is_refused() {
2622        // Bound first so every Location can name this server's own port.
2623        let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
2624        let port = listener.local_addr().unwrap().port();
2625        let hops: Vec<String> = (0..MAX_REDIRECTS + 2)
2626            .map(|i| redirect_to(&format!("http://redirect-loop.test:{port}/hop{i}")))
2627            .collect();
2628        let (addr, log) = spawn_http_on(listener, hops).await;
2629        test_host_override("redirect-loop.test", addr);
2630
2631        let err = guarded_get(
2632            &reqwest::Client::builder().build().unwrap(),
2633            &format!("http://redirect-loop.test:{port}/hop0"),
2634            &[],
2635        )
2636        .await
2637        .expect_err("an endless redirect chain was not refused");
2638        let rendered = format!("{err:#}");
2639        assert!(
2640            rendered.contains("too many redirects"),
2641            "refused for the wrong reason: {rendered}"
2642        );
2643        assert_eq!(
2644            log.lock().unwrap().len(),
2645            MAX_REDIRECTS + 1,
2646            "the guard made a different number of requests than its budget allows"
2647        );
2648    }
2649
2650    /// The content type must follow the body it describes. Because both come
2651    /// from the same value, a JSON body can never be labelled as a form.
2652    #[test]
2653    fn the_content_type_follows_the_body_kind() {
2654        assert_eq!(
2655            PostBody::Json(b"{}".to_vec()).content_type(),
2656            "application/json"
2657        );
2658        assert_eq!(
2659            PostBody::Form(&[("a", "b")]).content_type(),
2660            "application/x-www-form-urlencoded"
2661        );
2662        // And the bytes are encoded to match.
2663        assert_eq!(
2664            PostBody::Json(b"{\"a\":1}".to_vec()).into_bytes(),
2665            b"{\"a\":1}"
2666        );
2667        assert_eq!(PostBody::Form(&[("a", "b c")]).into_bytes(), b"a=b+c");
2668    }
2669
2670    /// Form encoding must percent-encode values; a value containing `&` or `=`
2671    /// must not be able to inject an extra parameter into the body.
2672    #[test]
2673    fn form_body_percent_encodes_and_cannot_inject_parameters() {
2674        let body = PostBody::Form(&[
2675            ("grant_type", "authorization_code"),
2676            ("code", "abc&scope=evil"),
2677            ("redirect_uri", "https://x.example/oauth/callback"),
2678        ])
2679        .into_bytes();
2680        let s = String::from_utf8(body).unwrap();
2681        assert!(s.contains("grant_type=authorization_code"));
2682        assert!(
2683            s.matches("scope=").count() == 0,
2684            "a `&` in a value injected a parameter: {s}"
2685        );
2686        assert!(s.contains("%26"), "the `&` was not encoded: {s}");
2687        assert!(s.contains("%3A%2F%2F"), "the `://` was not encoded: {s}");
2688    }
2689
2690    // ── SSRF ENFORCEMENT (not just the decision) ─────────────────────────────
2691    //
2692    // Everything below drives a REAL redirect through `guarded_get_inner`
2693    // against a real HTTP server. None of it was possible before the
2694    // `test_host_override` seam: the guard correctly refuses loopback, so a
2695    // local test server was unreachable through it, and three enforcement
2696    // properties had no coverage at all. Each had a mutation that left the
2697    // whole suite green.
2698
2699    /// A tiny HTTP server that replays canned responses and records every raw
2700    /// request it received. Returns its address and the request log.
2701    async fn spawn_http(
2702        responses: Vec<String>,
2703    ) -> (SocketAddr, std::sync::Arc<std::sync::Mutex<Vec<String>>>) {
2704        let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
2705        spawn_http_on(listener, responses).await
2706    }
2707
2708    /// [`spawn_http`] on a listener the caller already bound — for a test whose
2709    /// canned responses must name the server's own port (a redirect loop).
2710    async fn spawn_http_on(
2711        listener: tokio::net::TcpListener,
2712        responses: Vec<String>,
2713    ) -> (SocketAddr, std::sync::Arc<std::sync::Mutex<Vec<String>>>) {
2714        use tokio::io::{AsyncReadExt, AsyncWriteExt};
2715        let addr = listener.local_addr().unwrap();
2716        let log = std::sync::Arc::new(std::sync::Mutex::new(Vec::new()));
2717        let sink = std::sync::Arc::clone(&log);
2718        tokio::spawn(async move {
2719            let mut i = 0usize;
2720            loop {
2721                let Ok((mut sock, _)) = listener.accept().await else {
2722                    break;
2723                };
2724                // **The whole request, not the first 8 KB of it.**
2725                //
2726                // This used to be one `read` into a fixed buffer. Anything past
2727                // it was never captured, and the assertions over this log are
2728                // NEGATIVE — `!seen.contains("authorization:")` in the
2729                // cross-origin credential test — so a short capture satisfies
2730                // them exactly as well as a stripped header does. The two are
2731                // indistinguishable, and only one of them means the guard works.
2732                //
2733                // Same shape as `spawn_tls`, deliberately: three test servers
2734                // that read alike means the next one copied from any of them
2735                // starts correct. And the same limit applies — with no
2736                // `content-length` there is nothing to wait for, so a chunked
2737                // body stops after the head.
2738                let mut raw: Vec<u8> = Vec::new();
2739                let mut chunk = [0u8; 4096];
2740                loop {
2741                    // **A read error DISCARDS the connection rather than logging
2742                    // what arrived so far.** Breaking here and pushing the
2743                    // partial would put a truncated request in the log — the
2744                    // exact thing this change exists to stop, arriving by a
2745                    // different door. `spawn_tls` returns for the same reason.
2746                    let Ok(n) = sock.read(&mut chunk).await else {
2747                        return;
2748                    };
2749                    if n == 0 {
2750                        break;
2751                    }
2752                    raw.extend_from_slice(&chunk[..n]);
2753                    let Some(split) = raw.windows(4).position(|w| w == b"\r\n\r\n") else {
2754                        continue;
2755                    };
2756                    let (head, body) = raw.split_at(split + 4);
2757                    let want = String::from_utf8_lossy(head).lines().find_map(|l| {
2758                        let (k, v) = l.split_once(':')?;
2759                        k.eq_ignore_ascii_case("content-length")
2760                            .then(|| v.trim().parse::<usize>().ok())?
2761                    });
2762                    if want.is_none_or(|want| body.len() >= want) {
2763                        break;
2764                    }
2765                }
2766                if raw.is_empty() {
2767                    continue;
2768                }
2769                sink.lock()
2770                    .unwrap()
2771                    .push(String::from_utf8_lossy(&raw).to_string());
2772                let body = responses
2773                    .get(i)
2774                    .cloned()
2775                    .unwrap_or_else(|| responses.last().cloned().unwrap_or_default());
2776                i += 1;
2777                let _ = sock.write_all(body.as_bytes()).await;
2778                let _ = sock.flush().await;
2779            }
2780        });
2781        (addr, log)
2782    }
2783
2784    /// **The SSRF guard must ignore ambient proxy configuration.**
2785    ///
2786    /// reqwest defaults `auto_sys_proxy: true`. With `HTTP_PROXY` set in the
2787    /// process environment, a request is sent to the proxy in ABSOLUTE form —
2788    /// `GET http://host/path` — and over https as `CONNECT host:443`, for the
2789    /// PROXY to resolve the hostname.
2790    ///
2791    /// Precisely what breaks: `resolve_and_check` still runs its own lookup and
2792    /// still rejects forbidden IPs, so it is not that the check is skipped. It
2793    /// is that the checked address is no longer the address connected to — the
2794    /// proxy re-resolves the name on its own network, so the vetted result is
2795    /// decorative. DNS rebinding, split-horizon DNS and anything reachable from
2796    /// the proxy but not from here all come back. And the call returns 200, so
2797    /// it fails OPEN.
2798    ///
2799    /// Measured before `.no_proxy()` existed: vetted server 0 requests, proxy
2800    /// received `GET http://pin-vs-proxy.invalid/feed HTTP/1.1`, result
2801    /// `Ok(200)`.
2802    ///
2803    /// **Why this re-execs itself.** reqwest reads the proxy environment when
2804    /// the client is BUILT, so the variable has to be present before the
2805    /// builder runs. `set_var` is a data race against the ~39 `env::var` reads
2806    /// in this binary and is the one thing this codebase refuses to do in
2807    /// tests. So the parent owns both servers, and the child inherits
2808    /// `HTTP_PROXY` from birth — no mutation of a live environment anywhere.
2809    #[tokio::test]
2810    async fn the_pinned_client_ignores_ambient_proxy_configuration() {
2811        const VETTED: &str = "FR_AMBIENT_PROXY_VETTED";
2812        const HOST: &str = "ambient-proxy-probe.invalid";
2813
2814        if let Ok(vetted) = std::env::var(VETTED) {
2815            // ── child: HTTP_PROXY is already in our environment ──
2816            let addr: SocketAddr = vetted.parse().unwrap();
2817            let client = build_pinned_client(HOST, addr).expect("client");
2818            let _ = client.get(format!("http://{HOST}/probe")).send().await;
2819            return;
2820        }
2821
2822        // ── parent: owns both servers, so it can see who was contacted ──
2823        let (vetted_addr, vetted_log) = spawn_http(vec![ok_200()]).await;
2824        let (proxy_addr, proxy_log) = spawn_http(vec![ok_200()]).await;
2825
2826        // `tokio::process`, NOT `std::process`: a blocking `output()` here
2827        // would hold this single-threaded runtime, so the servers above could
2828        // never accept the child's connection — and the test would fail with
2829        // "did not reach the vetted address" for a reason that has nothing to
2830        // do with proxies. That exact false failure happened while writing it.
2831        let out = tokio::process::Command::new(std::env::current_exe().unwrap())
2832            // FULL path: `--exact` matches the whole test name including the
2833            // module. With the bare function name the child matched nothing,
2834            // ran zero tests, and exited 0 — so the parent saw a "successful"
2835            // child that had done nothing, and blamed the pin. The
2836            // `1 test` assertion below is there so that can never pass silently
2837            // again.
2838            .args([
2839                "net::tests::the_pinned_client_ignores_ambient_proxy_configuration",
2840                "--exact",
2841                "--test-threads=1",
2842            ])
2843            // **Clear the inherited proxy KILL-SWITCHES.**
2844            //
2845            // The child inherits this process's environment, and two inherited
2846            // values make the whole test vacuous — it passes with `.no_proxy()`
2847            // DELETED. Verified: `NO_PROXY='*'` and `REQUEST_METHOD=GET`
2848            // (hyper-util treats the latter as a CGI context and disables proxy
2849            // env entirely) each turn a genuine failure into `1 passed`.
2850            //
2851            // GitHub-hosted runners set none of these, so the gap was invisible
2852            // here — a self-hosted or corporate runner would have silently
2853            // neutered the regression test while it kept reporting success.
2854            .env_remove("NO_PROXY")
2855            .env_remove("no_proxy")
2856            .env_remove("REQUEST_METHOD")
2857            .env(VETTED, vetted_addr.to_string())
2858            .env("HTTP_PROXY", format!("http://{proxy_addr}"))
2859            .env("HTTPS_PROXY", format!("http://{proxy_addr}"))
2860            .env("ALL_PROXY", format!("http://{proxy_addr}"))
2861            .output()
2862            .await
2863            .expect("re-exec the test binary");
2864        let stdout = String::from_utf8_lossy(&out.stdout);
2865        assert!(
2866            out.status.success(),
2867            "child run failed: {}",
2868            String::from_utf8_lossy(&out.stderr)
2869        );
2870        assert!(
2871            stdout.contains("1 passed"),
2872            "the child ran no test, so this proves nothing about proxies — \
2873             check the --exact filter. Child stdout:\n{stdout}"
2874        );
2875
2876        let proxied = proxy_log.lock().unwrap().clone();
2877        let direct = vetted_log.lock().unwrap().len();
2878        assert!(
2879            proxied.is_empty(),
2880            "the pinned client used an ambient proxy, so the connect pin and \
2881             `is_forbidden_ip` were both bypassed — the proxy resolves the \
2882             hostname itself. Proxy saw: {proxied:?}"
2883        );
2884        assert_eq!(
2885            direct, 1,
2886            "the pinned client did not reach the vetted address it was pinned to",
2887        );
2888    }
2889
2890    fn ok_200() -> String {
2891        "HTTP/1.1 200 OK\r\nContent-Length: 2\r\nConnection: close\r\n\r\nhi".to_string()
2892    }
2893    fn redirect_to(loc: &str) -> String {
2894        format!("HTTP/1.1 302 Found\r\nLocation: {loc}\r\nContent-Length: 0\r\nConnection: close\r\n\r\n")
2895    }
2896    /// No `Location` and no body — which is what a 304 *is*, not a stub of one.
2897    fn not_modified_304() -> String {
2898        "HTTP/1.1 304 Not Modified\r\nETag: \"v1\"\r\nConnection: close\r\n\r\n".to_string()
2899    }
2900    /// `305 Use Proxy` — a relocating 3xx that names a PROXY, WITH a `Location`,
2901    /// so it would have been followed before the narrowing.
2902    fn use_proxy_305(loc: &str) -> String {
2903        format!("HTTP/1.1 305 Use Proxy\r\nLocation: {loc}\r\nContent-Length: 0\r\nConnection: close\r\n\r\n")
2904    }
2905
2906    /// **An unfollowable 3xx is refused, not handed back.**
2907    ///
2908    /// Narrowing the follow set to `301|302|303|307|308` left a choice for the
2909    /// rest: return them, or refuse. Returning is the worse one, because callers
2910    /// do not uniformly check the status — `web::resolve_feed_url` reads the body
2911    /// straight into feed autodiscovery, so a `305` whose error page carries a
2912    /// `<link rel="alternate">` would quietly become a subscription.
2913    ///
2914    /// A 305 also carries a `Location`, so before the narrowing it was FOLLOWED:
2915    /// the request went through a proxy the response chose. Refusing is the
2916    /// stricter behaviour in both directions.
2917    #[tokio::test]
2918    async fn an_unfollowable_3xx_is_refused_rather_than_returned() {
2919        let (addr, log) = spawn_http(vec![use_proxy_305("http://proxy.invalid:3128/")]).await;
2920        test_host_override("use-proxy.test", addr);
2921
2922        let err = guarded_get(
2923            &reqwest::Client::builder().build().unwrap(),
2924            &format!("http://use-proxy.test:{}/feed.xml", addr.port()),
2925            &[],
2926        )
2927        .await
2928        .expect_err("a 305 was returned to the caller instead of refused");
2929
2930        let msg = format!("{err:#}");
2931        assert!(
2932            msg.contains("305") && msg.contains("no single target"),
2933            "refused, but not as an unfollowable status: {msg}",
2934        );
2935        // Refused at the first hop: the proxy it named was never contacted.
2936        assert_eq!(log.lock().unwrap().len(), 1);
2937    }
2938
2939    /// **A `304 Not Modified` must reach the caller, not be read as a redirect.**
2940    ///
2941    /// `is_redirection()` is `300..=399`, so 304 — which carries no `Location`
2942    /// by definition — fell into the redirect branch and failed the whole fetch
2943    /// with "redirect response without a usable Location header". `feed.rs`
2944    /// sends `If-None-Match`/`If-Modified-Since` on every poll and has a correct
2945    /// 304 branch (`feed.rs`, `status == StatusCode::NOT_MODIFIED`) that could
2946    /// never be reached, so every feed answering "unchanged" was recorded as a
2947    /// failure and backed off exponentially.
2948    ///
2949    /// This was not theoretical: production logged it against `9to5mac.com`,
2950    /// `proton.me` and `kodi.tv`, and `/stats` reported 68 of 111 feeds failing
2951    /// while its own copy explained them away as "usually gone rather than
2952    /// flaky". Re-running the poller's conditional GET by hand returned
2953    /// `HTTP 304` with zero `Location` headers.
2954    ///
2955    /// The hop count is asserted too: a 304 must not provoke a second request.
2956    /// Returning the response but still looping would satisfy a status-only
2957    /// assertion while re-fetching every unchanged feed.
2958    #[tokio::test]
2959    async fn a_304_reaches_the_caller_instead_of_being_read_as_a_redirect() {
2960        let (addr, log) = spawn_http(vec![not_modified_304()]).await;
2961        test_host_override("not-modified.test", addr);
2962
2963        let resp = guarded_get(
2964            &reqwest::Client::builder().build().unwrap(),
2965            &format!("http://not-modified.test:{}/feed.xml", addr.port()),
2966            &[],
2967        )
2968        .await
2969        .expect("a 304 was treated as a redirect");
2970
2971        assert_eq!(
2972            resp.status(),
2973            reqwest::StatusCode::NOT_MODIFIED,
2974            "the 304 did not survive the guard intact",
2975        );
2976        assert_eq!(
2977            log.lock().unwrap().len(),
2978            1,
2979            "a 304 caused more than one request — it was followed, not returned",
2980        );
2981    }
2982
2983    /// **A real redirect is still followed** — the other half of the narrowing
2984    /// above, which would otherwise be satisfied by never following anything.
2985    #[tokio::test]
2986    async fn a_302_is_still_followed_after_the_304_narrowing() {
2987        let (b_addr, _b_log) = spawn_http(vec![ok_200()]).await;
2988        test_host_override("still-follows-b.test", b_addr);
2989        let (a_addr, a_log) = spawn_http(vec![redirect_to(&format!(
2990            "http://still-follows-b.test:{}/final",
2991            b_addr.port()
2992        ))])
2993        .await;
2994        test_host_override("still-follows-a.test", a_addr);
2995
2996        let resp = guarded_get(
2997            &reqwest::Client::builder().build().unwrap(),
2998            &format!("http://still-follows-a.test:{}/feed.xml", a_addr.port()),
2999            &[],
3000        )
3001        .await
3002        .expect("the 302 was not followed");
3003
3004        assert_eq!(resp.status(), reqwest::StatusCode::OK);
3005        assert_eq!(a_log.lock().unwrap().len(), 1);
3006    }
3007
3008    /// **A redirect to a forbidden address is refused — the marquee SSRF
3009    /// property, and until now it had no end-to-end test.**
3010    ///
3011    /// `guarded_get_refuses_private_redirect_target` builds a 302 stub and then
3012    /// throws it away (`let _ = addr;`), asserting on a private URL passed in
3013    /// directly. Nothing drove a redirect through the guard, and a mutation that
3014    /// validated only the first hop — resolving redirect targets with a bare
3015    /// `lookup_host` and no `is_forbidden_ip` — left all 679 tests passing.
3016    ///
3017    /// Here a real server really 302s to the cloud metadata endpoint. Only the
3018    /// test server's own host is exempted from the IP check; the redirect target
3019    /// is an IP literal and goes through the real one.
3020    #[tokio::test]
3021    async fn a_redirect_to_a_forbidden_address_is_refused() {
3022        let (addr, log) = spawn_http(vec![redirect_to(
3023            "http://169.254.169.254/latest/meta-data/",
3024        )])
3025        .await;
3026        test_host_override("hop-forbidden.test", addr);
3027
3028        let err = guarded_get(
3029            &reqwest::Client::builder().build().unwrap(),
3030            &format!("http://hop-forbidden.test:{}/feed.xml", addr.port()),
3031            &[],
3032        )
3033        .await
3034        .expect_err("a 302 to the metadata endpoint was followed");
3035
3036        let msg = format!("{err:#}");
3037        assert!(
3038            msg.contains("169.254.169.254") && msg.contains("forbidden"),
3039            "refused, but not by the address check: {msg}",
3040        );
3041        // The hop happened; the SECOND hop is what was stopped.
3042        assert_eq!(log.lock().unwrap().len(), 1);
3043    }
3044
3045    /// **Credentials do not follow a redirect off the original origin.**
3046    ///
3047    /// `hop_headers` is tested as a pure function; nothing tested that
3048    /// `guarded_get_inner` actually calls it. Swapping the call for a plain
3049    /// `extra_headers` — so a bearer token rides to whatever host an upstream
3050    /// names — left the whole suite green.
3051    ///
3052    /// Two real servers on two hosts. The first 302s to the second; the second
3053    /// records what it was sent.
3054    #[tokio::test]
3055    async fn credentials_are_dropped_when_a_redirect_leaves_the_origin() {
3056        let (b_addr, b_log) = spawn_http(vec![ok_200()]).await;
3057        test_host_override("cred-b.test", b_addr);
3058        let (a_addr, _a_log) = spawn_http(vec![redirect_to(&format!(
3059            "http://cred-b.test:{}/next",
3060            b_addr.port()
3061        ))])
3062        .await;
3063        test_host_override("cred-a.test", a_addr);
3064
3065        let resp = guarded_get(
3066            &reqwest::Client::builder().build().unwrap(),
3067            &format!("http://cred-a.test:{}/feed.xml", a_addr.port()),
3068            &[(
3069                HeaderName::from_static("authorization"),
3070                HeaderValue::from_static("Bearer super-secret"),
3071            )],
3072        )
3073        .await
3074        .expect("the cross-origin hop should still succeed, just without the token");
3075        assert!(resp.status().is_success());
3076
3077        let seen = b_log.lock().unwrap().join("\n").to_ascii_lowercase();
3078        // **Anchor the two negatives below.** `!contains` is satisfied by the
3079        // header being absent OR by the capture being short, and those are
3080        // indistinguishable from here. Asserting the hop was recorded at all
3081        // means an empty or truncated capture fails loudly instead of reading
3082        // as a pass — which, for a check about not leaking a bearer token
3083        // across origins, is the difference that matters.
3084        assert!(
3085            seen.contains("get /next"),
3086            "hop B recorded no request, so the assertions below prove nothing:\n{seen}",
3087        );
3088        assert!(
3089            !seen.contains("super-secret"),
3090            "the bearer token was forwarded across origins:\n{seen}",
3091        );
3092        assert!(
3093            !seen.contains("authorization:"),
3094            "the Authorization header survived a cross-origin redirect:\n{seen}",
3095        );
3096    }
3097
3098    /// **...but they DO survive a same-origin redirect.**
3099    ///
3100    /// The other direction, without which the test above is satisfied by a guard
3101    /// that strips every header always — which would quietly break every
3102    /// authenticated fetch in the app.
3103    ///
3104    /// A RELATIVE `Location` keeps the hop on the same origin without the
3105    /// response needing to know its own port.
3106    #[tokio::test]
3107    async fn credentials_survive_a_same_origin_redirect() {
3108        let (addr, log) = spawn_http(vec![redirect_to("/second"), ok_200()]).await;
3109        test_host_override("cred-same.test", addr);
3110
3111        let resp = guarded_get(
3112            &reqwest::Client::builder().build().unwrap(),
3113            &format!("http://cred-same.test:{}/feed.xml", addr.port()),
3114            &[(
3115                HeaderName::from_static("authorization"),
3116                HeaderValue::from_static("Bearer keep-me"),
3117            )],
3118        )
3119        .await
3120        .expect("a same-origin redirect should be followed");
3121        assert!(resp.status().is_success());
3122
3123        let reqs = log.lock().unwrap().clone();
3124        assert_eq!(reqs.len(), 2, "the redirect was not followed");
3125        assert!(
3126            reqs[1].to_ascii_lowercase().contains("keep-me"),
3127            "the token was stripped on a SAME-origin redirect — over-stripping \
3128             would break every authenticated fetch:\n{}",
3129            reqs[1],
3130        );
3131    }
3132
3133    /// **Every DNS answer is checked, not just the first.**
3134    ///
3135    /// A host that publishes `1.2.3.4, 127.0.0.1` must be rejected wholesale.
3136    /// Checking only the first answer left the suite green, because nothing
3137    /// exercised a multi-answer set — real DNS in a test cannot be made to
3138    /// return one.
3139    #[test]
3140    fn a_mixed_dns_answer_set_is_rejected_wholesale() {
3141        let public: SocketAddr = "1.2.3.4:80".parse().unwrap();
3142        let private: SocketAddr = "127.0.0.1:80".parse().unwrap();
3143        let link_local: SocketAddr = "169.254.169.254:80".parse().unwrap();
3144
3145        // All public: the first is pinned.
3146        assert_eq!(
3147            first_vetted(
3148                "ok.example",
3149                [public, "5.6.7.8:80".parse().unwrap()].into_iter()
3150            )
3151            .unwrap(),
3152            public,
3153        );
3154        // A forbidden answer ANYWHERE rejects the set — including last, which is
3155        // exactly what a first-answer-only check would miss.
3156        for bad in [private, link_local] {
3157            assert!(
3158                first_vetted("evil.example", [public, bad].into_iter()).is_err(),
3159                "{bad} in the answer set was accepted because a good answer came first",
3160            );
3161            assert!(first_vetted("evil.example", [bad, public].into_iter()).is_err());
3162        }
3163        // No answers at all is an error, not a silent pass.
3164        assert!(first_vetted("empty.example", std::iter::empty()).is_err());
3165    }
3166
3167    /// **The capture must hold the whole request, not the first 8 KB of it.**
3168    ///
3169    /// `spawn_http` recorded one `sock.read()` into a fixed 8 KB buffer and
3170    /// treated that as the request. Anything past it was never captured — and
3171    /// never seen by the assertions that read the capture.
3172    ///
3173    /// That matters because the assertions downstream are NEGATIVE:
3174    /// `credentials_are_dropped_when_a_redirect_leaves_the_origin` checks
3175    /// `!seen.contains("authorization:")`. A short capture satisfies that exactly
3176    /// as well as a stripped header does, and the two are indistinguishable.
3177    ///
3178    /// A body larger than the buffer makes the truncation deterministic rather
3179    /// than waiting on TCP segmentation, which is why this test can go red at
3180    /// all.
3181    #[tokio::test]
3182    async fn the_request_capture_is_not_truncated_at_the_buffer_size() {
3183        let (addr, log) = spawn_http(vec![ok_200()]).await;
3184        test_host_override("big-body.test", addr);
3185
3186        // Comfortably past the old 8 KB read, with a sentinel at the very end.
3187        let filler = "x".repeat(32 * 1024);
3188        let body = format!("{{\"pad\":\"{filler}\",\"tail\":\"THE-LAST-BYTES\"}}");
3189
3190        // Not `let _ =`: a refused POST leaves the capture empty, and
3191        // "captured no request at all" would be the only symptom with the
3192        // cause thrown away.
3193        guarded_post_json(
3194            &reqwest::Client::builder().build().unwrap(),
3195            &format!("http://big-body.test:{}/ingest", addr.port()),
3196            &[],
3197            body.into_bytes(),
3198        )
3199        .await
3200        .expect("the POST to the test server failed before anything was captured");
3201
3202        let seen = log.lock().unwrap().join("\n");
3203        // Positive anchor first: without it, the tail assertion below could pass
3204        // vacuously on an empty capture in some future refactor.
3205        assert!(
3206            seen.contains("POST /ingest"),
3207            "the server captured no request at all: {} bytes",
3208            seen.len()
3209        );
3210        assert!(
3211            seen.contains("THE-LAST-BYTES"),
3212            "the capture stops short of the request's end, so every negative \
3213             assertion over it — including the one about not leaking an \
3214             Authorization header across origins — can pass for the wrong \
3215             reason. captured {} bytes",
3216            seen.len()
3217        );
3218    }
3219
3220    // ── TLS test server ──────────────────────────────────────────────────────
3221
3222    /// How many times [`guarded_get_for_a_verdict`] asks, while the answer keeps
3223    /// being a timeout rather than a verdict about the certificate.
3224    ///
3225    /// **Nothing pins this number, and that is disclosed rather than implied.**
3226    /// Any value above one behaves identically on a healthy run, so pinning it
3227    /// would need a server that stalls past the 15 s per-read bound — a 15 s
3228    /// test. What IS pinned, in both directions, is the classifier the loop turns
3229    /// on: `a_timeout_is_recognised_as_a_timeout_and_not_a_certificate_verdict`
3230    /// for one, the assertion inside
3231    /// `the_test_ca_is_trusted_and_still_validates_hostnames` for the other.
3232    const VERDICT_ATTEMPTS: usize = 3;
3233
3234    /// Whether an error out of [`guarded_get`] is a TIMEOUT rather than a verdict
3235    /// about the certificate.
3236    ///
3237    /// Classified from `reqwest::Error::is_timeout` through the `with_context`
3238    /// layer `guarded_get_inner` adds, not from the message text — the message
3239    /// is not a contract, and matching on it is how this file's other
3240    /// error-shape assertion passed for the wrong reason once already.
3241    fn is_timeout(err: &anyhow::Error) -> bool {
3242        err.downcast_ref::<reqwest::Error>()
3243            .is_some_and(reqwest::Error::is_timeout)
3244    }
3245
3246    /// [`guarded_get`], asked again while the only answer is a timeout.
3247    ///
3248    /// **The certificate test is about whether the chain validates, and a
3249    /// timeout is not a verdict on that.** It was observed failing on the first
3250    /// HTTPS request in a freshly linked test binary — 11.7 s and 20.3 s
3251    /// measured on one macOS machine, against the 15 s per-read bound
3252    /// `build_pinned_client` sets. Nine later attempts on the same machine
3253    /// measured 8–17 ms, so the cause is NOT pinned; the leading candidate is
3254    /// CPU starvation with ~900 tests in flight, which no amount of warming
3255    /// would fix.
3256    ///
3257    /// Two earlier attempts at this are worth naming, because both were wrong in
3258    /// ways this one avoids. Widening `READ_TIMEOUT` changed a production
3259    /// constant to accommodate a test. Warming the platform verifier once per
3260    /// process rested on a claim that is simply false — reqwest builds
3261    /// `rustls_platform_verifier` whether or not an extra root is present
3262    /// (`reqwest-0.13/src/async_impl/client.rs`: both arms of
3263    /// `if config.root_certs.is_empty()`), so there was no test-only path to
3264    /// warm; it also failed silently, and issued a request into the caller's
3265    /// captured log.
3266    ///
3267    /// Retrying the timeout is insensitive to *which* cause it was, changes no
3268    /// production bound, and cannot mask a validation failure: a certificate
3269    /// verdict is returned on the first ask.
3270    async fn guarded_get_for_a_verdict(client: &Client, url: &str) -> Result<Response> {
3271        for _ in 1..VERDICT_ATTEMPTS {
3272            match guarded_get(client, url, &[]).await {
3273                Err(err) if is_timeout(&err) => {
3274                    eprintln!("asking {url} again after a timeout, not a verdict: {err:#}");
3275                }
3276                verdict => return verdict,
3277            }
3278        }
3279        guarded_get(client, url, &[]).await
3280    }
3281
3282    /// **The retry turns entirely on this classifier, so pin it.**
3283    ///
3284    /// A classifier that stops recognising timeouts leaves the certificate test
3285    /// exactly as latency-sensitive as it was, with nothing to say so. The
3286    /// opposite direction — a certificate error must NOT read as a timeout, or a
3287    /// genuine validation failure would be retried and then reported as one — is
3288    /// asserted where such an error already exists, in
3289    /// `the_test_ca_is_trusted_and_still_validates_hostnames`.
3290    ///
3291    /// Uses a real timeout against a socket that is accepted and never answered,
3292    /// through the same `with_context` wrapping `guarded_get_inner` applies, so
3293    /// the downcast is exercised through a context layer rather than on a bare
3294    /// error.
3295    #[tokio::test]
3296    async fn a_timeout_is_recognised_as_a_timeout_and_not_a_certificate_verdict() {
3297        let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
3298        let addr = listener.local_addr().unwrap();
3299        tokio::spawn(async move {
3300            // Accept and hold: answering nothing is the point, and dropping the
3301            // socket would end the request as a connection close instead.
3302            let mut held = Vec::new();
3303            while let Ok((sock, _)) = listener.accept().await {
3304                held.push(sock);
3305            }
3306        });
3307
3308        let client = Client::builder()
3309            .no_proxy()
3310            .timeout(Duration::from_millis(250))
3311            .build()
3312            .unwrap();
3313        let raw = client
3314            .get(format!("http://{addr}/never"))
3315            .send()
3316            .await
3317            .expect_err("a server that never answers must not produce a response");
3318        assert!(
3319            raw.is_timeout(),
3320            "the silent server ended the request some other way, so this test is \
3321             not exercising a timeout at all: {raw}",
3322        );
3323        let err = anyhow::Error::from(raw).context(format!("fetching http://{addr}/never"));
3324
3325        assert!(
3326            is_timeout(&err),
3327            "a real read timeout was not recognised as one, so the retry would \
3328             never retry and the certificate test stays latency-sensitive: {err:#}",
3329        );
3330    }
3331
3332    /// **The chain really validates — no invalid-cert acceptance anywhere.**
3333    ///
3334    /// The foundation every test below rests on. If this passed because
3335    /// validation were disabled rather than because the CA is trusted, none of
3336    /// the others would mean anything, so it is asserted directly: a host the
3337    /// leaf has NO SAN for must still fail.
3338    #[tokio::test]
3339    async fn the_test_ca_is_trusted_and_still_validates_hostnames() {
3340        let (addr, _log) = spawn_tls(|_| {
3341            let mut r = std::collections::HashMap::new();
3342            r.insert("/ok".to_string(), vec![TestResponse::json(200, "{}")]);
3343            r
3344        })
3345        .await;
3346        test_host_override("feed-tls.test", addr);
3347        // Registered, resolvable — but NOT in the leaf's SAN list.
3348        test_host_override("not-in-san.test", addr);
3349
3350        let client = reqwest::Client::builder().build().unwrap();
3351        // Asked again on a timeout — see `guarded_get_for_a_verdict`. A timeout
3352        // is not a verdict about this chain, and this test is only about the
3353        // verdict.
3354        let ok = guarded_get_for_a_verdict(
3355            &client,
3356            &format!("https://feed-tls.test:{}/ok", addr.port()),
3357        )
3358        .await
3359        .expect("a SAN-matching https host should be accepted");
3360        assert!(ok.status().is_success());
3361
3362        let err = guarded_get_for_a_verdict(
3363            &client,
3364            &format!("https://not-in-san.test:{}/ok", addr.port()),
3365        )
3366        .await
3367        .expect_err("a host with no SAN must still fail: validation is NOT disabled");
3368        // **The other half of the classifier, pinned where such an error exists.**
3369        // If a certificate verdict read as a timeout, this failure would be
3370        // retried `VERDICT_ATTEMPTS` times and then reported anyway — and the
3371        // retry would be masking exactly the failure it must never mask.
3372        assert!(
3373            !is_timeout(&err),
3374            "a certificate verdict was classified as a timeout, so the retry              would retry a genuine validation failure: {err:#}",
3375        );
3376        // **Assert the CERTIFICATE reason, not merely that it failed.**
3377        //
3378        // An earlier version accepted `msg.contains("name")`, which the DNS error
3379        // `nodename nor servname provided` also satisfies — so deleting the
3380        // `test_host_override` line above made this pass while proving nothing
3381        // about SAN validation. Verified: it did.
3382        let msg = format!("{err:#}").to_ascii_lowercase();
3383        assert!(
3384            msg.contains("notvalidforname") || msg.contains("invalid peer certificate"),
3385            "failed, but not because the certificate is invalid for this name — a \
3386             DNS or connect failure would prove nothing here: {msg}",
3387        );
3388    }
3389}