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