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