Skip to main content

feather_reader/
network.rs

1//! Read-only queries against the **public atproto network**.
2//!
3//! Where [`crate::atproto`] reads and writes *one user's own* PDS repo, this
4//! module asks the public relays a single question: **how many repos on the
5//! network hold a given collection?** Today that collection is
6//! `community.lexicon.rss.subscription`, and the answer is FeatherReader's
7//! adoption metric — the one number that says whether the portability claim on
8//! every page is being exercised by anybody but us.
9//!
10//! Three properties are load-bearing, and each is a rule rather than an
11//! intention (`design/NETWORK-SPEC.md` §4):
12//!
13//! 1. **It is never a source of truth.** The result is a projection: drop the
14//!    `network_stat` table and the next probe rebuilds it. Nothing in the reader
15//!    path reads it, and no reader surface may depend on it.
16//! 2. **It counts, it does not collect.** The relay answers with a list of DIDs;
17//!    we read `repos.len()` and drop the page. Persisting the DID list would
18//!    build a durable register of "accounts that use an RSS reader" on our disk,
19//!    for a feature whose entire output is an integer.
20//! 3. **The number is a LOWER BOUND, not a census.** Since sync v1.1 relays are
21//!    **non-archival**: a relay's index only covers hosts it actually crawls, so
22//!    a PDS no relay crawls is invisible to it. Two relays can therefore
23//!    disagree; we query every configured host, record each observation
24//!    separately, surface the max, and log the disagreement. Any copy derived
25//!    from this number must say "at least", never "exactly".
26//!
27//! **Why the SSRF guard, when the relay host is operator-configured?** Not
28//! because the operator is the threat — they can already edit the code. It earns
29//! its place for three other reasons. The shared `reqwest::Client` has *no*
30//! timeouts, so a hung relay would pin a background task forever, while
31//! [`crate::net`]'s per-hop pinned client bounds both the total request and the
32//! idle read. The shared client also follows up to ten redirects with no
33//! re-validation, whereas [`crate::net::guarded_get_no_privacy`] follows at most
34//! five and re-checks scheme + resolved IP on each — a relay that `302`s (DNS
35//! takeover, a misconfigured proxy, a captive portal on a self-hoster's network)
36//! is the realistic rebinding path. And it is the coherent choice: the same
37//! milestone that routes `PdsClient::list_records` through the guard should not
38//! open a fresh unguarded `http.get()` next door. It is free besides —
39//! [`crate::USER_AGENT`] and the body cap come along with it.
40
41use std::time::Duration;
42
43use reqwest::header::{HeaderMap, HeaderValue, ACCEPT, RETRY_AFTER};
44use reqwest::StatusCode;
45use serde::Deserialize;
46
47use crate::atproto::urlencode;
48
49/// The relay XRPC method that answers "which repos hold this collection?".
50pub const LIST_REPOS_BY_COLLECTION: &str = "com.atproto.sync.listReposByCollection";
51
52/// The public Bluesky relays queried by default, in order. Lives here (mirroring
53/// [`crate::atproto::DEFAULT_PLC_DIRECTORY`]) so [`crate::config`] imports the
54/// network fact rather than re-typing the hostnames.
55pub const DEFAULT_RELAY_HOSTS: [&str; 2] =
56    ["relay1.us-west.bsky.network", "relay1.us-east.bsky.network"];
57
58/// Repos requested per page. The relay's documented ceiling is 1000; 500 keeps a
59/// page body small (~30 KB) while making the whole current network one page.
60pub const DEFAULT_PAGE_LIMIT: u32 = 500;
61
62/// Hard cap on pages walked per host: 25 000 repos at [`DEFAULT_PAGE_LIMIT`],
63/// i.e. ~25 000× today's network. Hitting it means something went wrong (a
64/// cursor loop, a wrong collection) — the run is recorded as `truncated`, so the
65/// number reads as "at least N".
66pub const MAX_PAGES: usize = 50;
67
68/// Floor for the page limit. This is not defensive decoration: the live relay
69/// answers `limit=0` with **one** repo *and* a cursor, so a zero would burn the
70/// whole page budget one repo at a time.
71const MIN_PAGE_LIMIT: u32 = 1;
72
73/// Ceiling for the page limit (the relay's own maximum).
74const MAX_PAGE_LIMIT: u32 = 1000;
75
76/// Politeness delay between pages against one host.
77const DEFAULT_PAGE_DELAY: Duration = Duration::from_secs(1);
78
79/// Wall-clock budget for one host's whole walk. [`crate::net`] bounds each *hop*
80/// (30 s total, 15 s idle) but nothing bounds the walk: 50 pages × up to 5
81/// redirect hops × 30 s plus 49 s of inter-page sleep is a worst case near two
82/// hours. This is what keeps a slow relay from leaving a detached task alive
83/// across the next daily tick.
84///
85/// This is a **soft** budget, checked between pages by [`RelayClient::walk_pages`],
86/// which stops and returns what it has with `truncated = true`. It is
87/// deliberately NOT sized to let the full [`MAX_PAGES`] budget run: 50 pages ×
88/// ([`DEFAULT_PAGE_DELAY`] + [`crate::net::FETCH_TIMEOUT`]) is ~26 minutes, far
89/// too long to hold a background task for an optional metric. The two limits
90/// bound different things — `MAX_PAGES` bounds a *pathological* relay,
91/// this bounds a merely *slow* one — and whichever binds first, the result is a
92/// recorded lower bound rather than a lost run.
93const DEFAULT_HOST_BUDGET: Duration = Duration::from_secs(120);
94
95/// Slack added to [`DEFAULT_HOST_BUDGET`] + [`crate::net::WORST_CASE_REQUEST`]
96/// to form the **hard** backstop deadline. See [`RelayClient::hard_deadline`].
97///
98/// The backstop should never actually fire: `walk_pages` bounds each page by
99/// what is left of the budget, so the walk cannot outlive it. The margin exists
100/// only so that a future which somehow fails to observe that timeout is still
101/// bounded, and it is deliberately measured against a WORST-CASE request
102/// (redirect chain included) rather than a single `FETCH_TIMEOUT` — budgeting
103/// one hop was the bug in the first cut of this fix.
104const HARD_DEADLINE_SLACK: Duration = Duration::from_secs(10);
105
106/// How much of a non-2xx body is quoted back in the error (the atproto envelope
107/// `{"error","message"}` is short; a hostile body must not fill the log).
108const ERROR_SNIPPET_CHARS: usize = 200;
109
110// ---------------------------------------------------------------------------
111// Wire types
112// ---------------------------------------------------------------------------
113
114/// The `com.atproto.sync.listReposByCollection` response envelope.
115///
116/// Deliberately **not** `deny_unknown_fields`: the same forward-compatibility
117/// discipline [`crate::lexicon::FetchHint::Other`] applies, so a field added to
118/// the relay's response never breaks the probe.
119#[derive(Debug, Clone, Deserialize)]
120struct ListReposByCollectionOut {
121    /// One entry per repo the relay has indexed as holding the collection.
122    /// `#[serde(default)]` even though every observed response includes it.
123    #[serde(default)]
124    repos: Vec<RepoRef>,
125    /// The pagination cursor. MUST be optional: the final page omits the key
126    /// entirely rather than sending `null`.
127    #[serde(default)]
128    cursor: Option<String>,
129}
130
131/// One repo in a [`ListReposByCollectionOut`] page.
132///
133/// The DID is parsed but **never returned, accumulated, or persisted** — only
134/// `repos.len()` is read, and the page is dropped at the end of the iteration
135/// that produced it (see the module doc, rule 2).
136#[derive(Debug, Clone, Deserialize)]
137struct RepoRef {
138    #[allow(dead_code)]
139    #[serde(default)]
140    did: String,
141}
142
143// ---------------------------------------------------------------------------
144// Results
145// ---------------------------------------------------------------------------
146
147/// One relay's answer to "how many repos hold this collection?".
148#[derive(Debug, Clone, PartialEq, Eq)]
149pub struct AdoptionObservation {
150    /// The normalized relay base URL this number came from (the
151    /// `network_stat.source` column).
152    pub source: String,
153    /// The collection NSID that was counted.
154    pub collection: String,
155    /// Repos the relay has **indexed** as holding the collection — a lower
156    /// bound on network adoption, not a census (module doc, rule 3). Whether the
157    /// relay's index counts a repo that has since deleted its records is
158    /// unverified, which is why this says "indexed as holding" and not "holds".
159    pub repos: u64,
160    /// The walk stopped early, so `repos` is a floor of a lower bound: read it
161    /// as "at least N".
162    ///
163    /// Set by EITHER bound — the [`MAX_PAGES`] cap (a pathological relay) or the
164    /// `DEFAULT_HOST_BUDGET` wall-clock budget (a merely slow one). They mean
165    /// the same thing to a consumer, which is why they share one flag.
166    pub truncated: bool,
167    /// When the observation was taken (RFC3339, UTC, seconds precision —
168    /// the same shape every other TEXT timestamp in the DB uses).
169    pub observed_at: String,
170}
171
172/// One relay that did not answer, and why.
173#[derive(Debug, Clone, PartialEq, Eq)]
174pub struct RelayFailure {
175    /// The normalized relay base URL.
176    pub host: String,
177    /// The rendered [`RelayError`] — a log line, not a machine-readable code.
178    pub reason: String,
179}
180
181/// The outcome of one probe run across **every** configured relay.
182///
183/// Partial success is a real outcome (one relay up, one down), which a single
184/// `Result<AdoptionObservation>` cannot express — so this is not a `Result`, and
185/// the caller needs no error handling beyond [`succeeded`](Self::succeeded).
186#[derive(Debug, Clone, PartialEq, Eq)]
187pub struct AdoptionReport {
188    /// The collection that was counted.
189    pub collection: String,
190    /// One observation per host that answered, in configured order.
191    pub observations: Vec<AdoptionObservation>,
192    /// One entry per host that did not answer, in configured order.
193    pub failures: Vec<RelayFailure>,
194}
195
196impl AdoptionReport {
197    /// Whether any relay answered at all. A run where nothing answered leaves
198    /// the previous stored observation in place.
199    pub fn succeeded(&self) -> bool {
200        !self.observations.is_empty()
201    }
202
203    /// The observation to surface: the **highest** count, ties resolved to the
204    /// first configured host. Deterministic, so two instances with the same
205    /// config log the same number.
206    pub fn best(&self) -> Option<&AdoptionObservation> {
207        self.observations.iter().fold(None, |best, obs| match best {
208            Some(b) if b.repos >= obs.repos => Some(b),
209            _ => Some(obs),
210        })
211    }
212
213    /// Whether two relays disagree about the count — itself worth logging, since
214    /// non-archival relays index different host sets.
215    pub fn disagrees(&self) -> bool {
216        let mut counts = self.observations.iter().map(|o| o.repos);
217        match counts.next() {
218            Some(first) => counts.any(|n| n != first),
219            None => false,
220        }
221    }
222}
223
224// ---------------------------------------------------------------------------
225// Errors
226// ---------------------------------------------------------------------------
227
228/// Why a relay query failed. Every variant carries `host`, so a fan-out failure
229/// is attributable without the caller threading context back in.
230///
231/// Note the deliberate absence of a field named `source`: `thiserror` treats
232/// that name as `#[source]`, which requires `std::error::Error + 'static` —
233/// which `anyhow::Error` (what [`crate::net::guarded_get_no_privacy`] returns)
234/// does not implement. Transport failures are therefore flattened to a `String`
235/// built with anyhow's alternate formatter, which preserves the whole context
236/// chain in the one place it is consumed: the log line.
237#[derive(Debug, thiserror::Error)]
238pub enum RelayError {
239    /// The configured host is not a usable `http(s)` target.
240    #[error("relay host {host:?} is not a usable http(s) target: {reason}")]
241    BadHost {
242        /// The offending configured value.
243        host: String,
244        /// Why it was rejected.
245        reason: String,
246    },
247
248    /// The request never produced a response (DNS, TLS, connect, SSRF refusal).
249    #[error("relay request to {host:?} failed: {reason}")]
250    Transport {
251        /// The relay base URL.
252        host: String,
253        /// The flattened `anyhow` context chain.
254        reason: String,
255    },
256
257    /// The relay answered with a non-2xx status.
258    #[error("relay {host:?} returned HTTP {status}: {error}")]
259    Http {
260        /// The relay base URL.
261        host: String,
262        /// The HTTP status.
263        status: StatusCode,
264        /// A short snippet of the body (the atproto error envelope, usually).
265        error: String,
266    },
267
268    /// The relay rate-limited us. Distinct from [`RelayError::Http`] because the
269    /// required behaviour differs: abort the run, log `Retry-After`, never retry
270    /// tighter.
271    ///
272    /// `retry_after` is rendered INTO the `Display` string rather than only held
273    /// as a field. The scheduler logs `err.to_string()`, so a field the message
274    /// omits is unreachable — the value was parsed, stored, and then silently
275    /// dropped, leaving the operator with no idea how long to back off and
276    /// NETWORK-SPEC §4.6 ("`warn!` with `Retry-After` if present") unmet.
277    #[error("relay {host:?} rate-limited the probe (429){}",
278        match retry_after {
279            Some(secs) => format!(", Retry-After: {secs}s"),
280            None => String::new(),
281        })]
282    RateLimited {
283        /// The relay base URL.
284        host: String,
285        /// `Retry-After` in seconds, when the delta-seconds form was sent.
286        /// `None` covers both "absent" and "sent as an HTTP-date", which is
287        /// deliberately not parsed — see `retry_after_secs`.
288        retry_after: Option<u64>,
289    },
290
291    /// The body was not the expected JSON envelope.
292    #[error("relay {host:?} returned an unparseable body: {reason}")]
293    Malformed {
294        /// The relay base URL.
295        host: String,
296        /// The serde error.
297        reason: String,
298    },
299
300    /// The whole walk against one host blew its deadline.
301    #[error("relay {host:?} did not answer within {after:?}")]
302    TimedOut {
303        /// The relay base URL.
304        host: String,
305        /// The deadline that was exceeded.
306        after: Duration,
307    },
308}
309
310// ---------------------------------------------------------------------------
311// Host normalization
312// ---------------------------------------------------------------------------
313
314/// Normalize one configured relay host into an origin URL.
315///
316/// `FEATHERREADER_RELAY_HOSTS` is documented as a **bare-hostname** list
317/// (`relay1.us-west.bsky.network,…`) while the spec's prose example is a full
318/// `https://` URL — both must work, so a value with no `://` gets `https://`.
319/// Any path/query/fragment is stripped, so a stray `https://host/xrpc/` in
320/// config cannot produce `…/xrpc/xrpc/…`.
321pub fn normalize_relay_host(raw: &str) -> Result<String, RelayError> {
322    let trimmed = raw.trim();
323    if trimmed.is_empty() {
324        return Err(RelayError::BadHost {
325            host: raw.to_string(),
326            reason: "empty".to_string(),
327        });
328    }
329    let candidate = if trimmed.contains("://") {
330        trimmed.to_string()
331    } else {
332        format!("https://{trimmed}")
333    };
334    let url = url::Url::parse(&candidate).map_err(|err| RelayError::BadHost {
335        host: raw.to_string(),
336        reason: err.to_string(),
337    })?;
338    if !matches!(url.scheme(), "http" | "https") {
339        return Err(RelayError::BadHost {
340            host: raw.to_string(),
341            reason: format!("scheme {:?} is not http(s)", url.scheme()),
342        });
343    }
344    let host = url.host_str().ok_or_else(|| RelayError::BadHost {
345        host: raw.to_string(),
346        reason: "no host component".to_string(),
347    })?;
348    let mut origin = format!("{}://{}", url.scheme(), host);
349    if let Some(port) = url.port() {
350        origin.push_str(&format!(":{port}"));
351    }
352    Ok(origin)
353}
354
355// ---------------------------------------------------------------------------
356// The client
357// ---------------------------------------------------------------------------
358
359/// A thin, read-only client over the public relays.
360///
361/// Holds the shared `reqwest::Client` (never builds its own — one connection
362/// pool for the whole 512 MB box) and the normalized host list. Every request
363/// goes through [`crate::net::guarded_get_no_privacy`]; see the module doc for
364/// why.
365#[derive(Debug, Clone)]
366pub struct RelayClient {
367    http: reqwest::Client,
368    hosts: Vec<String>,
369    page_limit: u32,
370    max_pages: usize,
371    page_delay: Duration,
372    /// Soft, checked between pages; see [`DEFAULT_HOST_BUDGET`].
373    host_budget: Duration,
374}
375
376impl RelayClient {
377    /// Build a client over `hosts` (bare hostnames or full URLs; see
378    /// [`normalize_relay_host`]), deduped with configured order preserved.
379    ///
380    /// An **empty** host list is `Ok` and yields `is_enabled() == false` — that
381    /// is the operator's second off-switch, alongside an interval of `0`, and
382    /// must not be an error. A *malformed* host, by contrast, fails loud: a typo
383    /// should be visible, not silently probed forever.
384    pub fn new(http: reqwest::Client, hosts: &[String]) -> Result<Self, RelayError> {
385        let mut normalized: Vec<String> = Vec::with_capacity(hosts.len());
386        for raw in hosts {
387            let host = normalize_relay_host(raw)?;
388            if !normalized.contains(&host) {
389                normalized.push(host);
390            }
391        }
392        Ok(Self {
393            http,
394            hosts: normalized,
395            page_limit: DEFAULT_PAGE_LIMIT,
396            max_pages: MAX_PAGES,
397            page_delay: DEFAULT_PAGE_DELAY,
398            host_budget: DEFAULT_HOST_BUDGET,
399        })
400    }
401
402    /// Override the per-page repo limit, clamped into the relay's usable range.
403    /// The seam a future config knob would use; the clamp is what stops a `0`
404    /// (which the live relay answers with one repo *and* a cursor) from burning
405    /// the whole page budget.
406    pub fn with_page_limit(mut self, limit: u32) -> Self {
407        self.page_limit = limit.clamp(MIN_PAGE_LIMIT, MAX_PAGE_LIMIT);
408        self
409    }
410
411    /// The hard per-host deadline — a pure backstop.
412    ///
413    /// `walk_pages` already bounds itself: each page is wrapped in a timeout of
414    /// whatever is LEFT of [`Self::host_budget`], so the walk cannot outlive its
415    /// budget and the partial count is always reachable. This exists only to
416    /// catch a future that somehow fails to observe that timeout.
417    ///
418    /// Sized against [`crate::net::WORST_CASE_REQUEST`] — `(MAX_REDIRECTS + 1) ×
419    /// FETCH_TIMEOUT`, i.e. **180 s**, not one `FETCH_TIMEOUT`. Budgeting a
420    /// single request cost was the mistake in the first cut of this fix: at
421    /// `soft + 30 s + 10 s` the backstop sat *below* the cost of one redirecting
422    /// page, so a relay behind a captive portal or a `302` chain still had its
423    /// walk killed mid-page and still produced no observation. Derived, not
424    /// configured, so the two cannot be tuned into crossing.
425    fn hard_deadline(&self) -> Duration {
426        self.host_budget + crate::net::WORST_CASE_REQUEST + HARD_DEADLINE_SLACK
427    }
428
429    /// The normalized relay base URLs this client will query, in order.
430    pub fn hosts(&self) -> &[String] {
431        &self.hosts
432    }
433
434    /// Whether there is anything to query at all.
435    pub fn is_enabled(&self) -> bool {
436        !self.hosts.is_empty()
437    }
438
439    /// Count the repos holding `collection` on **every** configured relay.
440    ///
441    /// Hosts are queried **sequentially**, in configured order: the politeness
442    /// budget is per-run, both default relays are operated by the same party,
443    /// and sequential keeps peak memory at one page while making `observations`
444    /// deterministically ordered. A `429` from any host aborts the whole run
445    /// (both defaults sit behind one operator's rate limiter, so hammering the
446    /// next one after the first says "slow down" is the impolite reading).
447    ///
448    /// Never returns `Err`: a failure is data, recorded per-host on the report.
449    pub async fn count_repos_with_collection(&self, collection: &str) -> AdoptionReport {
450        let mut observations = Vec::new();
451        let mut failures = Vec::new();
452
453        for host in &self.hosts {
454            // The HARD backstop. Normally the soft budget inside `walk_pages`
455            // fires first and yields a truncated-but-RECORDED observation; this
456            // only catches a request still hanging past its own FETCH_TIMEOUT,
457            // where there is genuinely nothing to record.
458            match tokio::time::timeout(self.hard_deadline(), self.count_on_host(host, collection))
459                .await
460            {
461                Ok(Ok(obs)) => observations.push(obs),
462                Ok(Err(err)) => {
463                    let rate_limited = matches!(err, RelayError::RateLimited { .. });
464                    failures.push(RelayFailure {
465                        host: host.clone(),
466                        reason: err.to_string(),
467                    });
468                    if rate_limited {
469                        break;
470                    }
471                }
472                Err(_) => failures.push(RelayFailure {
473                    host: host.clone(),
474                    reason: RelayError::TimedOut {
475                        host: host.clone(),
476                        after: self.hard_deadline(),
477                    }
478                    .to_string(),
479                }),
480            }
481        }
482
483        AdoptionReport {
484            collection: collection.to_string(),
485            observations,
486            failures,
487        }
488    }
489
490    /// Walk one host's pages and return its observation.
491    async fn count_on_host(
492        &self,
493        host: &str,
494        collection: &str,
495    ) -> Result<AdoptionObservation, RelayError> {
496        let (repos, truncated) = self
497            .walk_pages(host, |cursor| self.fetch_page(host, collection, cursor))
498            .await?;
499        Ok(self.observe(host, collection, repos, truncated))
500    }
501
502    /// The pagination driver, generic over how a page is obtained so the
503    /// termination rules (absent cursor, empty page, page cap) are testable
504    /// without a socket — which matters because the SSRF guard rightly refuses a
505    /// loopback stub server, so an end-to-end fixture is not available.
506    ///
507    /// Returns `(repos, truncated)` — `truncated` is set by EITHER bound: the
508    /// [`MAX_PAGES`] cap or the [`DEFAULT_HOST_BUDGET`] wall-clock budget. Both
509    /// mean the same thing to the caller ("this count is a floor"), so they
510    /// share one flag and the `/about` copy reads "at least N" either way.
511    async fn walk_pages<F, Fut>(&self, host: &str, mut fetch: F) -> Result<(u64, bool), RelayError>
512    where
513        F: FnMut(Option<String>) -> Fut,
514        Fut: std::future::Future<Output = Result<ListReposByCollectionOut, RelayError>>,
515    {
516        let started = tokio::time::Instant::now();
517        let mut repos: u64 = 0;
518        let mut cursor: Option<String> = None;
519        for page_no in 0..self.max_pages {
520            // The SOFT budget. A slow relay used to blow the outer hard timeout
521            // and yield NO observation at all — every run, forever — because the
522            // timeout wrapped the whole walk and discarded its partial result.
523            // Stopping here keeps what we already counted and marks it truncated.
524            let remaining = self.host_budget.saturating_sub(started.elapsed());
525            if page_no > 0 && remaining.is_zero() {
526                return Ok((repos, true));
527            }
528            // Politeness: one second between pages against the same host.
529            if page_no > 0 && !self.page_delay.is_zero() {
530                tokio::time::sleep(self.page_delay).await;
531            }
532            let sent = cursor.take();
533            // Bound the page by what's LEFT of the budget, rather than trusting
534            // the outer deadline to be generous enough. One guarded request can
535            // cost `net::WORST_CASE_REQUEST` (redirects × FETCH_TIMEOUT), so a
536            // between-pages check alone can overshoot the budget by minutes and
537            // let the hard timeout kill the walk mid-page — silently restoring
538            // the very "no observation, ever" bug this exists to prevent. With
539            // the timeout scoped HERE, the walk cannot outlive its budget, so
540            // the partial count is always reachable and the outer deadline is a
541            // pure backstop rather than a number that has to be kept in sync.
542            let page = match tokio::time::timeout(remaining, fetch(sent.clone())).await {
543                Ok(res) => res?,
544                Err(_) => return Ok((repos, true)),
545            };
546            // `repos` is fed by an untrusted peer; never panic on overflow.
547            repos = repos.saturating_add(page.repos.len() as u64);
548            match advance(&page) {
549                // A relay that hands back the SAME cursor it was just given is
550                // not paginating. Following it would re-count the identical page
551                // up to `max_pages` times and publish the sum as "at least N" —
552                // inflating, in the UNSAFE direction, the one number this whole
553                // feature exists to state. Refuse the walk rather than record a
554                // count we already know is wrong.
555                PageStep::Continue(next) if Some(&next) == sent.as_ref() => {
556                    return Err(RelayError::Malformed {
557                        host: host.to_string(),
558                        reason: format!(
559                            "repeated cursor {next:?} at page {page_no} instead of advancing"
560                        ),
561                    })
562                }
563                PageStep::Continue(next) => cursor = Some(next),
564                PageStep::Done => return Ok((repos, false)),
565            }
566        }
567        // Fell out of the loop: the cap was hit, so the count is a floor.
568        Ok((repos, true))
569    }
570
571    /// Fetch and parse one page.
572    async fn fetch_page(
573        &self,
574        host: &str,
575        collection: &str,
576        cursor: Option<String>,
577    ) -> Result<ListReposByCollectionOut, RelayError> {
578        let url = self.page_url(host, collection, cursor.as_deref());
579        let resp = crate::net::guarded_get_no_privacy(
580            &self.http,
581            &url,
582            &[(ACCEPT, HeaderValue::from_static("application/json"))],
583        )
584        .await
585        .map_err(|err| RelayError::Transport {
586            host: host.to_string(),
587            reason: format!("{err:#}"),
588        })?;
589
590        // The 429 check MUST precede the generic non-2xx check, or RateLimited
591        // is unreachable.
592        if resp.status() == StatusCode::TOO_MANY_REQUESTS {
593            return Err(RelayError::RateLimited {
594                host: host.to_string(),
595                retry_after: retry_after_secs(resp.headers()),
596            });
597        }
598        if !resp.status().is_success() {
599            let status = resp.status();
600            let snippet = crate::net::read_capped(resp)
601                .await
602                .ok()
603                .map(|body| {
604                    String::from_utf8_lossy(&body)
605                        .chars()
606                        .take(ERROR_SNIPPET_CHARS)
607                        .collect::<String>()
608                })
609                .unwrap_or_default();
610            return Err(RelayError::Http {
611                host: host.to_string(),
612                status,
613                error: snippet,
614            });
615        }
616
617        // `read_capped` (streamed, 8 MiB, aborts mid-body) rather than
618        // `resp.json()`, which buffers a hostile body without bound. Peak memory
619        // for the whole probe is one page — O(1) in the size of the network.
620        let body = crate::net::read_capped(resp)
621            .await
622            .map_err(|err| RelayError::Transport {
623                host: host.to_string(),
624                reason: format!("{err:#}"),
625            })?;
626        serde_json::from_slice(&body).map_err(|err| RelayError::Malformed {
627            host: host.to_string(),
628            reason: err.to_string(),
629        })
630    }
631
632    /// The page URL. Pure — no I/O — so the query shape is unit-testable.
633    ///
634    /// Built with `format!` + [`urlencode`] because the declared `reqwest`
635    /// feature set is `rustls + gzip + json` (no `query` feature), exactly as
636    /// [`crate::atproto::resolve_handle`] does it. The `cursor` genuinely needs
637    /// encoding: live values are base64url, and a standard-base64 variant
638    /// carrying `+`, `/`, or `=` would silently corrupt the query.
639    fn page_url(&self, host: &str, collection: &str, cursor: Option<&str>) -> String {
640        let mut url = format!(
641            "{}/xrpc/{}?collection={}&limit={}",
642            host.trim_end_matches('/'),
643            LIST_REPOS_BY_COLLECTION,
644            urlencode(collection),
645            self.page_limit,
646        );
647        if let Some(cursor) = cursor {
648            url.push_str(&format!("&cursor={}", urlencode(cursor)));
649        }
650        url
651    }
652
653    /// Stamp an observation. One place mints the timestamp so the two exits of
654    /// [`walk_pages`](Self::walk_pages) cannot drift in format.
655    fn observe(
656        &self,
657        host: &str,
658        collection: &str,
659        repos: u64,
660        truncated: bool,
661    ) -> AdoptionObservation {
662        AdoptionObservation {
663            source: host.to_string(),
664            collection: collection.to_string(),
665            repos,
666            truncated,
667            observed_at: chrono::Utc::now().to_rfc3339_opts(chrono::SecondsFormat::Secs, true),
668        }
669    }
670}
671
672// ---------------------------------------------------------------------------
673// Pure helpers
674// ---------------------------------------------------------------------------
675
676/// What to do after a page.
677#[derive(Debug, PartialEq, Eq)]
678enum PageStep {
679    /// Continue with this cursor.
680    Continue(String),
681    /// Stop.
682    Done,
683}
684
685/// Termination rules, mirroring the guard already in
686/// [`crate::atproto::PdsClient::list_all_records`]: stop when the cursor is
687/// absent (the normal end — the final page omits the key), **and** stop when a
688/// cursor came back with an empty page, so a relay that echoes a cursor forever
689/// cannot spin the loop.
690///
691/// A malformed cursor is not an error case here: the live relay answers a
692/// garbage cursor with `200 {"repos":[]}` and no cursor, which terminates
693/// cleanly through the same rule.
694fn advance(page: &ListReposByCollectionOut) -> PageStep {
695    match &page.cursor {
696        Some(next) if !page.repos.is_empty() => PageStep::Continue(next.clone()),
697        _ => PageStep::Done,
698    }
699}
700
701/// `Retry-After` in seconds, for the log line. The HTTP-date form returns
702/// `None` rather than pulling a date parse in for a value we never compute with.
703fn retry_after_secs(headers: &HeaderMap) -> Option<u64> {
704    headers
705        .get(RETRY_AFTER)?
706        .to_str()
707        .ok()?
708        .trim()
709        .parse::<u64>()
710        .ok()
711}
712
713// ---------------------------------------------------------------------------
714// Tests — all pure. No socket: `net::guarded_get_no_privacy` refuses loopback by
715// design, so a local stub server cannot exercise the fetch path (net.rs's own
716// tests hit the same wall). Every decision is therefore factored into a pure
717// function and tested here; the fetch itself is `guarded_get_no_privacy` +
718// `read_capped`, both covered in `net.rs`.
719// ---------------------------------------------------------------------------
720
721#[cfg(test)]
722mod tests {
723    use super::*;
724
725    fn test_client() -> RelayClient {
726        RelayClient::new(
727            reqwest::Client::builder().build().unwrap(),
728            &[DEFAULT_RELAY_HOSTS[0].to_string()],
729        )
730        .unwrap()
731    }
732
733    /// A client whose page delay is zero, so the pagination tests don't sleep.
734    fn instant_client(max_pages: usize) -> RelayClient {
735        let mut c = test_client();
736        c.max_pages = max_pages;
737        c.page_delay = Duration::ZERO;
738        c
739    }
740
741    fn parse(body: &str) -> ListReposByCollectionOut {
742        serde_json::from_str(body).expect("relay body")
743    }
744
745    fn obs(source: &str, repos: u64) -> AdoptionObservation {
746        AdoptionObservation {
747            source: source.to_string(),
748            collection: crate::lexicon::nsid::SUBSCRIPTION.to_string(),
749            repos,
750            truncated: false,
751            observed_at: "2026-08-13T00:00:00Z".to_string(),
752        }
753    }
754
755    // -- host normalization ------------------------------------------------
756
757    #[test]
758    fn normalize_defaults_bare_host_to_https() {
759        assert_eq!(
760            normalize_relay_host("relay1.us-west.bsky.network").unwrap(),
761            "https://relay1.us-west.bsky.network"
762        );
763        assert_eq!(
764            normalize_relay_host("  relay1.us-east.bsky.network  ").unwrap(),
765            "https://relay1.us-east.bsky.network"
766        );
767    }
768
769    #[test]
770    fn normalize_strips_trailing_slash_and_path() {
771        assert_eq!(
772            normalize_relay_host("https://relay.example/").unwrap(),
773            "https://relay.example"
774        );
775        assert_eq!(
776            normalize_relay_host("https://relay.example/xrpc/whatever?x=1#f").unwrap(),
777            "https://relay.example"
778        );
779        // An explicit port survives (a self-hosted relay may use one).
780        assert_eq!(
781            normalize_relay_host("http://relay.example:8080/").unwrap(),
782            "http://relay.example:8080"
783        );
784    }
785
786    #[test]
787    fn normalize_rejects_bad_schemes_and_empties() {
788        for bad in ["wss://relay.example", "file:///etc/passwd", "", "   "] {
789            assert!(
790                normalize_relay_host(bad).is_err(),
791                "{bad:?} should be rejected"
792            );
793        }
794    }
795
796    // -- URL shape ---------------------------------------------------------
797
798    #[test]
799    fn page_url_without_cursor_is_exact() {
800        let c = test_client();
801        assert_eq!(
802            c.page_url(
803                "https://relay1.us-west.bsky.network",
804                crate::lexicon::nsid::SUBSCRIPTION,
805                None
806            ),
807            "https://relay1.us-west.bsky.network/xrpc/com.atproto.sync.listReposByCollection\
808             ?collection=community.lexicon.rss.subscription&limit=500"
809        );
810    }
811
812    #[test]
813    fn page_url_percent_encodes_the_cursor() {
814        let c = test_client();
815        let url = c.page_url("https://relay.example", "a.b.c", Some("aa+bb/cc=="));
816        assert!(
817            url.ends_with("&cursor=aa%2Bbb%2Fcc%3D%3D"),
818            "cursor must be percent-encoded, got: {url}"
819        );
820        // The NSID is unreserved, so it passes through readable.
821        assert!(url.contains("?collection=a.b.c&limit=500"), "{url}");
822    }
823
824    #[test]
825    fn page_limit_is_clamped() {
826        assert_eq!(test_client().with_page_limit(0).page_limit, MIN_PAGE_LIMIT);
827        assert_eq!(
828            test_client().with_page_limit(99_999).page_limit,
829            MAX_PAGE_LIMIT
830        );
831        assert_eq!(test_client().with_page_limit(200).page_limit, 200);
832    }
833
834    // -- wire parsing (bodies captured live from relay1.us-west) ------------
835
836    #[test]
837    fn parses_the_live_single_repo_terminal_page() {
838        let page = parse(r#"{"repos":[{"did":"did:plc:ohutz6x5acjmpuulp3x7wxxc"}]}"#);
839        assert_eq!(page.repos.len(), 1);
840        assert_eq!(page.cursor, None);
841    }
842
843    #[test]
844    fn parses_a_cursor_first_paged_body() {
845        let page = parse(
846            r#"{"cursor":"QQAAAGsAAAGVmhd3TGRpZDpwbGM6dGxkYW91amwzNzZ6dTV3ZXphem54ZmV2AA",
847                "repos":[{"did":"did:plc:qw4uaobncdi5ijsj4mthdboq"},
848                         {"did":"did:plc:tldaoujl376zu5wezaznxfev"}]}"#,
849        );
850        assert_eq!(page.repos.len(), 2);
851        assert!(page.cursor.is_some());
852    }
853
854    #[test]
855    fn parses_empty_absent_and_unknown_field_bodies() {
856        let empty = parse(r#"{"repos":[]}"#);
857        assert_eq!(empty.repos.len(), 0);
858        assert_eq!(empty.cursor, None);
859        // Missing `repos` entirely — `#[serde(default)]`.
860        assert_eq!(parse("{}").repos.len(), 0);
861        // Forward compatibility: an added field must not fail the parse.
862        let future = parse(r#"{"repos":[{"did":"did:web:lexicon.store","note":1}],"total":7}"#);
863        assert_eq!(future.repos.len(), 1);
864        assert_eq!(future.repos[0].did, "did:web:lexicon.store");
865    }
866
867    // -- pagination termination -------------------------------------------
868
869    #[test]
870    fn advance_stops_without_a_cursor() {
871        assert_eq!(
872            advance(&parse(r#"{"repos":[{"did":"a"}]}"#)),
873            PageStep::Done
874        );
875    }
876
877    #[test]
878    fn advance_continues_on_a_cursor_with_rows() {
879        assert_eq!(
880            advance(&parse(r#"{"repos":[{"did":"a"}],"cursor":"c1"}"#)),
881            PageStep::Continue("c1".to_string())
882        );
883    }
884
885    #[test]
886    fn advance_stops_on_an_empty_page_even_with_a_cursor() {
887        // The echo guard: a relay that returns a cursor forever cannot spin us.
888        assert_eq!(
889            advance(&parse(r#"{"repos":[],"cursor":"c1"}"#)),
890            PageStep::Done
891        );
892    }
893
894    #[tokio::test]
895    async fn walk_pages_follows_the_cursor_then_stops() {
896        let c = instant_client(MAX_PAGES);
897        let pages = [
898            r#"{"repos":[{"did":"a"},{"did":"b"}],"cursor":"c1"}"#,
899            r#"{"repos":[{"did":"c"}]}"#,
900        ];
901        let mut seen_cursors: Vec<Option<String>> = Vec::new();
902        let mut n = 0usize;
903        let (repos, truncated) = c
904            .walk_pages("https://relay.example", |cursor| {
905                seen_cursors.push(cursor);
906                let body = pages[n];
907                n += 1;
908                async move { Ok(parse(body)) }
909            })
910            .await
911            .unwrap();
912        assert_eq!(repos, 3);
913        assert!(!truncated);
914        assert_eq!(seen_cursors, vec![None, Some("c1".to_string())]);
915    }
916
917    #[tokio::test]
918    async fn walk_pages_stops_at_the_page_cap_and_marks_truncated() {
919        let c = instant_client(3);
920        // A relay that always returns a full-ish page and a genuinely FRESH
921        // cursor each time (a repeated one is refused — see the test below).
922        let mut n = 0usize;
923        let (repos, truncated) = c
924            .walk_pages("https://relay.example", |_| {
925                n += 1;
926                let body = format!(r#"{{"repos":[{{"did":"a"}},{{"did":"b"}}],"cursor":"c{n}"}}"#);
927                async move { Ok(parse(&body)) }
928            })
929            .await
930            .unwrap();
931        assert_eq!(repos, 6, "3 pages × 2 repos");
932        assert!(truncated, "hitting the cap makes the count a floor");
933    }
934
935    /// **Regression (v0.2.9 review).** A merely-SLOW relay used to produce no
936    /// observation at all: the hard `tokio::time::timeout` wrapped the entire
937    /// walk, so hitting it discarded every page already counted — and since the
938    /// deadline was reached the same way on every run, that host contributed
939    /// nothing, forever. The soft budget now stops between pages and keeps the
940    /// partial count, flagged `truncated` so it reads as "at least N".
941    ///
942    /// `start_paused` drives tokio's clock, so the simulated 6 s pages cost no
943    /// real time and the arithmetic is exact rather than timing-dependent.
944    #[tokio::test(start_paused = true)]
945    async fn walk_pages_keeps_a_partial_count_when_the_budget_runs_out() {
946        let mut c = instant_client(50);
947        c.host_budget = Duration::from_secs(10);
948
949        let mut n = 0usize;
950        let (repos, truncated) = c
951            .walk_pages("https://relay.example", |_| {
952                n += 1;
953                let body = format!(r#"{{"repos":[{{"did":"a"}},{{"did":"b"}}],"cursor":"c{n}"}}"#);
954                async move {
955                    // A slow relay: each page costs 6 s of the 10 s budget.
956                    tokio::time::sleep(Duration::from_secs(6)).await;
957                    Ok(parse(&body))
958                }
959            })
960            .await
961            .unwrap();
962
963        // Page 1 spends 6 s of the 10 s budget. Page 2 is then handed only the
964        // REMAINING 4 s and cannot finish in 6 s, so it is cancelled: the walk
965        // never overshoots its budget. That is what lets the outer deadline be a
966        // backstop rather than a number that has to be kept in sync with
967        // net.rs's redirect arithmetic.
968        assert_eq!(repos, 2, "page 1's count survives; page 2 never completed");
969        assert!(
970            truncated,
971            "a budget stop makes the count a floor, not a loss"
972        );
973        assert_eq!(n, 2, "page 2 was attempted, then cut off at the budget");
974    }
975
976    /// **Regression (v0.2.9 review, second pass).** The first cut of this fix
977    /// sized the backstop as `soft + FETCH_TIMEOUT + slack` = 160 s. But ONE
978    /// guarded request costs up to `(MAX_REDIRECTS + 1) × FETCH_TIMEOUT` = 180 s,
979    /// because `net::guarded_get_inner` re-validates each hop through a freshly
980    /// built client carrying its own full timeout. So the backstop sat *below*
981    /// the cost of a single redirecting page: a relay behind a captive portal or
982    /// a `302` chain still had its walk killed mid-page and still produced no
983    /// observation — the exact bug, narrowed to the redirect case.
984    ///
985    /// Asserting against the real worst case rather than restating
986    /// `hard_deadline()`'s own definition, which is what let it through: the old
987    /// assertion was `hard >= soft + FETCH_TIMEOUT` against a value *defined* as
988    /// `soft + FETCH_TIMEOUT + slack`, so it could not fail for any input.
989    #[test]
990    fn the_hard_deadline_clears_one_worst_case_request() {
991        let c = test_client();
992        let worst = crate::net::FETCH_TIMEOUT * (crate::net::MAX_REDIRECTS as u32 + 1);
993        assert_eq!(
994            crate::net::WORST_CASE_REQUEST,
995            worst,
996            "worst-case arithmetic"
997        );
998        assert!(
999            c.hard_deadline() >= c.host_budget + worst,
1000            "hard {:?} must clear soft {:?} + one worst-case request {:?}",
1001            c.hard_deadline(),
1002            c.host_budget,
1003            worst
1004        );
1005    }
1006
1007    /// The structural half of the same fix: a single page that outruns the
1008    /// remaining budget is cancelled and the partial count kept, rather than the
1009    /// walk running on and being killed by the outer deadline. This is what makes
1010    /// the backstop arithmetic a safety net instead of a load-bearing number.
1011    #[tokio::test(start_paused = true)]
1012    async fn a_page_that_outruns_the_budget_still_yields_what_was_counted() {
1013        let mut c = instant_client(50);
1014        c.host_budget = Duration::from_secs(10);
1015
1016        let mut n = 0usize;
1017        let (repos, truncated) = c
1018            .walk_pages("https://relay.example", |_| {
1019                n += 1;
1020                let first = n == 1;
1021                async move {
1022                    if first {
1023                        // Page 1 is fine.
1024                        return Ok(parse(r#"{"repos":[{"did":"a"}],"cursor":"c1"}"#));
1025                    }
1026                    // Page 2 hangs far past the whole budget — as a redirect
1027                    // chain would (up to 180 s for one request).
1028                    tokio::time::sleep(Duration::from_secs(600)).await;
1029                    Ok(parse(r#"{"repos":[{"did":"b"}],"cursor":"c2"}"#))
1030                }
1031            })
1032            .await
1033            .unwrap();
1034
1035        assert_eq!(repos, 1, "page 1's count survives the hung page 2");
1036        assert!(truncated, "and is reported as a floor");
1037    }
1038
1039    /// **Regression (v0.2.9 review).** `Retry-After` was parsed into the error
1040    /// and then unreachable: the scheduler logs `err.to_string()`, and the
1041    /// `Display` string omitted the field. NETWORK-SPEC §4.6 requires it be
1042    /// warned with, so it has to survive rendering, not just parsing.
1043    #[test]
1044    fn rate_limited_display_carries_retry_after() {
1045        let with = RelayError::RateLimited {
1046            host: "https://relay.example".to_string(),
1047            retry_after: Some(120),
1048        }
1049        .to_string();
1050        assert!(with.contains("Retry-After: 120s"), "{with}");
1051
1052        let without = RelayError::RateLimited {
1053            host: "https://relay.example".to_string(),
1054            retry_after: None,
1055        }
1056        .to_string();
1057        assert!(!without.contains("Retry-After"), "{without}");
1058        assert!(without.contains("rate-limited"), "{without}");
1059    }
1060
1061    #[tokio::test]
1062    async fn walk_pages_refuses_a_relay_that_repeats_its_cursor() {
1063        // The failure this guards: a relay stuck on one page would otherwise be
1064        // followed to the cap, summing the SAME page every pass and publishing
1065        // "at least 2 × max_pages" for what is really 2 repos — inflating the
1066        // number in the unsafe direction. No observation is better than a wrong
1067        // one, so the walk errors instead of returning a count.
1068        let c = instant_client(50);
1069        let err = c
1070            .walk_pages("https://relay.example", |_| async {
1071                Ok(parse(
1072                    r#"{"repos":[{"did":"a"},{"did":"b"}],"cursor":"stuck"}"#,
1073                ))
1074            })
1075            .await
1076            .unwrap_err();
1077        match err {
1078            RelayError::Malformed { host, reason } => {
1079                assert_eq!(host, "https://relay.example");
1080                assert!(reason.contains("repeated cursor"), "reason was {reason:?}");
1081            }
1082            other => panic!("expected Malformed, got {other:?}"),
1083        }
1084    }
1085
1086    #[tokio::test]
1087    async fn walk_pages_propagates_a_page_error() {
1088        let c = instant_client(MAX_PAGES);
1089        let err = c
1090            .walk_pages("https://relay.example", |_| async {
1091                Err(RelayError::Malformed {
1092                    host: "https://relay.example".to_string(),
1093                    reason: "expected value".to_string(),
1094                })
1095            })
1096            .await
1097            .unwrap_err();
1098        assert!(matches!(err, RelayError::Malformed { .. }));
1099    }
1100
1101    // -- Retry-After -------------------------------------------------------
1102
1103    #[test]
1104    fn retry_after_reads_delta_seconds_only() {
1105        let mut h = HeaderMap::new();
1106        assert_eq!(retry_after_secs(&h), None);
1107        h.insert(RETRY_AFTER, HeaderValue::from_static("120"));
1108        assert_eq!(retry_after_secs(&h), Some(120));
1109        h.insert(RETRY_AFTER, HeaderValue::from_static("  120 "));
1110        assert_eq!(retry_after_secs(&h), Some(120));
1111        // The HTTP-date form is deliberately not parsed.
1112        h.insert(
1113            RETRY_AFTER,
1114            HeaderValue::from_static("Wed, 21 Oct 2026 07:28:00 GMT"),
1115        );
1116        assert_eq!(retry_after_secs(&h), None);
1117    }
1118
1119    // -- client construction ----------------------------------------------
1120
1121    #[test]
1122    fn empty_host_list_is_ok_but_disabled() {
1123        let c = RelayClient::new(reqwest::Client::builder().build().unwrap(), &[]).unwrap();
1124        assert!(!c.is_enabled());
1125        assert!(c.hosts().is_empty());
1126    }
1127
1128    #[test]
1129    fn hosts_are_normalized_and_deduped_in_order() {
1130        let c = RelayClient::new(
1131            reqwest::Client::builder().build().unwrap(),
1132            &[
1133                "relay1.us-west.bsky.network".to_string(),
1134                "https://relay1.us-west.bsky.network/".to_string(),
1135                "relay1.us-east.bsky.network".to_string(),
1136            ],
1137        )
1138        .unwrap();
1139        assert_eq!(
1140            c.hosts(),
1141            [
1142                "https://relay1.us-west.bsky.network",
1143                "https://relay1.us-east.bsky.network"
1144            ]
1145        );
1146        assert!(c.is_enabled());
1147    }
1148
1149    #[test]
1150    fn one_bad_host_fails_loud() {
1151        let err = RelayClient::new(
1152            reqwest::Client::builder().build().unwrap(),
1153            &[
1154                "relay1.us-west.bsky.network".to_string(),
1155                "wss://relay1.us-east.bsky.network".to_string(),
1156            ],
1157        )
1158        .unwrap_err();
1159        assert!(matches!(err, RelayError::BadHost { .. }), "{err}");
1160    }
1161
1162    // -- report arithmetic -------------------------------------------------
1163
1164    #[test]
1165    fn report_surfaces_the_max_and_flags_disagreement() {
1166        let report = AdoptionReport {
1167            collection: crate::lexicon::nsid::SUBSCRIPTION.to_string(),
1168            observations: vec![obs("https://west", 2), obs("https://east", 40)],
1169            failures: Vec::new(),
1170        };
1171        assert!(report.succeeded());
1172        assert_eq!(report.best().unwrap().repos, 40);
1173        assert!(report.disagrees());
1174    }
1175
1176    #[test]
1177    fn report_ties_resolve_to_the_first_configured_host() {
1178        let report = AdoptionReport {
1179            collection: "c".to_string(),
1180            observations: vec![obs("https://west", 7), obs("https://east", 7)],
1181            failures: Vec::new(),
1182        };
1183        assert_eq!(report.best().unwrap().source, "https://west");
1184        assert!(!report.disagrees());
1185    }
1186
1187    #[test]
1188    fn report_with_only_failures_did_not_succeed() {
1189        let report = AdoptionReport {
1190            collection: "c".to_string(),
1191            observations: Vec::new(),
1192            failures: vec![RelayFailure {
1193                host: "https://west".to_string(),
1194                reason: "boom".to_string(),
1195            }],
1196        };
1197        assert!(!report.succeeded());
1198        assert!(report.best().is_none());
1199        assert!(!report.disagrees());
1200    }
1201
1202    // -- the guard is really in the path -----------------------------------
1203
1204    /// A relay host that resolves to an internal address is refused by the SSRF
1205    /// guard before any packet leaves the box (IP literal ⇒ no DNS, no connect).
1206    #[tokio::test]
1207    async fn internal_relay_host_is_refused_by_the_guard() {
1208        let c = RelayClient::new(
1209            reqwest::Client::builder().build().unwrap(),
1210            &["http://169.254.169.254".to_string()],
1211        )
1212        .unwrap();
1213        let report = c
1214            .count_repos_with_collection(crate::lexicon::nsid::SUBSCRIPTION)
1215            .await;
1216        assert!(!report.succeeded());
1217        assert_eq!(report.failures.len(), 1);
1218        let reason = &report.failures[0].reason;
1219        assert!(
1220            reason.contains("forbidden") || reason.contains("internal"),
1221            "expected an SSRF refusal, got: {reason}"
1222        );
1223    }
1224}