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}