feather_reader/feed.rs
1//! Feed fetch → parse → sanitize → store pipeline.
2//!
3//! This is the module that turns a feed URL into rows in the [`store`]. It is
4//! deliberately conservative on three axes, because a feed reader ingests
5//! **hostile, arbitrary web input**:
6//!
7//! 1. **Politeness** — fetches use a **conditional GET** (`If-None-Match` /
8//! `If-Modified-Since` from the stored `ETag` / `Last-Modified`), an
9//! identifiable [`crate::USER_AGENT`], a request timeout, and a simple
10//! exponential backoff hint on error. A `304 Not Modified` is a no-op:
11//! the feed is untouched apart from bumping its next-poll time.
12//! 2. **Safety** — every entry's HTML is run through `ammonia` before it is
13//! ever stored. Scripts, event handlers, `javascript:` URLs, tracking
14//! pixels' dangerous attributes, and other XSS vectors are stripped. Feeds
15//! carrying `<script>` is not hypothetical; treat all feed HTML as
16//! untrusted. The reader does not rely on this alone: it re-cleans the
17//! stored body with the same `sanitize_html` at render, through
18//! [`crate::sanitized_html::SanitizedHtml`] (#151).
19//! 3. **Robustness** — a malformed feed is **logged and skipped**, never a
20//! panic. One bad publisher must not take down the poller. All non-test
21//! paths use `Result`/`anyhow`; there are no `unwrap`/`expect`s.
22//!
23//! The normalized shape written to the store is the store's own
24//! [`store::NewFeed`] / [`store::NewEntry`]; dedup is by feed-native GUID via
25//! [`store::insert_entries`]'s `ON CONFLICT (feed_id, guid)` upsert.
26
27use std::time::Duration;
28
29use anyhow::{Context, Result};
30use chrono::{DateTime, SecondsFormat, Utc};
31use feed_rs::model::{Entry as RawEntry, Feed as RawFeed, Link as RawLink, Text};
32use reqwest::header::{ETAG, IF_MODIFIED_SINCE, IF_NONE_MATCH, LAST_MODIFIED};
33use reqwest::{Client, StatusCode};
34use sqlx::SqlitePool;
35use url::Url;
36
37use crate::store::{self, Feed, NewEntry, NewFeed};
38
39/// The privacy classification of a feed URL — the output of
40/// [`classify_feed_privacy`].
41///
42/// A **private** feed carries a secret (a token / key / auth credential) *in the
43/// URL itself* — a Substack `…/feed/private/<token>`, a Patreon `?auth=…` feed,
44/// a Ghost members `?uuid=` feed, a private-podcast token feed (Supercast,
45/// Supporting Cast, tokened Megaphone/Acast+), and so on. FeatherReader stores a
46/// user's subscriptions as records in their **public PDS** (unauthenticated
47/// `getRecord` / `listRecords` + the firehose, retained even after delete), so
48/// writing such a URL anywhere — the PDS *or* the server's own store — would risk
49/// leaking paid / members-only access.
50///
51/// **Decision (stopgap until atproto permissioned data ships): FeatherReader
52/// supports PUBLIC feeds only.** A feed classified [`FeedPrivacy::Private`] is
53/// *refused* at the add / import boundary — never fetched, never stored, never
54/// written to the PDS. There is no local-secret fallback and no override: the
55/// server holds NO private secret, ever, which keeps "your data lives in your
56/// public PDS" 100% honest.
57#[derive(Debug, Clone, PartialEq, Eq)]
58pub enum FeedPrivacy {
59 /// No secret detected in the URL; safe to add as a public feed.
60 Public,
61 /// A secret was detected in the URL. The `String` is a short, human-readable
62 /// reason (for logging / the skip report), e.g. `"substack private feed
63 /// path"`. The feed is refused — not fetched, stored, or written anywhere.
64 Private(String),
65}
66
67impl FeedPrivacy {
68 /// Whether this classification is [`FeedPrivacy::Private`].
69 pub fn is_private(&self) -> bool {
70 matches!(self, FeedPrivacy::Private(_))
71 }
72}
73
74/// Query-parameter *keys* that, when present with a long/opaque value, mark a URL
75/// as carrying a secret. Conservative and lowercase-compared; matched as a whole
76/// key (case-insensitive) so a benign `keyword=` does NOT trip `key`. This is the
77/// generic, provider-agnostic credential-in-query defence — it catches paid
78/// feeds from providers we've never heard of. Covers Patreon (`auth`), Ghost
79/// members (`uuid`), token-in-query feeds (`token`/`key`/`k`/`sig`/`hash`), and
80/// the long tail (`access`/`apikey`/`private`/`password`/`u`/`s`/`p`/…).
81const SECRET_QUERY_KEYS: &[&str] = &[
82 "token", "key", "auth", "secret", "k", "sig", "hash", "access", "apikey", "api_key", "uuid",
83 "id", "u", "s", "p", "private", "password", "pw",
84];
85
86/// Path *segments* / fragments that mark a private-feed URL shape. Matched as a
87/// case-insensitive substring of the (lowercased) path so `/feed/private/<tok>`,
88/// `/members/…`, `/subscriber/…` etc. all trip regardless of the token that
89/// follows. Provider-agnostic: many paid providers expose members-only feeds
90/// under one of these path conventions.
91const PRIVATE_PATH_MARKERS: &[&str] = &[
92 "/private/",
93 "/feed/private/",
94 "/rss/private/",
95 "/private-feed/",
96 "/members/",
97 "/member/",
98 "/subscriber/",
99];
100
101/// A KNOWN paid/private feed provider, matched by host substring + (optionally) a
102/// path/query marker specific to that provider. This is the **secondary**,
103/// precision layer on top of the generic heuristic — it names providers so the
104/// skip report can say *why* and so we catch provider-specific shapes that the
105/// generic pass might rate as borderline. Data-driven and easy to extend: add a
106/// row, don't touch the matcher.
107struct KnownProvider {
108 /// Substring that must appear in the URL host (lowercased), e.g.
109 /// `substack.com`.
110 host_contains: &'static str,
111 /// Optional lowercased substring that must appear in the path-or-query for a
112 /// match (a provider's private-feed marker). `None` = the host alone is
113 /// enough (used for hosts that ONLY serve private/tokened feeds).
114 marker: Option<&'static str>,
115 /// Human-readable reason for the skip report.
116 reason: &'static str,
117}
118
119/// The known-provider table. Covers paid NEWSLETTERS and private PODCASTS — an
120/// RSS reader ingests both. Kept intentionally verbose/commented so it's obvious
121/// what each row targets and safe to extend.
122const KNOWN_PROVIDERS: &[KnownProvider] = &[
123 // --- Paid newsletters -------------------------------------------------
124 // Substack private feed: author.substack.com/feed/private/<token>.
125 KnownProvider {
126 host_contains: "substack.com",
127 marker: Some("/feed/private/"),
128 reason: "Substack private feed",
129 },
130 // Patreon RSS carries the member token as ?auth=.
131 KnownProvider {
132 host_contains: "patreon.com",
133 marker: Some("auth="),
134 reason: "Patreon member feed",
135 },
136 // Ghost members feed: ?uuid=<member-uuid> (or a members token path).
137 KnownProvider {
138 host_contains: "ghost.io",
139 marker: Some("uuid="),
140 reason: "Ghost members feed",
141 },
142 // Buttondown paid RSS uses a per-subscriber token in the path/query.
143 KnownProvider {
144 host_contains: "buttondown.email",
145 marker: Some("token"),
146 reason: "Buttondown premium feed",
147 },
148 KnownProvider {
149 host_contains: "buttondown.com",
150 marker: Some("token"),
151 reason: "Buttondown premium feed",
152 },
153 // Beehiiv premium RSS carries a subscriber token.
154 KnownProvider {
155 host_contains: "beehiiv.com",
156 marker: Some("token"),
157 reason: "Beehiiv premium feed",
158 },
159 // Memberful-gated feeds (host or ?auth token).
160 KnownProvider {
161 host_contains: "memberful.com",
162 marker: None,
163 reason: "Memberful members feed",
164 },
165 // Pico / Steady member feeds.
166 KnownProvider {
167 host_contains: "pico.link",
168 marker: None,
169 reason: "Pico member feed",
170 },
171 KnownProvider {
172 host_contains: "steadyhq.com",
173 marker: None,
174 reason: "Steady member feed",
175 },
176 // --- Private podcasts -------------------------------------------------
177 // Supercast private podcast feeds (host serves tokened member feeds only).
178 KnownProvider {
179 host_contains: "supercast.com",
180 marker: None,
181 reason: "Supercast private podcast",
182 },
183 KnownProvider {
184 host_contains: "supercast.tech",
185 marker: None,
186 reason: "Supercast private podcast",
187 },
188 // Supporting Cast private podcast feeds (supportingcast.fm).
189 KnownProvider {
190 host_contains: "supportingcast.fm",
191 marker: None,
192 reason: "Supporting Cast private podcast",
193 },
194 // RedCircle private/exclusive feeds.
195 KnownProvider {
196 host_contains: "redcircle.com",
197 marker: Some("private"),
198 reason: "RedCircle private podcast",
199 },
200 // Private/tokened Megaphone, Acast+, and Omny feeds carry an access token.
201 KnownProvider {
202 host_contains: "megaphone.fm",
203 marker: Some("token"),
204 reason: "Megaphone private podcast",
205 },
206 KnownProvider {
207 host_contains: "acast.com",
208 marker: Some("token"),
209 reason: "Acast+ private podcast",
210 },
211 KnownProvider {
212 host_contains: "omny.fm",
213 marker: Some("token"),
214 reason: "Omny private podcast",
215 },
216 // Apple / Spotify subscriber podcast feeds carry a per-listener token.
217 KnownProvider {
218 host_contains: "podcasts.apple.com",
219 marker: Some("token"),
220 reason: "Apple subscriber podcast",
221 },
222 KnownProvider {
223 host_contains: "spotify.com",
224 marker: Some("token"),
225 reason: "Spotify subscriber podcast",
226 },
227];
228
229/// Whether a URL may be **stored or published** as a feed URL at all.
230///
231/// This is the storage-side twin of the scheme check `net::check_scheme` applies
232/// before fetching. The fetch side has always been safe, because nothing can
233/// reach the network except through `net.rs` — but "safe to fetch" and "safe to
234/// write down" are different questions, and only the first had an answer.
235///
236/// Two paths took a URL from outside and stored it with no validation at all:
237/// `resolve_subscriptions` (any atproto client can write a subscription record
238/// into a user's repo) and the OPML import (`xmlUrl` is whatever the file says).
239/// `classify_feed_privacy` does not cover this — it deliberately returns
240/// `Public` for an unparseable URL, on the stated assumption that "the add path
241/// will reject it as malformed regardless", and those two paths are the ones
242/// that never had an add path to do the rejecting.
243///
244/// Note that `javascript:alert(1)` and `file:///etc/passwd` both *parse* cleanly
245/// as URLs, so parsing is not the check — the scheme is.
246pub fn is_storable_feed_url(url: &str, allow_at_uri: bool) -> bool {
247 // **`at://` is checked BEFORE `Url::parse`, because `Url::parse` cannot read
248 // the form that matters.** `at://did:plc:…/…` fails to parse with *invalid
249 // port number* — the colons in the DID are taken as a port separator — while
250 // the handle form `at://alice.example.com/…` parses fine. So adding `"at"`
251 // to the `matches!` below would appear to work and silently reject every
252 // DID-based at-URI, which is all of them in practice.
253 if let Some(rest) = crate::atproto::strip_at_prefix(url) {
254 // **Recognised case-insensitively, stored canonically.** Schemes are
255 // case-insensitive, so `At://` names the same publication — but
256 // `feeds.url` is UNIQUE, so accepting both spellings is two rows for
257 // one publication, the hazard the canonical-handle rule exists for.
258 // Recognising it here rather than letting it fall through to the
259 // generic checks is what keeps it out of the poller: nothing can fetch
260 // it under any spelling.
261 if !url.starts_with(crate::atproto::AT_URI_PREFIX) {
262 return false;
263 }
264 // **This gates STORING only — polling is handled by exclusion**, by
265 // kind rather than by any re-description of the URL, and
266 // `FeedKind::POLLABLE` is the one place the why is written down.
267 return allow_at_uri && is_storable_publication_uri(rest);
268 }
269 match Url::parse(url) {
270 Ok(u) => {
271 // No host check: for http(s) the `url` crate refuses every hostless
272 // spelling at parse (`http://`, `https://?q`, `http:///` are all
273 // "empty host") and turns `https:///x` into host `x`. A conjunct
274 // requiring a non-empty host was unreachable — a test hunt listed
275 // it as untested, and the honest answer was that no input reaches
276 // it. The `Err` arm below is what refuses a hostless URL.
277 matches!(u.scheme(), "http" | "https")
278 }
279 Err(_) => false,
280 }
281}
282
283/// The body of an `at://` URI — `<did-or-handle>/<collection>/<rkey>` — judged
284/// as a **storable feed**.
285///
286/// An allowlist entry, not a loosening: exactly one foreign collection is
287/// accepted, `site.standard.publication`. The two paths this guard exists for
288/// (`resolve_subscriptions`, the OPML import) take records written by any
289/// atproto client, so "it is an at-URI" is not a reason to store it — only "it
290/// is a publication this reader knows how to poll" is.
291fn is_storable_publication_uri(rest: &str) -> bool {
292 let mut parts = rest.split('/');
293 let (Some(authority), Some(collection), Some(rkey)) =
294 (parts.next(), parts.next(), parts.next())
295 else {
296 return false;
297 };
298 parts.next().is_none()
299 && collection == crate::lexicon::nsid::STANDARD_PUBLICATION
300 // **The rkey is validated against atproto's rules, not a blacklist.**
301 //
302 // A blacklist was the first attempt and it leaked twice: `is_control()`
303 // is Unicode category Cc only, so a bidi override (Cf) passed — and it
304 // reordered both the manage page and `scheduler.rs`'s `%feed.url` log
305 // line. Worse, nothing stopped a query string or fragment living inside
306 // the rkey, which satisfies the three-segment check and is exactly what
307 // `classify_feed_privacy`'s `at://` exemption keys off: a token
308 // smuggled there would have been declared public.
309 //
310 // An allowlist cannot leak the next character class someone finds.
311 // The charset alone still admitted `.`, `..` and a 10 000-byte key;
312 // `is_valid_rkey` carries the length and reserved-name rules too.
313 && crate::atproto::is_valid_rkey(rkey)
314 && is_storable_at_authority(authority)
315}
316
317/// The DID form only. `did:plc:` identifiers are validated by
318/// [`crate::oauth::identity::is_atproto_did`] rather than a `did:` prefix check,
319/// which would accept `did:plc:TOOSHORT`.
320///
321/// **The handle form is not storable, for the reason the canonical-handle rule
322/// already gave:** `feeds.url` is UNIQUE, so `at://alice.example.com/…` beside
323/// `at://did:plc:…/…` is two rows — two sidebar entries, and two polled copies
324/// once the reader is wired — for one publication. A handle is a mutable name
325/// for a DID; the row is keyed on the identity. Resolving a pasted or imported
326/// handle to its DID is the reader's job (#165 already resolves DIDs to their
327/// PDS), and belongs at input, not in storage.
328fn is_storable_at_authority(authority: &str) -> bool {
329 crate::oauth::identity::is_atproto_did(authority)
330}
331
332/// Classify whether a feed URL carries a secret credential in the URL itself.
333///
334/// Returns [`FeedPrivacy::Private`] (with a reason) when the URL looks like it
335/// embeds a token / key / auth credential, else [`FeedPrivacy::Public`].
336///
337/// **Design — provider-agnostic first.** The primary defence is a generic
338/// credential-in-URL heuristic that catches paid feeds from *any* provider, not
339/// just the ones we've named; a secondary known-provider table adds precision
340/// (and a nicer reason) for the common paid newsletters and private podcasts. We
341/// deliberately **bias toward flagging**: a false-positive block of a public feed
342/// is low-harm (the user just can't add that one feed yet), whereas a false
343/// negative would leak a paid secret onto the public network — high-harm.
344///
345/// Detection (any one is sufficient):
346/// 1. **Userinfo** — `https://user:pass@host/…` embeds credentials directly.
347/// 2. **Known private-feed path markers** — `PRIVATE_PATH_MARKERS`
348/// (`/feed/private/`, `/members/`, `/subscriber/`, …).
349/// 3. **Credential query parameters** — a query key in `SECRET_QUERY_KEYS` with
350/// a long/opaque value (Patreon `?auth=`, Ghost `?uuid=`, `?token=`, …).
351/// 4. **High-entropy opaque token segments** — a long opaque blob (hex ≥ 16,
352/// base64url ≥ 16, or a UUID) anywhere in the path or a query value, even
353/// without a telltale name.
354/// 5. **Known providers** — `KNOWN_PROVIDERS` host (+ optional marker) match.
355///
356/// An unparseable URL is treated as [`FeedPrivacy::Public`]: the add path rejects
357/// a malformed URL downstream anyway, and we don't want a parse quirk to
358/// misclassify.
359pub fn classify_feed_privacy(url: &str) -> FeedPrivacy {
360 // **`at://` is classified deliberately, and NOT doing so refused real
361 // subscriptions.** An atproto rkey is a TID — 13 base32-sortable characters
362 // — which is exactly what the generic "high-entropy token in path"
363 // heuristic below is looking for. Measured: without this arm,
364 // `at://did:plc:…/site.standard.publication/3lab2c4d5e6f7g8h` is
365 // classified PRIVATE and the subscription refused, while the DID form slips
366 // through only because it fails to parse as a `Url` at all.
367 //
368 // `Public` is the right answer: a publication is a public record in a
369 // public repo and the rkey is a handle, not a secret, so there is no
370 // private/paid shape for this scheme to carry.
371 if let Some(rest) = crate::atproto::strip_at_prefix(url) {
372 // A non-canonical scheme spelling is an at-URI this reader will not
373 // store, not an unparseable string for the `Err(_) => Public` arm below
374 // to wave through. Fail closed.
375 if !url.starts_with(crate::atproto::AT_URI_PREFIX) {
376 return FeedPrivacy::Private("non-canonical at:// scheme spelling".to_string());
377 }
378 // **Only a WELL-FORMED publication URI is exempt.** The first version
379 // of this was a bare prefix match, which declared any attacker-chosen
380 // string starting `at://` safe to publish — skipping the userinfo
381 // check, the known-provider table, the private-path markers, the
382 // secret-query keys and the entropy heuristics all at once. That is a
383 // regression against every one of them, on a path
384 // (`rename_subscription`) where this function is the only gate and the
385 // value is written to the user's PUBLIC repo.
386 //
387 // Anything else falls through to the generic checks below, which is
388 // where a credential-bearing string belongs.
389 if is_storable_publication_uri(rest) {
390 return FeedPrivacy::Public;
391 }
392 // **Malformed `at://` is REFUSED, not passed through.** Falling through
393 // reaches `Url::parse`, which fails on the DID form and lands on the
394 // `Err(_) => Public` arm below — whose justification is "the add path
395 // will reject it as a malformed URL regardless".
396 //
397 // **Defence in depth, not a live gate.** This comment used to say the
398 // justification is false on the `rename_subscription` path, "where this
399 // function is the only gate". That stopped being true when storability
400 // moved ahead of privacy on the repoint: a review then found no
401 // production caller can reach this arm at all — add pre-checks the
402 // `at://` prefix, rename and OPML and `resolve_subscriptions` all
403 // refuse a non-storable URL first. It stays because a fail-closed
404 // branch is worth its keep for the next caller that arrives without
405 // one, and because deleting a guard on the grounds that nothing
406 // currently reaches it is how the next one gets it wrong. It is not
407 // load-bearing today, and saying so is the honest version.
408 return FeedPrivacy::Private("not a well-formed at:// publication URI".to_string());
409 }
410 let parsed = match Url::parse(url) {
411 Ok(u) => u,
412 // Can't parse => the add path will reject it as a malformed URL regardless.
413 Err(_) => return FeedPrivacy::Public,
414 };
415
416 // (1) Userinfo (`https://user:pass@host/…`) — credentials in the authority.
417 if !parsed.username().is_empty() || parsed.password().is_some() {
418 return FeedPrivacy::Private("credentials in URL userinfo".to_string());
419 }
420
421 let path_lower = parsed.path().to_ascii_lowercase();
422 let query_lower = parsed.query().unwrap_or("").to_ascii_lowercase();
423 let host_lower = parsed.host_str().unwrap_or("").to_ascii_lowercase();
424
425 // (0) Public-feed allowlist. A handful of large, fully-public feed shapes
426 // carry a high-entropy-looking id in the query that would otherwise trip the
427 // generic entropy heuristic. YouTube channel/playlist RSS
428 // (`youtube.com/feeds/videos.xml?channel_id=UC…` / `?playlist_id=PL…`) is the
429 // canonical way any reader subscribes to a channel — the id is a PUBLIC
430 // handle, not a secret. Allowlist it before the generic checks so we don't
431 // false-block it. (Userinfo / known-provider markers are checked below and
432 // still apply, so this can't be used to smuggle a credential.)
433 if is_public_youtube_feed(&host_lower, &path_lower, &parsed) {
434 return FeedPrivacy::Public;
435 }
436
437 // (5) Known-provider precision layer (checked early so its specific reason
438 // wins over a generic one). Host substring + optional path/query marker.
439 for kp in KNOWN_PROVIDERS {
440 if host_lower.contains(kp.host_contains) {
441 let marker_ok = match kp.marker {
442 None => true,
443 Some(m) => {
444 let m = m.to_ascii_lowercase();
445 path_lower.contains(&m) || query_lower.contains(&m)
446 }
447 };
448 if marker_ok {
449 return FeedPrivacy::Private(kp.reason.to_string());
450 }
451 }
452 }
453
454 // (2) Known private-feed path markers.
455 for marker in PRIVATE_PATH_MARKERS {
456 if path_lower.contains(marker) {
457 return FeedPrivacy::Private(format!("private feed path `{marker}`"));
458 }
459 }
460
461 // (3) Credential query parameters with a long/opaque value.
462 for (k, v) in parsed.query_pairs() {
463 let key = k.as_ref().to_ascii_lowercase();
464 if SECRET_QUERY_KEYS.iter().any(|sk| *sk == key) && value_is_opaque(v.as_ref()) {
465 return FeedPrivacy::Private(format!("credential query parameter `{key}`"));
466 }
467 }
468
469 // (4) High-entropy opaque token segments (an embedded key/token with no
470 // telltale name): hex ≥ 16, base64url ≥ 16, or a UUID, in path or query.
471 // The dominant real-world private-podcast shape delivers the token as a
472 // *filename* (`<token>.rss` / `<token>.xml`) or affixed inside a larger
473 // segment (`feed-<uuid>`), so [`segment_hides_secret`] strips a trailing feed
474 // extension AND scans dot/underscore/hyphen-delimited sub-parts, not just the
475 // whole segment.
476 for seg in parsed.path().split('/').filter(|s| !s.is_empty()) {
477 if segment_hides_secret(seg) {
478 return FeedPrivacy::Private("high-entropy token in path".to_string());
479 }
480 }
481 for (_, v) in parsed.query_pairs() {
482 if looks_like_embedded_secret(v.as_ref()) {
483 return FeedPrivacy::Private("high-entropy token in query".to_string());
484 }
485 }
486
487 FeedPrivacy::Public
488}
489
490/// Whether a *named* credential query value (`?token=<v>`) is long/opaque enough
491/// to count as a secret. A short value (e.g. an enum like `?token=none`) is not.
492/// We treat a UUID, or anything ≥ 8 chars that isn't an obvious plain word, as
493/// opaque — named credential keys already signal intent, so the length bar is
494/// low.
495fn value_is_opaque(v: &str) -> bool {
496 if v.is_empty() {
497 return false;
498 }
499 if is_uuid(v) {
500 return true;
501 }
502 v.len() >= 8
503}
504
505/// Heuristic: does `s` look like an embedded secret (an opaque high-entropy
506/// token), as opposed to an ordinary slug or word? Matches a UUID, a hex string
507/// ≥ 16 chars, or a base64url-ish blob ≥ 16 chars that mixes letters and digits
508/// and isn't a hyphen/dot slug. Deliberately strict so it only fires on things
509/// that really look like keys — the named-marker and known-provider checks cover
510/// the rest.
511fn looks_like_embedded_secret(s: &str) -> bool {
512 if is_uuid(s) {
513 return true;
514 }
515 // Hex string ≥ 16 chars (e.g. a 32-char MD5-ish token).
516 if s.len() >= 16 && s.chars().all(|c| c.is_ascii_hexdigit()) {
517 return true;
518 }
519 // base64url-ish opaque blob ≥ 16 chars.
520 if s.len() < 16 {
521 return false;
522 }
523 // Hyphen/dot-heavy slugs (`this-is-a-normal-post-title`) are not secrets.
524 let separators = s
525 .bytes()
526 .filter(|b| *b == b'-' || *b == b'.' || *b == b' ')
527 .count();
528 if separators >= 3 {
529 return false;
530 }
531 // Must be plausibly token-charset: base64url alphabet only. `=` is accepted
532 // as base64 padding (it only ever appears trailing on a real blob, so a
533 // padded base64url token like `…dnc=` still counts).
534 let token_chars = s
535 .chars()
536 .filter(|c| c.is_ascii_alphanumeric() || matches!(c, '_' | '-' | '='))
537 .count();
538 if token_chars < s.chars().count() {
539 return false;
540 }
541 let has_alpha = s.chars().any(|c| c.is_ascii_alphabetic());
542 let has_digit = s.chars().any(|c| c.is_ascii_digit());
543 if !(has_alpha && has_digit) {
544 // A token almost always mixes letters and digits; a pure-alpha long
545 // segment is far more likely to be a normal (if long) slug/word.
546 return false;
547 }
548 // Distinct-character ratio: real tokens use most of the alphabet, words
549 // repeat a small set. Require >= 10 distinct chars for a 16+ char blob.
550 let mut seen = std::collections::HashSet::new();
551 for c in s.chars() {
552 seen.insert(c.to_ascii_lowercase());
553 }
554 seen.len() >= 10
555}
556
557/// Known feed/file extensions a token filename may wear (`<token>.rss`,
558/// `<token>.xml`, …). Stripped before the whole-segment secret test so a
559/// tokened *filename* — the dominant private-podcast URL shape — is still caught.
560const FEED_EXTENSIONS: &[&str] = &["rss", "xml", "atom", "json", "rss20"];
561
562/// Whether a single path segment hides an embedded secret. Beyond the plain
563/// whole-segment [`looks_like_embedded_secret`] test, this also catches the two
564/// real-world shapes that wrap a token so the whole segment is no longer a clean
565/// blob:
566///
567/// 1. **Token-as-filename** — `<token>.rss` / `<token>.xml`: strip a trailing
568/// feed extension and re-test the stem.
569/// 2. **Token affixed inside a larger segment** — `feed-<uuid>`, `<token>.xml`,
570/// `pod_<hex32>`: split on `.`/`_`/`-` and test each sub-part, so a
571/// high-entropy blob delimited by an affix is still found.
572fn segment_hides_secret(seg: &str) -> bool {
573 if looks_like_embedded_secret(seg) {
574 return true;
575 }
576 // (1) Strip a trailing known feed extension and re-test the stem.
577 if let Some((stem, ext)) = seg.rsplit_once('.') {
578 if FEED_EXTENSIONS.contains(&ext.to_ascii_lowercase().as_str())
579 && looks_like_embedded_secret(stem)
580 {
581 return true;
582 }
583 }
584 // (2) A UUID embedded with an affix (`feed-<uuid>`, `<uuid>-audio`) — the
585 // `-` delimiters inside the UUID mean a naive split can't see it, so scan for
586 // a canonical UUID substring directly.
587 if contains_uuid(seg) {
588 return true;
589 }
590 // (3) Scan `.`/`_`/`-`-delimited sub-parts for a high-entropy blob affixed to
591 // an ordinary word (`pod_<hex32>`, `<hex32>.mp3`). Only fires on multi-part
592 // segments (a single-part segment was already covered by the whole-segment
593 // test above), so a plain `my-normal-post-slug` — whose parts are short
594 // dictionary words — can't trip it.
595 if seg.contains(['.', '_', '-']) {
596 for part in seg.split(['.', '_', '-']).filter(|p| !p.is_empty()) {
597 if looks_like_embedded_secret(part) {
598 return true;
599 }
600 }
601 }
602 false
603}
604
605/// Whether `s` contains a canonical 8-4-4-4-12 UUID as a substring (allowing an
606/// affix on either side, e.g. `feed-<uuid>` or `<uuid>-audio`). Slides a 36-char
607/// window over the string and tests each with [`is_uuid`].
608fn contains_uuid(s: &str) -> bool {
609 const UUID_LEN: usize = 36; // 8+4+4+4+12 + 4 hyphens.
610 let bytes = s.as_bytes();
611 if bytes.len() < UUID_LEN {
612 return false;
613 }
614 // ASCII-only window: a UUID is pure ASCII hex/hyphen, so byte indexing is
615 // safe here (a multi-byte char in the window just fails is_uuid).
616 (0..=bytes.len() - UUID_LEN).any(|i| s.get(i..i + UUID_LEN).map(is_uuid).unwrap_or(false))
617}
618
619/// Whether this is a PUBLIC YouTube channel/playlist RSS feed
620/// (`www.youtube.com/feeds/videos.xml?channel_id=UC…` or `?playlist_id=PL…`).
621/// The channel/playlist id is a public handle, not a credential, so these feeds
622/// must NOT be flagged by the generic entropy heuristic. We require the exact
623/// public host + feeds path + one of the two public id keys, so this narrow
624/// allowlist can't be abused to smuggle a `?token=` past classification.
625fn is_public_youtube_feed(host_lower: &str, path_lower: &str, parsed: &Url) -> bool {
626 let host_ok = host_lower == "youtube.com"
627 || host_lower == "www.youtube.com"
628 || host_lower.ends_with(".youtube.com");
629 if !host_ok || !path_lower.starts_with("/feeds/videos.xml") {
630 return false;
631 }
632 // Only the public id keys may appear; a `token`/`auth`/… key means treat it
633 // as a normal (potentially private) URL and let the checks below run.
634 parsed.query_pairs().all(|(k, _)| {
635 let k = k.as_ref().to_ascii_lowercase();
636 k == "channel_id" || k == "playlist_id" || k == "user"
637 })
638}
639
640/// Whether `s` is a canonical 8-4-4-4-12 hyphenated UUID (any hex case).
641fn is_uuid(s: &str) -> bool {
642 let groups = [8usize, 4, 4, 4, 12];
643 let parts: Vec<&str> = s.split('-').collect();
644 if parts.len() != groups.len() {
645 return false;
646 }
647 parts
648 .iter()
649 .zip(groups.iter())
650 .all(|(p, &n)| p.len() == n && p.chars().all(|c| c.is_ascii_hexdigit()))
651}
652
653/// How long a single feed fetch may take before we give up.
654const FETCH_TIMEOUT: Duration = Duration::from_secs(30);
655
656/// Per-read idle timeout: cap the wait for the *next* body chunk, so a server
657/// that trickles bytes forever can't tie up a fetch under the total timeout.
658const READ_TIMEOUT: Duration = Duration::from_secs(15);
659
660/// Base backoff applied after a failed poll; the caller multiplies this by the
661/// feed's consecutive-error count (with a ceiling) to space out retries.
662const BACKOFF_BASE: Duration = Duration::from_secs(300);
663
664/// Ceiling on backoff so a persistently broken feed still gets retried daily.
665const BACKOFF_MAX: Duration = Duration::from_secs(24 * 3600);
666
667/// The outcome of polling a single feed. Lets the scheduler decide how to
668/// reschedule (and lets tests assert what happened) without inspecting the DB.
669#[derive(Debug, Clone, PartialEq, Eq)]
670pub enum PollOutcome {
671 /// The feed was fetched, parsed, and stored. `new_entries` is the number of
672 /// entries inserted or updated by this poll.
673 Updated { new_entries: u64 },
674 /// The server returned `304 Not Modified` — nothing changed, nothing stored.
675 NotModified,
676 /// The feed was fetched, but **this instance** could not sanitize its
677 /// bodies now: every sanitize permit stayed held — by sanitizes other
678 /// feeds' polls gave up on — for the whole timeout (#226). Not the feed's
679 /// fault, so nothing is stored and nothing counts against it: no entries,
680 /// no validators, no error. It is polled again on its normal cadence.
681 Deferred,
682 /// The fetch or parse failed; the feed was left intact and skipped. Carries
683 /// the suggested backoff before the next attempt. Never a panic.
684 ///
685 /// **`kind` and `detail` are the reason, and they exist because their
686 /// absence cost a production investigation.** Until #159 the error was
687 /// logged here and discarded, so `feeds` recorded that a feed was failing
688 /// and never why — which is how sixty feeds broken by our own 304 handling
689 /// looked exactly like sixty dead blogs. `kind` is a small closed
690 /// vocabulary so failures can be counted by cause; `detail` is the message
691 /// for a human reading one row.
692 Failed {
693 backoff: Duration,
694 kind: FailureKind,
695 detail: String,
696 },
697}
698
699/// Why a poll failed, as a closed set.
700///
701/// Closed on purpose: the point is to *count* failures by cause, and a free-text
702/// kind cannot be counted. The detail string carries whatever else matters.
703#[derive(Debug, Clone, Copy, PartialEq, Eq)]
704pub enum FailureKind {
705 /// The request never produced a response — DNS, TLS, timeout, connection
706 /// refused, or a refusal by the SSRF guard.
707 Fetch,
708 /// A response arrived with a non-success status — or, from an atproto
709 /// PDS, a 2xx carrying an error envelope, which is the same refusal said
710 /// in the body.
711 Status,
712 /// The body was too large, or reading it failed part-way. For a
713 /// publication this includes a listing too large to walk: more pages,
714 /// bytes or records than the walk allows.
715 Body,
716 /// The body arrived and is not a feed this parser can read.
717 Parse,
718}
719
720/// What the poller does with a feed row.
721///
722/// **A column, not a predicate.** "Can this be fetched?" used to be
723/// `substr(url, 1, 5) = 'at://'` spliced into four statements, with a fifth
724/// reader that had already drifted from them. Deciding it once, in Rust, at
725/// insert — and storing the answer — means SQL cannot disagree with the
726/// fetcher, and wiring the standard.site reader becomes a change to this
727/// function plus a dispatch, rather than an edit to every statement that
728/// mentions a URL.
729#[derive(Debug, Clone, Copy, PartialEq, Eq)]
730pub enum FeedKind {
731 /// An RSS/Atom/JSON feed document fetched over http(s).
732 Rss,
733 /// An `at://…/site.standard.publication/…` record pair in somebody's PDS,
734 /// read by `standard_site` and pollable since 0.4.0.
735 Publication,
736 /// Any other `at://` row: another collection, or a spelling the storage
737 /// guard would refuse. Rows like this predate that guard. **Never
738 /// pollable** — handing one to the standard.site reader would fail its
739 /// collection check every tick and publish it as an unreachable publisher —
740 /// and counted as unpollable so the capacity it holds stays visible.
741 Unsupported,
742}
743
744impl FeedKind {
745 /// The kinds the scheduler may select: RSS, and since 0.4.0 standard.site
746 /// publications. [`FeedKind::Unsupported`] is the kind it never selects.
747 ///
748 /// **This is the canonical home of the exclusion; the other sites point
749 /// here.** It used to be a SQL string predicate in `store`, carrying
750 /// its own copy of the rule — which is how one reader (`count_feeds`) came
751 /// to drift from it unnoticed.
752 ///
753 /// Why an unpollable kind is skipped rather than failed: an `Unsupported`
754 /// row names nothing either reader can fetch — another collection, a handle,
755 /// a non-canonical spelling — so polling it could only fail.
756 /// Handing such a row to the poller does not leave the feature dormant — it
757 /// manufactures one permanent failure per row, which the public cause
758 /// histogram then reports as an unreachable publisher. Unsupported is not
759 /// broken, and telling those apart is the entire reason a failure cause is
760 /// recorded. (Rows like this exist: subscriptions written by other clients
761 /// before this reader refused the scheme.)
762 ///
763 /// **`store::count_feeds` is deliberately NOT filtered by this.** The
764 /// global ceiling bounds storage on a small box, and an unpollable row
765 /// occupies a row, so it counts against the cap. An earlier version of this
766 /// note listed the readers without naming the exception, which read as
767 /// completeness it did not have: a review found the ceiling consuming
768 /// capacity that appeared on no surface, since `/stats` measures the poller
769 /// and excludes these rows. `/admin/metrics` renders
770 /// `store::unpollable_feeds` for exactly that reason.
771 /// **Adding a kind here makes a population of rows due all at once.**
772 /// `store::due_feeds` orders `next_poll IS NOT NULL, next_poll ASC`, so a
773 /// row with no scheduled poll sorts ahead of every dated one. Rows that
774 /// were never pollable have no schedule, so the boot that reclassifies
775 /// them hands the poller a block of N rows that outrank every regular
776 /// feed until they drain — `ceil(N / batch)` ticks, measured, during which
777 /// `/stats` shows a climbing backlog and nothing logs why. Bounded and
778 /// harmless at ninety feeds; not at ten thousand. Whoever wires the next
779 /// kind should seed or stagger `next_poll` for the rows it admits.
780 pub const POLLABLE: &'static [FeedKind] = &[FeedKind::Rss, FeedKind::Publication];
781
782 /// The kinds whose entries the retention **window** applies to.
783 ///
784 /// **Age is the wrong retention policy for an archive, and that is a
785 /// measurement, not a preference.** Three real publications, read through
786 /// `standard_site::fetch` on 2026-09-27: Standard.site's newest document was
787 /// **131 days** old, Annotated's **109** (with its oldest at 373), and minus
788 /// listens' **241**. Against the instance default window of 14 days, every
789 /// one of them stored **zero** rows — a green poll, an empty feed, and an
790 /// info log as the only trace. Long-form publishing is not news-paced.
791 ///
792 /// So a publication is bounded by COUNT instead: `max_entries_per_feed` in
793 /// `insert_entries`, which caps a feed at the newest N plus up to N starred.
794 /// That is a real bound — it is what keeps this from being "retention off" —
795 /// and it is the one that suits a source whose value is its archive.
796 ///
797 /// A generous absolute ceiling still applies (`publication_retention_days`),
798 /// because "not aged out" must not mean "immortal": rows belonging to a feed
799 /// nobody polls any more would otherwise never be reaped at all, and the
800 /// per-feed trim only runs when a poll stores something.
801 pub const AGED: &'static [FeedKind] = &[FeedKind::Rss];
802
803 /// The column value. Stable — it is persisted.
804 pub fn as_str(self) -> &'static str {
805 match self {
806 FeedKind::Rss => "rss",
807 FeedKind::Publication => "publication",
808 FeedKind::Unsupported => "unsupported",
809 }
810 }
811
812 /// A closed vocabulary on the way back in: a kind written by a newer build
813 /// is not silently read as one this build knows.
814 pub fn parse(raw: &str) -> Option<Self> {
815 match raw {
816 "rss" => Some(FeedKind::Rss),
817 "publication" => Some(FeedKind::Publication),
818 "unsupported" => Some(FeedKind::Unsupported),
819 _ => None,
820 }
821 }
822
823 /// What a URL will be stored as. The only place the question is asked.
824 ///
825 /// `store::feeds.kind` is a cache of this function, re-derived from the URL
826 /// on every start and on every upsert, so a change here needs no migration
827 /// and SQL cannot hold an opinion of its own about what a row is.
828 ///
829 /// **Case-insensitive on purpose.** An earlier version was case-sensitive,
830 /// with a test pinning that a mixed-case `At://` row IS handed to the
831 /// poller — reasoning that if Rust does not call it an at-URI, neither
832 /// should anything else. That was wrong in the direction that matters: URL
833 /// schemes are case-insensitive, so `Url::parse` folds `At://` to scheme
834 /// `at` and `net::check_scheme` refuses it (the DID form does not parse at
835 /// all). The row could only fail, every tick, forever, and be published in
836 /// the `fetch` bucket as an unreachable publisher. Storing a non-canonical
837 /// spelling is separately refused, because `feeds.url` is UNIQUE.
838 pub fn of(url: &str) -> Self {
839 match crate::atproto::strip_at_prefix(url) {
840 // **Only a URI the storage guard would accept is a publication.**
841 // Since publications became pollable, "it is an at-URI" is no longer
842 // a safe enough reason: another collection would be polled, fail,
843 // and back off forever while reading as an unreachable publisher.
844 // The canonical lowercase scheme too, for the same reason: storage
845 // refuses `At://` (#183), and a legacy row spelled that way would
846 // otherwise be polled and fail its parse every tick.
847 Some(rest) if url.starts_with("at://") && is_storable_publication_uri(rest) => {
848 FeedKind::Publication
849 }
850 Some(_) => FeedKind::Unsupported,
851 None => FeedKind::Rss,
852 }
853 }
854}
855
856/// Apply a [`PollOutcome`] to the feed's row: settle the error columns AND
857/// reschedule it. **Both halves, always, from one place.**
858///
859/// The scheduler did this inline. `web::add_subscription` then copied only the
860/// first half, so a poll taken off the scheduler could clear a stale failure's
861/// COUNT while leaving the feed parked on its stale backoff HORIZON — reported
862/// healthy, not polled for up to 24h. And its failures fed `consecutive_errors`
863/// with no reschedule, so repeated Subscribe clicks drove a shared feed to the
864/// 24h ceiling for every subscriber. Two copies of a sequence drift; this is
865/// the one copy.
866///
867/// `cadence` is the interval to use on success; the scheduler derives it from
868/// the feed's hint, a direct caller passes the configured default. A store
869/// failure is logged and the reschedule still attempted, so a hiccup writing
870/// the count cannot strand the feed at a NULL `next_poll` that `due_feeds`
871/// would then re-poll every tick.
872pub async fn settle_poll(
873 pool: &sqlx::SqlitePool,
874 url: &str,
875 outcome: &PollOutcome,
876 cadence: Duration,
877) {
878 let delay = match outcome {
879 // A 304 is a healthy poll: it proves the fetch worked and nothing changed.
880 PollOutcome::Updated { .. } | PollOutcome::NotModified => {
881 if let Err(err) = crate::store::reset_feed_errors(pool, url).await {
882 tracing::warn!(feed = %url, %err, "failed to reset feed error count");
883 }
884 cadence
885 }
886 // Neither a success nor the feed's failure: leave its error count as
887 // it is, and keep its ordinary schedule.
888 PollOutcome::Deferred => cadence,
889 PollOutcome::Failed {
890 backoff,
891 kind,
892 detail,
893 } => {
894 // Recompute from the feed's REAL consecutive-error count so a
895 // persistently-broken feed climbs toward the ceiling instead of
896 // retrying at the floor forever; fall back to the outcome's floor.
897 match crate::store::bump_feed_errors(pool, url, *kind, detail).await {
898 Ok(count) => backoff_for(count.max(1) as u32),
899 Err(err) => {
900 tracing::warn!(feed = %url, %err, "failed to bump feed error count; using floor backoff");
901 *backoff
902 }
903 }
904 }
905 };
906 if let Err(err) = crate::store::set_next_poll(pool, url, delay).await {
907 tracing::error!(feed = %url, %err, "failed to persist next_poll");
908 }
909}
910
911/// Cap on a failure detail, applied where the string is BUILT.
912///
913/// It was originally applied only inside `store::bump_feed_errors`, which
914/// bounded the database row and nothing else — and widening
915/// [`PollOutcome::Failed`] with this field had quietly opened a second sink:
916/// `web::add_subscription` logs `?outcome` at INFO on a user-facing request
917/// path. Bounding at construction bounds every sink, including ones added
918/// later by someone who never reads this comment.
919pub const MAX_FAILURE_DETAIL_CHARS: usize = 300;
920
921/// Render an error chain into a bounded [`PollOutcome::Failed`] detail.
922///
923/// `{e:#}` — the anyhow CHAIN, not just the outermost context. "fetching
924/// https://…" alone says nothing; the cause is the part that would have named
925/// the 304 bug in #159.
926pub fn failure_detail(err: impl std::fmt::Display) -> String {
927 let s = err.to_string();
928 if s.chars().count() <= MAX_FAILURE_DETAIL_CHARS {
929 return s;
930 }
931 s.chars().take(MAX_FAILURE_DETAIL_CHARS).collect()
932}
933
934impl FailureKind {
935 /// The stable string stored in `feeds.last_error_kind` and aggregated on
936 /// `/stats`. Changing one of these silently rewrites history in the
937 /// aggregate, so they are spelled out rather than derived from the variant.
938 pub fn as_str(self) -> &'static str {
939 match self {
940 Self::Fetch => "fetch",
941 Self::Status => "status",
942 Self::Body => "body",
943 Self::Parse => "parse",
944 }
945 }
946
947 /// Read back a persisted `last_error_kind`. `None` for anything this
948 /// version does not know, so a row written by a newer build is not
949 /// silently attributed to a cause this one recognises — the same contract
950 /// [`crate::metrics::Backend::parse`] keeps for the same reason.
951 pub fn parse(raw: &str) -> Option<Self> {
952 match raw {
953 "fetch" => Some(Self::Fetch),
954 "status" => Some(Self::Status),
955 "body" => Some(Self::Body),
956 "parse" => Some(Self::Parse),
957 _ => None,
958 }
959 }
960
961 /// Every variant, so a test can assert over the whole set rather than a
962 /// list that drifts when a variant is added.
963 pub const ALL: [Self; 4] = [Self::Fetch, Self::Status, Self::Body, Self::Parse];
964}
965
966/// Build a `reqwest::Client` configured for polite **and safe** feed fetching.
967///
968/// Callers should build this **once** and share it (connection pooling), then
969/// hand a reference to [`poll_feed`]. Kept here so the fetch policy (UA,
970/// timeout, redirect behaviour) lives with the code that depends on it.
971///
972/// Auto-redirect is **disabled** on purpose: feed URLs are untrusted, so
973/// redirects are followed manually by [`crate::net::guarded_get`], which
974/// re-validates the scheme + resolved IP of every hop (SSRF defence). A client
975/// that silently followed redirects could be bounced onto `169.254.169.254` or
976/// `127.0.0.1` between the guard's check and the connect.
977pub fn build_client() -> Result<Client> {
978 Client::builder()
979 .user_agent(crate::USER_AGENT)
980 .timeout(FETCH_TIMEOUT)
981 .read_timeout(READ_TIMEOUT)
982 // Ignore ambient proxy configuration, for the same reason the pinned
983 // client does: a proxied request hands the hostname to the proxy to
984 // resolve, so `net`'s IP checks never see the address they are meant to
985 // vet. See `net::build_pinned_client`.
986 .no_proxy()
987 // No auto-redirect: net::guarded_get follows + re-validates each hop.
988 .redirect(reqwest::redirect::Policy::none())
989 .build()
990 .context("failed to build feed HTTP client")
991}
992
993/// Compute the backoff for the `n`th consecutive failure (1-based), clamped to
994/// `BACKOFF_MAX`. Exponential in the error count so transient blips retry soon
995/// while a durably-broken feed backs off toward daily.
996///
997/// The scheduler passes the feed's persisted `consecutive_errors` count (see
998/// [`crate::store::bump_feed_errors`]) so a feed that keeps failing actually
999/// climbs toward `BACKOFF_MAX` instead of retrying at the floor forever.
1000pub fn backoff_for(consecutive_errors: u32) -> Duration {
1001 let n = consecutive_errors.max(1);
1002 // Saturating shift: base * 2^(n-1), capped. Avoids overflow for large n.
1003 let factor = 1u64.checked_shl(n.saturating_sub(1)).unwrap_or(u64::MAX);
1004 let secs = BACKOFF_BASE
1005 .as_secs()
1006 .saturating_mul(factor)
1007 .min(BACKOFF_MAX.as_secs());
1008 Duration::from_secs(secs)
1009}
1010
1011/// Poll one feed by **what it is**: an RSS/Atom/JSON document over HTTP, or a
1012/// standard.site publication read from its author's PDS.
1013///
1014/// The single entry point the scheduler calls, so the choice of reader lives
1015/// with [`FeedKind`] and not in the poll loop. Same contract as [`poll_feed`]:
1016/// `Err` is a broken local store, never a misbehaving source.
1017pub async fn poll_feed_by_kind(
1018 pool: &SqlitePool,
1019 client: &Client,
1020 config: &crate::config::Config,
1021 feed: &Feed,
1022) -> Result<PollOutcome> {
1023 match FeedKind::of(&feed.url) {
1024 FeedKind::Rss => poll_feed(pool, client, feed, config.max_entries_per_feed).await,
1025 FeedKind::Publication => poll_publication(pool, client, config, feed).await,
1026 // Never selected: `due_feeds` reads only POLLABLE kinds. A defensive
1027 // failure rather than a panic if a caller hands one in anyway.
1028 FeedKind::Unsupported => Ok(PollOutcome::Failed {
1029 backoff: backoff_for(1),
1030 kind: FailureKind::Parse,
1031 detail: failure_detail(format!("{} is not a feed this reader can poll", feed.url)),
1032 }),
1033 }
1034}
1035
1036/// What a failed publication read is filed under in the cause histogram.
1037///
1038/// **`Fetch` means the request never produced a response**, so it is the
1039/// fallback, not the default answer: a PLC directory or PDS that answered —
1040/// with a 404 for a tombstoned DID, `RepoNotFound`, `RepoDeactivated` — is
1041/// `Status`, and an answer that was not what it claimed to be is `Parse`. Filed
1042/// as `Fetch`, a deleted account read as its server being down (found in
1043/// review). A listing refusal is typed for the same reason (#227): a 200 with
1044/// no `records` is `Parse`, a 200 error envelope is `Status`, and a listing or
1045/// body over its cap is `Body`.
1046pub(crate) fn publication_failure_kind(err: &anyhow::Error) -> FailureKind {
1047 use crate::atproto::{AtProtoError, DidResolutionCause, ListingTooLarge, UnreadableListing};
1048 for cause in err.chain() {
1049 if cause.is::<crate::standard_site::NotAPublication>() || cause.is::<serde_json::Error>() {
1050 return FailureKind::Parse;
1051 }
1052 if cause.is::<ListingTooLarge>() || cause.is::<crate::net::BodyTooLarge>() {
1053 return FailureKind::Body;
1054 }
1055 match cause.downcast_ref::<UnreadableListing>() {
1056 Some(UnreadableListing::NoRecords) => return FailureKind::Parse,
1057 Some(UnreadableListing::ErrorEnvelope { .. }) => return FailureKind::Status,
1058 None => {}
1059 }
1060 match cause.downcast_ref::<AtProtoError>() {
1061 Some(AtProtoError::Xrpc { .. }) => return FailureKind::Status,
1062 Some(AtProtoError::DidResolution { cause, .. }) => {
1063 return match cause {
1064 DidResolutionCause::Status => FailureKind::Status,
1065 DidResolutionCause::UnsupportedMethod | DidResolutionCause::NoPdsEndpoint => {
1066 FailureKind::Parse
1067 }
1068 DidResolutionCause::NotAPublicTarget => FailureKind::Fetch,
1069 }
1070 }
1071 _ => {}
1072 }
1073 }
1074 FailureKind::Fetch
1075}
1076
1077/// Read a standard.site publication from its author's PDS and store it.
1078///
1079/// A source failure — an unparseable URI, an unreachable PLC directory or
1080/// PDS, a walk that failed — is a [`PollOutcome::Failed`] with backoff, like an
1081/// RSS fetch failure, so `settle_poll` and `/stats` treat both kinds alike.
1082/// Only a broken local store is an `Err`, which `store_publication` decides.
1083async fn poll_publication(
1084 pool: &SqlitePool,
1085 client: &Client,
1086 config: &crate::config::Config,
1087 feed: &Feed,
1088) -> Result<PollOutcome> {
1089 poll_publication_group(pool, client, config, std::slice::from_ref(feed))
1090 .await
1091 .pop()
1092 .unwrap_or_else(|| Err(anyhow::anyhow!("no outcome for {}", feed.url)))
1093}
1094
1095/// Read several publications **from one repo** with one walk of its documents
1096/// (`standard_site::fetch_repo`), and store each. One outcome per feed, in
1097/// order. The caller groups by repo; a feed from another repo, or one whose
1098/// URL is not a publication URI, gets its own failure and does not affect the
1099/// rest.
1100///
1101/// **Why one walk:** cost scales with the repo, not the publication. Nine
1102/// publications in one repo were nine full walks of its documents.
1103///
1104/// A member the shared walk cut short is then re-read alone, within the same
1105/// read deadline (#229); see the comment above the re-read loop.
1106pub async fn poll_publication_group(
1107 pool: &SqlitePool,
1108 client: &Client,
1109 config: &crate::config::Config,
1110 feeds: &[Feed],
1111) -> Vec<Result<PollOutcome>> {
1112 poll_publication_group_with(
1113 pool,
1114 client,
1115 config,
1116 feeds,
1117 crate::atproto::MAX_LARGE_RECORDS,
1118 crate::atproto::MAX_LIST_BYTES,
1119 std::time::Instant::now(),
1120 )
1121 .await
1122}
1123
1124/// [`poll_publication_group`] with the walk's limits and the deadline's start
1125/// passed in, so the re-read of starved members is testable without 128 MiB
1126/// of documents or a clock that really runs out. Production passes
1127/// `MAX_LARGE_RECORDS`, `MAX_LIST_BYTES` and `Instant::now()`.
1128pub(crate) async fn poll_publication_group_with(
1129 pool: &SqlitePool,
1130 client: &Client,
1131 config: &crate::config::Config,
1132 feeds: &[Feed],
1133 per_site_cap: usize,
1134 budget_bytes: usize,
1135 deadline_started: std::time::Instant,
1136) -> Vec<Result<PollOutcome>> {
1137 let failed = |kind: FailureKind, detail: String| -> Result<PollOutcome> {
1138 Ok(PollOutcome::Failed {
1139 backoff: backoff_for(1),
1140 kind,
1141 detail: failure_detail(detail),
1142 })
1143 };
1144 let uris: Vec<Option<crate::standard_site::AtUri>> = feeds
1145 .iter()
1146 .map(|f| crate::standard_site::AtUri::parse(&f.url))
1147 .collect();
1148 let Some(did) = uris.iter().flatten().next().map(|u| u.authority.clone()) else {
1149 return feeds
1150 .iter()
1151 .map(|f| {
1152 failed(
1153 FailureKind::Parse,
1154 format!("{} is not a readable at:// URI", f.url),
1155 )
1156 })
1157 .collect();
1158 };
1159 // Only the feeds of that repo with a publication URI are read together.
1160 let readable: Vec<usize> = (0..feeds.len())
1161 .filter(|&i| {
1162 uris[i].as_ref().is_some_and(|u| {
1163 u.authority == did && u.collection == crate::lexicon::nsid::STANDARD_PUBLICATION
1164 })
1165 })
1166 .collect();
1167 let rkeys: Vec<String> = readable
1168 .iter()
1169 .map(|&i| uris[i].as_ref().unwrap().rkey.clone())
1170 .collect();
1171
1172 // **One deadline for the whole read.** It is otherwise bounded only per
1173 // request (FETCH_TIMEOUT x MAX_LIST_PAGES): hours, against a repo that
1174 // pages slowly, all of it holding up the publication loop.
1175 let fetched = tokio::time::timeout(
1176 config.publication_read_deadline,
1177 crate::standard_site::fetch_repo_capped(
1178 client,
1179 &config.oauth.plc_directory,
1180 &did,
1181 &rkeys,
1182 per_site_cap,
1183 budget_bytes,
1184 ),
1185 )
1186 .await;
1187 let mut reads: Vec<Option<anyhow::Result<crate::standard_site::PublicationRead>>> =
1188 feeds.iter().map(|_| None).collect();
1189 let repo_failure = match fetched {
1190 Ok(Ok(per)) => {
1191 for (slot, read) in readable.iter().zip(per) {
1192 reads[*slot] = Some(read);
1193 }
1194 None
1195 }
1196 Ok(Err(err)) => Some((publication_failure_kind(&err), format!("{err:#}"))),
1197 Err(_) => Some((
1198 FailureKind::Fetch,
1199 format!(
1200 "the read did not finish within {:?}",
1201 config.publication_read_deadline
1202 ),
1203 )),
1204 };
1205
1206 let (retention_days, retention_hard_days) = config.retention_for(FeedKind::Publication);
1207 let store = |url: String, read: crate::standard_site::PublicationRead| async move {
1208 crate::standard_site::store_publication(
1209 pool,
1210 &url,
1211 read,
1212 config.max_entries_per_feed,
1213 retention_days,
1214 retention_hard_days,
1215 )
1216 .await
1217 };
1218 let mut out = Vec::with_capacity(feeds.len());
1219 // The members the group's walk cut short, with how many entries the group
1220 // read gave each — taken before the read is moved into the store.
1221 let mut cut_short: Vec<(usize, usize)> = Vec::new();
1222 for (i, feed) in feeds.iter().enumerate() {
1223 let outcome = match (reads[i].take(), &repo_failure) {
1224 (Some(Ok(read)), _) => {
1225 if read.cut_short_by_group {
1226 cut_short.push((i, read.entries.len()));
1227 }
1228 store(feed.url.clone(), read).await
1229 }
1230 (Some(Err(err)), _) => failed(publication_failure_kind(&err), format!("{err:#}")),
1231 (None, Some((kind, detail))) if readable.contains(&i) => failed(*kind, detail.clone()),
1232 (None, _) => failed(
1233 FailureKind::Parse,
1234 format!("{} is not a publication in {did}", feed.url),
1235 ),
1236 };
1237 out.push(outcome);
1238 }
1239 drop(reads);
1240
1241 // **Re-read alone each member the group's walk cut short (#229).** One
1242 // repo's publications share one walk and one byte budget, so big siblings
1243 // can spend it and the walk ends before a quiet one's documents — every
1244 // poll. A static per-publication share fixed that and broke two worse
1245 // cases (a busy publication beside idle siblings was cut to a fraction;
1246 // one large document failed its publication every poll). Read alone, a
1247 // publication has the whole budget, so this is never worse than reading
1248 // it alone, by construction.
1249 //
1250 // **After every group read is stored and dropped, and one at a time**, so
1251 // at most one budget's worth of reads is held, as before. What the group
1252 // read stored stays: entries are upserted on `(feed_id, guid)`, so the
1253 // re-read's overlap updates rows rather than duplicating them.
1254 //
1255 // **Inside the same read deadline**, measured from the group read's start:
1256 // the re-reads must not stretch one group's poll past the bound the tick
1257 // was sized for. A member left over keeps its group outcome.
1258 if cut_short.is_empty() {
1259 return out;
1260 }
1261 // **Starved first.** The re-reads share one deadline and stop at the
1262 // first that overruns it, so the order decides who gets one. A big
1263 // sibling read alone can cost as much as the whole group read, and in
1264 // feed order it went first on every poll — leaving the quiet publication
1265 // this exists for `Failed` forever (review of #281). Fewest entries from
1266 // the group first, those with none ahead of all; ties keep feed order.
1267 cut_short.sort_by_key(|&(_, group_kept)| group_kept);
1268 let mut reread = 0usize;
1269 let mut improved = 0usize;
1270 for &(i, group_kept) in &cut_short {
1271 let feed = &feeds[i];
1272 let remaining = config
1273 .publication_read_deadline
1274 .saturating_sub(deadline_started.elapsed());
1275 if remaining.is_zero() {
1276 tracing::warn!(
1277 repo = %did,
1278 feed = %feed.url,
1279 left = cut_short.len() - reread,
1280 "the read deadline was spent before every publication its group cut short was re-read alone; the rest keep the group's outcome"
1281 );
1282 break;
1283 }
1284 let rkey = std::slice::from_ref(&uris[i].as_ref().expect("a read feed has a URI").rkey);
1285 let alone = tokio::time::timeout(
1286 remaining,
1287 crate::standard_site::fetch_repo_capped(
1288 client,
1289 &config.oauth.plc_directory,
1290 &did,
1291 rkey,
1292 per_site_cap,
1293 budget_bytes,
1294 ),
1295 )
1296 .await;
1297 // A re-read that fails keeps the group's outcome: the group read
1298 // already stored what it reached, and one failed attempt to read
1299 // further is not this feed's failure.
1300 let outcome = match alone {
1301 Err(_) => {
1302 tracing::warn!(
1303 repo = %did,
1304 feed = %feed.url,
1305 left = cut_short.len() - reread,
1306 "a re-read alone overran the read deadline; it and the rest keep the group's outcome"
1307 );
1308 break;
1309 }
1310 Ok(Err(err)) => {
1311 tracing::warn!(repo = %did, feed = %feed.url, err = %format!("{err:#}"),
1312 "a re-read alone failed; keeping the group's outcome");
1313 reread += 1;
1314 continue;
1315 }
1316 Ok(Ok(mut per)) => match per.pop() {
1317 Some(Ok(read)) => {
1318 let better = read.complete || read.entries.len() > group_kept;
1319 match store(feed.url.clone(), read).await {
1320 Ok(outcome) => {
1321 if better {
1322 improved += 1;
1323 }
1324 Ok(outcome)
1325 }
1326 // The group's rows are already stored; a failed second
1327 // store is not this feed's poll failing.
1328 Err(err) => {
1329 tracing::warn!(repo = %did, feed = %feed.url, err = %format!("{err:#}"),
1330 "storing a re-read failed; keeping the group's outcome");
1331 reread += 1;
1332 continue;
1333 }
1334 }
1335 }
1336 Some(Err(err)) => {
1337 tracing::warn!(repo = %did, feed = %feed.url, err = %format!("{err:#}"),
1338 "a re-read alone found no readable publication; keeping the group's outcome");
1339 reread += 1;
1340 continue;
1341 }
1342 None => {
1343 reread += 1;
1344 continue;
1345 }
1346 },
1347 };
1348 reread += 1;
1349 out[i] = outcome;
1350 }
1351 tracing::info!(
1352 repo = %did,
1353 cut_short = cut_short.len(),
1354 reread,
1355 improved,
1356 "re-read alone the publications their group's shared walk cut short"
1357 );
1358 out
1359}
1360
1361/// How long one entry's sanitize may run before the poll gives up on it (#226).
1362///
1363/// **Seconds, not milliseconds, on purpose.** Real bodies sanitize in about
1364/// 3.5 ms at most on an M-series Mac (production's largest body is 87 KB), so
1365/// 5 s is over 1,000x that: a much slower shared Fly vCPU, or one busy with
1366/// four polls at once, still never trips it on a real article. What it does
1367/// catch is the super-linear inputs ammonia has — 2 MiB of `&` runs take
1368/// ~2.4 s in release on that Mac, U+00A0 runs ~3.4 s, nested `<div>`s ~37 s,
1369/// and 8 MiB of `&` never finished in ten minutes — where the choice is
1370/// between waiting and giving up, and a timeout is a false positive only for
1371/// a body that is already nearly pathological.
1372pub(crate) const SANITIZE_TIMEOUT: Duration = Duration::from_secs(5);
1373
1374/// How many ingest sanitizes may run at once, process-wide (#226).
1375///
1376/// The poller's default `FEATHERREADER_POLL_CONCURRENCY` (4): in ordinary
1377/// operation every poll in flight gets a permit and none ever waits. What it
1378/// bounds is the bad case. A timed-out sanitize cannot be cancelled — ammonia
1379/// runs to completion on its blocking thread — and its permit is held until it
1380/// does, so hostile feeds can tie up at most this many blocking threads (and
1381/// CPUs) however many of them there are. Beyond that a poll waits for a permit
1382/// **asynchronously**, and for at most [`SANITIZE_TIMEOUT`] — see
1383/// [`sanitize_off_runtime`] for why that wait is bounded too.
1384pub(crate) const SANITIZE_CONCURRENCY: usize = 4;
1385
1386/// The permits behind [`SANITIZE_CONCURRENCY`].
1387static SANITIZE_PERMITS: tokio::sync::Semaphore =
1388 tokio::sync::Semaphore::const_new(SANITIZE_CONCURRENCY);
1389
1390/// The feeds with a sanitize still running, and how many of those a poll has
1391/// **given up on** (#226).
1392///
1393/// **Why:** a sanitize that timed out keeps running — ammonia cannot be
1394/// interrupted — and the feed's next poll is a full fetch (its validators were
1395/// not saved). Without this, every retry started another sanitize, and one
1396/// hostile feed could come to hold every permit and starve every other feed.
1397/// With it, a feed with a sanitize still running is not given another until
1398/// it returns: **one feed, at most one thread**, however often it is polled.
1399///
1400/// **Counted from the start**, by a guard the blocking closure owns, and
1401/// released only when ammonia returns or panics. Counting only once a poll
1402/// gave up let overlapping polls of one feed (the scheduler and a subscribe
1403/// POST, which polls inline) each start one before any timed out, and never
1404/// counted a sanitize whose poll was dropped mid-way (a client disconnecting).
1405/// A feed whose running sanitize was given up on is refused as its own
1406/// failure ([`SanitizeGaveUp::StillRunning`]); one whose sanitize is merely
1407/// in progress, or orphaned by a dropped poll, defers
1408/// ([`SanitizeGaveUp::Busy`]).
1409///
1410/// **By feed, not by body hash:** a hash let a feed dodge the refusal by
1411/// serving different bytes each fetch (a nonce, or another slow entry on
1412/// top), and it shared one feed's timeout with every other feed carrying the
1413/// same article.
1414pub(crate) struct InFlight(std::sync::Mutex<std::collections::BTreeMap<String, Slot>>);
1415
1416/// One feed's sanitizes in [`InFlight`]. Removed when both are zero.
1417#[derive(Default)]
1418pub(crate) struct Slot {
1419 /// Still running, given up on or not.
1420 running: usize,
1421 /// Of those, given up on by their poll.
1422 abandoned: usize,
1423}
1424
1425impl InFlight {
1426 pub(crate) const fn new() -> Self {
1427 InFlight(std::sync::Mutex::new(std::collections::BTreeMap::new()))
1428 }
1429
1430 fn lock(&self) -> std::sync::MutexGuard<'_, std::collections::BTreeMap<String, Slot>> {
1431 // Nothing here can leave the map inconsistent, so a poisoned lock is
1432 // still a usable one.
1433 self.0.lock().unwrap_or_else(|e| e.into_inner())
1434 }
1435}
1436
1437/// [`AbandonGuard::state`]: running, not given up on.
1438const RUNNING: u8 = 0;
1439/// Given up on by its poll, and counted as abandoned in [`InFlight`].
1440const ABANDONED: u8 = 1;
1441/// Returned (or panicked); never counted again.
1442const DONE: u8 = 2;
1443
1444/// One sanitize's standing in [`InFlight`], from registration until the
1445/// sanitize returns. The state is only read or changed with the map locked,
1446/// so "give up" and "return" cannot interleave: a sanitize that has already
1447/// returned is never counted as abandoned, and every count is taken back
1448/// exactly once.
1449struct AbandonGuard {
1450 set: &'static InFlight,
1451 feed: String,
1452 state: std::sync::Arc<std::sync::atomic::AtomicU8>,
1453}
1454
1455impl AbandonGuard {
1456 /// Register a sanitize for `feed`, unless one is already running: then
1457 /// why not — given up on (`StillRunning`) or merely in progress (`Busy`).
1458 fn register(set: &'static InFlight, feed: &str) -> Result<Self, SanitizeGaveUp> {
1459 let mut map = set.lock();
1460 let slot = map.entry(feed.to_owned()).or_default();
1461 if slot.abandoned > 0 {
1462 return Err(SanitizeGaveUp::StillRunning);
1463 }
1464 if slot.running > 0 {
1465 return Err(SanitizeGaveUp::Busy);
1466 }
1467 slot.running = 1;
1468 Ok(AbandonGuard {
1469 set,
1470 feed: feed.to_owned(),
1471 state: std::sync::Arc::new(std::sync::atomic::AtomicU8::new(RUNNING)),
1472 })
1473 }
1474
1475 /// The poll gave up on it: count it as abandoned, if still running.
1476 fn abandon(set: &'static InFlight, feed: &str, state: &std::sync::atomic::AtomicU8) {
1477 use std::sync::atomic::Ordering::SeqCst;
1478 let mut map = set.lock();
1479 if state.load(SeqCst) == RUNNING {
1480 state.store(ABANDONED, SeqCst);
1481 map.entry(feed.to_owned()).or_default().abandoned += 1;
1482 }
1483 }
1484}
1485
1486impl Drop for AbandonGuard {
1487 fn drop(&mut self) {
1488 use std::sync::atomic::Ordering::SeqCst;
1489 let mut map = self.set.lock();
1490 let was = self.state.swap(DONE, SeqCst);
1491 if let Some(slot) = map.get_mut(&self.feed) {
1492 slot.running = slot.running.saturating_sub(1);
1493 if was == ABANDONED {
1494 slot.abandoned = slot.abandoned.saturating_sub(1);
1495 }
1496 if slot.running == 0 && slot.abandoned == 0 {
1497 map.remove(&self.feed);
1498 }
1499 }
1500 }
1501}
1502
1503/// Whether ingest is being starved of sanitize permits, for `/stats` and the
1504/// logs (review of #274).
1505///
1506/// **Why:** one feed holds at most one thread (`InFlight`), but four hostile
1507/// feeds — four URL variants of one, say — can still hold all
1508/// `SANITIZE_CONCURRENCY` permits with sanitizes nobody can cancel. Every
1509/// other feed's poll then defers (`SanitizeGaveUp::NoPermit`): nothing is
1510/// stored and, by design, nothing is filed against those feeds — so without
1511/// this the instance would silently stop updating. Counted, not prevented:
1512/// the starvation lasts until those sanitizes finish.
1513pub struct Starvation {
1514 /// Polls deferred for want of a permit, since boot.
1515 deferred: std::sync::atomic::AtomicU64,
1516 /// Unix seconds of the first such deferral since a permit was last
1517 /// acquired; 0 while permits are coming free.
1518 since: std::sync::atomic::AtomicI64,
1519}
1520
1521/// After this long without a free permit, a deferral logs at error level.
1522pub const STARVED_ALERT_SECS: i64 = 5 * 60;
1523
1524impl Starvation {
1525 pub const fn new() -> Self {
1526 Starvation {
1527 deferred: std::sync::atomic::AtomicU64::new(0),
1528 since: std::sync::atomic::AtomicI64::new(0),
1529 }
1530 }
1531
1532 /// A poll got no permit at `now` (unix seconds). Returns how long the
1533 /// instance has gone without one.
1534 pub(crate) fn record_no_permit(&self, now: i64) -> i64 {
1535 use std::sync::atomic::Ordering::Relaxed;
1536 self.deferred.fetch_add(1, Relaxed);
1537 // From the exchange's own result, not a second load: a permit
1538 // acquired in between resets `since` to 0, and `now - 0` would log
1539 // decades of starvation.
1540 let since = match self.since.compare_exchange(0, now, Relaxed, Relaxed) {
1541 Ok(_) => now,
1542 Err(prev) => prev,
1543 };
1544 now - since
1545 }
1546
1547 /// A permit was acquired: not starved (any more).
1548 pub(crate) fn record_permit(&self) {
1549 self.since.store(0, std::sync::atomic::Ordering::Relaxed);
1550 }
1551
1552 /// Polls deferred since boot, and since when (unix seconds) no permit has
1553 /// come free, if that is now.
1554 pub fn snapshot(&self) -> (u64, Option<i64>) {
1555 use std::sync::atomic::Ordering::Relaxed;
1556 let since = self.since.load(Relaxed);
1557 (self.deferred.load(Relaxed), (since != 0).then_some(since))
1558 }
1559}
1560
1561impl Default for Starvation {
1562 fn default() -> Self {
1563 Self::new()
1564 }
1565}
1566
1567/// The production [`Starvation`], which `/stats` reads.
1568pub static SANITIZE_STARVATION: Starvation = Starvation::new();
1569
1570/// The production [`InFlight`].
1571static IN_FLIGHT: InFlight = InFlight::new();
1572
1573/// Where and for how long [`normalize_entries`] sanitizes. Production uses
1574/// [`SanitizeLimits::PRODUCTION`]; tests inject a short timeout and their own
1575/// permits and in-flight set, so an abandoned sanitize in one test cannot hold
1576/// up another.
1577#[derive(Clone, Copy)]
1578pub(crate) struct SanitizeLimits {
1579 pub(crate) permits: &'static tokio::sync::Semaphore,
1580 pub(crate) in_flight: &'static InFlight,
1581 pub(crate) timeout: Duration,
1582 /// Where permit starvation is recorded: [`SANITIZE_STARVATION`] in
1583 /// production.
1584 pub(crate) starvation: &'static Starvation,
1585 /// What runs on the blocking pool: [`sanitize_body`] in production. A seam
1586 /// for tests only — one that panics, or that takes the remaining permits
1587 /// mid-poll — so the paths those reach are tested deterministically.
1588 pub(crate) sanitize: fn(&str) -> String,
1589}
1590
1591impl SanitizeLimits {
1592 /// [`SANITIZE_PERMITS`], [`IN_FLIGHT`] and [`SANITIZE_TIMEOUT`].
1593 pub(crate) const PRODUCTION: SanitizeLimits = SanitizeLimits {
1594 permits: &SANITIZE_PERMITS,
1595 in_flight: &IN_FLIGHT,
1596 timeout: SANITIZE_TIMEOUT,
1597 starvation: &SANITIZE_STARVATION,
1598 sanitize: sanitize_body,
1599 };
1600}
1601
1602/// An entry body as it is stored: [`sanitize_html_bounded`] to
1603/// [`MAX_CONTENT_HTML_BYTES`].
1604fn sanitize_body(raw: &str) -> String {
1605 sanitize_html_bounded(raw, MAX_CONTENT_HTML_BYTES)
1606}
1607
1608/// Why [`sanitize_off_runtime`] produced no body.
1609#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1610enum SanitizeGaveUp {
1611 /// Every permit stayed taken for the whole timeout — by sanitizes other
1612 /// polls gave up on, so not the polled feed's fault.
1613 NoPermit,
1614 /// The sanitize did not finish within the timeout (or the runtime is
1615 /// shutting down and cancelled it before it started).
1616 TimedOut,
1617 /// A sanitize of this feed's was already given up on and is still
1618 /// running ([`InFlight`]); refused without starting another. The feed's
1619 /// fault exactly as [`SanitizeGaveUp::TimedOut`] is.
1620 StillRunning,
1621 /// A sanitize of this feed's is still running but nobody gave up on it:
1622 /// an overlapping poll, or one dropped mid-sanitize ([`InFlight`]). Not
1623 /// started, and not the feed's fault: the poll defers, as for
1624 /// [`SanitizeGaveUp::NoPermit`].
1625 Busy,
1626}
1627
1628/// [`sanitize_html_bounded`] on the blocking pool, under a permit and a
1629/// timeout.
1630///
1631/// **Never on an async worker**: ammonia is super-linear on some inputs (#226),
1632/// and inline it stalled a tokio worker — and the poller with it — for as long
1633/// as it ran. The permit is acquired asynchronously and moved INTO the blocking
1634/// closure, so a sanitize the caller has given up on keeps it until ammonia
1635/// returns: the bound counts threads actually busy, not polls still waiting.
1636///
1637/// **The wait for a permit is bounded by the same timeout**, separately from
1638/// the sanitize. Permits held by abandoned sanitizes can stay taken for
1639/// minutes, and an unbounded wait would stall every poll behind them — and a
1640/// graceful shutdown, which drains the polls in flight, past Fly's
1641/// `kill_timeout`. So one entry costs a poll at most two timeouts.
1642///
1643/// **A feed with a sanitize still running is refused before any of that** —
1644/// no permit, no thread — until that sanitize finishes; see [`InFlight`].
1645async fn sanitize_off_runtime(
1646 feed: &str,
1647 raw: String,
1648 limits: SanitizeLimits,
1649) -> Result<String, SanitizeGaveUp> {
1650 // Registered before the permit wait, so overlapping polls of one feed
1651 // cannot both get past here. If the wait fails, or this future is
1652 // dropped before the spawn, the guard drops and the count goes with it.
1653 let guard = AbandonGuard::register(limits.in_flight, feed)?;
1654 let state = std::sync::Arc::clone(&guard.state);
1655 let permit = match tokio::time::timeout(limits.timeout, limits.permits.acquire()).await {
1656 Ok(Ok(permit)) => {
1657 limits.starvation.record_permit();
1658 permit
1659 }
1660 // `Ok(Err(_))` is a closed semaphore, which this never does.
1661 Ok(Err(_)) | Err(_) => {
1662 let starved = limits.starvation.record_no_permit(Utc::now().timestamp());
1663 if starved >= STARVED_ALERT_SECS {
1664 tracing::error!(
1665 starved_secs = starved,
1666 "ingest is stalled: no sanitize permit has come free for {}m, so every \
1667 poll is deferred and nothing is stored (#226: slow bodies' abandoned \
1668 sanitizes hold them all)",
1669 starved / 60
1670 );
1671 }
1672 return Err(SanitizeGaveUp::NoPermit);
1673 }
1674 };
1675 let sanitize = limits.sanitize;
1676 let task = tokio::task::spawn_blocking(move || {
1677 // Both released when ammonia returns — or panics — not when the poll
1678 // gives up.
1679 let _permit = permit;
1680 let _guard = guard;
1681 sanitize(&raw)
1682 });
1683 match tokio::time::timeout(limits.timeout, task).await {
1684 Ok(Ok(html)) => Ok(html),
1685 // A panic in the sanitizer propagates as it did when this ran inline.
1686 Ok(Err(err)) if err.is_panic() => std::panic::resume_unwind(err.into_panic()),
1687 Ok(Err(_)) | Err(_) => {
1688 // From now until it returns, this feed is refused.
1689 AbandonGuard::abandon(limits.in_flight, feed, &state);
1690 Err(SanitizeGaveUp::TimedOut)
1691 }
1692 }
1693}
1694
1695/// Normalize a poll's entries, sanitizing each body off the async runtime
1696/// ([`sanitize_off_runtime`]). Returns the entries and whether a sanitize
1697/// timed out — or `None` when the poll should be **deferred**.
1698///
1699/// **After the first timeout no further body is sanitized in this poll**, so
1700/// one hostile feed costs at most one abandoned blocking thread per poll. What
1701/// stops the later entries is the per-feed refusal ([`InFlight`]): the
1702/// abandoned sanitize is still running, so [`sanitize_off_runtime`] refuses
1703/// each of them at once. The `timed_out` guard here is only a backup, for the
1704/// race in which that sanitize finishes between being given up on and the
1705/// next entry — without it, that entry would be sanitized after all. The
1706/// timed-out entry and every later entry with a body are returned with
1707/// `content_html: None` and `keep_stored_content: true`: an entry already
1708/// stored keeps the body it has, and a new one is stored without a body (the
1709/// reader shows its title and a link to the original). Nothing degraded is
1710/// ever built or stored in its place — no excerpt, no cut of the raw input.
1711///
1712/// **No permit is not a timeout.** It means other feeds' abandoned sanitizes
1713/// hold every permit, which says nothing about this feed, so the whole poll
1714/// is deferred (`None`): storing its new entries without bodies, and filing a
1715/// failure against it, is what let one hostile feed degrade every other one.
1716/// Nothing is stored, not even the entries sanitized before it — the next
1717/// poll is a full fetch and stores them all, with bodies.
1718async fn normalize_entries(
1719 feed_url: &str,
1720 raw: &[RawEntry],
1721 limits: SanitizeLimits,
1722) -> Option<(Vec<NewEntry>, bool)> {
1723 let mut timed_out = false;
1724 let mut out = Vec::with_capacity(raw.len());
1725 for e in raw {
1726 let mut entry = entry_without_body(e);
1727 if let Some(body) = entry_raw_body(e) {
1728 if timed_out {
1729 // Normally `sanitize_off_runtime` would refuse it anyway (the
1730 // abandoned sanitize is still in flight); this covers the
1731 // race where that sanitize has already finished.
1732 entry.keep_stored_content = true;
1733 } else {
1734 match sanitize_off_runtime(feed_url, body.to_owned(), limits).await {
1735 Ok(html) => entry.content_html = Some(html),
1736 Err(SanitizeGaveUp::NoPermit) => {
1737 tracing::warn!(
1738 feed = %feed_url,
1739 timeout = ?limits.timeout,
1740 "no sanitize permit came free in time (all held by abandoned \
1741 sanitizes); deferring this feed's poll without storing it"
1742 );
1743 return None;
1744 }
1745 Err(SanitizeGaveUp::Busy) => {
1746 tracing::info!(
1747 feed = %feed_url,
1748 "a sanitize for this feed is already running (an overlapping \
1749 or dropped poll); deferring this poll without storing it"
1750 );
1751 return None;
1752 }
1753 Err(why) => {
1754 tracing::warn!(
1755 feed = %feed_url,
1756 entry = %entry.guid,
1757 raw_bytes = body.len(),
1758 timeout = ?limits.timeout,
1759 ?why,
1760 "an entry body was not sanitized in time; this poll stores no \
1761 further bodies for the feed and keeps any already stored"
1762 );
1763 timed_out = true;
1764 entry.keep_stored_content = true;
1765 }
1766 }
1767 }
1768 }
1769 out.push(entry);
1770 }
1771 Some((out, timed_out))
1772}
1773
1774/// Fetch, parse, sanitize, normalize, and store a single feed.
1775///
1776/// Performs a conditional GET using the feed's stored `ETag` / `Last-Modified`.
1777/// On `304` it returns [`PollOutcome::NotModified`] without touching entries. On
1778/// `200` it parses with `feed-rs`, sanitizes every entry's HTML with `ammonia`,
1779/// upserts the feed row (carrying the fresh validators) and inserts new entries
1780/// (deduped by GUID). Any fetch/parse error is logged and returned as
1781/// [`PollOutcome::Failed`] — it never panics and never propagates as `Err` for
1782/// a merely-broken feed, so one bad publisher can't stall the scheduler.
1783///
1784/// `Err` is reserved for *store* failures (a broken local DB is a real error the
1785/// caller should see), not for feed misbehaviour.
1786///
1787/// `max_entries_per_feed` caps how many entries this feed retains after insert
1788/// (newest N by published date); `<= 0` disables the per-feed trim.
1789///
1790/// Bodies are sanitized off the async runtime with a per-entry timeout
1791/// (`SANITIZE_TIMEOUT`; see `normalize_entries`). **A poll in which a sanitize
1792/// timed out does not save the response's `ETag` / `Last-Modified`** — the
1793/// next poll must get a full `200` to store the bodies this one could not, and
1794/// a saved validator would answer it with `304`s until the feed next changed.
1795/// Such a poll is reported as a [`FailureKind::Body`] failure, after storing
1796/// what it has, so the feed backs off and shows why in `/stats`: the body was
1797/// the feed's. A poll that could not get a sanitize permit at all — every one
1798/// held by OTHER feeds' abandoned sanitizes — stores nothing and returns
1799/// [`PollOutcome::Deferred`], which counts against no one.
1800pub async fn poll_feed(
1801 pool: &SqlitePool,
1802 client: &Client,
1803 feed: &Feed,
1804 max_entries_per_feed: i64,
1805) -> Result<PollOutcome> {
1806 poll_feed_with(
1807 pool,
1808 client,
1809 feed,
1810 max_entries_per_feed,
1811 SanitizeLimits::PRODUCTION,
1812 )
1813 .await
1814}
1815
1816/// [`poll_feed`] with the sanitize limits injected.
1817async fn poll_feed_with(
1818 pool: &SqlitePool,
1819 client: &Client,
1820 feed: &Feed,
1821 max_entries_per_feed: i64,
1822 limits: SanitizeLimits,
1823) -> Result<PollOutcome> {
1824 // --- conditional GET (through the SSRF guard) ----------------------------
1825 // The guard re-validates the scheme + resolved IP of the target and of every
1826 // redirect hop, so a subscribed feed can't bounce the poller onto an
1827 // internal address (cloud metadata / loopback). Conditional-GET validators
1828 // ride along as extra headers.
1829 let mut extra: Vec<(reqwest::header::HeaderName, reqwest::header::HeaderValue)> = Vec::new();
1830 if let Some(etag) = feed.etag.as_deref() {
1831 if let Ok(v) = reqwest::header::HeaderValue::from_str(etag) {
1832 extra.push((IF_NONE_MATCH, v));
1833 }
1834 }
1835 if let Some(lm) = feed.last_modified.as_deref() {
1836 if let Ok(v) = reqwest::header::HeaderValue::from_str(lm) {
1837 extra.push((IF_MODIFIED_SINCE, v));
1838 }
1839 }
1840
1841 let resp = match crate::net::guarded_get(client, &feed.url, &extra).await {
1842 Ok(r) => r,
1843 Err(e) => {
1844 tracing::warn!(feed = %feed.url, error = %e, "feed fetch failed (or blocked by SSRF guard)");
1845 return Ok(PollOutcome::Failed {
1846 backoff: backoff_for(1),
1847 kind: FailureKind::Fetch,
1848 detail: failure_detail(format!("{e:#}")),
1849 });
1850 }
1851 };
1852
1853 let status = resp.status();
1854 if status == StatusCode::NOT_MODIFIED {
1855 tracing::debug!(feed = %feed.url, "feed not modified (304)");
1856 // Bump last_polled/next_poll only; leave validators + entries untouched.
1857 touch_polled(pool, &feed.url, None, None)
1858 .await
1859 .with_context(|| format!("touch_polled after 304 for {}", feed.url))?;
1860 return Ok(PollOutcome::NotModified);
1861 }
1862 if !status.is_success() {
1863 tracing::warn!(feed = %feed.url, %status, "feed returned non-success status");
1864 return Ok(PollOutcome::Failed {
1865 backoff: backoff_for(1),
1866 kind: FailureKind::Status,
1867 detail: failure_detail(status),
1868 });
1869 }
1870
1871 // Capture validators for the *next* conditional GET before consuming body.
1872 let new_etag = header_str(resp.headers().get(ETAG));
1873 let new_last_modified = header_str(resp.headers().get(LAST_MODIFIED));
1874
1875 // Stream the body with a hard byte cap, aborting mid-stream if it exceeds
1876 // it. We never trust Content-Length: reqwest's gzip layer strips it, so a
1877 // small gzip bomb could otherwise inflate to GBs before any size check.
1878 let body = match crate::net::read_capped(resp).await {
1879 Ok(b) => b,
1880 Err(e) => {
1881 tracing::warn!(feed = %feed.url, error = %e, "feed body rejected (too large / read error)");
1882 return Ok(PollOutcome::Failed {
1883 backoff: backoff_for(1),
1884 kind: FailureKind::Body,
1885 detail: failure_detail(format!("{e:#}")),
1886 });
1887 }
1888 };
1889
1890 // --- parse (malformed feed => log + skip, never panic) -------------------
1891 let parsed = match parse_feed(&body[..]) {
1892 Ok(f) => f,
1893 Err(e) => {
1894 tracing::warn!(feed = %feed.url, error = %e, "malformed feed; skipping");
1895 return Ok(PollOutcome::Failed {
1896 backoff: backoff_for(1),
1897 kind: FailureKind::Parse,
1898 detail: failure_detail(format!("{e:#}")),
1899 });
1900 }
1901 };
1902
1903 // --- normalize + sanitize ------------------------------------------------
1904 let Some((entries, sanitize_timed_out)) =
1905 normalize_entries(&feed.url, &parsed.entries, limits).await
1906 else {
1907 // Every permit is held by other feeds' abandoned sanitizes: store
1908 // nothing from this poll (not even its validators) and blame nothing.
1909 return Ok(PollOutcome::Deferred);
1910 };
1911
1912 let (title, site_url) = feed_metadata(&parsed);
1913 // A timed-out poll keeps the validators already stored (`None` is "keep
1914 // current" to the upsert), so the next poll is a full fetch.
1915 let (etag, last_modified) = if sanitize_timed_out {
1916 (None, None)
1917 } else {
1918 (new_etag, new_last_modified)
1919 };
1920 let new_feed = NewFeed {
1921 url: feed.url.clone(),
1922 title,
1923 site_url,
1924 etag,
1925 last_modified,
1926 last_polled: Some(now_rfc3339()),
1927 next_poll: None, // the scheduler owns cadence; leave it to set next_poll.
1928 };
1929
1930 // --- store (a store failure IS a real error) -----------------------------
1931 let feed_id = store::upsert_feed(pool, &new_feed)
1932 .await
1933 .with_context(|| format!("upsert_feed for {}", feed.url))?;
1934 let n = store::insert_entries(pool, feed_id, &entries, max_entries_per_feed)
1935 .await
1936 .with_context(|| format!("insert_entries for {}", feed.url))?;
1937
1938 tracing::info!(feed = %feed.url, entries = n, "feed polled");
1939 if sanitize_timed_out {
1940 return Ok(PollOutcome::Failed {
1941 backoff: backoff_for(1),
1942 kind: FailureKind::Body,
1943 detail: failure_detail(format!(
1944 "an entry body was not sanitized within {:?}; entries were stored \
1945 without new bodies",
1946 limits.timeout
1947 )),
1948 });
1949 }
1950 Ok(PollOutcome::Updated { new_entries: n })
1951}
1952
1953/// Bump `last_polled` (and optionally validators) without changing entries —
1954/// used on the `304 Not Modified` path.
1955async fn touch_polled(
1956 pool: &SqlitePool,
1957 url: &str,
1958 etag: Option<String>,
1959 last_modified: Option<String>,
1960) -> Result<()> {
1961 // `None` means "keep current" — upsert_feed COALESCEs the validators, so a
1962 // 304 that repeats no headers leaves the stored ones untouched. This used to
1963 // re-read the row and re-supply them by hand because the upsert clobbered
1964 // unconditionally; the read-modify-write is gone now that the upsert is
1965 // honest, and with it a race where a concurrent poll's validators could be
1966 // read here and written back stale.
1967 let nf = NewFeed {
1968 url: url.to_string(),
1969 etag,
1970 last_modified,
1971 last_polled: Some(now_rfc3339()),
1972 ..Default::default()
1973 };
1974 store::upsert_feed(pool, &nf).await?;
1975 Ok(())
1976}
1977
1978/// Extract `(title, site_url)` from a parsed feed. `site_url` prefers an
1979/// `alternate`/no-rel HTML link over the feed's self link.
1980fn feed_metadata(parsed: &RawFeed) -> (Option<String>, Option<String>) {
1981 let title = parsed
1982 .title
1983 .as_ref()
1984 .map(|t| bound_text(text_plain(t), MAX_TITLE_BYTES));
1985 let site_url = parsed
1986 .links
1987 .iter()
1988 // Prefer an explicit human-facing page: rel="alternate" or no rel at all.
1989 .find(|l| {
1990 l.rel.as_deref() == Some("alternate")
1991 || (l.rel.is_none()
1992 && l.media_type.as_deref() != Some("application/rss+xml")
1993 && l.media_type.as_deref() != Some("application/atom+xml"))
1994 })
1995 .or_else(|| {
1996 parsed
1997 .links
1998 .iter()
1999 .find(|l| l.rel.as_deref() != Some("self"))
2000 })
2001 .or_else(|| parsed.links.first())
2002 .map(|l| bound_text(l.href.clone(), MAX_URL_BYTES));
2003 (title, site_url)
2004}
2005
2006/// Parse a feed body with feed-rs, generating ids for id-less entries the way
2007/// feed-rs 2.4 did.
2008///
2009/// An entry's id is its dedup key, so how feed-rs fills a missing one is part
2010/// of FeatherReader's storage contract. feed-rs 3.0 also puts `<comments>` and
2011/// `wfw:commentRss` URLs in `entry.links`, and its default generator hashes the
2012/// FIRST link with the title — so an id-less item listing its comments link
2013/// before `<link>` got a new id on upgrade, and was stored a second time. This
2014/// generator hashes the first link that is not a comments link, which is the
2015/// link 2.4 hashed, through the same public 2.4/3.0 function.
2016///
2017/// With no such link, it returns an empty id rather than feed-rs's fallback (a
2018/// random UUID, which made the item a new row on every poll); `normalize_entry`
2019/// then derives [`stable_guid`]. `parse` has no base URI, so feed-rs's
2020/// uri+title branch was never reached and nothing else is lost.
2021///
2022/// **A panic inside feed-rs is returned as an error.** feed-rs 3.0 panics on
2023/// some hostile-but-plausible input — an author address with a multi-byte
2024/// character beside it, e.g. `jose@example.com(José)` (its name/address
2025/// splitter slices on a byte that is not a character boundary). Feed bodies are
2026/// arbitrary web input, so any such panic is a malformed feed: it takes the
2027/// ordinary parse-failure path (logged, `FailureKind::Parse`, backed off) rather
2028/// than unwinding out of the poll task, which recorded nothing and left the
2029/// feed silently un-polled. Unwinding is safe here: the parser is built and
2030/// dropped inside the closure and touches no state of ours.
2031fn parse_feed(body: &[u8]) -> Result<RawFeed> {
2032 let parsed = std::panic::catch_unwind(|| {
2033 feed_rs::parser::Builder::new()
2034 .id_generator(entry_id)
2035 .build()
2036 .parse(body)
2037 });
2038 match parsed {
2039 Ok(result) => Ok(result?),
2040 Err(panic) => {
2041 let why = panic
2042 .downcast_ref::<String>()
2043 .map(String::as_str)
2044 .or_else(|| panic.downcast_ref::<&str>().copied())
2045 .unwrap_or("non-string panic payload");
2046 anyhow::bail!("feed parser panicked: {why}")
2047 }
2048 }
2049}
2050
2051/// The id generator [`parse_feed`] installs; see there.
2052fn entry_id(links: &[RawLink], title: &Option<Text>, _uri: Option<&str>) -> String {
2053 match links.iter().find(|l| is_primary_link(l)) {
2054 Some(link) => feed_rs::parser::generate_id_from_link_and_title(link, title),
2055 None => String::new(),
2056 }
2057}
2058
2059/// Whether a link is one of the entry's own links rather than a pointer to its
2060/// comments (feed-rs 3.0's `<comments>` / `wfw:commentRss`, marked by
2061/// `target`). Only these are candidates for the permalink and the id, which is
2062/// all feed-rs 2.4 ever put in `entry.links`.
2063fn is_primary_link(l: &RawLink) -> bool {
2064 l.target.is_none()
2065}
2066
2067/// Turn a parsed [`RawEntry`] into the store's [`NewEntry`], sanitizing HTML
2068/// inline. **Tests only**: the poller sanitizes off the async runtime, through
2069/// [`normalize_entries`]. Both are [`entry_without_body`] plus
2070/// [`sanitize_html_bounded`] of [`entry_raw_body`], so they store the same bytes.
2071#[cfg(test)]
2072fn normalize_entry(e: &RawEntry) -> NewEntry {
2073 NewEntry {
2074 content_html: entry_raw_body(e)
2075 .map(|raw| sanitize_html_bounded(raw, MAX_CONTENT_HTML_BYTES)),
2076 ..entry_without_body(e)
2077 }
2078}
2079
2080/// The markup an entry's stored body is sanitized from: the full `content`
2081/// body, else the `summary`. Whichever is chosen is **always** passed through
2082/// [`sanitize_html_bounded`] before storage.
2083fn entry_raw_body(e: &RawEntry) -> Option<&str> {
2084 e.content
2085 .as_ref()
2086 .and_then(|c| c.body.as_deref())
2087 .or_else(|| e.summary.as_ref().map(|t| t.content.as_str()))
2088}
2089
2090/// Everything about a [`RawEntry`] except its body, which is left `None` for
2091/// the caller to sanitize. GUID falls back to the entry link, then to a stable
2092/// hash of title+link, so an entry missing an `id` still deduplicates instead
2093/// of being re-inserted forever.
2094fn entry_without_body(e: &RawEntry) -> NewEntry {
2095 let url = entry_link(e);
2096
2097 // GUID may use the raw link (dedup key only, never rendered), so prefer the
2098 // entry's first raw link for identity even when it's not a safe href.
2099 let guid = if !e.id.trim().is_empty() {
2100 e.id.trim().to_string()
2101 } else if let Some(link) = raw_entry_link(e) {
2102 link
2103 } else {
2104 // Last resort: derive a stable id so re-fetches dedup rather than dupe.
2105 stable_guid(e)
2106 };
2107
2108 NewEntry {
2109 guid: bound_guid(guid),
2110 url: url.map(|u| bound_text(u, MAX_URL_BYTES)),
2111 title: e
2112 .title
2113 .as_ref()
2114 .map(|t| bound_text(text_plain(t), MAX_TITLE_BYTES)),
2115 author: entry_author(e).map(|a| bound_text(a, MAX_AUTHOR_BYTES)),
2116 published: entry_time(e),
2117 content_html: None,
2118 fetched_at: None, // store defaults to "now".
2119 keep_stored_content: false,
2120 }
2121}
2122
2123/// The raw best-permalink URL for an entry (no scheme filtering) — used only as
2124/// a dedup GUID, never rendered as an href.
2125///
2126/// Comments links are never candidates (see [`is_primary_link`]).
2127fn raw_entry_link(e: &RawEntry) -> Option<String> {
2128 let mut links = e.links.iter().filter(|l| is_primary_link(l));
2129 links
2130 .clone()
2131 .find(|l| l.rel.as_deref() == Some("alternate") || l.rel.is_none())
2132 .or_else(|| links.next())
2133 .map(|l| l.href.clone())
2134}
2135
2136/// The best display/permalink URL for an entry, **scheme-allow-listed** so it is
2137/// safe to render as an `href`: prefer `rel="alternate"` or a no-rel link, else
2138/// the first link — but only if it is an `http`/`https` URL. A `javascript:` or
2139/// `data:` permalink (a stored-XSS vector that survives HTML escaping, since it
2140/// carries no HTML-special characters) is dropped here at ingest, before it can
2141/// ever reach the store or a template.
2142fn entry_link(e: &RawEntry) -> Option<String> {
2143 raw_entry_link(e).and_then(|href| crate::net::safe_link(&href))
2144}
2145
2146/// The first author's name, if it has one.
2147///
2148/// feed-rs 3.0 splits RSS's `address (Name)` into an email and a name, but
2149/// keeps the parentheses: `(Name)`. They are stripped here. An author given
2150/// only as an address has no name and no byline — the address is not shown.
2151/// (feed-rs 2.4 named every RSS `<author>` "author", and an Atom author with no
2152/// name "unknown"; those placeholders were stored as bylines.)
2153///
2154/// The byline is the first author that has a name once stripped: an item may
2155/// give `<author>` as a bare address and the name in `<dc:creator>`.
2156fn entry_author(e: &RawEntry) -> Option<String> {
2157 e.authors.iter().find_map(|p| {
2158 let name = p.name.as_deref()?.trim();
2159 let name = [('(', ')'), ('<', '>'), ('[', ']')]
2160 .iter()
2161 .find_map(|&(open, close)| name.strip_prefix(open)?.strip_suffix(close))
2162 .unwrap_or(name)
2163 .trim();
2164 (!name.is_empty()).then(|| name.to_string())
2165 })
2166}
2167
2168/// Best CREDIBLE publication time (published, else updated) as an RFC3339
2169/// string, or `None` when neither is credible.
2170///
2171/// **A date in the future is discarded, not clamped, and not stored.** There was
2172/// no upper bound here, so an item dated in the year 2999 was stored verbatim
2173/// and became permanent: both retention sweeps test
2174/// `COALESCE(published, fetched_at) < cutoff` and a future date is never less
2175/// than either, the per-feed keep-set orders on the same expression `DESC` where
2176/// it is rank one forever, and every list view puts it at the top. A publisher
2177/// with a broken clock does that by accident; anyone wanting a permanent slot at
2178/// the top of a reader's list does it on purpose.
2179///
2180/// Discarded rather than clamped to now because the entry upsert refreshes
2181/// `published` on every poll while stamping `fetched_at` once — so a value
2182/// derived from the current clock is rewritten every cycle and the row can never
2183/// age at all. Clamping relocates the defect. Undated is the honest answer, and
2184/// `fetched_at` then dates the row and holds still. That is the rule the
2185/// publication path already follows; see `standard_site::entries_from_records`.
2186///
2187/// **Each candidate is judged separately**, so a bogus `<published>` beside a
2188/// credible `<updated>` keeps the good date. That helps Atom, and RSS 2 only
2189/// when the item carries an `<atom:updated>` (read since feed-rs 3.0):
2190/// otherwise `feed-rs` copies `published` into `updated` when `updated` is
2191/// absent (`parser/rss2/mod.rs`), so the second candidate holds the same value
2192/// and the fall-through is a no-op. Worth keeping where two independent dates
2193/// exist; worth not overstating where they do not.
2194///
2195/// **The ceiling is [`MAX_FUTURE_PUBLISHED_DAYS`], NOT the publication path's
2196/// clock-skew grace, and the asymmetry is deliberate.** There, a refused
2197/// `publishedAt` falls back to the record key's TID — the real write time, a
2198/// credible date — so a five-minute bound costs almost nothing. Here there is no
2199/// such fallback: refusing leaves the entry undated and the reader sees no date
2200/// at all. Five minutes is sized for skew between two clocks, while the ordinary
2201/// cause of a future `pubDate` is a local time stamped `+0000` (up to 14 hours
2202/// out, the widest real UTC offset) or a post scheduled a little ahead. Those are
2203/// dates worth keeping, and they stop being future on their own.
2204///
2205/// What the bound must prevent is a date that can never become past, because that
2206/// is what makes a row permanently unsweepable, un-evictable and first in the
2207/// list.
2208fn entry_time(e: &RawEntry) -> Option<String> {
2209 let ceiling = Utc::now() + chrono::Duration::days(MAX_FUTURE_PUBLISHED_DAYS);
2210 e.published
2211 .filter(|d| *d <= ceiling)
2212 .or_else(|| e.updated.filter(|d| *d <= ceiling))
2213 .map(fmt_time)
2214}
2215
2216/// How far ahead of now a feed may date an entry before [`entry_time`] refuses
2217/// the date and lets `fetched_at` stand in.
2218///
2219/// Two days: the widest real UTC offset is +14:00, so a local time mislabelled as
2220/// UTC lands inside this, as does a post scheduled slightly ahead. Both are dates
2221/// worth keeping, and both stop being future without help. Anything further is
2222/// refused, because a date that never becomes past is what makes a row
2223/// permanently unsweepable and permanently first in the reading list.
2224pub(crate) const MAX_FUTURE_PUBLISHED_DAYS: i64 = 2;
2225
2226/// Extract the plain string content of a feed [`Text`] node.
2227fn text_plain(t: &Text) -> String {
2228 t.content.trim().to_string()
2229}
2230
2231/// Sanitize hostile feed **HTML** with ammonia's whitelist cleaner. Applied to
2232/// every RSS/Atom entry body unconditionally, because every one of them is
2233/// markup.
2234///
2235/// **Not for plain text.** An earlier version of this comment claimed it was
2236/// "safe on plain text too (it will simply escape/strip as needed)". It is
2237/// not: `clean` PARSES its input, so a bare `<` in prose swallows the rest —
2238/// `"if x<y then z"` comes back as `"if x"`. That sentence is how a
2239/// plain-text field got run through here once already. Use
2240/// [`plain_text_to_html`].
2241pub(crate) fn sanitize_html(raw: &str) -> String {
2242 ammonia::clean(raw)
2243}
2244
2245/// Render **plain text** into the HTML the `content_html` column holds.
2246///
2247/// **Not [`sanitize_html`].** `ammonia::clean` parses its input as markup, so a
2248/// bare `<` in prose swallows the rest: measured here, `"if x<y then z"` comes
2249/// back as `"if x"`. That is correct for an RSS body, which IS markup, and
2250/// silent data loss for a field a lexicon defines as text. Escape first, then
2251/// add the only markup this needs — line breaks, which the column's consumer
2252/// renders as HTML and would otherwise collapse.
2253pub(crate) fn plain_text_to_html(raw: &str) -> String {
2254 let escaped = raw
2255 .replace('&', "&")
2256 .replace('<', "<")
2257 .replace('>', ">");
2258 // Safe by order: every `<` from the input is already `<` before this
2259 // adds a real tag.
2260 escaped.replace('\n', "<br>")
2261}
2262
2263/// **What one entry, or one feed row, may store per field (#205).**
2264///
2265/// Every field below arrives from someone else's server — an RSS document or a
2266/// publisher's `site.standard.document` record — and nothing between the wire
2267/// and SQLite used to shorten it. A title could be megabytes, then indexed, read
2268/// back and rendered into every list view of that feed.
2269///
2270/// Chosen from measurement, not guessed. Production's 4,389 entries on
2271/// 2026-10-03: title max 253 bytes (p99 148), url max 235 (p99 178), author max
2272/// 26, content_html max 86,969 (p99 17,535). Each bound is at least 8x the
2273/// largest real value, so no ordinary article is touched; `MAX_TITLE_BYTES` is
2274/// `site.standard.document`'s own `title.maxLength`.
2275///
2276/// **Truncated, never refused.** Dropping an article because one field is long
2277/// is the failure mode the retention work was careful to avoid.
2278pub(crate) const MAX_TITLE_BYTES: usize = 5_000;
2279/// See [`MAX_TITLE_BYTES`].
2280pub(crate) const MAX_AUTHOR_BYTES: usize = 1_000;
2281/// See [`MAX_TITLE_BYTES`]. A truncated URL is a broken link, which was judged
2282/// better than no link; at 35x the longest real one it should never happen.
2283pub(crate) const MAX_URL_BYTES: usize = 8_192;
2284/// See [`MAX_TITLE_BYTES`]. Applies to the STORED HTML, after sanitizing or
2285/// escaping — see [`sanitize_html_bounded`] for why that is the bound that matters.
2286pub(crate) const MAX_CONTENT_HTML_BYTES: usize = 2 * 1024 * 1024;
2287/// An entry id longer than this is replaced by a stable hash of the whole id
2288/// (see [`bound_guid`]): it is the dedup key, under a UNIQUE index.
2289pub(crate) const MAX_GUID_BYTES: usize = 2_048;
2290
2291/// The largest index `<= at` that is a character boundary of `s`.
2292fn floor_char_boundary(s: &str, at: usize) -> usize {
2293 if at >= s.len() {
2294 return s.len();
2295 }
2296 (0..=at).rev().find(|&i| s.is_char_boundary(i)).unwrap_or(0)
2297}
2298
2299/// `s` cut to at most `max` bytes, on a character boundary. Plain-text fields
2300/// only: cutting markup or an escaped string here could split a tag or an
2301/// entity, which is what [`sanitize_html_bounded`] and [`plain_text_to_html_bounded`] exist for.
2302pub(crate) fn bound_text(mut s: String, max: usize) -> String {
2303 let cut = floor_char_boundary(&s, max);
2304 s.truncate(cut);
2305 s
2306}
2307
2308/// Plain text escaped into `content_html`, cut so the **output** fits `max`.
2309///
2310/// **Exact, in one pass.** Escaping is a fixed size per character (`&` is
2311/// five bytes, `<` and `>` four, a newline `<br>` four, anything else its UTF-8
2312/// length), so the longest prefix whose escaped form fits is found by adding
2313/// those up — no rendering, no search. Three rounds of review found bugs in a
2314/// generic re-render search that this replaces (#224).
2315pub(crate) fn plain_text_to_html_bounded(raw: &str, max: usize) -> String {
2316 let mut size = 0usize;
2317 let mut cut = raw.len();
2318 for (i, c) in raw.char_indices() {
2319 let escaped = match c {
2320 '&' => 5,
2321 '<' | '>' | '\n' => 4,
2322 c => c.len_utf8(),
2323 };
2324 if size + escaped > max {
2325 cut = i;
2326 break;
2327 }
2328 size += escaped;
2329 }
2330 plain_text_to_html(&raw[..cut])
2331}
2332
2333/// Feed HTML sanitized into `content_html`, cut so the **output** fits `max`.
2334///
2335/// **Sanitize once, then cut the sanitized output, not the input.** Cutting
2336/// the input made the result depend on how much of it the sanitizer would
2337/// strip — a large `data:` image, `<style>` or unterminated comment — and the
2338/// searches that tried to account for that kept nothing, or ran for hours, in
2339/// review (#224). Sanitized HTML is already clean: cutting it on a character
2340/// boundary and sanitizing that prefix again only closes the tags the cut left
2341/// open, so the second pass usually grows it by little. The first cut leaves a
2342/// margin for that growth; deeply nested markup, whose closers can outgrow any
2343/// margin, falls through to a bounded bisection.
2344///
2345/// Cost: the first sanitize is the one every body always had; the bounded
2346/// passes — one, or at most 1 + [`SANITIZE_BOUND_ATTEMPTS`] — run on at most
2347/// `max` bytes of already-clean HTML, and only for a body over the bound.
2348pub(crate) fn sanitize_html_bounded(raw: &str, max: usize) -> String {
2349 let clean = sanitize_html(raw);
2350 if clean.len() <= max {
2351 return clean;
2352 }
2353 // First, the cut that almost always works: just under the bound, with a
2354 // margin for the closers the cut leaves open. When it fits, that is the
2355 // answer — within 1/64 of the bound, in one extra pass.
2356 let margin = (max / 64).max(64);
2357 let first = floor_char_boundary(&clean, max.saturating_sub(margin));
2358 let again = sanitize_html(&clean[..first]);
2359 if again.len() <= max {
2360 return again;
2361 }
2362 // **Then a bounded bisection, not a widening margin.** A closing tag is
2363 // longer than the tag it opens, so a cut through deeply nested markup can
2364 // grow past the bound by more than any fixed margin; widening the margin
2365 // 4x a round reached a cut of 0 and stored nothing where nearly all of it
2366 // fit (found in review). `lo` always fits (the empty prefix does), `hi`
2367 // never does, and the next probe is taken AFTER they move.
2368 let (mut lo, mut hi) = (0usize, first);
2369 let mut best = String::new();
2370 let resolution = (max / 1024).max(1);
2371 for _ in 0..SANITIZE_BOUND_ATTEMPTS {
2372 if hi - lo <= resolution {
2373 break;
2374 }
2375 let mut cut = floor_char_boundary(&clean, lo + (hi - lo) / 2);
2376 if cut <= lo {
2377 // A multi-byte character straddles the midpoint: step past it
2378 // rather than give up.
2379 cut = ceil_char_boundary(&clean, lo + 1);
2380 if cut >= hi {
2381 break;
2382 }
2383 }
2384 let out = sanitize_html(&clean[..cut]);
2385 if out.len() <= max {
2386 lo = cut;
2387 best = out;
2388 } else {
2389 hi = cut;
2390 }
2391 }
2392 best
2393}
2394
2395/// The smallest index `>= at` that is a character boundary of `s`.
2396fn ceil_char_boundary(s: &str, at: usize) -> usize {
2397 (at..=s.len())
2398 .find(|&i| s.is_char_boundary(i))
2399 .unwrap_or(s.len())
2400}
2401
2402/// The most bisection passes [`sanitize_html_bounded`] spends after its first
2403/// cut: enough to resolve a 2 MiB bound to about 1/1024 of it.
2404const SANITIZE_BOUND_ATTEMPTS: usize = 14;
2405
2406/// An entry id, or a stable stand-in for one too long to index.
2407///
2408/// The id is the dedup key under `UNIQUE (feed_id, guid)`, so truncating it
2409/// would merge distinct entries that share a long prefix. A hash of the WHOLE
2410/// id keeps them apart and keeps the same entry deduplicating across polls —
2411/// the same construction as [`stable_guid`].
2412pub(crate) fn bound_guid(guid: String) -> String {
2413 if guid.len() <= MAX_GUID_BYTES {
2414 return guid;
2415 }
2416 use std::hash::{Hash, Hasher};
2417 let mut h = dedup_hasher();
2418 guid.hash(&mut h);
2419 format!("featherreader:long-guid:{:016x}", h.finish())
2420}
2421
2422/// The hasher behind the stored dedup keys of [`bound_guid`] and
2423/// [`stable_guid`]: SipHash-1-3 with zero keys, from the `siphasher` crate.
2424///
2425/// **Not `std`'s `DefaultHasher`**, whose algorithm the standard library
2426/// documents as unspecified and free to change between Rust releases: a
2427/// toolchain bump could re-key every such entry, and each would be stored a
2428/// second time. `DefaultHasher` is SipHash-1-3 with zero keys today, so this
2429/// produces the values already stored; tests pin them.
2430fn dedup_hasher() -> siphasher::sip::SipHasher13 {
2431 siphasher::sip::SipHasher13::new_with_keys(0, 0)
2432}
2433
2434/// Format a chrono timestamp as RFC3339 (UTC, seconds precision) to match the
2435/// store's string columns.
2436pub(crate) fn fmt_time(dt: DateTime<Utc>) -> String {
2437 dt.to_rfc3339_opts(SecondsFormat::Secs, true)
2438}
2439
2440/// "Now" in the store's RFC3339 shape.
2441fn now_rfc3339() -> String {
2442 Utc::now().to_rfc3339_opts(SecondsFormat::Secs, true)
2443}
2444
2445/// A stable GUID derived from an entry's title + first link, for feeds that
2446/// supply neither an id nor a usable link id. Deterministic so re-fetches dedup.
2447fn stable_guid(e: &RawEntry) -> String {
2448 use std::hash::{Hash, Hasher};
2449 let mut h = dedup_hasher();
2450 e.title.as_ref().map(|t| t.content.as_str()).hash(&mut h);
2451 e.links
2452 .iter()
2453 .find(|l| is_primary_link(l))
2454 .map(|l| l.href.as_str())
2455 .hash(&mut h);
2456 e.summary.as_ref().map(|s| s.content.as_str()).hash(&mut h);
2457 format!("featherreader:synthetic:{:016x}", h.finish())
2458}
2459
2460/// Decode an HTTP header value to an owned `String`, dropping non-UTF-8 values.
2461fn header_str(v: Option<&reqwest::header::HeaderValue>) -> Option<String> {
2462 v.and_then(|h| h.to_str().ok()).map(str::to_string)
2463}
2464
2465/// Discover a feed URL from a site's HTML via
2466/// `<link rel="alternate" type="application/rss+xml|atom+xml" href="…">`.
2467///
2468/// Returns the first RSS/Atom autodiscovery link found, resolved against the
2469/// page URL if the `href` is relative. This is what lets a user paste a *site*
2470/// URL and have FeatherReader find the actual feed ("subscribe by URL").
2471/// Returns `None` if the HTML carries no autodiscovery link.
2472///
2473/// The `base` is the URL the HTML was fetched from, used to resolve relative
2474/// `href`s. Pass `None` to only accept absolute hrefs.
2475pub fn discover_feed(site_html: &str, base: Option<&Url>) -> Option<Url> {
2476 // Parse the HTML with html5ever (via ammonia's dependency graph is separate;
2477 // use a light hand-rolled scan over <link> tags to avoid a new dependency).
2478 // We look for <link ...> elements whose rel contains "alternate" and whose
2479 // type is an RSS/Atom feed media type, and take the href.
2480 for tag in link_tags(site_html) {
2481 let rel = attr(&tag, "rel").unwrap_or_default().to_ascii_lowercase();
2482 let typ = attr(&tag, "type").unwrap_or_default().to_ascii_lowercase();
2483 let is_feed_type = typ.contains("application/rss+xml")
2484 || typ.contains("application/atom+xml")
2485 || typ.contains("application/feed+json")
2486 || typ.contains("application/json");
2487 // rel="alternate" is the standard; be lenient and also accept a bare
2488 // feed type with any rel, but require the feed media type either way.
2489 let rel_ok = rel.split_whitespace().any(|r| r == "alternate") || rel.is_empty();
2490 if is_feed_type && rel_ok {
2491 if let Some(href) = attr(&tag, "href") {
2492 let href = href.trim();
2493 if href.is_empty() {
2494 continue;
2495 }
2496 // Absolute URL wins directly; otherwise resolve against `base`.
2497 // Either way, only http(s): the href is publisher-controlled and
2498 // `Url::parse` accepts any scheme, so this is where an `at://`
2499 // (or `file:`, `javascript:`) alternate would otherwise become
2500 // the URL the add path stores — after its input gate has run.
2501 // Skip, don't stop: a later real feed link still wins.
2502 let resolved = match Url::parse(href) {
2503 Ok(u) => Some(u),
2504 Err(_) => base.and_then(|b| b.join(href).ok()),
2505 };
2506 match resolved {
2507 Some(u) if matches!(u.scheme(), "http" | "https") => return Some(u),
2508 _ => continue,
2509 }
2510 }
2511 }
2512 }
2513 None
2514}
2515
2516/// Extract the raw text of every `<link ...>` tag (self-closing or not) from an
2517/// HTML string. A deliberately small, allocation-light scan — feed
2518/// autodiscovery does not need a full DOM, and avoiding one keeps the dependency
2519/// surface minimal (design bias: boring, small-dependency).
2520fn link_tags(html: &str) -> Vec<String> {
2521 let mut out = Vec::new();
2522 let bytes = html.as_bytes();
2523 let lower = html.to_ascii_lowercase();
2524 let mut search_from = 0usize;
2525 while let Some(rel_idx) = lower[search_from..].find("<link") {
2526 let start = search_from + rel_idx;
2527 // Ensure it's a tag boundary ("<link" followed by whitespace, '>' or '/').
2528 let after = bytes.get(start + 5).copied();
2529 let boundary = matches!(after, Some(b) if b == b' ' || b == b'\t' || b == b'\n' || b == b'\r' || b == b'>' || b == b'/');
2530 if !boundary {
2531 search_from = start + 5;
2532 continue;
2533 }
2534 // Find the closing '>' for this tag.
2535 if let Some(end_rel) = html[start..].find('>') {
2536 let end = start + end_rel;
2537 out.push(html[start..=end].to_string());
2538 search_from = end + 1;
2539 } else {
2540 break;
2541 }
2542 }
2543 out
2544}
2545
2546/// Read an attribute value from a single tag string, handling both single- and
2547/// double-quoted values. Case-insensitive attribute name match.
2548fn attr(tag: &str, name: &str) -> Option<String> {
2549 let lower = tag.to_ascii_lowercase();
2550 let needle = format!("{name}=");
2551 let mut from = 0usize;
2552 while let Some(rel) = lower[from..].find(&needle) {
2553 let name_start = from + rel;
2554 // Guard against matching a suffix of a longer attribute name
2555 // (e.g. matching "type=" inside "mytype="): the char before must be a
2556 // tag/whitespace boundary.
2557 let ok_prefix = name_start == 0
2558 || matches!(
2559 tag.as_bytes().get(name_start - 1),
2560 Some(b' ') | Some(b'\t') | Some(b'\n') | Some(b'\r') | Some(b'<')
2561 );
2562 let val_start = name_start + needle.len();
2563 if !ok_prefix {
2564 from = val_start;
2565 continue;
2566 }
2567 let rest = &tag[val_start..];
2568 let quote = rest.chars().next();
2569 let value = match quote {
2570 Some('"') => rest[1..].split('"').next(),
2571 Some('\'') => rest[1..].split('\'').next(),
2572 // Unquoted: read up to whitespace, '>' or '/'.
2573 _ => rest
2574 .split(|c: char| c.is_whitespace() || c == '>' || c == '/')
2575 .next(),
2576 };
2577 return value.map(str::to_string);
2578 }
2579 None
2580}
2581
2582#[cfg(test)]
2583mod tests {
2584 use super::*;
2585
2586 const RSS_SAMPLE: &str = r#"<?xml version="1.0" encoding="UTF-8"?>
2587<rss version="2.0">
2588 <channel>
2589 <title>Example RSS Feed</title>
2590 <link>https://example.com/</link>
2591 <description>An example feed for tests</description>
2592 <item>
2593 <title>First post</title>
2594 <link>https://example.com/first</link>
2595 <guid>https://example.com/first</guid>
2596 <author>alice@example.com (Alice)</author>
2597 <pubDate>Fri, 10 Jul 2026 08:00:00 GMT</pubDate>
2598 <description><![CDATA[<p>Hello <b>world</b>.</p><script>alert('xss')</script><img src="x" onerror="alert(1)">]]></description>
2599 </item>
2600 <item>
2601 <title>Second post</title>
2602 <link>https://example.com/second</link>
2603 <guid>guid-second</guid>
2604 <pubDate>Sat, 11 Jul 2026 08:00:00 GMT</pubDate>
2605 <description><![CDATA[<a href="javascript:alert(1)">click</a><a href="https://ok.example/">ok</a>]]></description>
2606 </item>
2607 </channel>
2608</rss>"#;
2609
2610 const ATOM_SAMPLE: &str = r#"<?xml version="1.0" encoding="utf-8"?>
2611<feed xmlns="http://www.w3.org/2005/Atom">
2612 <title>Example Atom Feed</title>
2613 <link rel="alternate" href="https://atom.example.com/"/>
2614 <link rel="self" href="https://atom.example.com/feed.xml"/>
2615 <id>urn:uuid:feed-1</id>
2616 <updated>2026-07-11T08:00:00Z</updated>
2617 <entry>
2618 <title>Atom entry</title>
2619 <id>urn:uuid:entry-1</id>
2620 <link rel="alternate" href="https://atom.example.com/a"/>
2621 <author><name>Bob</name></author>
2622 <updated>2026-07-11T08:00:00Z</updated>
2623 <content type="html"><![CDATA[<p>Safe <em>text</em>.</p><script>steal()</script><iframe src="evil"></iframe>]]></content>
2624 </entry>
2625</feed>"#;
2626
2627 /// [`ATOM_SAMPLE`] with the `self` link before the `alternate` one.
2628 const ATOM_SELF_FIRST: &str = r#"<?xml version="1.0" encoding="utf-8"?>
2629<feed xmlns="http://www.w3.org/2005/Atom">
2630 <title>Example Atom Feed</title>
2631 <link rel="self" href="https://atom.example.com/feed.xml"/>
2632 <link rel="alternate" href="https://atom.example.com/"/>
2633 <id>urn:uuid:feed-1</id>
2634 <updated>2026-07-11T08:00:00Z</updated>
2635 <entry>
2636 <title>Atom entry</title>
2637 <id>urn:uuid:entry-1</id>
2638 <link rel="alternate" href="https://atom.example.com/a"/>
2639 <author><name>Bob</name></author>
2640 <updated>2026-07-11T08:00:00Z</updated>
2641 <content type="html"><![CDATA[<p>Safe <em>text</em>.</p><script>steal()</script><iframe src="evil"></iframe>]]></content>
2642 </entry>
2643</feed>"#;
2644
2645 // ---- #205: what one entry may store, per field ----------------------
2646
2647 /// An RSS document carrying one item with exactly these fields.
2648 fn rss_with_fields(title: &str, link: &str, author: &str, body: &str, guid: &str) -> String {
2649 format!(
2650 r#"<?xml version="1.0"?><rss version="2.0" xmlns:dc="http://purl.org/dc/elements/1.1/"><channel>
2651<title>{title}</title><link>https://example.com/{link}</link>
2652<item><title>{title}</title><link>https://example.com/{link}</link><guid>{guid}</guid>
2653<dc:creator>{author}</dc:creator><description><![CDATA[{body}]]></description></item>
2654</channel></rss>"#
2655 )
2656 }
2657
2658 #[test]
2659 fn bound_text_cuts_on_a_character_boundary() {
2660 // 'é' is two bytes, so an odd limit lands mid-character.
2661 let cut = bound_text("é".repeat(100), 51);
2662 assert!(cut.len() <= 51, "not bounded: {} bytes", cut.len());
2663 assert_eq!(cut, "é".repeat(25), "cut too short or mid-character");
2664 assert_eq!(
2665 bound_text("short".into(), 51),
2666 "short",
2667 "a short value changed"
2668 );
2669 }
2670
2671 #[test]
2672 fn rendered_content_fits_even_when_rendering_grows_it() {
2673 // Escaping turns each `&` into `&` — five bytes from one.
2674 let out = plain_text_to_html_bounded(&"&".repeat(1_000), 100);
2675 assert!(
2676 out.len() <= 100,
2677 "escaped output not bounded: {} bytes",
2678 out.len()
2679 );
2680 assert!(!out.is_empty(), "bounded to nothing");
2681 assert_eq!(
2682 out.matches("&").count() * 5,
2683 out.len(),
2684 "cut mid-entity: {out}"
2685 );
2686 }
2687
2688 #[test]
2689 fn an_rss_items_text_fields_are_bounded() {
2690 let big = "x".repeat(100_000);
2691 let xml = rss_with_fields(&big, &big, &big, "body", "id-1");
2692 let parsed = parse_feed(xml.as_bytes()).unwrap();
2693 let e = normalize_entry(&parsed.entries[0]);
2694 assert!(e.title.as_ref().unwrap().len() <= MAX_TITLE_BYTES, "title");
2695 assert!(e.url.as_ref().unwrap().len() <= MAX_URL_BYTES, "url");
2696 assert!(
2697 e.author.as_ref().unwrap().len() <= MAX_AUTHOR_BYTES,
2698 "author"
2699 );
2700 let (title, site) = feed_metadata(&parsed);
2701 assert!(title.unwrap().len() <= MAX_TITLE_BYTES, "feed title");
2702 assert!(site.unwrap().len() <= MAX_URL_BYTES, "feed site url");
2703 }
2704
2705 #[test]
2706 fn an_rss_body_is_bounded_and_still_well_formed() {
2707 // Long enough to need cutting, with markup straddling the cut.
2708 let body = format!("<p>{}<b>tail</b></p>", "a".repeat(MAX_CONTENT_HTML_BYTES));
2709 let xml = rss_with_fields("t", "l", "a", &body, "id-2");
2710 let parsed = parse_feed(xml.as_bytes()).unwrap();
2711 let html = normalize_entry(&parsed.entries[0]).content_html.unwrap();
2712 assert!(
2713 html.len() <= MAX_CONTENT_HTML_BYTES,
2714 "body not bounded: {}",
2715 html.len()
2716 );
2717 assert_eq!(
2718 sanitize_html(&html),
2719 html,
2720 "the stored body is not well-formed sanitized HTML"
2721 );
2722 }
2723
2724 /// Review of #224: cutting the INPUT first threw away content that would
2725 /// have fit. A large inline `data:` image, which ammonia strips anyway,
2726 /// ahead of the article left the cut ending inside the image tag, and the
2727 /// article was stored as an empty body.
2728 #[test]
2729 fn content_the_sanitizer_strips_does_not_count_against_the_bound() {
2730 const MAX: usize = 64 * 1024;
2731 let body = format!(
2732 r#"<p><img src="data:image/png;base64,{}"></p><p>the article</p>"#,
2733 "A".repeat(MAX + 16 * 1024)
2734 );
2735 let html = sanitize_html_bounded(&body, MAX);
2736 assert!(
2737 html.contains("the article"),
2738 "the article was cut away: {} bytes kept",
2739 html.len()
2740 );
2741 assert!(html.len() <= MAX);
2742 }
2743
2744 /// Review of #224 (second round): the re-cut loop shrank its input by a
2745 /// few bytes a round when the bytes beyond the cut were ones the sanitizer
2746 /// strips anyway — measured at ~32 bytes/round, hours for one entry, on the
2747 /// poller's async task. The number of renders must be bounded.
2748 /// Review of #224: a cut into bytes the sanitizer strips anyway (here an
2749 /// unterminated comment) crept a few bytes a round, for hours. Bounding
2750 /// works on the SANITIZED output now, so stripped input costs nothing.
2751 #[test]
2752 fn content_cut_beside_stripped_bytes_still_keeps_what_fits() {
2753 const MAX: usize = 64 * 1024;
2754 let body = format!("<p>{}</p><!--{}", "a".repeat(MAX + 8), "x".repeat(16 * MAX));
2755 let html = sanitize_html_bounded(&body, MAX);
2756 assert!(html.len() <= MAX);
2757 assert!(
2758 html.len() > MAX - MAX / 16,
2759 "kept far less than fits: {}",
2760 html.len()
2761 );
2762 }
2763
2764 /// Also from that review: a cut landing inside a stripped prefix stored an
2765 /// empty body when an article that fits followed it.
2766 /// Also from review: a cut landing inside a stripped prefix stored an
2767 /// empty body when an article that fits followed it.
2768 #[test]
2769 fn a_stripped_prefix_does_not_leave_an_empty_body() {
2770 const MAX: usize = 64 * 1024;
2771 let body = format!(
2772 r#"<p><img src="data:image/png;base64,{}"></p><p>{}</p>"#,
2773 "A".repeat(3 * MAX),
2774 "&".repeat(MAX)
2775 );
2776 let html = sanitize_html_bounded(&body, MAX);
2777 assert!(html.len() <= MAX);
2778 assert!(
2779 html.len() > MAX - MAX / 16,
2780 "nothing like what fits was kept: {} bytes",
2781 html.len()
2782 );
2783 }
2784
2785 /// Third review: a multi-byte character where the search's stale probe
2786 /// landed ended it early, keeping 0 bytes where ~2 MiB fit.
2787 #[test]
2788 fn a_stripped_multibyte_prefix_does_not_leave_an_empty_body() {
2789 const MAX: usize = 64 * 1024;
2790 let body = format!(
2791 "<!--{}--><p>{}</p>",
2792 "漢".repeat(MAX),
2793 "&".repeat(MAX * 3 / 10)
2794 );
2795 let html = sanitize_html_bounded(&body, MAX);
2796 assert!(html.len() <= MAX);
2797 assert!(
2798 html.len() > MAX - MAX / 16,
2799 "kept {} of ~{MAX} that fits",
2800 html.len()
2801 );
2802 }
2803
2804 /// Fourth review of #224: a closing tag is longer than the tag it closes,
2805 /// so a cut through nested markup grew past the bound on re-sanitizing,
2806 /// and the widening margin jumped straight to a cut of 0 — an empty body
2807 /// where nearly all of it fit.
2808 #[test]
2809 fn nested_markup_is_cut_not_emptied() {
2810 const MAX: usize = 64 * 1024;
2811 let html = sanitize_html_bounded(&"<span>".repeat(16_384), MAX);
2812 assert!(html.len() <= MAX);
2813 assert!(
2814 html.len() > MAX / 3,
2815 "kept {} of ~{MAX} that fits",
2816 html.len()
2817 );
2818
2819 let mixed = format!("{}{}", "t".repeat(MAX * 3 / 4), "<span>".repeat(MAX / 8));
2820 let html = sanitize_html_bounded(&mixed, MAX);
2821 assert!(html.len() <= MAX);
2822 assert!(
2823 html.len() > MAX - MAX / 16,
2824 "kept {} of ~{MAX} that fits",
2825 html.len()
2826 );
2827 }
2828
2829 /// The plain-text bound is exact: escaping is linear, so the longest
2830 /// fitting prefix is found in one pass, never by search.
2831 #[test]
2832 fn the_plain_text_bound_is_exact() {
2833 const MAX: usize = 64 * 1024;
2834 let raw = format!("{}{}", "漢".repeat(MAX / 4), "&".repeat(MAX / 10));
2835 let out = plain_text_to_html_bounded(&raw, MAX);
2836 assert!(out.len() <= MAX);
2837 // The next character would not have fitted: '&' escapes to 5 bytes.
2838 assert!(out.len() > MAX - 5, "kept {} of {MAX}", out.len());
2839 assert!(
2840 raw.starts_with(&out.replace("&", "&")),
2841 "not a prefix of the input"
2842 );
2843 }
2844
2845 #[test]
2846 fn an_overlong_rss_guid_becomes_a_stable_short_one() {
2847 let long = "g".repeat(10_000);
2848 let xml = rss_with_fields("t", "l", "a", "b", &long);
2849 let parsed = parse_feed(xml.as_bytes()).unwrap();
2850 let a = normalize_entry(&parsed.entries[0]).guid;
2851 let b = normalize_entry(&parsed.entries[0]).guid;
2852 assert!(a.len() <= MAX_GUID_BYTES, "guid not bounded: {}", a.len());
2853 assert_eq!(
2854 a, b,
2855 "the stand-in is not stable, so the entry would duplicate"
2856 );
2857 let other = rss_with_fields("t", "l", "a", "b", &format!("{long}h"));
2858 let other = parse_feed(other.as_bytes()).unwrap();
2859 assert_ne!(
2860 a,
2861 normalize_entry(&other.entries[0]).guid,
2862 "two ids collapsed into one"
2863 );
2864 }
2865
2866 #[test]
2867 fn an_ordinary_long_article_is_untouched() {
2868 // The largest body production held on 2026-10-03 was 86,969 bytes.
2869 let body = format!("<p>{}</p>", "word ".repeat(18_000));
2870 let xml = rss_with_fields(
2871 "A normal title",
2872 "post",
2873 "Author",
2874 &body,
2875 "https://example.com/post",
2876 );
2877 let parsed = parse_feed(xml.as_bytes()).unwrap();
2878 let e = normalize_entry(&parsed.entries[0]);
2879 assert_eq!(e.content_html.unwrap(), sanitize_html(&body));
2880 assert_eq!(e.title.as_deref(), Some("A normal title"));
2881 assert_eq!(e.guid, "https://example.com/post");
2882 }
2883
2884 /// **A stated date in the future is discarded, not stored and not clamped.**
2885 ///
2886 /// `entry_time` was `e.published.or(e.updated)` with no ceiling, so an item
2887 /// dated in the year 2999 was stored verbatim and then became permanent:
2888 /// both retention sweeps test `COALESCE(published, fetched_at) < cutoff` and
2889 /// a future date is never less than either; the per-feed keep-set orders on
2890 /// the same expression `DESC`, where it is rank one forever; and every list
2891 /// view orders on `published DESC`, where it sits at the top. One item in one
2892 /// feed, there for good. A publisher with a broken clock does this by
2893 /// accident.
2894 ///
2895 /// Discarded rather than clamped to now, which is the rule `#186`
2896 /// established on the atproto side and the reasoning transfers exactly: the
2897 /// entry upsert refreshes `published` on every poll but stamps `fetched_at`
2898 /// once, so a value derived from the current clock is rewritten every cycle
2899 /// and the row can never age at all. Clamping moves the defect. Falling back
2900 /// to undated lets `fetched_at` date it, and that holds still.
2901 ///
2902 /// The ceiling is [`MAX_FUTURE_PUBLISHED_DAYS`] (two days), deliberately
2903 /// looser than the publication path's five-minute grace: a publication can
2904 /// fall back to its record key's TID, and a feed has no such fallback.
2905 #[test]
2906 fn a_future_dated_rss_item_is_stored_undated_rather_than_dated_in_2999() {
2907 let future = r#"<?xml version="1.0"?>
2908<rss version="2.0"><channel><title>Clock</title><link>https://clock.example/</link>
2909<item><title>From the future</title><link>https://clock.example/1</link>
2910<guid>https://clock.example/1</guid>
2911<pubDate>Sat, 01 Jan 2999 00:00:00 GMT</pubDate></item>
2912</channel></rss>"#;
2913 let parsed = parse_feed(future.as_bytes()).expect("should parse");
2914 // The fixture is only meaningful if feed-rs actually read the date.
2915 assert!(
2916 parsed.entries[0].published.is_some(),
2917 "the fixture's pubDate did not parse, so this test proves nothing",
2918 );
2919
2920 let e = normalize_entry(&parsed.entries[0]);
2921 assert_eq!(
2922 e.published, None,
2923 "a year-2999 date was stored, which makes the row unsweepable, \
2924 un-evictable and permanently first in the reading list",
2925 );
2926
2927 // And the other direction: an ordinary past date must survive, or this
2928 // would be satisfied by discarding every date.
2929 let past = future.replace("01 Jan 2999", "01 Jan 2020");
2930 let parsed = parse_feed(past.as_bytes()).expect("should parse");
2931 let e = normalize_entry(&parsed.entries[0]);
2932 assert!(
2933 e.published
2934 .as_deref()
2935 .is_some_and(|p| p.starts_with("2020")),
2936 "an ordinary past date was discarded: {:?}",
2937 e.published,
2938 );
2939
2940 // **A merely MISLABELLED date must survive.** The ordinary cause of a
2941 // future `pubDate` is a local time stamped `+0000` — up to 14 hours out,
2942 // not a clock a few minutes fast. Refusing those would leave real
2943 // articles undated and dateless on screen, which is why the bound is two
2944 // days rather than the publication path's five-minute skew grace.
2945 let soon = (Utc::now() + chrono::Duration::hours(14)).to_rfc2822();
2946 let near = future.replace("Sat, 01 Jan 2999 00:00:00 GMT", &soon);
2947 let parsed = parse_feed(near.as_bytes()).expect("should parse");
2948 assert!(
2949 parsed.entries[0].published.is_some(),
2950 "the mislabelled-date fixture did not parse",
2951 );
2952 let e = normalize_entry(&parsed.entries[0]);
2953 assert!(
2954 e.published.is_some(),
2955 "a date 14 hours ahead — the widest real UTC offset — was refused, \
2956 so a timezone-mislabelled article loses its date entirely",
2957 );
2958
2959 // And the bound still bounds: a month out is refused.
2960 let far = (Utc::now() + chrono::Duration::days(30)).to_rfc2822();
2961 let month = future.replace("Sat, 01 Jan 2999 00:00:00 GMT", &far);
2962 let parsed = parse_feed(month.as_bytes()).expect("should parse");
2963 let e = normalize_entry(&parsed.entries[0]);
2964 assert_eq!(
2965 e.published, None,
2966 "a date a month ahead was kept, so the row leads the list for a month",
2967 );
2968
2969 // **Each candidate is judged separately, not the winner of `or`.**
2970 //
2971 // Atom specifically: for RSS 2 `feed-rs` copies `published` into
2972 // `updated` when `updated` is absent, so the second candidate holds the
2973 // same bogus value and the fall-through cannot help. Only a format
2974 // carrying two independent dates exercises this.
2975 let both = r#"<?xml version="1.0"?>
2976<feed xmlns="http://www.w3.org/2005/Atom"><title>Clock</title>
2977<entry><title>Mixed</title><id>https://clock.example/2</id>
2978<link href="https://clock.example/2"/>
2979<published>2999-01-01T00:00:00Z</published>
2980<updated>2020-06-01T00:00:00Z</updated></entry></feed>"#;
2981 let parsed = parse_feed(both.as_bytes()).expect("should parse");
2982 assert!(
2983 parsed.entries[0].published.is_some() && parsed.entries[0].updated.is_some(),
2984 "the fixture needs BOTH dates parsed for this case to mean anything",
2985 );
2986 let e = normalize_entry(&parsed.entries[0]);
2987 assert!(
2988 e.published
2989 .as_deref()
2990 .is_some_and(|p| p.starts_with("2020")),
2991 "a credible `updated` was discarded along with a bogus `published`, \
2992 leaving the entry undated: {:?}",
2993 e.published,
2994 );
2995 }
2996
2997 /// Parse a static RSS sample through feed-rs + our normalize/sanitize path
2998 /// (no network) and assert the entries come out sanitized and well-shaped.
2999 #[test]
3000 fn rss_parses_and_sanitizes() {
3001 let parsed = parse_feed(RSS_SAMPLE.as_bytes()).expect("RSS should parse");
3002 assert_eq!(
3003 parsed.title.as_ref().map(text_plain).as_deref(),
3004 Some("Example RSS Feed")
3005 );
3006 assert_eq!(parsed.entries.len(), 2);
3007
3008 let (title, site) = feed_metadata(&parsed);
3009 assert_eq!(title.as_deref(), Some("Example RSS Feed"));
3010 assert_eq!(site.as_deref(), Some("https://example.com/"));
3011
3012 let e0 = normalize_entry(&parsed.entries[0]);
3013 assert_eq!(e0.guid, "https://example.com/first");
3014 assert_eq!(e0.title.as_deref(), Some("First post"));
3015 assert_eq!(e0.url.as_deref(), Some("https://example.com/first"));
3016 assert!(e0.published.is_some());
3017 let html0 = e0.content_html.expect("content present");
3018 // Sanitized: benign markup kept, script + onerror stripped.
3019 assert!(html0.contains("Hello"));
3020 assert!(html0.contains("<b>world</b>") || html0.contains("<b>"));
3021 assert!(!html0.to_ascii_lowercase().contains("<script"));
3022 assert!(!html0.to_ascii_lowercase().contains("onerror"));
3023 assert!(!html0.to_ascii_lowercase().contains("alert"));
3024
3025 // Second entry: javascript: URL scrubbed, safe link kept.
3026 let e1 = normalize_entry(&parsed.entries[1]);
3027 assert_eq!(e1.guid, "guid-second");
3028 let html1 = e1.content_html.expect("content present");
3029 assert!(!html1.to_ascii_lowercase().contains("javascript:"));
3030 assert!(html1.contains("https://ok.example/"));
3031 }
3032
3033 /// Same, for an Atom sample: alternate link is the site URL, dangerous
3034 /// elements are stripped from entry content.
3035 /// **`rel="alternate"` is preferred over a `rel="self"` listed FIRST.** In
3036 /// `ATOM_SAMPLE` the alternate link is already first, so "prefer alternate"
3037 /// and "take the first link" were indistinguishable; the selection could
3038 /// be replaced by `links.first()` with the suite green. A feed listing
3039 /// `self` first — very common in Atom — would store the feed XML URL as
3040 /// the subscription's `siteUrl`, published to the reader's PDS.
3041 #[test]
3042 fn atom_prefers_alternate_over_a_self_link_listed_first() {
3043 let parsed = parse_feed(ATOM_SELF_FIRST.as_bytes()).expect("Atom should parse");
3044 let (title, site) = feed_metadata(&parsed);
3045 assert_eq!(title.as_deref(), Some("Example Atom Feed"));
3046 // alternate link preferred over rel="self".
3047 assert_eq!(site.as_deref(), Some("https://atom.example.com/"));
3048 }
3049
3050 #[test]
3051 fn atom_parses_and_sanitizes() {
3052 let parsed = parse_feed(ATOM_SAMPLE.as_bytes()).expect("Atom should parse");
3053 let (title, site) = feed_metadata(&parsed);
3054 assert_eq!(title.as_deref(), Some("Example Atom Feed"));
3055 // alternate link preferred over rel="self".
3056 assert_eq!(site.as_deref(), Some("https://atom.example.com/"));
3057
3058 assert_eq!(parsed.entries.len(), 1);
3059 let e = normalize_entry(&parsed.entries[0]);
3060 assert_eq!(e.guid, "urn:uuid:entry-1");
3061 assert_eq!(e.title.as_deref(), Some("Atom entry"));
3062 assert_eq!(e.author.as_deref(), Some("Bob"));
3063 assert_eq!(e.url.as_deref(), Some("https://atom.example.com/a"));
3064 let html = e.content_html.expect("content present");
3065 assert!(html.contains("Safe"));
3066 assert!(!html.to_ascii_lowercase().contains("<script"));
3067 assert!(!html.to_ascii_lowercase().contains("<iframe"));
3068 }
3069
3070 /// **The rkey obeys all of atproto's record-key rules, not just the
3071 /// charset.** `.` and `..` are reserved and the length is 1..=512; the
3072 /// charset alone admitted both and a 10 000-character key, into a UNIQUE
3073 /// column and the user's public PDS. The rule is the one the repo's own
3074 /// TID tests already state.
3075 #[test]
3076 fn an_rkey_must_obey_atprotos_length_and_dot_rules() {
3077 let uri = |rkey: &str| {
3078 format!("at://did:plc:ohutz6x5acjmpuulp3x7wxxc/site.standard.publication/{rkey}")
3079 };
3080 assert!(
3081 !is_storable_feed_url(&uri("."), true),
3082 "`.` is a reserved rkey"
3083 );
3084 assert!(
3085 !is_storable_feed_url(&uri(".."), true),
3086 "`..` is a reserved rkey"
3087 );
3088 assert!(
3089 !is_storable_feed_url(&uri(&"a".repeat(513)), true),
3090 "an rkey over 512 bytes was accepted"
3091 );
3092 assert!(
3093 is_storable_feed_url(&uri(&"a".repeat(512)), true),
3094 "an rkey of exactly 512 bytes is valid"
3095 );
3096 assert!(is_storable_feed_url(&uri("3lab2c4d5e6f7g8h"), true));
3097 }
3098
3099 /// **Every `KNOWN_PROVIDERS` row is pinned by its own reason.**
3100 ///
3101 /// The provider test used URLs that the generic heuristics catch on their
3102 /// own — Substack's `/feed/private/` is also a path marker, Patreon's
3103 /// `?auth=…` an opaque secret key — and never asserted the reason. With
3104 /// 17 of 18 rows deleted, the suite stayed green. The provider layer runs
3105 /// FIRST so its specific reason wins; asserting the reason pins each row
3106 /// even where a generic rule would still refuse the URL. Values are kept
3107 /// short and plain so the generic query rule (`value_is_opaque`) does not
3108 /// fire — most of these are refused by the provider row alone.
3109 #[test]
3110 fn each_known_provider_is_caught_by_its_own_row() {
3111 for (url, reason) in [
3112 (
3113 "https://author.substack.com/feed/private/x",
3114 "Substack private feed",
3115 ),
3116 (
3117 "https://www.patreon.com/rss/creator?auth=ab",
3118 "Patreon member feed",
3119 ),
3120 ("https://blog.ghost.io/rss/?uuid=x", "Ghost members feed"),
3121 (
3122 "https://buttondown.email/me/rss?token=x",
3123 "Buttondown premium feed",
3124 ),
3125 (
3126 "https://buttondown.com/me/rss?token=x",
3127 "Buttondown premium feed",
3128 ),
3129 (
3130 "https://rss.beehiiv.com/feeds/x.xml?token=x",
3131 "Beehiiv premium feed",
3132 ),
3133 (
3134 "https://example.memberful.com/feed",
3135 "Memberful members feed",
3136 ),
3137 ("https://example.pico.link/feed", "Pico member feed"),
3138 ("https://steadyhq.com/rss/example", "Steady member feed"),
3139 (
3140 "https://example.supercast.com/feed",
3141 "Supercast private podcast",
3142 ),
3143 (
3144 "https://example.supercast.tech/feed",
3145 "Supercast private podcast",
3146 ),
3147 (
3148 "https://example.supportingcast.fm/feed",
3149 "Supporting Cast private podcast",
3150 ),
3151 (
3152 "https://feeds.redcircle.com/x?private=1",
3153 "RedCircle private podcast",
3154 ),
3155 (
3156 "https://feeds.megaphone.fm/x?token=x",
3157 "Megaphone private podcast",
3158 ),
3159 (
3160 "https://feeds.acast.com/public/shows/x?token=x",
3161 "Acast+ private podcast",
3162 ),
3163 (
3164 "https://omny.fm/shows/x/playlists/podcast.rss?token=x",
3165 "Omny private podcast",
3166 ),
3167 (
3168 "https://podcasts.apple.com/feed/x?token=x",
3169 "Apple subscriber podcast",
3170 ),
3171 (
3172 "https://anchor.spotify.com/s/x/podcast/rss?token=x",
3173 "Spotify subscriber podcast",
3174 ),
3175 ] {
3176 match classify_feed_privacy(url) {
3177 FeedPrivacy::Private(r) => {
3178 assert_eq!(r, reason, "{url} was refused by another rule")
3179 }
3180 FeedPrivacy::Public => panic!("{url} was not refused at all"),
3181 }
3182 }
3183 }
3184
3185 /// **Plain text is escaped, not sanitised.** `ammonia::clean` parses its
3186 /// input as markup, so a `<` in prose swallows everything after it:
3187 /// measured in this tree, `"if x<y then z"` becomes `"if x"`. That is the
3188 /// right function for an RSS body (which IS markup) and exactly the wrong
3189 /// one for a field the lexicon defines as plain text — it silently deletes
3190 /// the reader's content.
3191 #[test]
3192 fn plain_text_is_escaped_rather_than_swallowed() {
3193 assert_eq!(
3194 plain_text_to_html("Vec<String> is a type"),
3195 "Vec<String> is a type"
3196 );
3197 assert_eq!(plain_text_to_html("if x<y then z"), "if x<y then z");
3198 assert_eq!(plain_text_to_html("a & b"), "a & b");
3199 // Line structure survives into a field rendered as HTML.
3200 assert_eq!(plain_text_to_html("one\ntwo"), "one<br>two");
3201 // And it is still safe: the escaping happens before any markup is added.
3202 let hostile = plain_text_to_html("<script>alert(1)</script>");
3203 assert!(!hostile.contains("<script"), "{hostile}");
3204 }
3205
3206 #[test]
3207 fn discover_finds_rss_link() {
3208 let html = r#"<!doctype html><html><head>
3209 <title>Blog</title>
3210 <link rel="stylesheet" href="/style.css">
3211 <link rel="alternate" type="application/rss+xml" title="RSS" href="/feed.xml">
3212 </head><body>hi</body></html>"#;
3213 let base = Url::parse("https://blog.example.com/").unwrap();
3214 let found = discover_feed(html, Some(&base)).expect("should discover feed");
3215 assert_eq!(found.as_str(), "https://blog.example.com/feed.xml");
3216 }
3217
3218 #[test]
3219 fn discover_finds_atom_absolute_link() {
3220 let html = r#"<head><link rel="alternate" type="application/atom+xml" href="https://x.example/atom"></head>"#;
3221 let found = discover_feed(html, None).expect("should discover absolute feed");
3222 assert_eq!(found.as_str(), "https://x.example/atom");
3223 }
3224
3225 /// **Autodiscovery only ever yields an http(s) URL.**
3226 ///
3227 /// The href is publisher-controlled and `Url::parse` accepts any scheme, so
3228 /// a page could hand the add path an `at://` publication URI (or anything
3229 /// else) that the user never typed — and the add path's input gate has
3230 /// already run by then. A non-http(s) alternate is skipped, not returned,
3231 /// so a later real feed link still wins.
3232 #[test]
3233 fn discover_skips_a_non_http_alternate() {
3234 let at_link = r#"<link rel="alternate" type="application/rss+xml" href="at://did:plc:ohutz6x5acjmpuulp3x7wxxc/site.standard.publication/3lab2c4d5e6f7g8h">"#;
3235 assert!(
3236 discover_feed(&format!("<head>{at_link}</head>"), None).is_none(),
3237 "an at:// alternate was handed back as a feed URL"
3238 );
3239 let ftp_link =
3240 r#"<link rel="alternate" type="application/atom+xml" href="ftp://x.example/atom">"#;
3241 assert!(discover_feed(&format!("<head>{ftp_link}</head>"), None).is_none());
3242
3243 let real =
3244 r#"<link rel="alternate" type="application/atom+xml" href="https://x.example/atom">"#;
3245 let found = discover_feed(&format!("<head>{at_link}{real}</head>"), None)
3246 .expect("the http(s) link after a skipped one must still be found");
3247 assert_eq!(found.as_str(), "https://x.example/atom");
3248 }
3249
3250 #[test]
3251 fn discover_returns_none_without_feed_link() {
3252 let html =
3253 r#"<head><link rel="stylesheet" href="/s.css"><link rel="icon" href="/f.ico"></head>"#;
3254 assert!(discover_feed(html, None).is_none());
3255 }
3256
3257 #[test]
3258 fn synthetic_guid_is_stable_and_dedups() {
3259 // An item with neither guid nor link: `parse_feed` leaves its id empty
3260 // (feed-rs's own fallback is a random UUID). Clearing id and links
3261 // anyway keeps this a test of *our* synthetic fallback alone.
3262 let xml = r#"<?xml version="1.0"?><rss version="2.0"><channel>
3263 <title>t</title>
3264 <item><title>only a title</title><description>body</description></item>
3265 </channel></rss>"#;
3266 let mut parsed = parse_feed(xml.as_bytes()).expect("parse");
3267 parsed.entries[0].id.clear();
3268 parsed.entries[0].links.clear();
3269 let g1 = normalize_entry(&parsed.entries[0]).guid;
3270 let g2 = normalize_entry(&parsed.entries[0]).guid;
3271 assert_eq!(g1, g2);
3272 assert!(g1.starts_with("featherreader:synthetic:"));
3273 }
3274
3275 // ---- feed-rs 3.0: entry identity and links held to their 2.4 values ----
3276 //
3277 // The guid is the dedup key under `UNIQUE (feed_id, guid)`. If a parser
3278 // upgrade changes it, every reader gets every live entry of that feed a
3279 // second time. The pinned ids below are what feed-rs 2.4.0 generated for
3280 // these exact bytes.
3281
3282 /// An id-less item with an ordinary link keeps the id 2.4 generated.
3283 #[test]
3284 fn a_generated_entry_id_is_the_one_feed_rs_2_4_produced() {
3285 let xml = r#"<?xml version="1.0"?><rss version="2.0"><channel><title>t</title><item><title>No guid here</title><link>https://n.example/post</link></item></channel></rss>"#;
3286 let e = normalize_entry(&parse_feed(xml.as_bytes()).expect("parse").entries[0]);
3287 assert_eq!(e.guid, "5813b43a0512aaef2750311bf4d978a");
3288 assert_eq!(e.url.as_deref(), Some("https://n.example/post"));
3289 }
3290
3291 /// **feed-rs 3.0 adds `<comments>` and `wfw:commentRss` to `entry.links`.**
3292 /// Listed before `<link>`, the comments URL became the entry's permalink,
3293 /// and for an id-less item it was hashed into the generated id, so the
3294 /// same item got a new guid after the upgrade: a duplicate in every
3295 /// reader's list.
3296 #[test]
3297 fn a_comments_link_listed_first_is_neither_the_permalink_nor_the_id() {
3298 let xml = r#"<?xml version="1.0"?>
3299<rss version="2.0" xmlns:wfw="http://wellformedweb.org/CommentAPI/"><channel><title>c</title><link>https://c.example/</link>
3300<item><title>Comments listed first, no guid</title><comments>https://c.example/1#comments</comments><link>https://c.example/1</link><wfw:commentRss>https://c.example/1/feed</wfw:commentRss></item>
3301<item><title>Comments first, guid present</title><comments>https://c.example/6#comments</comments><link>https://c.example/6</link><guid>c6</guid></item>
3302</channel></rss>"#;
3303 let parsed = parse_feed(xml.as_bytes()).expect("parse");
3304 let e0 = normalize_entry(&parsed.entries[0]);
3305 assert_eq!(
3306 e0.guid, "cd0017f2746ee934cf45ca0125100796",
3307 "the generated id moved, so this item would be stored twice"
3308 );
3309 assert_eq!(e0.url.as_deref(), Some("https://c.example/1"));
3310 let e1 = normalize_entry(&parsed.entries[1]);
3311 assert_eq!(e1.guid, "c6");
3312 assert_eq!(e1.url.as_deref(), Some("https://c.example/6"));
3313 }
3314
3315 /// The same for Atom, where 3.0 also reads `wfw:commentRss`.
3316 #[test]
3317 fn an_atom_comment_feed_link_is_neither_the_permalink_nor_the_id() {
3318 let xml = r#"<?xml version="1.0" encoding="utf-8"?>
3319<feed xmlns="http://www.w3.org/2005/Atom" xmlns:wfw="http://wellformedweb.org/CommentAPI/"><title>a</title>
3320<entry><title>commentRss before link, no id</title><wfw:commentRss>https://a.example/1/feed</wfw:commentRss><link href="https://a.example/1"/><updated>2020-06-01T00:00:00Z</updated></entry>
3321</feed>"#;
3322 let e = normalize_entry(&parse_feed(xml.as_bytes()).expect("parse").entries[0]);
3323 assert_eq!(e.guid, "5148a3d11836efc42f82a7b25a38383d");
3324 assert_eq!(e.url.as_deref(), Some("https://a.example/1"));
3325 }
3326
3327 /// **An item with no guid and no permalink gets a guid that holds still.**
3328 /// feed-rs's default generator falls back to a random UUID when an entry
3329 /// has no link, so such an item was a new row on every poll (in 2.4 as in
3330 /// 3.0). It now falls through to [`stable_guid`]. A comments link is not a
3331 /// permalink, so an item carrying only one is in the same position — and
3332 /// is not given the comments page as its URL.
3333 #[test]
3334 fn an_item_without_guid_or_permalink_dedups_across_polls() {
3335 let xml = r#"<?xml version="1.0"?><rss version="2.0"><channel><title>t</title>
3336<item><title>only a title</title><description>body</description></item>
3337<item><title>Only a comments link</title><comments>https://c.example/2#comments</comments></item>
3338</channel></rss>"#;
3339 let first = parse_feed(xml.as_bytes()).expect("parse");
3340 let second = parse_feed(xml.as_bytes()).expect("parse");
3341 for i in 0..2 {
3342 let a = normalize_entry(&first.entries[i]);
3343 let b = normalize_entry(&second.entries[i]);
3344 assert_eq!(
3345 a.guid, b.guid,
3346 "entry {i}'s guid changed between two parses"
3347 );
3348 assert!(a.guid.starts_with("featherreader:synthetic:"), "{}", a.guid);
3349 assert_eq!(a.url, None, "entry {i}");
3350 }
3351 }
3352
3353 /// **The author is the person's name.** feed-rs 2.4 named every RSS
3354 /// `<author>` "author" (the element name; the text went to `email`) and an
3355 /// Atom author with an empty `<name>` "unknown", and FeatherReader stored
3356 /// those words as the byline. 3.0 splits name from address but leaves the
3357 /// `(Name)` of the RSS `address (Name)` form in its parentheses.
3358 #[test]
3359 fn the_author_is_the_name_not_the_element_or_the_address() {
3360 let xml = r#"<?xml version="1.0"?>
3361<rss version="2.0" xmlns:dc="http://purl.org/dc/elements/1.1/"><channel><title>p</title>
3362<item><title>a1</title><guid>a1</guid><author>alice@example.com (Alice Example)</author></item>
3363<item><title>a2</title><guid>a2</guid><author>bob@example.com</author></item>
3364<item><title>a3</title><guid>a3</guid><author>Carol</author></item>
3365<item><title>a4</title><guid>a4</guid><dc:creator>Dave <dave@example.com></dc:creator></item>
3366<item><title>a5</title><guid>a5</guid><author>erin@example.com ()</author></item>
3367</channel></rss>"#;
3368 let parsed = parse_feed(xml.as_bytes()).expect("parse");
3369 let authors: Vec<Option<String>> = parsed
3370 .entries
3371 .iter()
3372 .map(|e| normalize_entry(e).author)
3373 .collect();
3374 assert_eq!(
3375 authors,
3376 vec![
3377 Some("Alice Example".to_string()),
3378 None,
3379 Some("Carol".to_string()),
3380 Some("Dave".to_string()),
3381 None,
3382 ]
3383 );
3384
3385 let atom = r#"<?xml version="1.0" encoding="utf-8"?>
3386<feed xmlns="http://www.w3.org/2005/Atom"><title>a</title>
3387<entry><id>x1</id><title>t</title><author><name></name></author></entry>
3388<entry><id>x2</id><title>t</title><author><name>Bob</name></author></entry>
3389</feed>"#;
3390 let parsed = parse_feed(atom.as_bytes()).expect("parse");
3391 assert_eq!(normalize_entry(&parsed.entries[0]).author, None);
3392 assert_eq!(
3393 normalize_entry(&parsed.entries[1]).author.as_deref(),
3394 Some("Bob")
3395 );
3396 }
3397
3398 /// **The byline is the first author that HAS a name.** An item giving
3399 /// `<author>` as a bare address and the name in `<dc:creator>` has two
3400 /// people, and the first has no name; taking only the first lost the byline.
3401 #[test]
3402 fn the_byline_is_the_first_author_with_a_name() {
3403 let xml = r#"<?xml version="1.0"?>
3404<rss version="2.0" xmlns:dc="http://purl.org/dc/elements/1.1/"><channel><title>p</title>
3405<item><title>b1</title><guid>b1</guid><author>bob@example.com</author><dc:creator>Bob Smith</dc:creator></item>
3406<item><title>b2</title><guid>b2</guid><author>erin@example.com ()</author><dc:creator>Erin</dc:creator></item>
3407<item><title>b3</title><guid>b3</guid><author>only@example.com</author></item>
3408</channel></rss>"#;
3409 let parsed = parse_feed(xml.as_bytes()).expect("parse");
3410 let authors: Vec<Option<String>> = parsed
3411 .entries
3412 .iter()
3413 .map(|e| normalize_entry(e).author)
3414 .collect();
3415 assert_eq!(
3416 authors,
3417 vec![
3418 Some("Bob Smith".to_string()),
3419 Some("Erin".to_string()),
3420 None
3421 ]
3422 );
3423 }
3424
3425 /// Author strings that make feed-rs 3.0 panic: its name/address splitter
3426 /// (`parser/util/mod.rs`, `parse_person_name_email`) slices one byte either
3427 /// side of the address, which is not a character boundary when a multi-byte
3428 /// character touches it.
3429 const PANICKING_AUTHORS: [&str; 3] = [
3430 "jose@example.com(José)",
3431 "«zoe@example.com» Zoë Long Name Here",
3432 "Zoë Long Name «zoe@example.com»",
3433 ];
3434
3435 /// Every place feed-rs runs that splitter, as a whole document around `v`.
3436 fn documents_with_author(v: &str) -> Vec<(&'static str, String)> {
3437 vec![
3438 (
3439 "rss author",
3440 format!(
3441 r#"<?xml version="1.0"?><rss version="2.0"><channel><title>t</title><item><title>x</title><guid>g</guid><author>{v}</author></item></channel></rss>"#
3442 ),
3443 ),
3444 (
3445 "rss dc:creator",
3446 format!(
3447 r#"<?xml version="1.0"?><rss version="2.0" xmlns:dc="http://purl.org/dc/elements/1.1/"><channel><title>t</title><item><title>x</title><guid>g</guid><dc:creator>{v}</dc:creator></item></channel></rss>"#
3448 ),
3449 ),
3450 (
3451 "rss managingEditor",
3452 format!(
3453 r#"<?xml version="1.0"?><rss version="2.0"><channel><title>t</title><managingEditor>{v}</managingEditor><item><title>x</title><guid>g</guid></item></channel></rss>"#
3454 ),
3455 ),
3456 (
3457 "rss webMaster",
3458 format!(
3459 r#"<?xml version="1.0"?><rss version="2.0"><channel><title>t</title><webMaster>{v}</webMaster><item><title>x</title><guid>g</guid></item></channel></rss>"#
3460 ),
3461 ),
3462 (
3463 "rss1 dc:creator",
3464 format!(
3465 r#"<?xml version="1.0"?><rdf:RDF xmlns:rdf="http://www.w3.org/1999/02/22-rdf-syntax-ns#" xmlns="http://purl.org/rss/1.0/" xmlns:dc="http://purl.org/dc/elements/1.1/"><channel><title>t</title></channel><item><title>x</title><link>https://e.example/1</link><dc:creator>{v}</dc:creator></item></rdf:RDF>"#
3466 ),
3467 ),
3468 (
3469 "json author",
3470 format!(
3471 r#"{{"version":"https://jsonfeed.org/version/1.1","title":"t","items":[{{"id":"1","content_text":"x","authors":[{{"name":"{v}"}}]}}]}}"#
3472 ),
3473 ),
3474 ]
3475 }
3476
3477 /// **A parser panic is a parse failure, not a crash.** feed-rs 3.0 panics
3478 /// on [`PANICKING_AUTHORS`]. Uncaught, the panic escaped `poll_feed`, killed
3479 /// the poll task after the scheduler had already moved `next_poll`, and
3480 /// recorded no failure — the feed stopped updating with nothing to say why.
3481 #[test]
3482 fn a_feed_rs_panic_is_returned_as_an_error() {
3483 for v in PANICKING_AUTHORS {
3484 for (slot, doc) in documents_with_author(v) {
3485 let err = match parse_feed(doc.as_bytes()) {
3486 Ok(_) => {
3487 panic!("{slot} {v:?}: parsed; the fixture no longer reaches the panic")
3488 }
3489 Err(e) => format!("{e:#}"),
3490 };
3491 assert!(err.contains("panicked"), "{slot} {v:?}: {err}");
3492 }
3493 }
3494 // The same slots with an ASCII neighbour parse normally.
3495 for (slot, doc) in documents_with_author("jose@example.com (Jose)") {
3496 assert!(parse_feed(doc.as_bytes()).is_ok(), "{slot}");
3497 }
3498 }
3499
3500 /// End to end: a poll of a feed that panics the parser reports
3501 /// `FailureKind::Parse`, so it is backed off and its `last_error` says why.
3502 #[tokio::test]
3503 async fn a_poll_of_a_feed_that_panics_the_parser_is_a_parse_failure() {
3504 let doc = &documents_with_author(PANICKING_AUTHORS[0])[0].1;
3505 let base = crate::net::tests::serve_body(doc.as_bytes().to_vec()).await;
3506 let port: u16 = base
3507 .trim_end_matches('/')
3508 .rsplit(':')
3509 .next()
3510 .unwrap()
3511 .parse()
3512 .unwrap();
3513 crate::net::test_host_override(
3514 "panicking-author.test",
3515 std::net::SocketAddr::from(([127, 0, 0, 1], port)),
3516 );
3517 let url = format!("http://panicking-author.test:{port}/feed.xml");
3518 let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
3519 crate::store::upsert_feed(
3520 &pool,
3521 &crate::store::NewFeed {
3522 url: url.clone(),
3523 ..Default::default()
3524 },
3525 )
3526 .await
3527 .unwrap();
3528 let feed = crate::store::get_feed_by_url(&pool, &url)
3529 .await
3530 .unwrap()
3531 .unwrap();
3532 let client = build_client().unwrap();
3533 let outcome = poll_feed(&pool, &client, &feed, 0)
3534 .await
3535 .expect("a parser panic surfaced as a store error");
3536 match outcome {
3537 PollOutcome::Failed {
3538 kind: FailureKind::Parse,
3539 detail,
3540 ..
3541 } => assert!(detail.contains("panicked"), "{detail}"),
3542 other => panic!("expected a parse failure, got {other:?}"),
3543 }
3544 }
3545
3546 // ---- dedup-key hashes are fixed across toolchains -------------------
3547 //
3548 // `stable_guid` and `bound_guid` are stored dedup keys. They were built on
3549 // `std`'s `DefaultHasher`, whose algorithm the standard library documents
3550 // as unspecified and subject to change between releases; a toolchain bump
3551 // could have re-keyed every such entry and duplicated it. The values below
3552 // were produced by that `DefaultHasher` code, so they also prove the
3553 // replacement hashes identically today.
3554
3555 #[test]
3556 fn a_synthetic_guid_is_the_value_already_stored() {
3557 let xml = r#"<?xml version="1.0"?><rss version="2.0"><channel><title>t</title>
3558<item><title>only a title</title><description>body</description></item>
3559<item><description>no title either</description></item>
3560<item><title>Tïtle wíth ünïcode</title></item>
3561</channel></rss>"#;
3562 let parsed = parse_feed(xml.as_bytes()).expect("parse");
3563 let guids: Vec<String> = parsed
3564 .entries
3565 .iter()
3566 .map(|e| normalize_entry(e).guid)
3567 .collect();
3568 assert_eq!(
3569 guids,
3570 vec![
3571 "featherreader:synthetic:964034cf9a24b551".to_string(),
3572 "featherreader:synthetic:baa7c19198c76b30".into(),
3573 "featherreader:synthetic:34b72cc7fa0a59ac".into(),
3574 ]
3575 );
3576 }
3577
3578 #[test]
3579 fn a_long_guid_is_the_value_already_stored() {
3580 let a = bound_guid("g".repeat(MAX_GUID_BYTES + 1));
3581 let b = bound_guid(format!(
3582 "https://long.example/{}",
3583 "ü".repeat(MAX_GUID_BYTES)
3584 ));
3585 assert_eq!(
3586 (a.as_str(), b.as_str()),
3587 (
3588 "featherreader:long-guid:966404db20ef8f05",
3589 "featherreader:long-guid:a1d631602455ace4",
3590 )
3591 );
3592 }
3593
3594 #[test]
3595 fn entry_link_scheme_allowlist_neutralizes_javascript() {
3596 // An entry whose only link is a javascript: URL must yield no href.
3597 let xml = r#"<?xml version="1.0"?><rss version="2.0"><channel>
3598 <title>t</title>
3599 <item>
3600 <title>evil</title>
3601 <link>javascript:alert(document.domain)</link>
3602 <guid>evil-1</guid>
3603 </item>
3604 </channel></rss>"#;
3605 let parsed = parse_feed(xml.as_bytes()).expect("parse");
3606 let e = normalize_entry(&parsed.entries[0]);
3607 // url is dropped (not a safe http(s) link)…
3608 assert_eq!(e.url, None);
3609 // …but the entry still dedups (guid preserved from <guid>).
3610 assert_eq!(e.guid, "evil-1");
3611
3612 // A data: URL is likewise dropped.
3613 let xml2 = r#"<?xml version="1.0"?><rss version="2.0"><channel>
3614 <title>t</title>
3615 <item><title>d</title><link>data:text/html,<script>1</script></link><guid>d1</guid></item>
3616 </channel></rss>"#;
3617 let parsed2 = parse_feed(xml2.as_bytes()).expect("parse");
3618 let e2 = normalize_entry(&parsed2.entries[0]);
3619 assert_eq!(e2.url, None);
3620
3621 // A normal https link survives.
3622 let xml3 = r#"<?xml version="1.0"?><rss version="2.0"><channel>
3623 <title>t</title>
3624 <item><title>ok</title><link>https://ok.example/post</link><guid>ok1</guid></item>
3625 </channel></rss>"#;
3626 let parsed3 = parse_feed(xml3.as_bytes()).expect("parse");
3627 let e3 = normalize_entry(&parsed3.entries[0]);
3628 assert_eq!(e3.url.as_deref(), Some("https://ok.example/post"));
3629 }
3630
3631 #[test]
3632 fn classify_privacy_flags_secret_urls_across_providers() {
3633 // --- Known providers: newsletters ---
3634 // Substack private feed path.
3635 assert!(
3636 classify_feed_privacy("https://author.substack.com/feed/private/deadbeefcafe1234")
3637 .is_private()
3638 );
3639 // Patreon ?auth= member feed.
3640 assert!(classify_feed_privacy(
3641 "https://www.patreon.com/rss/author?auth=Zm9vYmFyc2VjcmV0dG9rZW4"
3642 )
3643 .is_private());
3644 // Ghost members feed via ?uuid=.
3645 assert!(classify_feed_privacy(
3646 "https://blog.ghost.io/rss/?uuid=1f2e3d4c-5b6a-7089-90ab-cdef01234567"
3647 )
3648 .is_private());
3649
3650 // --- Known providers: private podcasts ---
3651 // Supporting Cast tokened podcast feed.
3652 assert!(classify_feed_privacy(
3653 "https://feeds.supportingcast.fm/show/abcdef0123456789abcdef01"
3654 )
3655 .is_private());
3656 // Supercast private podcast (host alone is enough).
3657 assert!(classify_feed_privacy("https://feeds.supercast.com/12345/rss").is_private());
3658
3659 // --- Generic, provider-agnostic heuristic ---
3660 // Named credential query params with an opaque value.
3661 assert!(
3662 classify_feed_privacy("https://example.com/feed?token=Zm9vYmFyc2VjcmV0").is_private()
3663 );
3664 assert!(
3665 classify_feed_privacy("https://example.com/feed?key=Zm9vYmFyc2VjcmV0").is_private()
3666 );
3667 assert!(
3668 classify_feed_privacy("https://example.com/feed?secret=Zm9vYmFyc2VjcmV0").is_private()
3669 );
3670 // Userinfo credentials in the authority.
3671 assert!(classify_feed_privacy("https://user:pass@example.com/feed").is_private());
3672 // A `/private/` path segment on an unknown host.
3673 assert!(classify_feed_privacy("https://blog.example.com/private/rss").is_private());
3674 // `/members/` path convention.
3675 assert!(classify_feed_privacy("https://news.example.com/members/feed.xml").is_private());
3676 // A high-entropy opaque token embedded in the path with no telltale name.
3677 assert!(
3678 classify_feed_privacy("https://feeds.example.com/aB3xK9zQ7mP2rT5wL8nD4vF6")
3679 .is_private()
3680 );
3681 // A bare UUID path segment (many tokened feeds).
3682 assert!(classify_feed_privacy(
3683 "https://feeds.example.com/1f2e3d4c-5b6a-7089-90ab-cdef01234567"
3684 )
3685 .is_private());
3686 }
3687
3688 /// The dominant real-world private-podcast shape delivers the token as a
3689 /// FILENAME (`<token>.rss` / `<token>.xml`) or affixed inside a larger
3690 /// segment (`feed-<uuid>`). Named providers are caught by their host rule;
3691 /// these are UNKNOWN-provider CDNs that must still be caught by the generic
3692 /// backstop, so the secret is never fetched or stored.
3693 #[test]
3694 fn classify_privacy_catches_tokened_filenames_on_unknown_hosts() {
3695 // **A stem only the extension-strip branch can see.** Every other case
3696 // here is also caught by the sub-part scan (branch 3) or the UUID scan,
3697 // so deleting the strip left the suite green — `FEED_EXTENSIONS` was
3698 // effectively dead. This stem's `-`-separated parts are each too short
3699 // to look like a secret on their own, and the whole segment fails on
3700 // the `.` — only stripping `.rss` and re-testing the 26-char stem sees
3701 // it. That is exactly the shape a hyphen-bearing base64url token
3702 // filename takes.
3703 assert!(
3704 classify_feed_privacy("https://cdn.example/feeds/aB3xK9pQ-7mZ2vN8w-Qr5tYuW.rss")
3705 .is_private(),
3706 "a token stem visible only after stripping the extension was not caught"
3707 );
3708 // hex-32 token as an .xml filename.
3709 assert!(classify_feed_privacy(
3710 "https://cdn.somepod.io/f/a1b2c3d4e5f60718293a4b5c6d7e8f90.xml"
3711 )
3712 .is_private());
3713 // hex-32 token as a .rss filename on an unknown CDN.
3714 assert!(classify_feed_privacy(
3715 "https://dcs.megaphone.example/network/a1b2c3d4e5f60718293a4b5c6d7e8f90.rss"
3716 )
3717 .is_private());
3718 // UUID + .xml filename.
3719 assert!(classify_feed_privacy(
3720 "https://brandnew.example/feed/1f2e3d4c-5b6a-7089-90ab-cdef01234567.xml"
3721 )
3722 .is_private());
3723 // UUID affixed with a prefix (`feed-<uuid>`) — split can't see it, the
3724 // UUID-substring scan must.
3725 assert!(classify_feed_privacy(
3726 "https://x.example/feed-1f2e3d4c-5b6a-7089-90ab-cdef01234567"
3727 )
3728 .is_private());
3729 // UUID + .rss suffix.
3730 assert!(classify_feed_privacy(
3731 "https://x.example/1f2e3d4c-5b6a-7089-90ab-cdef01234567.rss"
3732 )
3733 .is_private());
3734 // hex-16 token as an .xml filename.
3735 assert!(classify_feed_privacy("https://x.example/feed/9f8e7d6c5b4a3928.xml").is_private());
3736 // A base64url token with `=` padding as a clean path segment.
3737 assert!(
3738 classify_feed_privacy("https://cdn.pod.io/f/YWJjZGVmZ2hpamtsbW5vcHFyc3R1dnc=")
3739 .is_private()
3740 );
3741 }
3742
3743 /// YouTube channel/playlist RSS feeds are FULLY PUBLIC (the id is a public
3744 /// handle, not a secret) and are the standard way to subscribe to a channel —
3745 /// they must NOT be false-blocked by the generic entropy heuristic.
3746 #[test]
3747 fn classify_privacy_allows_public_youtube_feeds() {
3748 assert_eq!(
3749 classify_feed_privacy(
3750 "https://www.youtube.com/feeds/videos.xml?channel_id=UC-lHJZR3Gqxm24_Vd_AJ5Yw"
3751 ),
3752 FeedPrivacy::Public
3753 );
3754 assert_eq!(
3755 classify_feed_privacy(
3756 "https://www.youtube.com/feeds/videos.xml?playlist_id=PLFgquLnL59alCl_2TQvOiD5Vgm1hCaGSI"
3757 ),
3758 FeedPrivacy::Public
3759 );
3760 // Bare host form too.
3761 assert_eq!(
3762 classify_feed_privacy(
3763 "https://youtube.com/feeds/videos.xml?channel_id=UC-lHJZR3Gqxm24_Vd_AJ5Yw"
3764 ),
3765 FeedPrivacy::Public
3766 );
3767 // The allowlist is narrow: a `token=` on the YouTube feeds path still
3768 // classifies private (can't smuggle a credential through the allowlist).
3769 assert!(classify_feed_privacy(
3770 "https://www.youtube.com/feeds/videos.xml?token=Zm9vYmFyc2VjcmV0dG9rZW4"
3771 )
3772 .is_private());
3773 }
3774
3775 #[test]
3776 fn classify_privacy_leaves_normal_public_feeds_public() {
3777 // Plain feed documents.
3778 assert_eq!(
3779 classify_feed_privacy("https://example.com/feed.xml"),
3780 FeedPrivacy::Public
3781 );
3782 assert_eq!(
3783 classify_feed_privacy("https://blog.example.com/rss"),
3784 FeedPrivacy::Public
3785 );
3786 assert_eq!(
3787 classify_feed_privacy("https://blog.example.com/rss.xml"),
3788 FeedPrivacy::Public
3789 );
3790 // A Substack PUBLIC feed (`/feed`, not `/feed/private/`) stays public.
3791 assert_eq!(
3792 classify_feed_privacy("https://author.substack.com/feed"),
3793 FeedPrivacy::Public
3794 );
3795 // A WordPress `/feed` endpoint.
3796 assert_eq!(
3797 classify_feed_privacy("https://wordpress.example.com/feed/"),
3798 FeedPrivacy::Public
3799 );
3800 // A plain Atom feed.
3801 assert_eq!(
3802 classify_feed_privacy("https://example.org/atom.xml"),
3803 FeedPrivacy::Public
3804 );
3805 // A long, hyphenated slug must NOT be mistaken for an embedded secret.
3806 assert_eq!(
3807 classify_feed_privacy("https://example.com/2026/07/my-first-long-blog-post-title/feed"),
3808 FeedPrivacy::Public
3809 );
3810 // A benign query key that merely contains "key" as a substring is fine.
3811 assert_eq!(
3812 classify_feed_privacy("https://example.com/feed?keyword=rust"),
3813 FeedPrivacy::Public
3814 );
3815 // A short, non-opaque value on a named key (e.g. an enum) is not a secret.
3816 assert_eq!(
3817 classify_feed_privacy("https://example.com/feed?p=2"),
3818 FeedPrivacy::Public
3819 );
3820 // An empty credential value is not a secret.
3821 assert_eq!(
3822 classify_feed_privacy("https://example.com/feed?token="),
3823 FeedPrivacy::Public
3824 );
3825 // A hyphenated slug ending in a feed extension must NOT be seen as a
3826 // tokened filename (the stem is short dictionary words, not a blob).
3827 assert_eq!(
3828 classify_feed_privacy("https://example.com/my-first-long-blog-post.xml"),
3829 FeedPrivacy::Public
3830 );
3831 // A short hex episode id in an .xml filename (< 16 chars) is not a secret.
3832 assert_eq!(
3833 classify_feed_privacy("https://example.com/episodes/ab12cd.xml"),
3834 FeedPrivacy::Public
3835 );
3836 // A dotted host-style filename slug stays public.
3837 assert_eq!(
3838 classify_feed_privacy("https://example.com/category/tech-news/feed.xml"),
3839 FeedPrivacy::Public
3840 );
3841 // Unparseable URL: treated as Public (add path rejects it downstream).
3842 assert_eq!(classify_feed_privacy("not a url"), FeedPrivacy::Public);
3843 }
3844
3845 /// **The detail is bounded where it is CONSTRUCTED, not only where it is
3846 /// stored.**
3847 ///
3848 /// Review found that widening `PollOutcome::Failed` with this field opened a
3849 /// second sink nobody looked at: `web.rs`'s `add_subscription` logs
3850 /// `?outcome` at INFO on a user-facing request path, so the whole
3851 /// untruncated anyhow chain — redirect-hop URLs, the SSRF guard's refusal
3852 /// text naming a resolved internal address — went to the access log.
3853 ///
3854 /// Bounding inside `bump_feed_errors` protected the database and nothing
3855 /// else. Bounding at construction protects every sink, including the ones
3856 /// added later.
3857 #[test]
3858 fn a_failure_detail_is_bounded_at_construction() {
3859 let huge = "x".repeat(10_000);
3860 let outcome = PollOutcome::Failed {
3861 backoff: BACKOFF_BASE,
3862 kind: FailureKind::Fetch,
3863 detail: failure_detail(&huge),
3864 };
3865 let PollOutcome::Failed { detail, .. } = &outcome else {
3866 panic!("wrong variant");
3867 };
3868 assert!(
3869 detail.chars().count() <= MAX_FAILURE_DETAIL_CHARS,
3870 "detail was {} chars",
3871 detail.chars().count(),
3872 );
3873 // And the Debug rendering — which is what actually reached the log — is
3874 // bounded with it.
3875 assert!(format!("{outcome:?}").len() < 1_000);
3876 }
3877
3878 /// **Every failure kind has its own label, and they round-trip.**
3879 ///
3880 /// Review found that collapsing all four `as_str` arms to `"fetch"` left
3881 /// the whole suite green: every test of these columns passed string
3882 /// literals, so nothing tied a variant to its label. A histogram whose
3883 /// buckets all say the same thing is worse than no histogram — it reports a
3884 /// single confident cause for four different failures.
3885 ///
3886 /// Asserted over `ALL` rather than a hand-written list, so adding a variant
3887 /// without a label fails here instead of silently sharing one.
3888 #[test]
3889 fn every_failure_kind_has_a_distinct_round_tripping_label() {
3890 let mut seen = std::collections::BTreeSet::new();
3891 for kind in FailureKind::ALL {
3892 let label = kind.as_str();
3893 assert!(
3894 seen.insert(label),
3895 "{label:?} is used by more than one FailureKind",
3896 );
3897 assert_eq!(
3898 FailureKind::parse(label),
3899 Some(kind),
3900 "{label:?} does not read back as the kind that wrote it",
3901 );
3902 }
3903 assert_eq!(seen.len(), FailureKind::ALL.len());
3904 // A label from a newer build is not attributed to a cause this one
3905 // knows — the `metrics::Backend::parse` contract.
3906 assert_eq!(FailureKind::parse("quota"), None);
3907 }
3908
3909 /// **Escalation reaches `settle_poll`.** `backoff_for` grows with the
3910 /// count and is tested alone; nothing asserted that the poll path passes
3911 /// the COUNT in. `backoff_for(1)` in its place left the whole suite green
3912 /// — a permanently dead feed retrying forever at the first-failure floor,
3913 /// which the comment on that line says must not happen.
3914 #[tokio::test]
3915 async fn backoff_escalates_with_consecutive_failures() -> anyhow::Result<()> {
3916 let pool = crate::store::init_url("sqlite::memory:").await?;
3917 let url = "https://dead.example/feed.xml";
3918 crate::store::upsert_feed(
3919 &pool,
3920 &crate::store::NewFeed {
3921 url: url.to_string(),
3922 ..Default::default()
3923 },
3924 )
3925 .await?;
3926 for _ in 0..5 {
3927 crate::store::bump_feed_errors(&pool, url, FailureKind::Fetch, "down").await?;
3928 }
3929 let before = chrono::Utc::now();
3930 settle_poll(
3931 &pool,
3932 url,
3933 &PollOutcome::Failed {
3934 backoff: Duration::from_secs(300),
3935 kind: FailureKind::Fetch,
3936 detail: "still down".to_string(),
3937 },
3938 Duration::from_secs(3600),
3939 )
3940 .await;
3941 let next: String = sqlx::query_scalar("SELECT next_poll FROM feeds WHERE url = ?1")
3942 .bind(url)
3943 .fetch_one(&pool)
3944 .await?;
3945 let next = chrono::DateTime::parse_from_rfc3339(&next)?.with_timezone(&chrono::Utc);
3946 let delay = (next - before).num_seconds();
3947 let expected = backoff_for(6).as_secs() as i64;
3948 assert!(
3949 (delay - expected).abs() <= 60,
3950 "sixth failure scheduled {delay}s out; escalation says {expected}s"
3951 );
3952 assert!(
3953 delay > backoff_for(1).as_secs() as i64 + 60,
3954 "the sixth failure landed on the first-failure floor"
3955 );
3956 Ok(())
3957 }
3958
3959 #[test]
3960 fn backoff_grows_and_is_capped() {
3961 assert_eq!(backoff_for(1), BACKOFF_BASE);
3962 assert!(backoff_for(2) > backoff_for(1));
3963 assert_eq!(backoff_for(100), BACKOFF_MAX);
3964 }
3965
3966 /// **Storable and pollable are ONE decision.**
3967 ///
3968 /// Review found the sequencing error this closes: making `at://` storable
3969 /// while nothing can poll it does not leave the feature dormant, it creates
3970 /// permanent failures that the cause histogram then publishes as
3971 /// unreachable publishers — the exact conflation it exists to end.
3972 #[test]
3973 fn an_at_uri_is_not_storable_while_standard_site_is_off() {
3974 let uri = "at://did:plc:ohutz6x5acjmpuulp3x7wxxc/site.standard.publication/3lab";
3975 assert!(
3976 !is_storable_feed_url(uri, false),
3977 "stored a feed nothing can poll"
3978 );
3979 assert!(is_storable_feed_url(uri, true));
3980 assert!(is_storable_feed_url("https://example.com/feed.xml", false));
3981 assert!(is_storable_feed_url("https://example.com/feed.xml", true));
3982 }
3983
3984 /// **A non-canonical scheme spelling is recognised and refused.** URL
3985 /// schemes are case-insensitive, so `At://` names the same thing as
3986 /// `at://` — but `feeds.url` is UNIQUE, so accepting both is two rows for
3987 /// one publication. Recognised (not passed through to the generic checks
3988 /// as if it were an ordinary URL), then refused for the spelling.
3989 /// Since publications are polled, only a URI the storage guard accepts is
3990 /// a publication. Any other at-URI is Unsupported, which no poller reads.
3991 #[test]
3992 fn only_a_storable_publication_uri_is_a_publication() {
3993 let did = "did:plc:ohutz6x5acjmpuulp3x7wxxc";
3994 let pubn = crate::lexicon::nsid::STANDARD_PUBLICATION;
3995 assert_eq!(
3996 FeedKind::of(&format!("at://{did}/{pubn}/3lab")),
3997 FeedKind::Publication
3998 );
3999 for unsupported in [
4000 format!("at://{did}/app.bsky.feed.post/3lab"),
4001 format!("At://{did}/{pubn}/3lab"),
4002 format!("at://alice.example.com/{pubn}/3lab"),
4003 format!("at://did:plc:short/{pubn}/3lab"),
4004 ] {
4005 assert_eq!(
4006 FeedKind::of(&unsupported),
4007 FeedKind::Unsupported,
4008 "{unsupported}"
4009 );
4010 }
4011 assert_eq!(FeedKind::of("https://example.com/feed.xml"), FeedKind::Rss);
4012 assert!(!FeedKind::POLLABLE.contains(&FeedKind::Unsupported));
4013 assert_eq!(FeedKind::parse("unsupported"), Some(FeedKind::Unsupported));
4014 }
4015
4016 #[test]
4017 fn a_non_canonical_at_uri_spelling_is_recognised_and_refused() {
4018 for odd in [
4019 "At://did:plc:ohutz6x5acjmpuulp3x7wxxc/site.standard.publication/3lab",
4020 "AT://did:plc:ohutz6x5acjmpuulp3x7wxxc/site.standard.publication/3lab",
4021 ] {
4022 assert!(!is_storable_feed_url(odd, true), "stored {odd:?}");
4023 // Fails CLOSED: it is an at-URI this reader will not store, not an
4024 // unparseable string that the `Err(_) => Public` arm waves through.
4025 assert!(
4026 classify_feed_privacy(odd).is_private(),
4027 "{odd:?} was declared publishable"
4028 );
4029 }
4030 }
4031
4032 /// **The handle form is not storable — the DID form is the identity.**
4033 ///
4034 /// `feeds.url` is UNIQUE; a handle and its DID would be two rows for one
4035 /// publication, and a handle can change hands. Every spelling is refused,
4036 /// canonical or not; resolving one to a DID is the input path's job.
4037 #[test]
4038 fn a_handle_form_publication_uri_is_not_storable() {
4039 for authority in [
4040 "alice.example.com",
4041 "EXAMPLE.COM",
4042 "169.254.169.254",
4043 "pds.internal",
4044 "printer.local",
4045 "host:8080",
4046 "-.-",
4047 "a b.c",
4048 ] {
4049 let uri = format!("at://{authority}/site.standard.publication/3lab");
4050 assert!(
4051 !is_storable_feed_url(&uri, true),
4052 "accepted authority {authority:?}"
4053 );
4054 }
4055 }
4056
4057 /// A control character or space in the at-URI is refused: `scheduler.rs`
4058 /// logs `%feed.url` with Display, and `feeds.url` is UNIQUE.
4059 #[test]
4060 fn an_at_uri_with_control_characters_is_not_storable() {
4061 for bad in [
4062 "at://did:plc:ohutz6x5acjmpuulp3x7wxxc/site.standard.publication/3lab\n",
4063 "at://did:plc:ohutz6x5acjmpuulp3x7wxxc/site.standard.publication/3lab ",
4064 "at://did:plc:ohutz6x5acjmpuulp3x7wxxc/site.standard.publication/3l\tab",
4065 ] {
4066 assert!(!is_storable_feed_url(bad, true), "accepted {bad:?}");
4067 }
4068 }
4069
4070 /// **The `at://` exemption is a REGRESSION unless it is narrow.**
4071 ///
4072 /// Fan-out review found the arm I added was a bare prefix match, so *any*
4073 /// attacker-chosen string starting `at://` was declared safe to publish —
4074 /// skipping the userinfo check, the known-provider table, the private-path
4075 /// markers, the secret-query keys and the entropy heuristics. Measured
4076 /// against `main`, these three went from `Private` to `Public`.
4077 ///
4078 /// That matters because `rename_subscription` caches the URL AND rewrites
4079 /// the user's PUBLIC PDS record, with `classify_feed_privacy` as its only
4080 /// gate.
4081 #[test]
4082 fn a_credential_bearing_at_uri_is_still_private() {
4083 for hostile in [
4084 "at://user:pass@private.example.com/feed/private/TOKEN?apikey=deadbeefdeadbeef",
4085 "at://patreon.com/rss/12345?auth=deadbeefdeadbeefdeadbeef",
4086 "at://did:plc:ohutz6x5acjmpuulp3x7wxxc/site.standard.publication/3lab?apikey=sekrit",
4087 ] {
4088 assert!(
4089 matches!(classify_feed_privacy(hostile), FeedPrivacy::Private(_)),
4090 "declared public: {hostile}"
4091 );
4092 }
4093 }
4094
4095 /// An rkey is `[A-Za-z0-9._:~-]` per atproto. Without that, a query string
4096 /// or path fragment smuggled into the rkey satisfies the three-segment
4097 /// check — which is what the exemption above keys off.
4098 #[test]
4099 fn an_rkey_outside_the_atproto_charset_is_not_storable() {
4100 for bad in [
4101 "at://did:plc:ohutz6x5acjmpuulp3x7wxxc/site.standard.publication/3lab?apikey=sekrit",
4102 "at://did:plc:ohutz6x5acjmpuulp3x7wxxc/site.standard.publication/3lab#frag",
4103 "at://did:plc:ohutz6x5acjmpuulp3x7wxxc/site.standard.publication/3lab%2Fevil",
4104 "at://did:plc:ohutz6x5acjmpuulp3x7wxxc/site.standard.publication/caf\u{e9}",
4105 "at://did:plc:ohutz6x5acjmpuulp3x7wxxc/site.standard.publication/3lab\u{202e}x",
4106 ] {
4107 assert!(!is_storable_feed_url(bad, true), "accepted rkey in {bad:?}");
4108 }
4109 // The legitimate charset still passes.
4110 for good in [
4111 "at://did:plc:ohutz6x5acjmpuulp3x7wxxc/site.standard.publication/3lab2c4d5e6f7g8h",
4112 "at://did:plc:ohutz6x5acjmpuulp3x7wxxc/site.standard.publication/a.b_c~d-e",
4113 ] {
4114 assert!(is_storable_feed_url(good, true), "refused {good:?}");
4115 }
4116 }
4117
4118 /// **A real rkey is a TID, and a TID looks exactly like a secret.**
4119 ///
4120 /// Without an explicit `at://` arm, `classify_feed_privacy` runs the generic
4121 /// "high-entropy token in path" heuristic over the rkey. Measured: a
4122 /// realistic 16-char rkey on a handle-form at-URI is classified PRIVATE and
4123 /// the subscription REFUSED. The DID form escaped only because it fails to
4124 /// parse as a `Url` at all — so the bug was invisible from that side.
4125 ///
4126 /// The first version of this test used the rkey `3lab`, which is too short
4127 /// to trip the heuristic, so it passed with and without the fix.
4128 #[test]
4129 fn a_realistic_at_uri_rkey_is_not_mistaken_for_a_secret() {
4130 for uri in [
4131 "at://did:plc:ohutz6x5acjmpuulp3x7wxxc/site.standard.publication/3lab2c4d5e6f7g8h",
4132 "at://did:plc:ohutz6x5acjmpuulp3x7wxxc/site.standard.publication/aB3xK9pQ7mZ2vN8w",
4133 "at://did:plc:ohutz6x5acjmpuulp3x7wxxc/site.standard.publication/3lab2c4d5e6f7g8h",
4134 ] {
4135 assert_eq!(
4136 classify_feed_privacy(uri),
4137 FeedPrivacy::Public,
4138 "a publication rkey was mistaken for a credential: {uri}"
4139 );
4140 }
4141 }
4142
4143 /// **The DID form is the one that matters, and the one `Url::parse` cannot
4144 /// read.**
4145 ///
4146 /// `Url::parse("at://did:plc:…/…")` fails with *invalid port number* — the
4147 /// colons in the DID are taken as a port separator. So the obvious
4148 /// implementation, adding `"at"` to the `matches!` on `u.scheme()`, silently
4149 /// rejects every DID-based at-URI while appearing to work: the handle form
4150 /// (`at://alice.example.com/…`) parses fine and would pass such a test.
4151 ///
4152 /// All 19 at-URI rows in production are the DID form. A test written with a
4153 /// handle would have passed against an implementation that cannot store a
4154 /// single one of them.
4155 #[test]
4156 fn a_did_form_publication_uri_is_storable() {
4157 assert!(is_storable_feed_url(
4158 "at://did:plc:ohutz6x5acjmpuulp3x7wxxc/site.standard.publication/3lab",
4159 true
4160 ));
4161 }
4162
4163 /// **An allowlist entry, not a loosening.** `at://` is accepted for exactly
4164 /// one foreign collection. Any other collection is somebody else's lexicon
4165 /// arriving through a path (`resolve_subscriptions`, OPML import) that takes
4166 /// records from outside with no add-path to reject them.
4167 #[test]
4168 fn an_at_uri_for_another_collection_is_not_storable() {
4169 assert!(!is_storable_feed_url(
4170 "at://did:plc:ohutz6x5acjmpuulp3x7wxxc/community.lexicon.rss.subscription/3lab",
4171 true
4172 ));
4173 assert!(!is_storable_feed_url(
4174 "at://did:plc:ohutz6x5acjmpuulp3x7wxxc/app.bsky.feed.post/3lab",
4175 true
4176 ));
4177 }
4178
4179 /// The malformed shapes, each of which a naive `split('/')` would accept.
4180 #[test]
4181 fn a_malformed_at_uri_is_not_storable() {
4182 for bad in [
4183 "at://",
4184 "at://did:plc:ohutz6x5acjmpuulp3x7wxxc",
4185 "at://did:plc:ohutz6x5acjmpuulp3x7wxxc/site.standard.publication",
4186 "at://did:plc:ohutz6x5acjmpuulp3x7wxxc/site.standard.publication/",
4187 "at:///site.standard.publication/3lab",
4188 "at://did:plc:ohutz6x5acjmpuulp3x7wxxc/site.standard.publication/3lab/extra",
4189 "at://not-a-did-or-handle/site.standard.publication/3lab",
4190 "at://did:plc:TOOSHORT/site.standard.publication/3lab",
4191 ] {
4192 assert!(!is_storable_feed_url(bad, true), "accepted {bad:?}");
4193 }
4194 }
4195
4196 /// **The reason the function exists, unchanged.** Mutating the new branch to
4197 /// accept any scheme makes this fail while the at-URI tests keep passing —
4198 /// that asymmetry is what says the change was an allowlist entry.
4199 #[test]
4200 fn the_refused_schemes_are_still_refused() {
4201 for bad in [
4202 "javascript:alert(1)",
4203 "file:///etc/passwd",
4204 "data:text/html,<script>",
4205 "ftp://example.com/feed.xml",
4206 "at:did:plc:ohutz6x5acjmpuulp3x7wxxc/site.standard.publication/3lab",
4207 ] {
4208 assert!(!is_storable_feed_url(bad, true), "accepted {bad:?}");
4209 }
4210 // A hostless http(s) URL is a PARSE error, not a parsed URL with no
4211 // host — `https:///feed.xml` even parses as host `feed.xml`. What
4212 // refuses these is the `Err` arm, so that is what this pins.
4213 for hostless in ["http://", "https://?q=1", "http:///"] {
4214 assert!(
4215 !is_storable_feed_url(hostless, true),
4216 "a hostless URL {hostless:?} was storable"
4217 );
4218 }
4219 assert!(is_storable_feed_url("https://example.com/feed.xml", true));
4220 assert!(is_storable_feed_url("http://example.com/feed.xml", true));
4221 }
4222
4223 /// **What re-cleaning a stored body at render costs (#151).** Not a test: a
4224 /// measurement, kept so the numbers in the PR and CHANGELOG can be
4225 /// reproduced.
4226 ///
4227 /// ```text
4228 /// FEATHER_BENCH_FEEDS=/path/to/dir/of/feed/files \
4229 /// cargo test --release --lib render_reclean_cost -- --ignored --nocapture
4230 /// ```
4231 ///
4232 /// 1. Typical bodies: real feed documents from that directory, put through
4233 /// `normalize_entry` so each is exactly what ingest would store, then
4234 /// re-cleaned with `SanitizedHtml::clean`. Also counts bodies the
4235 /// re-clean changed.
4236 /// 2. Pathological bodies: #226's quadratic inputs in their stored
4237 /// (fixed-point) form, up to the 2 MiB stored bound, through the bare
4238 /// sanitizer. A hostile feed can make ingest store these (#226); they
4239 /// are what the renderer's permits, single flight and cache are sized
4240 /// against.
4241 /// `FEATHER_BENCH_UNCAPPED_MAX=<bytes>` skips the larger sizes (2 MiB of
4242 /// nesting takes ~37 s).
4243 #[test]
4244 #[ignore = "benchmark; see the doc comment"]
4245 fn render_reclean_cost() {
4246 use crate::sanitized_html::SanitizedHtml;
4247 use std::time::{Duration, Instant};
4248
4249 fn pct(sorted: &[Duration], p: f64) -> Duration {
4250 sorted[((sorted.len() as f64 - 1.0) * p).round() as usize]
4251 }
4252
4253 if let Ok(dir) = std::env::var("FEATHER_BENCH_FEEDS") {
4254 let mut times = Vec::new();
4255 let mut sizes = Vec::new();
4256 let (mut changed, mut files) = (0usize, 0usize);
4257 for path in std::fs::read_dir(&dir).unwrap() {
4258 let bytes = std::fs::read(path.unwrap().path()).unwrap();
4259 let Ok(parsed) = parse_feed(&bytes) else {
4260 continue;
4261 };
4262 files += 1;
4263 for e in &parsed.entries {
4264 let Some(stored) = normalize_entry(e).content_html else {
4265 continue;
4266 };
4267 // Best of three, to take scheduler noise out of a small body.
4268 let mut best = Duration::MAX;
4269 let mut out = None;
4270 for _ in 0..3 {
4271 let t = Instant::now();
4272 out = Some(SanitizedHtml::clean(&stored));
4273 best = best.min(t.elapsed());
4274 }
4275 let out = out.unwrap();
4276 changed += usize::from(out.as_str() != stored);
4277 times.push(best);
4278 sizes.push(stored.len());
4279 }
4280 }
4281 times.sort();
4282 sizes.sort();
4283 println!(
4284 "real feeds: {files} files, {} bodies; size p50 {} B, p99 {} B, max {} B",
4285 times.len(),
4286 sizes[sizes.len() / 2],
4287 sizes[((sizes.len() - 1) as f64 * 0.99).round() as usize],
4288 sizes[sizes.len() - 1],
4289 );
4290 println!(
4291 " re-clean p50 {:?}, p99 {:?}, max {:?}; changed by re-cleaning: {changed}",
4292 pct(×, 0.5),
4293 pct(×, 0.99),
4294 times[times.len() - 1],
4295 );
4296 }
4297
4298 let bound = MAX_CONTENT_HTML_BYTES;
4299 let uncapped_max = std::env::var("FEATHER_BENCH_UNCAPPED_MAX")
4300 .ok()
4301 .and_then(|v| v.parse::<usize>().ok())
4302 .unwrap_or(bound);
4303 for size in [bound / 8, bound / 4, bound / 2, bound] {
4304 if size > uncapped_max {
4305 continue;
4306 }
4307 let depth = size / 11;
4308 for (name, stored) in [
4309 (
4310 "'&' run",
4311 format!("<p>{}</p>", "&".repeat((size - 7) / 5)),
4312 ),
4313 (
4314 "U+00A0 run",
4315 format!("<p>{}</p>", " ".repeat((size - 7) / 6)),
4316 ),
4317 (
4318 "nested <div>",
4319 format!("{}{}", "<div>".repeat(depth), "</div>".repeat(depth)),
4320 ),
4321 ("plain text", format!("<p>{}</p>", "a".repeat(size - 7))),
4322 ] {
4323 let t = Instant::now();
4324 let out = sanitize_html(&stored);
4325 println!(
4326 "pathological {name:>13} {:>8} B: {:?} (fixed point: {})",
4327 stored.len(),
4328 t.elapsed(),
4329 out == stored
4330 );
4331 }
4332 }
4333 }
4334}
4335
4336/// #226: ingest sanitizing runs off the async runtime, under a permit and a
4337/// per-entry timeout, and a timed-out poll neither loses a stored body nor
4338/// saves the validators that would turn the next poll into a `304`.
4339#[cfg(test)]
4340mod sanitize_off_runtime_tests {
4341 use super::*;
4342 use std::sync::atomic::{AtomicBool, Ordering};
4343 use std::sync::{Arc, Mutex};
4344 use std::time::Instant;
4345 use tokio::sync::Semaphore;
4346
4347 /// A body ammonia is super-linear on: one text node of `&`. 48 Ki of them
4348 /// take ~2.5 s to sanitize in a debug build on an M-series Mac (32 Ki
4349 /// measured 1.1 s, 64 Ki 4.5 s) — far past [`SHORT`], short enough that the
4350 /// abandoned thread finishes soon after the test.
4351 fn slow_body() -> String {
4352 format!("<p>{}</p>", "&".repeat(48 * 1024))
4353 }
4354
4355 /// The injected timeout for a poll that must give up on [`slow_body`].
4356 const SHORT: Duration = Duration::from_millis(100);
4357
4358 /// Generous enough that [`slow_body`] always finishes in a debug build.
4359 const GENEROUS: Duration = Duration::from_secs(120);
4360
4361 const ETAG_V1: &str = "\"v1\"";
4362 const LAST_MODIFIED_V1: &str = "Tue, 06 Oct 2026 08:00:00 GMT";
4363
4364 /// Permits of the test's own, so an abandoned sanitize here never holds up
4365 /// another test's poll (the production semaphore is process-wide).
4366 fn permits(n: usize) -> &'static Semaphore {
4367 Box::leak(Box::new(Semaphore::new(n)))
4368 }
4369
4370 /// `permits` and `timeout`, with an in-flight set of the call's own.
4371 fn limits(permits: &'static Semaphore, timeout: Duration) -> SanitizeLimits {
4372 SanitizeLimits {
4373 permits,
4374 in_flight: Box::leak(Box::new(InFlight::new())),
4375 timeout,
4376 starvation: Box::leak(Box::new(Starvation::new())),
4377 sanitize: sanitize_body,
4378 }
4379 }
4380
4381 const ORDINARY_A: &str = "<p>Before the slow one <script>x()</script><em>kept</em>.</p>";
4382 const ORDINARY_B: &str = "<p>After the slow one, <b>bold</b>.</p>";
4383
4384 /// An RSS document: `a` (ordinary), `slow` (the given body), `b` (ordinary),
4385 /// `bare` (no body at all).
4386 fn rss(slow: &str) -> String {
4387 let item = |guid: &str, body: Option<&str>| {
4388 let desc = body
4389 .map(|b| format!("<description><![CDATA[{b}]]></description>"))
4390 .unwrap_or_default();
4391 format!(
4392 "<item><title>{guid}</title><link>https://slow.example/{guid}</link>\
4393 <guid>urn:{guid}</guid>{desc}</item>"
4394 )
4395 };
4396 format!(
4397 r#"<?xml version="1.0"?><rss version="2.0"><channel><title>Slow</title>
4398<link>https://slow.example/</link>{}{}{}{}</channel></rss>"#,
4399 item("a", Some(ORDINARY_A)),
4400 item("slow", Some(slow)),
4401 item("b", Some(ORDINARY_B)),
4402 item("bare", None),
4403 )
4404 }
4405
4406 /// Serve `body` with `ETag` / `Last-Modified`, and answer `304` to a
4407 /// request that presents the `ETag` — as a real server would, which is what
4408 /// makes a wrongly saved validator visible: the next poll gets no bodies.
4409 async fn serve_with_validators(body: Arc<Mutex<Vec<u8>>>) -> u16 {
4410 use tokio::io::{AsyncReadExt, AsyncWriteExt};
4411 let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
4412 let port = listener.local_addr().unwrap().port();
4413 tokio::spawn(async move {
4414 while let Ok((mut sock, _)) = listener.accept().await {
4415 let body = body.lock().unwrap().clone();
4416 tokio::spawn(async move {
4417 let mut req = Vec::new();
4418 let mut buf = [0u8; 1024];
4419 while !req.windows(4).any(|w| w == b"\r\n\r\n") {
4420 match sock.read(&mut buf).await {
4421 Ok(0) | Err(_) => return,
4422 Ok(n) => req.extend_from_slice(&buf[..n]),
4423 }
4424 }
4425 let head = String::from_utf8_lossy(&req).to_ascii_lowercase();
4426 let not_modified =
4427 head.contains(&format!("if-none-match: {}", ETAG_V1.to_lowercase()));
4428 let reply = if not_modified {
4429 format!("HTTP/1.1 304 Not Modified\r\nETag: {ETAG_V1}\r\nConnection: close\r\n\r\n")
4430 .into_bytes()
4431 } else {
4432 let mut r = format!(
4433 "HTTP/1.1 200 OK\r\nContent-Type: application/rss+xml\r\nETag: {ETAG_V1}\r\n\
4434 Last-Modified: {LAST_MODIFIED_V1}\r\nContent-Length: {}\r\n\
4435 Connection: close\r\n\r\n",
4436 body.len()
4437 )
4438 .into_bytes();
4439 r.extend_from_slice(&body);
4440 r
4441 };
4442 let _ = sock.write_all(&reply).await;
4443 let _ = sock.flush().await;
4444 });
4445 }
4446 });
4447 port
4448 }
4449
4450 /// A store with one feed row for `host`, served by a fixture with `body`.
4451 async fn fixture(host: &str, body: String) -> (SqlitePool, Feed, Arc<Mutex<Vec<u8>>>) {
4452 let body = Arc::new(Mutex::new(body.into_bytes()));
4453 let port = serve_with_validators(Arc::clone(&body)).await;
4454 crate::net::test_host_override(host, std::net::SocketAddr::from(([127, 0, 0, 1], port)));
4455 let url = format!("http://{host}:{port}/feed.xml");
4456 let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
4457 crate::store::upsert_feed(
4458 &pool,
4459 &NewFeed {
4460 url: url.clone(),
4461 ..Default::default()
4462 },
4463 )
4464 .await
4465 .unwrap();
4466 let feed = crate::store::get_feed_by_url(&pool, &url)
4467 .await
4468 .unwrap()
4469 .unwrap();
4470 (pool, feed, body)
4471 }
4472
4473 async fn bodies(pool: &SqlitePool) -> Vec<(String, Option<String>)> {
4474 sqlx::query_as("SELECT guid, content_html FROM entries ORDER BY guid")
4475 .fetch_all(pool)
4476 .await
4477 .unwrap()
4478 }
4479
4480 async fn validators(pool: &SqlitePool, url: &str) -> (Option<String>, Option<String>) {
4481 let f = crate::store::get_feed_by_url(pool, url)
4482 .await
4483 .unwrap()
4484 .unwrap();
4485 (f.etag, f.last_modified)
4486 }
4487
4488 fn clean(raw: &str) -> Option<String> {
4489 Some(sanitize_html_bounded(raw, MAX_CONTENT_HTML_BYTES))
4490 }
4491
4492 /// The red test for #226. On a current-thread runtime a sanitize run
4493 /// inline stops every timer until it returns; off the runtime, a ticker
4494 /// keeps firing throughout the poll, and the poll ends at the timeout.
4495 #[tokio::test(flavor = "current_thread")]
4496 async fn a_slow_body_is_sanitized_off_the_runtime_and_the_poll_gives_up_on_it() {
4497 let (pool, feed, _) = fixture("slow-body.test", rss(&slow_body())).await;
4498
4499 // The largest gap between consecutive ticks while the poll runs.
4500 let done = Arc::new(AtomicBool::new(false));
4501 let ticker = {
4502 let done = Arc::clone(&done);
4503 tokio::spawn(async move {
4504 let mut last = Instant::now();
4505 let mut worst = Duration::ZERO;
4506 while !done.load(Ordering::SeqCst) {
4507 tokio::time::sleep(Duration::from_millis(5)).await;
4508 worst = worst.max(last.elapsed());
4509 last = Instant::now();
4510 }
4511 worst
4512 })
4513 };
4514 tokio::task::yield_now().await;
4515
4516 let client = build_client().unwrap();
4517 let started = Instant::now();
4518 let outcome = poll_feed_with(&pool, &client, &feed, 0, limits(permits(4), SHORT))
4519 .await
4520 .unwrap();
4521 let took = started.elapsed();
4522 done.store(true, Ordering::SeqCst);
4523 let worst_gap = ticker.await.unwrap();
4524
4525 assert!(
4526 worst_gap < Duration::from_millis(700),
4527 "the runtime stalled for {worst_gap:?} during the poll: the sanitize ran on it"
4528 );
4529 // Tight enough that a sanitize timeout a few times `SHORT` fails it
4530 // (vacuous-test hunt of #274: `< 1 s` let an 8x timeout through;
4531 // 6x leaves a loaded CI runner room).
4532 assert!(
4533 took < SHORT * 6,
4534 "the poll took {took:?}; it should end at the {SHORT:?} timeout"
4535 );
4536 assert!(
4537 matches!(
4538 outcome,
4539 PollOutcome::Failed {
4540 kind: FailureKind::Body,
4541 ..
4542 }
4543 ),
4544 "{outcome:?}"
4545 );
4546 assert_eq!(
4547 bodies(&pool).await,
4548 vec![
4549 // Before the slow entry: sanitized and stored as always.
4550 ("urn:a".to_string(), clean(ORDINARY_A)),
4551 // After it: refused while the slow entry's abandoned sanitize
4552 // runs (the per-feed refusal; `timed_out` backs it up), so
4553 // stored without a body.
4554 ("urn:b".to_string(), None),
4555 ("urn:bare".to_string(), None),
4556 // The slow entry itself: no body, and nothing degraded instead.
4557 ("urn:slow".to_string(), None),
4558 ]
4559 );
4560 assert_eq!(
4561 validators(&pool, &feed.url).await,
4562 (None, None),
4563 "a timed-out poll saved its validators: the next poll would be a 304"
4564 );
4565 }
4566
4567 /// A re-poll that times out keeps every body already stored, and saves no
4568 /// validator; the next poll that can sanitize stores the real bodies
4569 /// (rather than being answered `304` forever).
4570 #[tokio::test]
4571 async fn a_timed_out_repoll_keeps_stored_bodies_and_the_next_poll_stores_real_ones() {
4572 let slow = slow_body();
4573 let (pool, feed, _) = fixture("slow-repoll.test", rss(&slow)).await;
4574 let feed_id = feed.id;
4575 let stored = |guid: &str, body: &str| NewEntry {
4576 guid: guid.to_string(),
4577 content_html: Some(body.to_string()),
4578 ..Default::default()
4579 };
4580 crate::store::insert_entries(
4581 &pool,
4582 feed_id,
4583 &[
4584 stored("urn:slow", "<p>old slow</p>"),
4585 stored("urn:b", "<p>old b</p>"),
4586 ],
4587 0,
4588 )
4589 .await
4590 .unwrap();
4591
4592 let client = build_client().unwrap();
4593 let sem = permits(2);
4594 poll_feed_with(&pool, &client, &feed, 0, limits(sem, SHORT))
4595 .await
4596 .unwrap();
4597 assert_eq!(
4598 bodies(&pool).await,
4599 vec![
4600 ("urn:a".to_string(), clean(ORDINARY_A)),
4601 ("urn:b".to_string(), Some("<p>old b</p>".to_string())),
4602 ("urn:bare".to_string(), None),
4603 ("urn:slow".to_string(), Some("<p>old slow</p>".to_string())),
4604 ],
4605 "a timed-out poll changed a stored body"
4606 );
4607 assert_eq!(validators(&pool, &feed.url).await, (None, None));
4608
4609 // The next poll — same server, same document — with time to finish.
4610 let feed = crate::store::get_feed_by_url(&pool, &feed.url)
4611 .await
4612 .unwrap()
4613 .unwrap();
4614 let outcome = poll_feed_with(&pool, &client, &feed, 0, limits(sem, GENEROUS))
4615 .await
4616 .unwrap();
4617 assert!(
4618 matches!(outcome, PollOutcome::Updated { .. }),
4619 "{outcome:?}"
4620 );
4621 assert_eq!(
4622 bodies(&pool).await,
4623 vec![
4624 ("urn:a".to_string(), clean(ORDINARY_A)),
4625 ("urn:b".to_string(), clean(ORDINARY_B)),
4626 ("urn:bare".to_string(), None),
4627 ("urn:slow".to_string(), clean(&slow)),
4628 ]
4629 );
4630 assert_eq!(
4631 validators(&pool, &feed.url).await,
4632 (
4633 Some(ETAG_V1.to_string()),
4634 Some(LAST_MODIFIED_V1.to_string())
4635 )
4636 );
4637
4638 // And the fixture does honour them, so the check above is meaningful.
4639 let feed = crate::store::get_feed_by_url(&pool, &feed.url)
4640 .await
4641 .unwrap()
4642 .unwrap();
4643 let outcome = poll_feed_with(&pool, &client, &feed, 0, limits(sem, GENEROUS))
4644 .await
4645 .unwrap();
4646 assert!(matches!(outcome, PollOutcome::NotModified), "{outcome:?}");
4647 }
4648
4649 /// Ordinary bodies are stored exactly as the inline path (#224's
4650 /// `sanitize_html_bounded`, which `normalize_entry` and its tests pin)
4651 /// produces them.
4652 #[tokio::test]
4653 async fn ordinary_bodies_are_stored_byte_identical_to_the_inline_path() {
4654 let big = format!("<p>{}</p>", "word ".repeat(18_000));
4655 let doc = rss(&big);
4656 let (pool, feed, _) = fixture("ordinary-bodies.test", doc.clone()).await;
4657 let client = build_client().unwrap();
4658 let outcome = poll_feed_with(&pool, &client, &feed, 0, limits(permits(4), GENEROUS))
4659 .await
4660 .unwrap();
4661 assert!(
4662 matches!(outcome, PollOutcome::Updated { .. }),
4663 "{outcome:?}"
4664 );
4665
4666 let mut inline: Vec<(String, Option<String>)> = parse_feed(doc.as_bytes())
4667 .unwrap()
4668 .entries
4669 .iter()
4670 .map(normalize_entry)
4671 .map(|e| (e.guid, e.content_html))
4672 .collect();
4673 inline.sort();
4674 assert_eq!(bodies(&pool).await, inline);
4675 assert_eq!(
4676 validators(&pool, &feed.url).await,
4677 (
4678 Some(ETAG_V1.to_string()),
4679 Some(LAST_MODIFIED_V1.to_string())
4680 )
4681 );
4682 }
4683
4684 /// With every permit held, a poll waits for one **asynchronously**: other
4685 /// work on the same current-thread runtime goes on, and the poll finishes
4686 /// once a permit is free.
4687 ///
4688 /// A wait that blocks the runtime would never let the test release the
4689 /// permit, so it would hang rather than fail; a timeout on that same
4690 /// runtime could never fire. So the runtime runs on a thread of its own,
4691 /// and the test fails if it has not finished in 10 s (vacuous-test hunt of
4692 /// #274).
4693 #[test]
4694 fn a_poll_waits_for_a_permit_without_blocking_the_runtime() {
4695 use std::sync::mpsc::RecvTimeoutError;
4696 let (done_tx, done_rx) = std::sync::mpsc::channel();
4697 let runtime = std::thread::spawn(move || {
4698 tokio::runtime::Builder::new_current_thread()
4699 .enable_all()
4700 .build()
4701 .unwrap()
4702 .block_on(poll_waits_for_a_permit());
4703 let _ = done_tx.send(());
4704 });
4705 match done_rx.recv_timeout(Duration::from_secs(10)) {
4706 Ok(()) => runtime.join().unwrap(),
4707 // The body panicked: fail with its panic.
4708 Err(RecvTimeoutError::Disconnected) => {
4709 std::panic::resume_unwind(runtime.join().unwrap_err())
4710 }
4711 Err(RecvTimeoutError::Timeout) => {
4712 panic!("the poll blocked the runtime while waiting for a permit")
4713 }
4714 }
4715 }
4716
4717 async fn poll_waits_for_a_permit() {
4718 let (pool, feed, _) = fixture("permit-wait.test", rss(ORDINARY_B)).await;
4719 let sem = permits(1);
4720 let held = sem.acquire().await.unwrap();
4721
4722 let poll = tokio::spawn(async move {
4723 let client = build_client().unwrap();
4724 let outcome = poll_feed_with(&pool, &client, &feed, 0, limits(sem, GENEROUS))
4725 .await
4726 .unwrap();
4727 (outcome, bodies(&pool).await)
4728 });
4729 // Timers still fire on this runtime while the poll waits.
4730 let started = Instant::now();
4731 for _ in 0..20 {
4732 tokio::time::sleep(Duration::from_millis(10)).await;
4733 }
4734 assert!(started.elapsed() < Duration::from_secs(2));
4735 assert!(!poll.is_finished(), "the poll did not wait for a permit");
4736
4737 drop(held);
4738 let (outcome, stored) = poll.await.unwrap();
4739 assert!(
4740 matches!(outcome, PollOutcome::Updated { .. }),
4741 "{outcome:?}"
4742 );
4743 assert!(stored.contains(&("urn:slow".to_string(), clean(ORDINARY_B))));
4744 }
4745
4746 /// **Permits held elsewhere are not this feed's fault** (review of #274).
4747 ///
4748 /// When every permit is held — by sanitizes other feeds' polls gave up on
4749 /// — a poll waits at most the timeout and then defers: nothing stored (no
4750 /// bodiless new entries, no validators), no failure recorded, the error
4751 /// count untouched, the next poll on the ordinary cadence. Filed as a
4752 /// `Body` failure, one hostile feed had turned every healthy feed's poll
4753 /// into a failure, backing them off toward 24 h.
4754 #[tokio::test]
4755 async fn a_poll_starved_of_permits_defers_without_blaming_the_feed() {
4756 let (pool, feed, _) = fixture("permit-starved.test", rss(ORDINARY_B)).await;
4757 // An earlier, unrelated failure, which a deferral must leave as it is.
4758 crate::store::bump_feed_errors(&pool, &feed.url, FailureKind::Fetch, "earlier")
4759 .await
4760 .unwrap();
4761 let sem = permits(1);
4762 let _held = sem.acquire().await.unwrap();
4763
4764 let client = build_client().unwrap();
4765 let started = Instant::now();
4766 let outcome = poll_feed_with(&pool, &client, &feed, 0, limits(sem, SHORT))
4767 .await
4768 .unwrap();
4769 // It waited the timeout for a permit, and not much more (vacuous-test
4770 // hunt of #274: `< 2 s` let a 15x permit wait through; 10x leaves a
4771 // loaded CI runner room).
4772 let took = started.elapsed();
4773 assert!(took >= SHORT && took < SHORT * 10, "{took:?}");
4774 assert_eq!(outcome, PollOutcome::Deferred);
4775 let cadence = Duration::from_secs(3600);
4776 settle_poll(&pool, &feed.url, &outcome, cadence).await;
4777
4778 assert_eq!(bodies(&pool).await, vec![], "stored entries without bodies");
4779 assert_eq!(validators(&pool, &feed.url).await, (None, None));
4780 let after = crate::store::get_feed_by_url(&pool, &feed.url)
4781 .await
4782 .unwrap()
4783 .unwrap();
4784 assert_eq!(
4785 after.consecutive_errors, 1,
4786 "the deferral changed the error count"
4787 );
4788 let next = DateTime::parse_from_rfc3339(after.next_poll.as_deref().unwrap()).unwrap();
4789 let in_secs = (next.with_timezone(&Utc) - Utc::now()).num_seconds();
4790 assert!(
4791 (3500..=3600).contains(&in_secs),
4792 "next poll in {in_secs}s, not on the {cadence:?} cadence"
4793 );
4794 }
4795
4796 /// Wait (bounded) until every one of `sem`'s `n` permits is free again —
4797 /// i.e. every abandoned sanitize holding one has finished.
4798 async fn wait_for_permits(sem: &Semaphore, n: usize) {
4799 let deadline = Instant::now() + Duration::from_secs(120);
4800 while sem.available_permits() < n {
4801 assert!(
4802 Instant::now() < deadline,
4803 "an abandoned sanitize never finished"
4804 );
4805 tokio::time::sleep(Duration::from_millis(20)).await;
4806 }
4807 }
4808
4809 /// **One hostile feed, at most one thread** (review of #274). A re-poll
4810 /// while the feed's abandoned sanitize is still running starts no second
4811 /// sanitize and takes no permit, whatever body it now serves: it is
4812 /// refused at once, as timed out. Once that sanitize finishes, the feed is
4813 /// released and its body is sanitized again. (Renamed in the vacuous-test
4814 /// hunt of #274: it dated from when the refusal was keyed by body hash.)
4815 #[tokio::test]
4816 async fn a_feed_is_refused_until_its_abandoned_sanitize_finishes() {
4817 let slow = slow_body();
4818 let (pool, feed, served) = fixture("abandoned-body.test", rss(&slow)).await;
4819 let client = build_client().unwrap();
4820 let sem = permits(4);
4821 let set: &'static InFlight = Box::leak(Box::new(InFlight::new()));
4822 let short = SanitizeLimits {
4823 permits: sem,
4824 in_flight: set,
4825 timeout: SHORT,
4826 starvation: Box::leak(Box::new(Starvation::new())),
4827 sanitize: sanitize_body,
4828 };
4829
4830 poll_feed_with(&pool, &client, &feed, 0, short)
4831 .await
4832 .unwrap();
4833 assert_eq!(
4834 sem.available_permits(),
4835 3,
4836 "the first poll abandoned one sanitize"
4837 );
4838
4839 // Re-polls while it runs, each serving a changed body: each refused
4840 // at once, no permit taken.
4841 for nonce in 0..3 {
4842 *served.lock().unwrap() = rss(&format!("{slow}<!-- {nonce} -->")).into_bytes();
4843 let started = Instant::now();
4844 let outcome = poll_feed_with(&pool, &client, &feed, 0, short)
4845 .await
4846 .unwrap();
4847 assert!(
4848 matches!(
4849 outcome,
4850 PollOutcome::Failed {
4851 kind: FailureKind::Body,
4852 ..
4853 }
4854 ),
4855 "{outcome:?}"
4856 );
4857 assert!(
4858 started.elapsed() < Duration::from_secs(1),
4859 "{:?}",
4860 started.elapsed()
4861 );
4862 }
4863 assert_eq!(
4864 sem.available_permits(),
4865 3,
4866 "a re-poll started another sanitize for a feed already abandoned"
4867 );
4868
4869 // Once the abandoned sanitize returns, the feed is no longer refused.
4870 *served.lock().unwrap() = rss(&slow).into_bytes();
4871 wait_for_permits(sem, 4).await;
4872 let outcome = poll_feed_with(
4873 &pool,
4874 &client,
4875 &feed,
4876 0,
4877 SanitizeLimits {
4878 timeout: GENEROUS,
4879 ..short
4880 },
4881 )
4882 .await
4883 .unwrap();
4884 assert!(
4885 matches!(outcome, PollOutcome::Updated { .. }),
4886 "{outcome:?}"
4887 );
4888 assert!(bodies(&pool)
4889 .await
4890 .contains(&("urn:slow".to_string(), clean(&slow))));
4891 }
4892
4893 /// **One hostile feed, at most one thread — whatever bytes it serves**
4894 /// (second review of #274). Tracking by body hash let a feed dodge the
4895 /// refusal by changing its slow body each fetch (a nonce, or another slow
4896 /// entry on top), leaving one more uncancellable sanitize behind per poll
4897 /// until it held every permit. Tracked by feed, any body from a feed with
4898 /// an abandoned sanitize still running is refused, without a permit.
4899 #[tokio::test]
4900 async fn a_feed_cannot_dodge_the_refusal_by_changing_its_body() {
4901 let sem = permits(4);
4902 let short = limits(sem, SHORT);
4903 let slow = slow_body();
4904 let first = sanitize_off_runtime("https://hostile.test/feed", slow.clone(), short).await;
4905 assert_eq!(first, Err(SanitizeGaveUp::TimedOut));
4906 for nonce in 0..3 {
4907 let changed = format!("{slow}<!-- {nonce} -->");
4908 let started = Instant::now();
4909 let again = sanitize_off_runtime("https://hostile.test/feed", changed, short).await;
4910 assert_eq!(again, Err(SanitizeGaveUp::StillRunning));
4911 assert!(started.elapsed() < Duration::from_secs(1));
4912 }
4913 assert_eq!(
4914 sem.available_permits(),
4915 3,
4916 "a changed body from the same feed started another sanitize"
4917 );
4918 wait_for_permits(sem, 4).await;
4919 }
4920
4921 /// **Another feed's timeout is not this feed's** (second review of #274).
4922 /// Shared by body hash, one feed's abandoned sanitize refused every other
4923 /// feed carrying the same article, as a `Body` failure of their own. Now
4924 /// another feed is sanitized as usual, and a feed with an ordinary body
4925 /// is never held up by a hostile one.
4926 #[tokio::test]
4927 async fn one_feeds_abandoned_sanitize_does_not_refuse_another_feed() {
4928 let sem = permits(4);
4929 let short = limits(sem, SHORT);
4930 let slow = slow_body();
4931 assert_eq!(
4932 sanitize_off_runtime("https://hostile.test/feed", slow.clone(), short).await,
4933 Err(SanitizeGaveUp::TimedOut)
4934 );
4935 // The same article in another feed: its own attempt, not a refusal.
4936 assert_eq!(
4937 sanitize_off_runtime("https://mirror.test/feed", slow, short).await,
4938 Err(SanitizeGaveUp::TimedOut)
4939 );
4940 assert_eq!(
4941 sanitize_off_runtime("https://ordinary.test/feed", ORDINARY_B.to_string(), short).await,
4942 Ok(clean(ORDINARY_B).unwrap())
4943 );
4944 wait_for_permits(sem, 4).await;
4945 }
4946
4947 /// The **content** body is stored, not the summary, and a body that
4948 /// sanitizes past [`MAX_CONTENT_HTML_BYTES`] is stored within it — with
4949 /// literal expectations, not ones built from the functions under test
4950 /// (vacuous-test hunt of #274: the poll path could have read the summary
4951 /// only, or dropped the bound, and the byte-identical test still passed).
4952 #[tokio::test]
4953 async fn the_content_body_is_stored_over_the_summary_and_within_the_bound() {
4954 let huge = format!("<p>{}<b>tail</b></p>", "a".repeat(MAX_CONTENT_HTML_BYTES));
4955 let item = |guid: &str, content: &str, summary: &str| {
4956 format!(
4957 "<item><title>{guid}</title><link>https://content.example/{guid}</link>\
4958 <guid>urn:{guid}</guid>\
4959 <description><![CDATA[{summary}]]></description>\
4960 <content:encoded><![CDATA[{content}]]></content:encoded></item>"
4961 )
4962 };
4963 let doc = format!(
4964 r#"<?xml version="1.0"?><rss version="2.0"
4965xmlns:content="http://purl.org/rss/1.0/modules/content/"><channel><title>Content</title>
4966<link>https://content.example/</link>{}{}</channel></rss>"#,
4967 item(
4968 "full",
4969 "<p>The <b>full</b> body.<script>x()</script></p>",
4970 "<p>Only the teaser.</p>"
4971 ),
4972 item("huge", &huge, "<p>Short.</p>"),
4973 );
4974 let (pool, feed, _) = fixture("content-body.test", doc).await;
4975 let client = build_client().unwrap();
4976 let outcome = poll_feed_with(&pool, &client, &feed, 0, limits(permits(4), GENEROUS))
4977 .await
4978 .unwrap();
4979 assert!(
4980 matches!(outcome, PollOutcome::Updated { .. }),
4981 "{outcome:?}"
4982 );
4983 let stored = bodies(&pool).await;
4984 assert_eq!(
4985 stored[0],
4986 (
4987 "urn:full".to_string(),
4988 Some("<p>The <b>full</b> body.</p>".to_string())
4989 )
4990 );
4991 let (guid, body) = &stored[1];
4992 assert_eq!(guid, "urn:huge");
4993 let len = body.as_deref().map_or(0, str::len);
4994 assert!(
4995 (MAX_CONTENT_HTML_BYTES / 2..=MAX_CONTENT_HTML_BYTES).contains(&len),
4996 "stored {len} bytes against a bound of {MAX_CONTENT_HTML_BYTES}"
4997 );
4998 }
4999
5000 /// Two permits, emptied mid-poll by [`sanitize_then_starve`].
5001 static STARVE: Semaphore = Semaphore::const_new(2);
5002 static STARVE_RAN: AtomicBool = AtomicBool::new(false);
5003
5004 /// [`sanitize_body`], which first — while its own sanitize holds one of
5005 /// [`STARVE`]'s permits — queues a task for both, so it takes the other
5006 /// now and this one as soon as it is released, and keeps them. The poll's
5007 /// next entry then finds no permit, deterministically.
5008 fn sanitize_then_starve(raw: &str) -> String {
5009 tokio::runtime::Handle::current()
5010 .spawn(async { STARVE.acquire_many(2).await.unwrap().forget() });
5011 let deadline = Instant::now() + Duration::from_secs(10);
5012 while STARVE.available_permits() > 0 {
5013 assert!(Instant::now() < deadline, "the starving task never queued");
5014 std::thread::sleep(Duration::from_millis(1));
5015 }
5016 STARVE_RAN.store(true, Ordering::SeqCst);
5017 sanitize_body(raw)
5018 }
5019
5020 /// **No permit after earlier entries were sanitized still defers the
5021 /// whole poll**: the bodies already sanitized are not stored, no entry is
5022 /// stored without its body, no validator is saved and no failure is
5023 /// filed (vacuous-test hunt of #274: every starved-poll test starved the
5024 /// first entry, so storing the part already done went unnoticed).
5025 #[tokio::test]
5026 async fn a_poll_starved_after_its_first_entry_still_stores_nothing() {
5027 let (pool, feed, _) = fixture("starved-midway.test", rss(ORDINARY_B)).await;
5028 let client = build_client().unwrap();
5029 let outcome = poll_feed_with(
5030 &pool,
5031 &client,
5032 &feed,
5033 0,
5034 SanitizeLimits {
5035 sanitize: sanitize_then_starve,
5036 ..limits(&STARVE, SHORT)
5037 },
5038 )
5039 .await
5040 .unwrap();
5041 assert!(
5042 STARVE_RAN.load(Ordering::SeqCst),
5043 "the first entry was not sanitized"
5044 );
5045 assert_eq!(outcome, PollOutcome::Deferred);
5046 assert_eq!(
5047 bodies(&pool).await,
5048 vec![],
5049 "stored part of a deferred poll"
5050 );
5051 assert_eq!(validators(&pool, &feed.url).await, (None, None));
5052 }
5053
5054 fn sanitize_panics(_: &str) -> String {
5055 panic!("the sanitizer panicked")
5056 }
5057
5058 /// A panic in the sanitizer propagates to the poll, as it did when the
5059 /// sanitize ran inline — not swallowed as a timeout (vacuous-test hunt of
5060 /// #274: the `resume_unwind` arm had no test).
5061 #[tokio::test]
5062 async fn a_sanitizer_panic_propagates_to_the_poll() {
5063 let limits = SanitizeLimits {
5064 sanitize: sanitize_panics,
5065 ..limits(permits(1), GENEROUS)
5066 };
5067 let joined = tokio::spawn(async move {
5068 sanitize_off_runtime("https://panics.test/feed", ORDINARY_B.to_string(), limits).await
5069 })
5070 .await;
5071 let err = joined.expect_err("the panic did not propagate");
5072 assert_eq!(
5073 err.into_panic().downcast_ref::<&str>(),
5074 Some(&"the sanitizer panicked")
5075 );
5076 }
5077
5078 /// **Counted from the start, not from the timeout** (third review of
5079 /// #274). A sanitize counted against its feed only once a poll gave up
5080 /// on it left two ways for one feed to take every permit: overlapping
5081 /// polls (the scheduler, and a subscribe POST, which polls inline) all
5082 /// passed the check before any of them timed out; and a poll dropped
5083 /// mid-sanitize (a client disconnecting from the subscribe request)
5084 /// never counted its sanitize at all. Now a feed with any sanitize still
5085 /// running is not given another: a poll that finds one in flight defers
5086 /// (`Busy`), without a permit.
5087 #[tokio::test]
5088 async fn a_feed_with_a_sanitize_in_flight_is_not_given_another() {
5089 let sem = permits(4);
5090 let set: &'static InFlight = Box::leak(Box::new(InFlight::new()));
5091 let generous = SanitizeLimits {
5092 permits: sem,
5093 in_flight: set,
5094 timeout: GENEROUS,
5095 starvation: Box::leak(Box::new(Starvation::new())),
5096 sanitize: sanitize_body,
5097 };
5098 let feed = "https://hostile.test/feed";
5099
5100 // Overlapping: one poll's sanitize is running, not yet given up on.
5101 let first = tokio::spawn(sanitize_off_runtime(feed, slow_body(), generous));
5102 while sem.available_permits() == 4 {
5103 tokio::time::sleep(Duration::from_millis(5)).await;
5104 }
5105 let started = Instant::now();
5106 assert_eq!(
5107 sanitize_off_runtime(feed, slow_body(), generous).await,
5108 Err(SanitizeGaveUp::Busy)
5109 );
5110 assert!(started.elapsed() < Duration::from_secs(1));
5111 assert_eq!(
5112 sem.available_permits(),
5113 3,
5114 "an overlapping poll of the same feed started a second sanitize"
5115 );
5116
5117 // Dropped: the poll goes away mid-sanitize; the sanitize runs on and
5118 // still counts.
5119 first.abort();
5120 let _ = first.await;
5121 assert_eq!(
5122 sanitize_off_runtime(feed, slow_body(), generous).await,
5123 Err(SanitizeGaveUp::Busy)
5124 );
5125 assert_eq!(
5126 sem.available_permits(),
5127 3,
5128 "a dropped poll's sanitize was not counted against its feed"
5129 );
5130
5131 // Once it returns, the feed is sanitized as usual.
5132 wait_for_permits(sem, 4).await;
5133 assert_eq!(
5134 sanitize_off_runtime(feed, ORDINARY_B.to_string(), generous).await,
5135 Ok(clean(ORDINARY_B).unwrap())
5136 );
5137 assert!(set.lock().is_empty(), "the registry kept a finished feed");
5138 }
5139
5140 /// **A starved instance is counted, not silent** (third review of #274).
5141 /// Four hostile feeds can hold every permit; the polls that then defer
5142 /// store nothing and file nothing against their feeds, so the starvation
5143 /// itself is recorded: each deferral counts, the first marks since when,
5144 /// and the next permit acquired clears that (the count stays).
5145 #[tokio::test]
5146 async fn permit_starvation_is_recorded_and_cleared() {
5147 let sem = permits(1);
5148 let short = limits(sem, SHORT);
5149 assert_eq!(short.starvation.snapshot(), (0, None));
5150 let held = sem.acquire().await.unwrap();
5151 let before = Utc::now().timestamp();
5152 for feed in ["https://one.test/feed", "https://two.test/feed"] {
5153 assert_eq!(
5154 sanitize_off_runtime(feed, ORDINARY_B.to_string(), short).await,
5155 Err(SanitizeGaveUp::NoPermit)
5156 );
5157 }
5158 let (deferred, since) = short.starvation.snapshot();
5159 assert_eq!(
5160 deferred, 2,
5161 "a deferral for want of a permit went uncounted"
5162 );
5163 let since = since.expect("a starved instance did not record since when");
5164 assert!(since >= before && since <= Utc::now().timestamp());
5165
5166 drop(held);
5167 assert_eq!(
5168 sanitize_off_runtime("https://one.test/feed", ORDINARY_B.to_string(), short).await,
5169 Ok(clean(ORDINARY_B).unwrap())
5170 );
5171 assert_eq!(
5172 short.starvation.snapshot(),
5173 (2, None),
5174 "a permit came free but the instance still reads as starved"
5175 );
5176 }
5177}