Skip to main content

feather_reader/
standard_site.rs

1//! Reading `standard.site` publications as feeds.
2//!
3//! A publication is not a feed document — it is a record in somebody's atproto
4//! repo, and its "entries" are separate records in the same repo. So this reads
5//! two collections rather than fetching one URL:
6//!
7//! ```text
8//! at://<did>/site.standard.publication/<rkey>
9//!   ├─ resolve <did> → PDS
10//!   ├─ listRecords site.standard.publication  → name, url (every publication
11//!   │                                           in the repo, matched on rkey)
12//!   └─ listRecords site.standard.document     → paged, filtered on `site`
13//! ```
14//!
15//! **Unauthenticated throughout.** This reads *someone else's* repo with no
16//! session, which is why it cannot reuse [`crate::oauth::xrpc::Repo`]: that type
17//! takes its base URL from the session's PDS, hardcodes `repo` to
18//! `session.sub`, and DPoP-signs every send. None of that survives contact with
19//! "read a stranger's repo".
20//!
21//! **Only `textContent` and `description` are read; `content` is ignored.**
22//! `content` is an open union — measured across 449 real documents it carried
23//! six different wrappers and twenty-two block types from five vendor
24//! namespaces, growing with every platform that adopts the lexicon, and it
25//! would drag an HTML-sanitisation surface over foreign input. A document with
26//! neither field still yields an entry: title, date and a link is what an RSS
27//! reader shows for a title-only feed, and is not a failure state.
28
29use serde::Deserialize;
30
31use crate::lexicon::nsid;
32
33/// Documents per `listRecords` page.
34///
35/// Smaller than the protocol default of 100 because a `site.standard.document`
36/// carries the whole article — ~17 KB measured, and the `content` union this
37/// module ignores is still in the wire bytes — while
38/// [`crate::net::read_capped`] bounds a response at 8 MB. 100 long-form
39/// articles per page can exceed that and fail the walk outright.
40const DOCUMENT_PAGE_SIZE: u32 = 25;
41
42/// A parsed `at://` URI: `at://<authority>/<collection>/<rkey>`.
43///
44/// Parsed by hand rather than with `url::Url`, which **cannot read the form that
45/// matters**: `at://did:plc:…/…` fails with *invalid port number*, because the
46/// colons in the DID are taken as a port separator. The handle form parses
47/// fine, which is what makes the failure easy to miss.
48#[derive(Debug, Clone, PartialEq, Eq)]
49pub struct AtUri {
50    pub authority: String,
51    pub collection: String,
52    pub rkey: String,
53}
54
55impl AtUri {
56    /// Parse, or `None` if this is not a well-formed three-segment at-URI.
57    pub fn parse(uri: &str) -> Option<Self> {
58        let rest = uri.strip_prefix(crate::atproto::AT_URI_PREFIX)?;
59        let mut parts = rest.split('/');
60        let (authority, collection, rkey) = (parts.next()?, parts.next()?, parts.next()?);
61        if parts.next().is_some()
62            || authority.is_empty()
63            || collection.is_empty()
64            || rkey.is_empty()
65        {
66            return None;
67        }
68        Some(Self {
69            authority: authority.to_string(),
70            collection: collection.to_string(),
71            rkey: rkey.to_string(),
72        })
73    }
74}
75
76impl std::fmt::Display for AtUri {
77    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
78        write!(
79            f,
80            "at://{}/{}/{}",
81            self.authority, self.collection, self.rkey
82        )
83    }
84}
85
86/// What one read of a publication produced.
87///
88/// **`complete` is the fact the store cannot recover afterwards.** A truncated
89/// walk and a short publication return the same entries; the difference decides
90/// whether an empty result is a quiet blog or a read that gave up, and the
91/// caller has no other way to tell. `fetch` used to drop it on the floor after
92/// logging a warning.
93#[derive(Debug, Clone)]
94pub struct PublicationRead {
95    pub publication: Publication,
96    pub entries: Vec<Entry>,
97    pub complete: bool,
98}
99
100/// The publication record — a pointer, not a feed. Supplies the title and the
101/// base URL that document `path`s are joined onto.
102#[derive(Debug, Clone, PartialEq, Eq)]
103pub struct Publication {
104    pub name: Option<String>,
105    pub url: String,
106}
107
108/// One document, mapped onto the shape the feed pipeline already stores.
109#[derive(Debug, Clone, PartialEq, Eq)]
110pub struct Entry {
111    /// The document's own `at://` URI.
112    ///
113    /// **Not derived from `path`.** `path` is mutable — a publisher who moves an
114    /// article would duplicate their whole archive on the next poll, because
115    /// dedup is `UNIQUE (feed_id, guid)`.
116    pub guid: String,
117    pub title: String,
118    /// When the document was published, re-spelled by `crate::feed::fmt_time`
119    /// into the store's one RFC3339 shape. The reading order sorts on this
120    /// column as a string, so a publisher's spelling cannot go in verbatim.
121    ///
122    /// Taken from `publishedAt` when it parses AND is not in the future,
123    /// otherwise from the TID in the record key, otherwise `None`. See
124    /// `entries_from_records` for why a future date is discarded rather than
125    /// clamped, and why an undated entry is left undated.
126    pub published: Option<String>,
127    /// The joined, **scheme-vetted** permalink — `None` when the document's
128    /// `path` does not resolve to a safe href on the publication's origin.
129    /// `Option` because the guarantee cannot be met unconditionally and the
130    /// store's column is optional too; a title-only entry is not a failure.
131    pub url: Option<String>,
132    /// `description`, else `textContent`, escaped by
133    /// `crate::feed::plain_text_to_html` — both are plain text in the
134    /// lexicon, and the column they land in is rendered as HTML.
135    pub summary: Option<String>,
136}
137
138impl From<Entry> for crate::store::NewEntry {
139    /// The shape the poller stores. Kept here, next to the fields it maps,
140    /// so wiring the reader to the scheduler has nothing left to decide.
141    fn from(e: Entry) -> Self {
142        crate::store::NewEntry {
143            // The document's URI, which the publisher's PDS chose: bounded like
144            // any entry id before it reaches the UNIQUE index (found in review).
145            guid: crate::feed::bound_guid(e.guid),
146            url: e.url,
147            title: Some(e.title),
148            author: None,
149            published: e.published,
150            content_html: e.summary,
151            fetched_at: None,
152            keep_stored_content: false,
153        }
154    }
155}
156
157#[derive(Debug, Deserialize)]
158struct PublicationValue {
159    name: Option<String>,
160    url: String,
161}
162
163#[derive(Debug, Deserialize)]
164struct DocumentValue {
165    title: String,
166    /// Optional so a document without one is an entry with no date, the same
167    /// answer a garbage one gets — the strictness ran the other way, making
168    /// the field this module is willing to DISCARD the one whose absence was
169    /// fatal to the whole record.
170    #[serde(rename = "publishedAt")]
171    published_at: Option<String>,
172    /// Optional for the same reason as `publishedAt`: a document with no
173    /// `path` keeps its title, date and summary rather than vanishing from
174    /// the feed entirely. `Entry.url` is already `Option`.
175    path: Option<String>,
176    /// The at-URI of the publication this document belongs to.
177    ///
178    /// **Load-bearing.** A repo can hold several publications — measured, some
179    /// do — so documents must be filtered by this rather than assumed to belong
180    /// to the one being polled.
181    site: String,
182    #[serde(rename = "textContent")]
183    text_content: Option<String>,
184    description: Option<String>,
185}
186
187/// Find the publication named by `rkey` among a repo's publication records.
188///
189/// Returns its **canonical** at-URI — the one the PDS itself minted — alongside
190/// the record. That canonical URI is what documents reference in their `site`
191/// field, so it is the key the filter uses, rather than the string the reader
192/// subscribed with. Storage is DID-form only (#164), so today the two agree;
193/// taking the PDS's spelling keeps them agreeing if they ever stop.
194pub fn publication_from_records(
195    rkey: &str,
196    records: &[crate::atproto::RecordEntry],
197) -> Option<(String, Publication)> {
198    let entry = records
199        .iter()
200        .find(|r| AtUri::parse(&r.uri).is_some_and(|u| u.rkey == rkey))?;
201    let value: PublicationValue = serde_json::from_value(entry.value.clone()).ok()?;
202    // **The base of every Entry.url, so it is vetted as the href it becomes.**
203    // `net::safe_link` is the same check the RSS entry pipeline applies at
204    // `feed.rs`, and the reason `safe_link.rs` exists as a type at all: the
205    // procedural version of this guarantee was deleted once with a green suite.
206    let url = crate::net::safe_link(&value.url)?;
207    Some((
208        entry.uri.clone(),
209        Publication {
210            name: value
211                .name
212                .map(|n| crate::feed::bound_text(n, crate::feed::MAX_TITLE_BYTES)),
213            url: crate::feed::bound_text(url, crate::feed::MAX_URL_BYTES),
214        },
215    ))
216}
217
218/// Map a repo's document records onto entries, keeping only those belonging to
219/// `canonical_site`.
220pub fn entries_from_records(
221    canonical_site: &str,
222    publication: &Publication,
223    records: &[crate::atproto::RecordEntry],
224) -> Vec<Entry> {
225    // **Normalised to a directory.** `Url::join` is RFC-3986: against a base of
226    // `https://example.com/blog`, a relative `posts/a` resolves to
227    // `/posts/a`, silently dropping the subpath every permalink needs. A
228    // trailing slash makes the base a directory, which is what a publication
229    // URL means.
230    let base = url::Url::parse(&publication.url).ok().map(|mut u| {
231        if !u.path().ends_with('/') {
232            u.set_path(&format!("{}/", u.path()));
233        }
234        u
235    });
236    // One "now" for the whole batch, so two documents of the same poll are
237    // judged against the same instant rather than one being called credible and
238    // an identical one not. (`tid_timestamp` reads the clock again for its own
239    // ceiling; that only ever tightens a bound five minutes away, so it needs
240    // no such agreement.)
241    let now = chrono::Utc::now();
242    // The SAME allowance the record-key branch gets. A publisher's clock runs
243    // ahead of ours as readily as a PDS's does, and judging a stated date
244    // against a bare `now` while the rkey below gets five minutes discarded a
245    // perfectly good date and left the newest post undated — dated, then, by
246    // when we first saw it rather than when it was written.
247    let ceiling = now + chrono::Duration::seconds(crate::atproto::CLOCK_SKEW_GRACE_SECS);
248    records
249        .iter()
250        .filter_map(|record| {
251            // A document that does not deserialise is SKIPPED, not fatal: one
252            // malformed record must not cost a publisher its whole feed.
253            let doc: DocumentValue = serde_json::from_value(record.value.clone()).ok()?;
254            if doc.site != canonical_site {
255                return None;
256            }
257            Some(Entry {
258                guid: record.uri.clone(),
259                title: crate::feed::bound_text(doc.title, crate::feed::MAX_TITLE_BYTES),
260                // **Three rules, in order: what the publisher credibly said,
261                // then when the record was written, then nothing.**
262                //
263                // A stated date in the future is DISCARDED, not clamped to now.
264                // The store refreshes `published` on every poll but stamps
265                // `fetched_at` only once, so a date derived from the current
266                // clock is rewritten hourly: the row stays the newest thing in
267                // the feed forever, is never swept, is never evicted by the
268                // per-feed cap, and sits at the top of the reading list re-dated
269                // to today. Clamping moved the defect rather than fixing it.
270                // Every fall-through here holds still instead — the TID is the
271                // record's real write time, and an undated row is dated by
272                // `fetched_at`, which never moves. An rkey that is not a TID
273                // leaves the entry undated rather than inventing a date from a
274                // slug.
275                published: doc
276                    .published_at
277                    .as_deref()
278                    .and_then(|raw| chrono::DateTime::parse_from_rfc3339(raw).ok())
279                    .map(|d| d.with_timezone(&chrono::Utc))
280                    .filter(|d| *d <= ceiling)
281                    .or_else(|| {
282                        AtUri::parse(&record.uri)
283                            .and_then(|uri| crate::atproto::tid_timestamp(&uri.rkey))
284                    })
285                    .map(crate::feed::fmt_time),
286                // `non_blank` for the same reason the summary uses it: a blank
287                // path joins to the publication's own base, so a handful of
288                // documents with an empty `path` became a handful of entries
289                // all linking to the site root.
290                url: non_blank(doc.path)
291                    .as_deref()
292                    .and_then(|path| join_path(base.as_ref(), path))
293                    .map(|u| crate::feed::bound_text(u, crate::feed::MAX_URL_BYTES)),
294                // `description` first — the authored summary — but only when it
295                // actually says something: a blank one must not shadow the body.
296                // Then ESCAPED, not sanitised: both fields are plain text.
297                summary: non_blank(doc.description)
298                    .or_else(|| non_blank(doc.text_content))
299                    .map(|raw| {
300                        crate::feed::plain_text_to_html_bounded(
301                            &raw,
302                            crate::feed::MAX_CONTENT_HTML_BYTES,
303                        )
304                    }),
305            })
306        })
307        .collect()
308}
309
310/// The ingest floor: the oldest `published` an entry may carry and still be worth
311/// storing, or `None` for "store anything".
312///
313/// **The floor is the window the SWEEP would use, which is not the smaller of the
314/// two.** `store::prune_old_entries` honours the hard ceiling only when it is
315/// strictly older than the rolling window, or when there is no window at all
316/// (`hard_days > 0 && (days <= 0 || hard_days > days)`); a ceiling inside the
317/// window is logged and ignored there, because the hard delete spares nothing and
318/// would otherwise delete exactly the rows the soft delete exists to spare.
319///
320/// So the earliest thing that can delete a row is `retention_days` when there is
321/// a window, and the ceiling only when there is not.
322///
323/// **The pair comes from [`crate::config::Config::retention_for`], and for a
324/// publication it is not the RSS window.** That function is the one home for which
325/// window applies to which kind, precisely so this and the sweep cannot drift: a
326/// publication gets `(0, publication_retention_days)` — no rolling window, and a
327/// generous archive ceiling — because a 14-day window stored ZERO rows from every
328/// real publication measured, their newest documents being 109 to 241 days old.
329/// Fed that pair, the `days == 0` branch below falls through to the ceiling, which
330/// is exactly the sweep pass that can delete such a row.
331///
332/// `a_publications_floor_is_its_archive_ceiling_not_the_rss_window` asserts the
333/// two halves agree; it is the test that fails if either side is changed alone.
334///
335/// Two wrong versions, both worth naming. Keying on `retention_days` alone left
336/// `retention_days = 0` unfloored while the ceiling still deleted at
337/// `retention_hard_days` — the exact cycle the floor exists to prevent, and a
338/// worse one, since the ceiling spares nothing and a starred entry came back
339/// unstarred rather than merely unread. Reaching for the MINIMUM then
340/// over-corrected: at `days = 180, hard = 30` the sweep ignores the ceiling and
341/// deletes nothing before 180 days, while `min` floored ingest at 30 and silently
342/// discarded five months of a publisher's archive that nothing would have deleted.
343///
344/// **An unrepresentable window is no floor, said directly.** `RETENTION_DAYS`
345/// parses into a `u32` with no upper bound, and `u32::MAX` days is an operator
346/// saying "keep everything"; both the duration and the subtraction can fail, and
347/// either failure means `None` here. The previous shape fell back to a sentinel
348/// instant (`MIN_UTC`) and left the outcome resting on the fact that the row
349/// comparison is LEXICOGRAPHIC: `fmt_time` renders an out-of-range year with a
350/// `+`/`-` sign, which sorts either side of a 4-digit year by ASCII accident
351/// rather than by date. Verified: swapping that fallback to `MAX_UTC` — a floor
352/// of the year 262143, which should discard every entry in existence — changed
353/// nothing, because `"2026…" > "+262143…"`. With `None` the ambiguity is gone,
354/// and the comparison never sees a date outside the range it can order.
355///
356/// `now` is a parameter so this is testable without the clock.
357fn ingest_floor(
358    retention_days: u32,
359    retention_hard_days: u32,
360    now: chrono::DateTime<chrono::Utc>,
361) -> Option<String> {
362    let days = if retention_days > 0 {
363        retention_days
364    } else if retention_hard_days > 0 {
365        retention_hard_days
366    } else {
367        return None;
368    };
369    let window = chrono::Duration::try_days(days.into())?;
370    now.checked_sub_signed(window).map(crate::feed::fmt_time)
371}
372
373/// Persist one publication read, and say what the poll amounted to.
374///
375/// `retention_days` is the ingest floor: an entry whose stated date is already
376/// older than the window is not stored, because storing it means the next sweep
377/// deletes it and the next poll re-inserts it with a new row id — which loses its
378/// read state and arrives unread, on that cycle, forever.
379///
380/// **For a publication the caller passes `Config::retention_for`'s pair, which is
381/// the archive ceiling rather than the 14-day window** — so in practice almost
382/// nothing is floored out here, which is the point: measured, a 14-day floor
383/// dropped every document of every real publication tried. See `ingest_floor`.
384///
385/// **An entry with no date is stored anyway.** That is a decision, not an
386/// oversight: it is dated by `fetched_at` instead, which does not move, and the
387/// alternative is discarding an article the reader can never see.
388///
389/// Three consequences, stated because an earlier version of this comment named
390/// only the first and called it "accepted":
391///
392/// * It resurrects once per retention window. `fetched_at` is stamped at the
393///   first insert and never refreshed, so the row ages out, is swept once the
394///   reader has read it, and is re-inserted unread by the next poll.
395/// * Under the hard ceiling it comes back **unstarred as well**. The soft sweep
396///   spares `starred = 1 OR read = 0`; the ceiling spares nothing.
397/// * Until then it OUTRANKS the publication's real articles. Both the sweep and
398///   the per-feed cap order on `COALESCE(published, fetched_at)`, so an undated
399///   entry sorts by the moment it was fetched — that is, as the newest thing in
400///   the feed — and a publication whose documents are mostly undated can push
401///   dated articles out of `max_entries_per_feed`.
402///
403/// **`new_entries` is an upper bound, not a count of rows that survived.**
404/// `insert_entries` reports what it inserted, and the per-feed cap then trims
405/// within the same call, so a poll that inserted 30 rows into a feed capped at 10
406/// can report 30.
407pub async fn store_publication(
408    pool: &sqlx::SqlitePool,
409    url: &str,
410    read: PublicationRead,
411    max_entries_per_feed: i64,
412    retention_days: u32,
413    retention_hard_days: u32,
414) -> anyhow::Result<crate::feed::PollOutcome> {
415    let offered = read.entries.len();
416
417    // **The failure check runs before anything is written.** Stamping the feed
418    // and then returning `Failed` is the shape that makes a broken publication
419    // read as freshly polled on `/stats`, and it is easy to write by accident
420    // because the upsert is the natural first step.
421    //
422    // **Keyed on what was OFFERED, not on what survives the floor.** A truncated
423    // read that produced nothing means the walk gave up before its first record.
424    // A truncated read whose entries are merely older than the window is a
425    // healthy poll of an old publication, and calling that a failure puts it into
426    // backoff that widens forever.
427    if !read.complete && offered == 0 {
428        return Ok(crate::feed::PollOutcome::Failed {
429            backoff: crate::feed::backoff_for(1),
430            kind: crate::feed::FailureKind::Body,
431            detail: crate::feed::failure_detail(
432                "the publication read stopped before its first document",
433            ),
434        });
435    }
436
437    // See [`ingest_floor`] for which of the two windows this is, and why.
438    let floor = ingest_floor(retention_days, retention_hard_days, chrono::Utc::now());
439    let rows: Vec<crate::store::NewEntry> = read
440        .entries
441        .into_iter()
442        .filter(|e| match (&floor, &e.published) {
443            // Lexicographic on the store's one RFC3339 spelling, which sorts
444            // chronologically by construction.
445            (Some(floor), Some(published)) => published.as_str() >= floor.as_str(),
446            // No floor, or no date. The undated case is the decision the doc
447            // above records: kept, dated by `fetched_at`, and it resurrects once
448            // per window.
449            _ => true,
450        })
451        .map(Into::into)
452        .collect();
453
454    // **Say when the floor emptied the read.** A complete read of three
455    // year-old posts and a complete read of an empty publication both return
456    // `Updated { new_entries: 0 }`, stamp the feed, and look green — so a
457    // subscriber to an archived blog gets a blank feed and nothing anywhere says
458    // why. `offered` is already in hand.
459    if offered > rows.len() {
460        tracing::info!(
461            feed = %url,
462            offered,
463            stored = rows.len(),
464            "the retention floor dropped documents older than the window"
465        );
466    }
467
468    let feed_id = crate::store::upsert_feed(
469        pool,
470        &crate::store::NewFeed {
471            url: url.to_string(),
472            title: read.publication.name.clone(),
473            site_url: Some(read.publication.url.clone()),
474            last_polled: Some(crate::feed::fmt_time(chrono::Utc::now())),
475            ..Default::default()
476        },
477    )
478    .await?;
479
480    let new_entries =
481        crate::store::insert_entries(pool, feed_id, &rows, max_entries_per_feed).await?;
482    Ok(crate::feed::PollOutcome::Updated { new_entries })
483}
484
485fn non_blank(s: Option<String>) -> Option<String> {
486    s.filter(|v| !v.trim().is_empty())
487}
488
489/// Join a document `path` onto the publication's base URL.
490///
491/// **`Url::join`, not concatenation.** Concatenating produced
492/// `https://x.com/https://evil.example/a` for an absolute path and buried the
493/// path inside the query for a base carrying one. `join` also keeps the result
494/// on the publication's own origin for a relative path, which is the only shape
495/// the lexicon describes.
496fn join_path(base: Option<&url::Url>, path: &str) -> Option<String> {
497    // **`safe_link` on the way out, not only on the base.** The scheme
498    // guarantee used to live solely in `publication_from_records`; this
499    // function and `Publication` are both `pub`, so a caller that built a
500    // `Publication` some other way (the step-3 poller, from a stored row) gave
501    // an unparseable base — and the no-base branch then returned the
502    // document's `path` verbatim, putting `javascript:` into an entry link.
503    // **No base, no URL.** This branch used to return `safe_link(path)`, which
504    // vets the scheme but NOT the origin — so a caller holding a `Publication`
505    // it did not build through `publication_from_records` (the step-3 poller,
506    // from a stored row whose `site_url` is NULL or malformed) would publish a
507    // publisher-controlled `https://evil.example/x` as a permalink under that
508    // publication's name. The two branches agree now: off-origin is `None`, and
509    // "no origin to be off" is also `None`.
510    let base = base?;
511    match base.join(path) {
512        // A path that resolves off the publication's origin is not a path, it
513        // is a redirect the publisher smuggled into a field we render as theirs.
514        Ok(joined) if joined.origin() == base.origin() => crate::net::safe_link(joined.as_str()),
515        // **No URL, not the homepage.** Falling back to the base gave every
516        // affected entry the same href pointing at the site root — which is
517        // what a publication on an apex domain whose documents live on `www.`
518        // or a CDN would produce for its whole archive, with nothing to say
519        // anything had been dropped.
520        _ => None,
521    }
522}
523
524/// What the document walk should do with one record.
525///
526/// A pure decision so it can be tested without a PDS: the walk itself is a
527/// closure over the network.
528#[derive(Debug, Clone, Copy, PartialEq, Eq)]
529enum DocumentFate {
530    /// Belongs to the requested publication at this index.
531    Keep(usize),
532    /// Belongs to another publication **in this repo** — normal, and the
533    /// reason the `site` filter exists. Not a signal of anything.
534    Sibling,
535    /// References a publication this repo does not have: a `site` spelling
536    /// nothing can ever match. Indistinguishable from a quiet blog without
537    /// saying so, which is why it is counted.
538    Orphan,
539    /// Not a document this reader understands.
540    Malformed,
541}
542
543fn classify_document(
544    record: &crate::atproto::RecordEntry,
545    wanted: &std::collections::HashMap<String, usize>,
546    known: &std::collections::HashSet<&str>,
547) -> DocumentFate {
548    match serde_json::from_value::<DocumentValue>(record.value.clone()) {
549        Ok(doc) if wanted.contains_key(&doc.site) => DocumentFate::Keep(wanted[&doc.site]),
550        Ok(doc) if known.contains(doc.site.as_str()) => DocumentFate::Sibling,
551        Ok(_) => DocumentFate::Orphan,
552        Err(_) => DocumentFate::Malformed,
553    }
554}
555
556/// Read a publication and its documents through the hardened anonymous client.
557///
558/// **Deliberately thin.** Everything that makes this fetch safe already exists
559/// in [`crate::atproto`] and is reused rather than rebuilt:
560///
561/// - [`crate::atproto::resolve_did_to_pds`] runs
562///   [`crate::net::assert_public_target`] on the `serviceEndpoint`, which is a
563///   stranger's string;
564/// - every read goes through [`crate::net::guarded_get_no_privacy`], re-vetting
565///   per hop and pinning the connection, which closes the rebinding window and
566///   supplies the `User-Agent` that 4 of 19 measured endpoints demand;
567/// - [`crate::net::read_capped`] bounds each response;
568/// - [`crate::atproto::PdsClient::list_all_records`] bounds the page count AND
569///   detects a repeated or absent cursor — the trap this module's first draft
570///   walked into, already solved there;
571/// - an XRPC error envelope surfaces as [`crate::atproto::AtProtoError::Xrpc`] rather
572///   than deserialising into an empty page.
573///
574/// The first draft of this module reimplemented all of that, worse. The only
575/// logic left here is the part that is genuinely about standard.site.
576pub async fn fetch(
577    http: &reqwest::Client,
578    plc_directory: &str,
579    uri: &AtUri,
580) -> anyhow::Result<PublicationRead> {
581    // The collection is part of the identity of what was subscribed to, and
582    // this function is `pub`: without the check it lists publications and
583    // matches on rkey alone, so `at://did/app.bsky.feed.post/<rkey>` would be
584    // "read as a publication" whenever a publication shares that rkey.
585    if uri.collection != nsid::STANDARD_PUBLICATION {
586        return Err(
587            NotAPublication(format!("{uri} is not a {} URI", nsid::STANDARD_PUBLICATION)).into(),
588        );
589    }
590    fetch_repo(
591        http,
592        plc_directory,
593        &uri.authority,
594        std::slice::from_ref(&uri.rkey),
595    )
596    .await?
597    .pop()
598    .unwrap_or_else(|| Err(NotAPublication(format!("{uri} was not read")).into()))
599}
600
601/// **The repo answered, and what it holds is not this publication** — the
602/// record was deleted, never existed, or is unreadable, or the URI names
603/// another collection.
604///
605/// Typed so the poller can tell it from a network failure: filed as `Fetch`,
606/// an author deleting their publication read as their server being down, in
607/// the public cause histogram (found in review).
608#[derive(Debug)]
609pub struct NotAPublication(pub String);
610
611impl std::fmt::Display for NotAPublication {
612    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
613        f.write_str(&self.0)
614    }
615}
616
617impl std::error::Error for NotAPublication {}
618
619/// Read several publications in ONE repo with one walk of its documents.
620pub async fn fetch_repo(
621    http: &reqwest::Client,
622    plc_directory: &str,
623    did: &str,
624    rkeys: &[String],
625) -> anyhow::Result<Vec<anyhow::Result<PublicationRead>>> {
626    fetch_repo_capped(
627        http,
628        plc_directory,
629        did,
630        rkeys,
631        crate::atproto::MAX_LARGE_RECORDS,
632        crate::atproto::MAX_LIST_BYTES,
633    )
634    .await
635}
636
637/// [`fetch_repo`] with the per-publication document cap passed in, so the
638/// fairness between siblings is testable without 2,000 documents.
639pub(crate) async fn fetch_repo_capped(
640    http: &reqwest::Client,
641    plc_directory: &str,
642    did: &str,
643    requested: &[String],
644    per_site_cap: usize,
645    budget_bytes: usize,
646) -> anyhow::Result<Vec<anyhow::Result<PublicationRead>>> {
647    use anyhow::Context;
648
649    // Nothing asked for, nothing fetched (found in review: an empty request
650    // still resolved the DID and listed the repo).
651    if requested.is_empty() {
652        return Ok(Vec::new());
653    }
654    // Each rkey read once; a repeat gets a copy of its read. Indexing by the
655    // raw list made the first of two equal rkeys an empty, "complete" read.
656    let mut rkeys: Vec<String> = Vec::new();
657    for r in requested {
658        if !rkeys.contains(r) {
659            rkeys.push(r.clone());
660        }
661    }
662    let rkeys = rkeys.as_slice();
663
664    let pds = crate::atproto::resolve_did_to_pds(http, plc_directory, did)
665        .await
666        .with_context(|| format!("resolving the PDS for {did}"))?;
667    let client = crate::atproto::PdsClient::anonymous(http.clone(), pds, did.to_string());
668
669    // **One budget for both walks, because they nest.** The publications are
670    // still held when the documents walk runs, so two ceilings would let this
671    // read hold twice the bound the box was sized for.
672    let mut budget = crate::atproto::ByteBudget::new(budget_bytes);
673    let (publications, skipped_publications) = client
674        .list_all_records_skipping_within(nsid::STANDARD_PUBLICATION, &mut budget)
675        .await
676        .with_context(|| format!("listing publications for {did}"))?;
677    if skipped_publications > 0 {
678        tracing::warn!(
679            repo = %did,
680            skipped = skipped_publications,
681            "skipped malformed publication records in this repo"
682        );
683    }
684
685    // Each requested publication, if this repo holds it.
686    let wanted: Vec<Option<(String, Publication)>> = rkeys
687        .iter()
688        .map(|rkey| publication_from_records(rkey, &publications))
689        .collect();
690    let index_of: std::collections::HashMap<String, usize> = wanted
691        .iter()
692        .enumerate()
693        .filter_map(|(i, w)| w.as_ref().map(|(site, _)| (site.clone(), i)))
694        .collect();
695    let not_found = |rkey: &str| {
696        anyhow::Error::new(NotAPublication(format!(
697            "at://{did}/{}/{rkey} is not a readable site.standard.publication",
698            nsid::STANDARD_PUBLICATION
699        )))
700    };
701    if index_of.is_empty() {
702        return Ok(requested.iter().map(|r| Err(not_found(r))).collect());
703    }
704
705    // **One walk of the documents for every requested publication, with a
706    // cap PER publication.** Filtered inside the walk, so each cap counts
707    // that publication's own documents: a shared cap would let a busy
708    // publication fill the window and starve a quiet sibling — permanently,
709    // and worse with every post the busy one makes.
710    let known: std::collections::HashSet<&str> =
711        publications.iter().map(|p| p.uri.as_str()).collect();
712    let mut kept_per = vec![0usize; rkeys.len()];
713    let mut capped = vec![false; rkeys.len()];
714    let mut orphaned = 0usize;
715    let documents = client
716        .list_recent_matching_within(
717            nsid::STANDARD_DOCUMENT,
718            per_site_cap.saturating_mul(index_of.len()),
719            &mut budget,
720            DOCUMENT_PAGE_SIZE,
721            |record| match classify_document(record, &index_of, &known) {
722                DocumentFate::Keep(i) if kept_per[i] < per_site_cap => {
723                    kept_per[i] += 1;
724                    true
725                }
726                DocumentFate::Keep(i) => {
727                    capped[i] = true;
728                    false
729                }
730                // Orphans are counted while WALKING: a `site` spelling nothing
731                // in this repo matches looks exactly like an empty publication
732                // otherwise.
733                DocumentFate::Orphan => {
734                    orphaned += 1;
735                    false
736                }
737                DocumentFate::Sibling | DocumentFate::Malformed => false,
738            },
739        )
740        .await
741        .with_context(|| format!("listing documents for {did}"))?;
742
743    // Split the one walk back into one read per publication.
744    let mut per_site: Vec<Vec<crate::atproto::RecordEntry>> = vec![Vec::new(); rkeys.len()];
745    for record in documents.records {
746        let site = serde_json::from_value::<DocumentValue>(record.value.clone()).map(|d| d.site);
747        if let Some(&i) = site.ok().as_deref().and_then(|s| index_of.get(s)) {
748            per_site[i].push(record);
749        }
750    }
751    let reads: Vec<anyhow::Result<PublicationRead>> = wanted
752        .into_iter()
753        .zip(per_site)
754        .enumerate()
755        .map(|(i, (w, records))| match w {
756            None => Err(not_found(&rkeys[i])),
757            Some((site, publication)) => {
758                let walk = crate::atproto::RecordWalk {
759                    records,
760                    complete: documents.complete && !capped[i],
761                    malformed: documents.malformed,
762                };
763                Ok(read_from(publication, &site, walk, orphaned))
764            }
765        })
766        .collect();
767    // Moved out, not cloned: one copy of every read is the bound the shared
768    // budget was sized for. Only a repeated rkey gets a copy.
769    let mut reads: Vec<Option<anyhow::Result<PublicationRead>>> =
770        reads.into_iter().map(Some).collect();
771    let positions: Vec<usize> = requested
772        .iter()
773        .map(|r| {
774            rkeys
775                .iter()
776                .position(|k| k == r)
777                .expect("every requested rkey is in rkeys")
778        })
779        .collect();
780    Ok(positions
781        .iter()
782        .enumerate()
783        .map(|(n, &u)| {
784            let repeated_later = positions[n + 1..].contains(&u);
785            match (&reads[u], repeated_later) {
786                (Some(Ok(read)), true) => Ok(read.clone()),
787                (Some(Err(_)), _) | (None, _) => Err(not_found(&requested[n])),
788                (Some(Ok(_)), false) => match reads[u].take() {
789                    Some(Ok(read)) => Ok(read),
790                    _ => Err(not_found(&requested[n])),
791                },
792            }
793        })
794        .collect())
795}
796
797/// Assemble a [`PublicationRead`] from a finished walk.
798///
799/// **Extracted so `complete` is testable.** It is the one fact the store cannot
800/// recover afterwards, and reaching the `false` case through `fetch` needs a PLC
801/// directory, a PDS, and a walk that truncates — three mocked hosts for one
802/// boolean. Reaching it here needs a struct.
803fn read_from(
804    publication: Publication,
805    canonical_site: &str,
806    documents: crate::atproto::RecordWalk,
807    orphaned: usize,
808) -> PublicationRead {
809    let entries = entries_from_records(canonical_site, &publication, &documents.records);
810
811    // **A walk that stopped early is not a short archive.** Reading part of a
812    // publication is acceptable; reporting it as the whole of one is not, and
813    // when the part is empty — a quiet publication whose busy sibling fills
814    // every page this reader will fetch — the feed looks healthy and stays
815    // empty forever.
816    // Logged AND returned. Logging it alone is how the fact that the walk gave
817    // up reached nobody who could act on it: a truncated read and a short
818    // publication produce the same entries, and only the caller can tell the
819    // difference between "quiet blog" and "we stopped early".
820    if !documents.complete {
821        tracing::warn!(
822            site = %canonical_site,
823            kept = entries.len(),
824            "stopped reading this publication before its documents ran out"
825        );
826    }
827    // **A spelling mismatch on `site` looks exactly like an empty
828    // publication.** A publication with no documents is normal, so the poller
829    // would call this healthy forever; if the repo HAD documents and none
830    // matched, say so, because that is the shape of a bug rather than of a
831    // quiet blog.
832    // **Skipped records are said out loud** (#177), for the same reason as
833    // orphans below: a document we could not read looks exactly like one that
834    // was never written.
835    if documents.malformed > 0 {
836        tracing::warn!(
837            site = %canonical_site,
838            skipped = documents.malformed,
839            "skipped malformed document records for this publication"
840        );
841    }
842    if orphaned > 0 {
843        tracing::warn!(
844            site = %canonical_site,
845            orphaned,
846            "documents in this repo reference no publication in it — a `site` spelling nothing matches"
847        );
848    }
849    PublicationRead {
850        publication,
851        entries,
852        complete: documents.complete,
853    }
854}
855
856#[cfg(test)]
857pub(crate) mod tests {
858    use super::*;
859    use crate::atproto::RecordEntry;
860    use serde_json::json;
861
862    const DID: &str = "did:plc:ohutz6x5acjmpuulp3x7wxxc";
863
864    /// **A walk that gave up must not report itself as a whole publication.**
865    ///
866    /// `fetch` used to log that and return only the entries, so the fact reached
867    /// nobody who could act on it — and a truncated read and a quiet blog produce
868    /// exactly the same entries.
869    #[test]
870    fn a_read_carries_whether_the_walk_finished() {
871        let publication = Publication {
872            name: Some("Scan's Lab".to_string()),
873            url: "https://example.com/blog/".to_string(),
874        };
875        for complete in [true, false] {
876            let walk = crate::atproto::RecordWalk {
877                records: Vec::new(),
878                complete,
879                malformed: 0,
880            };
881            let read = read_from(publication.clone(), "at://d/c/r", walk, 0);
882            assert_eq!(
883                read.complete, complete,
884                "the walk said complete={complete} and the read said {}",
885                read.complete
886            );
887        }
888    }
889
890    // -- storing a publication read -----------------------------------------
891
892    fn read_of(entries: Vec<Entry>, complete: bool) -> PublicationRead {
893        PublicationRead {
894            publication: Publication {
895                name: Some("Scan's Lab".to_string()),
896                url: "https://example.com/blog/".to_string(),
897            },
898            entries,
899            complete,
900        }
901    }
902
903    fn entry_dated(guid: &str, published: Option<&str>) -> Entry {
904        Entry {
905            guid: guid.to_string(),
906            title: "T".to_string(),
907            published: published.map(str::to_string),
908            url: Some("https://example.com/blog/a".to_string()),
909            summary: None,
910        }
911    }
912
913    fn days_ago(n: i64) -> String {
914        crate::feed::fmt_time(chrono::Utc::now() - chrono::Duration::days(n))
915    }
916
917    const PUB_URL: &str = "at://did:plc:ohutz6x5acjmpuulp3x7wxxc/site.standard.publication/3lab";
918
919    async fn pool() -> sqlx::SqlitePool {
920        crate::store::init_url("sqlite::memory:").await.unwrap()
921    }
922
923    #[tokio::test]
924    async fn a_complete_read_stores_its_entries_and_reports_them() {
925        let pool = pool().await;
926        let read = read_of(
927            vec![
928                entry_dated("at://d/c/1", Some(&days_ago(1))),
929                entry_dated("at://d/c/2", Some(&days_ago(2))),
930            ],
931            true,
932        );
933        let outcome = store_publication(&pool, PUB_URL, read, 0, 14, 180)
934            .await
935            .unwrap();
936        assert!(
937            matches!(
938                outcome,
939                crate::feed::PollOutcome::Updated { new_entries: 2 }
940            ),
941            "expected two new entries, got {outcome:?}"
942        );
943        let n: i64 = sqlx::query_scalar("SELECT COUNT(*) FROM entries")
944            .fetch_one(&pool)
945            .await
946            .unwrap();
947        assert_eq!(n, 2, "the entries were not stored");
948    }
949
950    /// **Starvation is keyed on what was OFFERED, not on what was stored.**
951    ///
952    /// A truncated read that produced nothing is a failure: the walk gave up
953    /// before the first record. But a truncated read whose entries the retention
954    /// floor dropped is not — the read worked, the entries are simply older than
955    /// the window, and calling that a failure puts a healthy publication into
956    /// backoff forever.
957    #[tokio::test]
958    async fn an_incomplete_read_that_offered_nothing_is_a_failure() {
959        let pool = pool().await;
960        let outcome = store_publication(&pool, PUB_URL, read_of(vec![], false), 0, 14, 180)
961            .await
962            .unwrap();
963        assert!(
964            matches!(outcome, crate::feed::PollOutcome::Failed { .. }),
965            "a truncated read that produced nothing is not a healthy poll: {outcome:?}"
966        );
967    }
968
969    #[tokio::test]
970    async fn an_incomplete_read_whose_entries_the_floor_dropped_is_not_a_failure() {
971        let pool = pool().await;
972        let read = read_of(
973            vec![entry_dated("at://d/c/old", Some(&days_ago(900)))],
974            false,
975        );
976        let outcome = store_publication(&pool, PUB_URL, read, 0, 14, 180)
977            .await
978            .unwrap();
979        assert!(
980            !matches!(outcome, crate::feed::PollOutcome::Failed { .. }),
981            "the read offered an entry; the floor dropping it is not a failed poll: {outcome:?}"
982        );
983    }
984
985    /// A failed poll must not look like a successful one on `/stats`.
986    #[tokio::test]
987    async fn a_failed_read_does_not_stamp_last_polled() {
988        let pool = pool().await;
989        let outcome = store_publication(&pool, PUB_URL, read_of(vec![], false), 0, 14, 180)
990            .await
991            .unwrap();
992        assert!(matches!(outcome, crate::feed::PollOutcome::Failed { .. }));
993        let stamped: Option<String> =
994            sqlx::query_scalar("SELECT last_polled FROM feeds WHERE url = ?1")
995                .bind(PUB_URL)
996                .fetch_optional(&pool)
997                .await
998                .unwrap()
999                .flatten();
1000        assert_eq!(
1001            stamped, None,
1002            "a failed poll stamped last_polled, so the feed reads as freshly polled"
1003        );
1004    }
1005
1006    #[tokio::test]
1007    async fn an_entry_already_older_than_the_window_is_not_stored() {
1008        let pool = pool().await;
1009        let read = read_of(
1010            vec![
1011                entry_dated("at://d/c/fresh", Some(&days_ago(1))),
1012                entry_dated("at://d/c/ancient", Some(&days_ago(900))),
1013            ],
1014            true,
1015        );
1016        store_publication(&pool, PUB_URL, read, 0, 14, 180)
1017            .await
1018            .unwrap();
1019        let guids: Vec<String> = sqlx::query_scalar("SELECT guid FROM entries ORDER BY guid")
1020            .fetch_all(&pool)
1021            .await
1022            .unwrap();
1023        assert_eq!(
1024            guids,
1025            vec!["at://d/c/fresh".to_string()],
1026            "an entry the next sweep would delete was stored anyway"
1027        );
1028    }
1029
1030    /// **The floor follows whichever window actually deletes, not the rolling one.**
1031    ///
1032    /// `retention_days = 0` is a supported configuration meaning "no rolling
1033    /// window", and the hard ceiling stays alive independently. Keying the ingest
1034    /// floor on the rolling window alone let a 900-day-old entry in, which the
1035    /// ceiling then deleted and the next poll re-inserted — and the ceiling spares
1036    /// nothing, so a starred entry came back unstarred.
1037    #[tokio::test]
1038    async fn the_floor_follows_the_hard_ceiling_when_the_window_is_disabled() {
1039        let pool = pool().await;
1040        let read = read_of(
1041            vec![entry_dated("at://d/c/ancient", Some(&days_ago(900)))],
1042            true,
1043        );
1044        store_publication(&pool, PUB_URL, read, 0, 0, 180)
1045            .await
1046            .unwrap();
1047        let n: i64 = sqlx::query_scalar("SELECT COUNT(*) FROM entries")
1048            .fetch_one(&pool)
1049            .await
1050            .unwrap();
1051        assert_eq!(
1052            n, 0,
1053            "an entry the hard ceiling will delete was stored, so it will resurrect"
1054        );
1055    }
1056
1057    /// **The WINDOW is the floor when there is one, because it deletes first.**
1058    ///
1059    /// With a 14-day rolling window and a 180-day ceiling, an entry 100 days old
1060    /// is inside the ceiling and outside the window — so the sweep takes it once
1061    /// it has been read and the next poll puts it back. Taking the longer of the
1062    /// two would store it.
1063    #[tokio::test]
1064    async fn the_floor_follows_the_window_when_both_are_set() {
1065        let pool = pool().await;
1066        let read = read_of(
1067            vec![entry_dated("at://d/c/hundred", Some(&days_ago(100)))],
1068            true,
1069        );
1070        store_publication(&pool, PUB_URL, read, 0, 14, 180)
1071            .await
1072            .unwrap();
1073        let n: i64 = sqlx::query_scalar("SELECT COUNT(*) FROM entries")
1074            .fetch_one(&pool)
1075            .await
1076            .unwrap();
1077        assert_eq!(
1078            n, 0,
1079            "an entry inside the ceiling but outside the window was stored, so it will cycle"
1080        );
1081    }
1082
1083    /// With both windows off there is no floor, because nothing will delete it.
1084    /// **The two halves of the retention policy, asserted against each other.**
1085    ///
1086    /// `Config::retention_for` decides which window a kind gets; `ingest_floor`
1087    /// decides what is worth storing; `store::prune_old_entries` decides what is
1088    /// deleted. The first two are asserted here against the third's rule, because
1089    /// a drift between them is the resurrection cycle — a row the store keeps, the
1090    /// sweep deletes, and the next poll re-inserts unread — and nothing else in the
1091    /// suite would notice one side changing alone.
1092    ///
1093    /// The publication case is the one that matters: fed the RSS window a
1094    /// publication stores nothing at all, because a real one's newest document is
1095    /// months old.
1096    #[test]
1097    fn a_publications_floor_is_its_archive_ceiling_not_the_rss_window() {
1098        let config = crate::config::Config::default();
1099        let now = chrono::DateTime::parse_from_rfc3339("2026-09-29T12:00:00Z")
1100            .unwrap()
1101            .with_timezone(&chrono::Utc);
1102
1103        let (days, hard) = config.retention_for(crate::feed::FeedKind::Publication);
1104        let floor = ingest_floor(days, hard, now).expect("a publication has an ingest floor");
1105        let expected = crate::feed::fmt_time(
1106            now - chrono::Duration::days(config.publication_retention_days.into()),
1107        );
1108        assert_eq!(
1109            floor, expected,
1110            "a publication's floor must be its archive ceiling ({} days), because \
1111             that is the only sweep pass that can delete its rows",
1112            config.publication_retention_days,
1113        );
1114
1115        // And the number this rules OUT: the RSS window would drop every document
1116        // of every real publication measured (newest 109 to 241 days old).
1117        let rss_floor = ingest_floor(config.retention_days, config.retention_hard_days, now)
1118            .expect("an RSS feed has an ingest floor");
1119        assert!(
1120            floor < rss_floor,
1121            "the publication floor ({floor}) is no older than the RSS one \
1122             ({rss_floor}), so a months-old document would still be dropped",
1123        );
1124        let a_real_publications_newest_document =
1125            crate::feed::fmt_time(now - chrono::Duration::days(109));
1126        assert!(
1127            a_real_publications_newest_document.as_str() >= floor.as_str(),
1128            "the newest document a real publication offered would be refused at \
1129             ingest: {a_real_publications_newest_document} against a floor of {floor}",
1130        );
1131        assert!(
1132            a_real_publications_newest_document.as_str() < rss_floor.as_str(),
1133            "this assertion is only meaningful while the RSS window WOULD have \
1134             dropped it, and it no longer does",
1135        );
1136    }
1137
1138    /// **The floor, all four configurations, without a database or the clock.**
1139    ///
1140    /// The end-to-end tests below observe the floor through what gets stored,
1141    /// which cannot see the difference between "no floor" and "a floor the
1142    /// lexicographic row comparison happens to sort past". This asserts the value
1143    /// itself.
1144    #[test]
1145    fn the_ingest_floor_is_the_window_the_sweep_would_use() {
1146        let now = chrono::DateTime::parse_from_rfc3339("2026-09-26T12:00:00Z")
1147            .unwrap()
1148            .with_timezone(&chrono::Utc);
1149
1150        assert_eq!(
1151            ingest_floor(14, 180, now).as_deref(),
1152            Some("2026-09-12T12:00:00Z"),
1153            "with both set, the floor is the WINDOW — the thing that deletes first",
1154        );
1155        assert_eq!(
1156            ingest_floor(180, 30, now).as_deref(),
1157            Some("2026-03-30T12:00:00Z"),
1158            "a ceiling INSIDE the window is one `prune_old_entries` ignores, so it \
1159             must not lower the floor — `min` here discarded five months of archive \
1160             that nothing would have deleted",
1161        );
1162        assert_eq!(
1163            ingest_floor(0, 30, now).as_deref(),
1164            Some("2026-08-27T12:00:00Z"),
1165            "with no window the ceiling stands alone, and it still deletes",
1166        );
1167        assert_eq!(
1168            ingest_floor(0, 0, now),
1169            None,
1170            "with no retention at all there is nothing to floor against",
1171        );
1172        // `u32::MAX` days is an operator saying "keep everything". Both the
1173        // duration and the subtraction fail there, and the answer is the ABSENCE
1174        // of a floor — not a sentinel instant whose rendering the row comparison
1175        // then has to sort correctly by accident.
1176        assert_eq!(
1177            ingest_floor(u32::MAX, u32::MAX, now),
1178            None,
1179            "an unrepresentable window produced a floor, so the comparison is \
1180             resting on how `fmt_time` renders an out-of-range year",
1181        );
1182    }
1183
1184    /// **A ceiling INSIDE the window is one the sweep ignores, so the floor must
1185    /// ignore it too.**
1186    ///
1187    /// `prune_old_entries` honours `retention_hard_days` only when it is strictly
1188    /// older than `retention_days` (or when there is no window at all): a ceiling
1189    /// inside the window is logged and dropped, because the hard delete spares
1190    /// nothing and would otherwise delete exactly the rows the soft delete exists
1191    /// to spare. So at `days = 180, hard = 30` nothing is deleted before 180 days.
1192    ///
1193    /// An earlier version took the MINIMUM of the two, floored ingest at 30 days,
1194    /// and silently discarded five months of a publisher's archive that nothing
1195    /// would ever have deleted. The older test above passes under both rules —
1196    /// `min(14, 180)` and "the window" are both 14 — which is why this case is
1197    /// the one that had to be written.
1198    #[tokio::test]
1199    async fn a_ceiling_inside_the_window_does_not_lower_the_floor() {
1200        let pool = pool().await;
1201        let read = read_of(
1202            vec![entry_dated("at://d/c/hundred", Some(&days_ago(100)))],
1203            true,
1204        );
1205        store_publication(&pool, PUB_URL, read, 0, 180, 30)
1206            .await
1207            .unwrap();
1208        let n: i64 = sqlx::query_scalar("SELECT COUNT(*) FROM entries")
1209            .fetch_one(&pool)
1210            .await
1211            .unwrap();
1212        assert_eq!(
1213            n, 1,
1214            "an entry inside the 180-day window was dropped because of a 30-day \
1215             ceiling the sweep ignores — five months of archive discarded at ingest \
1216             that nothing would have deleted",
1217        );
1218    }
1219
1220    /// **A complete read of an empty publication is a healthy poll, not a
1221    /// failure.** The failure branch is keyed on `!complete && offered == 0`, and
1222    /// dropping either half of that makes a brand-new or fully-archived
1223    /// publication go into backoff that widens forever.
1224    #[tokio::test]
1225    async fn a_complete_read_of_nothing_is_not_a_failure() {
1226        let pool = pool().await;
1227        let outcome = store_publication(&pool, PUB_URL, read_of(vec![], true), 0, 14, 180)
1228            .await
1229            .unwrap();
1230        assert!(
1231            matches!(
1232                outcome,
1233                crate::feed::PollOutcome::Updated { new_entries: 0 }
1234            ),
1235            "a complete read of an empty publication was not a healthy poll: {outcome:?}"
1236        );
1237    }
1238
1239    /// **The other half of `a_failed_read_does_not_stamp_last_polled`.** That test
1240    /// alone is satisfied by never stamping at all, which would leave every
1241    /// publication permanently due and `/stats` permanently wrong.
1242    #[tokio::test]
1243    async fn a_successful_read_stamps_last_polled() {
1244        let pool = pool().await;
1245        let read = read_of(vec![entry_dated("at://d/c/1", Some(&days_ago(1)))], true);
1246        store_publication(&pool, PUB_URL, read, 0, 14, 180)
1247            .await
1248            .unwrap();
1249        let stamped: Option<String> =
1250            sqlx::query_scalar("SELECT last_polled FROM feeds WHERE url = ?1")
1251                .bind(PUB_URL)
1252                .fetch_one(&pool)
1253                .await
1254                .unwrap();
1255        assert!(
1256            stamped.is_some(),
1257            "a successful read left `last_polled` NULL, so the publication stays \
1258             due forever and `/stats` never shows it as polled",
1259        );
1260    }
1261
1262    /// **The saturating window must saturate in the KEEPING direction.**
1263    ///
1264    /// `an_absurd_retention_window_does_not_panic` only asserts that it returns.
1265    /// A saturation that produced `DateTime::MAX` instead — or a floor of "now" —
1266    /// would satisfy it while silently discarding the publisher's entire archive:
1267    /// `u32::MAX` days is an operator saying "keep everything".
1268    #[tokio::test]
1269    async fn an_absurd_retention_window_keeps_everything_rather_than_nothing() {
1270        let pool = pool().await;
1271        let read = read_of(
1272            vec![entry_dated("at://d/c/ancient", Some(&days_ago(10_000)))],
1273            true,
1274        );
1275        store_publication(&pool, PUB_URL, read, 0, u32::MAX, u32::MAX)
1276            .await
1277            .expect("a huge window is a wide floor, not a crash");
1278        let n: i64 = sqlx::query_scalar("SELECT COUNT(*) FROM entries")
1279            .fetch_one(&pool)
1280            .await
1281            .unwrap();
1282        assert_eq!(
1283            n, 1,
1284            "a window of u32::MAX days dropped a 27-year-old entry, so the \
1285             saturation went the wrong way",
1286        );
1287    }
1288
1289    /// **`max_entries_per_feed` is passed through, and every test above passes
1290    /// `0`.** So a call that dropped the argument, or passed a constant, would
1291    /// have gone unnoticed — and this is the cap that decides which of a
1292    /// publisher's articles a reader keeps.
1293    #[tokio::test]
1294    async fn the_per_feed_cap_is_the_one_the_caller_passed() {
1295        let pool = pool().await;
1296        let read = read_of(
1297            (0..5)
1298                .map(|i| entry_dated(&format!("at://d/c/{i}"), Some(&days_ago(i + 1))))
1299                .collect(),
1300            true,
1301        );
1302        store_publication(&pool, PUB_URL, read, 2, 0, 0)
1303            .await
1304            .unwrap();
1305        let n: i64 = sqlx::query_scalar("SELECT COUNT(*) FROM entries")
1306            .fetch_one(&pool)
1307            .await
1308            .unwrap();
1309        assert_eq!(
1310            n, 2,
1311            "five entries under a cap of two left {n} rows, so the caller's cap is \
1312             not the one being applied",
1313        );
1314    }
1315
1316    #[tokio::test]
1317    async fn no_retention_at_all_means_no_ingest_floor() {
1318        let pool = pool().await;
1319        let read = read_of(
1320            vec![entry_dated("at://d/c/ancient", Some(&days_ago(900)))],
1321            true,
1322        );
1323        store_publication(&pool, PUB_URL, read, 0, 0, 0)
1324            .await
1325            .unwrap();
1326        let n: i64 = sqlx::query_scalar("SELECT COUNT(*) FROM entries")
1327            .fetch_one(&pool)
1328            .await
1329            .unwrap();
1330        assert_eq!(
1331            n, 1,
1332            "with nothing deleting it, an old entry is worth keeping"
1333        );
1334    }
1335
1336    /// A retention window near `u32::MAX` must not panic the poller.
1337    #[tokio::test]
1338    async fn an_absurd_retention_window_does_not_panic() {
1339        let pool = pool().await;
1340        let read = read_of(vec![entry_dated("at://d/c/x", Some(&days_ago(1)))], true);
1341        store_publication(&pool, PUB_URL, read, 0, u32::MAX, u32::MAX)
1342            .await
1343            .expect("a huge window is a wide floor, not a crash");
1344    }
1345
1346    /// The feed row learns the publication's name and homepage — the only reason
1347    /// beyond the timestamp that the upsert is there at all. `upsert_feed`
1348    /// COALESCEs both, so dropping either is silent.
1349    #[tokio::test]
1350    async fn the_feed_row_learns_the_publications_name_and_site() {
1351        let pool = pool().await;
1352        let read = read_of(vec![entry_dated("at://d/c/1", Some(&days_ago(1)))], true);
1353        store_publication(&pool, PUB_URL, read, 0, 14, 180)
1354            .await
1355            .unwrap();
1356        let (title, site): (Option<String>, Option<String>) =
1357            sqlx::query_as("SELECT title, site_url FROM feeds WHERE url = ?1")
1358                .bind(PUB_URL)
1359                .fetch_one(&pool)
1360                .await
1361                .unwrap();
1362        assert_eq!(
1363            title.as_deref(),
1364            Some("Scan's Lab"),
1365            "the name never reached the row"
1366        );
1367        assert_eq!(
1368            site.as_deref(),
1369            Some("https://example.com/blog/"),
1370            "the homepage never reached the row"
1371        );
1372    }
1373
1374    /// **The undated entry is kept, deliberately.** It is dated by `fetched_at`,
1375    /// which holds still, and the alternative is discarding an article the reader
1376    /// can never see. It resurrects once per retention window; that is accepted.
1377    #[tokio::test]
1378    async fn an_undated_entry_is_stored_rather_than_dropped() {
1379        let pool = pool().await;
1380        let read = read_of(vec![entry_dated("at://d/c/undated", None)], true);
1381        store_publication(&pool, PUB_URL, read, 0, 14, 180)
1382            .await
1383            .unwrap();
1384        let n: i64 = sqlx::query_scalar("SELECT COUNT(*) FROM entries")
1385            .fetch_one(&pool)
1386            .await
1387            .unwrap();
1388        assert_eq!(n, 1, "an entry with no date was discarded");
1389    }
1390
1391    /// A record key from a real atproto repo, and the instant it decodes to.
1392    ///
1393    /// **Fixed, not minted.** A test that mints a TID expects "about now", and
1394    /// "about now" is satisfied by any source of the current time — including a
1395    /// fallback that has had the rkey ripped out of it and returns `Utc::now()`
1396    /// instead. Both tests named for the rkey fallback used to survive exactly
1397    /// that mutation. A historical key with a stated answer cannot.
1398    const PAST_TID: &str = "3jzfcijpj2z2a";
1399    const PAST_TID_WRITTEN_AT: &str = "2023-06-30T15:03:01Z";
1400
1401    fn rec(collection: &str, rkey: &str, value: serde_json::Value) -> RecordEntry {
1402        RecordEntry {
1403            uri: format!("at://{DID}/{collection}/{rkey}"),
1404            cid: None,
1405            value,
1406        }
1407    }
1408
1409    fn publication(rkey: &str, url: &str) -> RecordEntry {
1410        rec(
1411            nsid::STANDARD_PUBLICATION,
1412            rkey,
1413            json!({ "name": "Scan's Lab", "url": url }),
1414        )
1415    }
1416
1417    fn document(rkey: &str, site: &str, title: &str, path: &str) -> RecordEntry {
1418        rec(
1419            nsid::STANDARD_DOCUMENT,
1420            rkey,
1421            json!({
1422                "title": title,
1423                "publishedAt": "2026-07-11T00:00:00Z",
1424                "path": path,
1425                "site": site,
1426                "textContent": "body",
1427            }),
1428        )
1429    }
1430
1431    /// A PLC directory and the author's PDS, both on one loopback port, for a
1432    /// repo holding `records` (`(collection, rkey, value)`; an empty rkey serves
1433    /// a malformed envelope with no `uri`). Each collection is
1434    /// served as one full page carrying a cursor, then an empty page, because a
1435    /// real PDS returns a cursor on its final page and the walk must stop on the
1436    /// empty one. Returns the PLC base URL and a count of `listRecords` calls.
1437    pub(crate) async fn serve_repo(
1438        did: &'static str,
1439        records: Vec<(&'static str, &'static str, serde_json::Value)>,
1440    ) -> (String, std::sync::Arc<std::sync::atomic::AtomicUsize>) {
1441        use axum::extract::{Query, Request};
1442        use std::collections::HashMap;
1443        use std::sync::atomic::{AtomicUsize, Ordering};
1444        use std::sync::Arc;
1445
1446        let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
1447        let addr = listener.local_addr().unwrap();
1448        let port = addr.port();
1449        // Unique per server: the override table is process-wide and tests run
1450        // in parallel, so a shared name would let one test reach another's.
1451        let (plc_host, pds_host) = (
1452            format!("plc-{port}.repo.test"),
1453            format!("pds-{port}.repo.test"),
1454        );
1455        for host in [&plc_host, &pds_host] {
1456            crate::net::test_host_override(host, addr);
1457        }
1458        let hits = Arc::new(AtomicUsize::new(0));
1459        let counter = Arc::clone(&hits);
1460        let records = Arc::new(records);
1461        let app = axum::Router::new().fallback(
1462            move |Query(q): Query<HashMap<String, String>>, req: Request| {
1463                let records = Arc::clone(&records);
1464                let counter = Arc::clone(&counter);
1465                async move {
1466                    let path = req.uri().path().to_string();
1467                    if path == format!("/{did}") {
1468                        return axum::Json(json!({
1469                            "id": did,
1470                            "service": [{
1471                                "id": "#atproto_pds",
1472                                "type": "AtprotoPersonalDataServer",
1473                                "serviceEndpoint": format!("http://{pds_host}:{port}"),
1474                            }],
1475                        }));
1476                    }
1477                    assert_eq!(path, "/xrpc/com.atproto.repo.listRecords", "unexpected request");
1478                    counter.fetch_add(1, Ordering::SeqCst);
1479                    let collection = q.get("collection").cloned().unwrap_or_default();
1480                    // Pages honour `limit` and an offset cursor, and a cursor
1481                    // comes back on the final page too, as a real PDS's does.
1482                    let limit: usize = q.get("limit").and_then(|l| l.parse().ok()).unwrap_or(50);
1483                    let offset: usize = q.get("cursor").and_then(|c| c.parse().ok()).unwrap_or(0);
1484                    let page: Vec<_> = records
1485                        .iter()
1486                        .filter(|(c, _, _)| *c == collection)
1487                        .skip(offset)
1488                        .take(limit)
1489                        .map(|(c, rkey, value)| {
1490                            // An empty rkey serves the #177 shape: an envelope
1491                            // with no `uri`, which no record can be read from.
1492                            if rkey.is_empty() {
1493                                return json!({ "cid": "bafy", "value": value });
1494                            }
1495                            json!({ "uri": format!("at://{did}/{c}/{rkey}"), "cid": "bafy", "value": value })
1496                        })
1497                        .collect();
1498                    let next = (offset + page.len()).to_string();
1499                    axum::Json(json!({ "records": page, "cursor": next }))
1500                }
1501            },
1502        );
1503        tokio::spawn(async move { axum::serve(listener, app).await.unwrap() });
1504        (format!("http://{plc_host}:{port}"), hits)
1505    }
1506
1507    /// **The mocked repo reads through the real fetch path.** Pins that
1508    /// [`serve_repo`] is a faithful enough PDS for A0 and the polling tests to
1509    /// mean something: PLC resolution, both walks, the `site` filter and the
1510    /// empty-page stop all run for real.
1511    #[tokio::test]
1512    async fn fetch_reads_a_publication_from_a_mocked_repo() {
1513        const OWN: &str = "did:plc:fetchmock";
1514        let site = format!("at://{OWN}/{}/mine", nsid::STANDARD_PUBLICATION);
1515        let sibling = format!("at://{OWN}/{}/other", nsid::STANDARD_PUBLICATION);
1516        let doc = |site: &str, title: &str| {
1517            json!({ "title": title, "publishedAt": "2026-07-11T00:00:00Z",
1518                    "path": "/p", "site": site })
1519        };
1520        let (plc, hits) = serve_repo(
1521            OWN,
1522            vec![
1523                (
1524                    nsid::STANDARD_PUBLICATION,
1525                    "mine",
1526                    json!({ "name": "Mine", "url": "https://mine.example" }),
1527                ),
1528                (
1529                    nsid::STANDARD_PUBLICATION,
1530                    "other",
1531                    json!({ "name": "Other", "url": "https://other.example" }),
1532                ),
1533                (nsid::STANDARD_DOCUMENT, "3l2fmaaaaaa2a", doc(&site, "kept")),
1534                (
1535                    nsid::STANDARD_DOCUMENT,
1536                    "3l2fmaaaaaa2b",
1537                    doc(&sibling, "sibling's"),
1538                ),
1539            ],
1540        )
1541        .await;
1542        let client = crate::feed::build_client().unwrap();
1543        let read = fetch(&client, &plc, &AtUri::parse(&site).unwrap())
1544            .await
1545            .unwrap();
1546        assert!(read.complete, "the walk did not finish");
1547        assert_eq!(read.publication.name.as_deref(), Some("Mine"));
1548        let titles: Vec<&str> = read.entries.iter().map(|e| e.title.as_str()).collect();
1549        assert_eq!(
1550            titles,
1551            vec!["kept"],
1552            "the site filter let a sibling through"
1553        );
1554        assert_eq!(
1555            hits.load(std::sync::atomic::Ordering::SeqCst),
1556            4,
1557            "two collections, each a full page then an empty one"
1558        );
1559    }
1560
1561    /// A read that fails — here, a DID whose PLC directory cannot be reached —
1562    /// is a poll FAILURE with backoff, never an `Err`: `Err` from the poll seam
1563    /// means the local store is broken.
1564    #[tokio::test]
1565    async fn a_failed_publication_read_is_a_poll_failure() {
1566        let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
1567        // A VALID DID, so the row is a Publication and the failure is the
1568        // network's — an invalid one is Unsupported and fails for another reason.
1569        let url = format!(
1570            "at://did:plc:unreachableaaaaaaaaaaaaa/{}/x",
1571            nsid::STANDARD_PUBLICATION
1572        );
1573        crate::store::upsert_feed(
1574            &pool,
1575            &crate::store::NewFeed {
1576                url: url.clone(),
1577                ..Default::default()
1578            },
1579        )
1580        .await
1581        .unwrap();
1582        let feed = crate::store::get_feed_by_url(&pool, &url)
1583            .await
1584            .unwrap()
1585            .unwrap();
1586        let mut config = crate::config::Config::default();
1587        config.oauth.plc_directory = "http://plc.nowhere.invalid".into();
1588        let client = crate::feed::build_client().unwrap();
1589        let outcome = crate::feed::poll_feed_by_kind(&pool, &client, &config, &feed)
1590            .await
1591            .expect("a source failure surfaced as a store error");
1592        assert!(
1593            matches!(
1594                outcome,
1595                crate::feed::PollOutcome::Failed {
1596                    kind: crate::feed::FailureKind::Fetch,
1597                    ..
1598                }
1599            ),
1600            "an unreachable publication was not a fetch failure: {outcome:?}"
1601        );
1602    }
1603
1604    /// Review of #225: an author who deletes their publication record is not
1605    /// an unreachable publisher. Every fetch error used to be filed as
1606    /// `Fetch`, putting it in the public histogram's network bucket.
1607    #[tokio::test]
1608    async fn a_deleted_publication_record_is_not_an_unreachable_publisher() {
1609        const GONE: &str = "did:plc:goneaaaaaaaaaaaaaaaaaaaa";
1610        let site = format!("at://{GONE}/{}/deleted", nsid::STANDARD_PUBLICATION);
1611        // The repo answers, but holds no record with that rkey.
1612        let (plc, _) = serve_repo(
1613            GONE,
1614            vec![(
1615                nsid::STANDARD_PUBLICATION,
1616                "another",
1617                json!({ "name": "Other", "url": "https://o.example" }),
1618            )],
1619        )
1620        .await;
1621        let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
1622        crate::store::upsert_feed(
1623            &pool,
1624            &crate::store::NewFeed {
1625                url: site.clone(),
1626                ..Default::default()
1627            },
1628        )
1629        .await
1630        .unwrap();
1631        let feed = crate::store::get_feed_by_url(&pool, &site)
1632            .await
1633            .unwrap()
1634            .unwrap();
1635        let mut config = crate::config::Config::default();
1636        config.oauth.plc_directory = plc;
1637        let client = crate::feed::build_client().unwrap();
1638        let outcome = crate::feed::poll_feed_by_kind(&pool, &client, &config, &feed)
1639            .await
1640            .unwrap();
1641        assert!(
1642            matches!(
1643                outcome,
1644                crate::feed::PollOutcome::Failed {
1645                    kind: crate::feed::FailureKind::Parse,
1646                    ..
1647                }
1648            ),
1649            "a deleted publication was filed as a network failure: {outcome:?}"
1650        );
1651    }
1652
1653    /// A PLC directory and PDS on one port that answer however the test says:
1654    /// the PLC lookup with `plc_status` (200 serves a DID document), and every
1655    /// `listRecords` with `(status, body)` after `delay` — or, when `endless`,
1656    /// a fresh page of one sibling document, forever.
1657    async fn serve_answering(
1658        did: &'static str,
1659        plc_status: u16,
1660        pds: (u16, &'static str),
1661        delay: std::time::Duration,
1662        endless: bool,
1663    ) -> String {
1664        use std::sync::atomic::{AtomicUsize, Ordering};
1665        use std::sync::Arc;
1666        let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
1667        let addr = listener.local_addr().unwrap();
1668        let port = addr.port();
1669        let (plc_host, pds_host) = (
1670            format!("plc-{port}.answer.test"),
1671            format!("pds-{port}.answer.test"),
1672        );
1673        crate::net::test_host_override(&plc_host, addr);
1674        crate::net::test_host_override(&pds_host, addr);
1675        let pages = Arc::new(AtomicUsize::new(0));
1676        let endpoint = format!("http://{pds_host}:{port}");
1677        let app = axum::Router::new().fallback(move |req: axum::extract::Request| {
1678            let pages = Arc::clone(&pages);
1679            let endpoint = endpoint.clone();
1680            async move {
1681                use axum::response::IntoResponse;
1682                if req.uri().path() == format!("/{did}") {
1683                    let status = axum::http::StatusCode::from_u16(plc_status).unwrap();
1684                    let doc = json!({ "id": did, "service": [{ "id": "#atproto_pds",
1685                        "type": "AtprotoPersonalDataServer", "serviceEndpoint": endpoint }] });
1686                    return (status, axum::Json(doc)).into_response();
1687                }
1688                tokio::time::sleep(delay).await;
1689                if endless {
1690                    let n = pages.fetch_add(1, Ordering::SeqCst);
1691                    let body = json!({ "records": [{
1692                        "uri": format!("at://{did}/{}/3lend{n:08}", nsid::STANDARD_DOCUMENT),
1693                        "cid": "b",
1694                        "value": { "title": "x", "path": "/x", "publishedAt": "2026-07-11T00:00:00Z",
1695                                   "site": format!("at://{did}/{}/other", nsid::STANDARD_PUBLICATION) } }],
1696                        "cursor": format!("c{n}") });
1697                    return axum::Json(body).into_response();
1698                }
1699                let status = axum::http::StatusCode::from_u16(pds.0).unwrap();
1700                (status, [("content-type", "application/json")], pds.1).into_response()
1701            }
1702        });
1703        tokio::spawn(async move { axum::serve(listener, app).await.unwrap() });
1704        format!("http://{plc_host}:{port}")
1705    }
1706
1707    async fn poll_publication_at(
1708        did: &str,
1709        plc: String,
1710        deadline: Option<std::time::Duration>,
1711    ) -> crate::feed::PollOutcome {
1712        let site = format!("at://{did}/{}/mine", nsid::STANDARD_PUBLICATION);
1713        let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
1714        crate::store::upsert_feed(
1715            &pool,
1716            &crate::store::NewFeed {
1717                url: site.clone(),
1718                ..Default::default()
1719            },
1720        )
1721        .await
1722        .unwrap();
1723        let feed = crate::store::get_feed_by_url(&pool, &site)
1724            .await
1725            .unwrap()
1726            .unwrap();
1727        let mut config = crate::config::Config::default();
1728        config.oauth.plc_directory = plc;
1729        if let Some(d) = deadline {
1730            config.publication_read_deadline = d;
1731        }
1732        let client = crate::feed::build_client().unwrap();
1733        crate::feed::poll_feed_by_kind(&pool, &client, &config, &feed)
1734            .await
1735            .unwrap()
1736    }
1737
1738    fn kind_of(outcome: &crate::feed::PollOutcome) -> Option<crate::feed::FailureKind> {
1739        match outcome {
1740            crate::feed::PollOutcome::Failed { kind, .. } => Some(*kind),
1741            _ => None,
1742        }
1743    }
1744
1745    /// Review of #225: one publication read had no overall deadline — only
1746    /// FETCH_TIMEOUT per request x MAX_LIST_PAGES — and the tick waited for it.
1747    #[tokio::test]
1748    async fn a_publication_read_has_an_overall_deadline() {
1749        const DID: &str = "did:plc:slowrepoaaaaaaaaaaaaaaaa";
1750        let plc = serve_answering(
1751            DID,
1752            200,
1753            (200, ""),
1754            std::time::Duration::from_millis(50),
1755            true,
1756        )
1757        .await;
1758        let started = std::time::Instant::now();
1759        let outcome =
1760            poll_publication_at(DID, plc, Some(std::time::Duration::from_millis(300))).await;
1761        assert!(
1762            started.elapsed() < std::time::Duration::from_secs(3),
1763            "the read ran {:?}",
1764            started.elapsed()
1765        );
1766        assert_eq!(
1767            kind_of(&outcome),
1768            Some(crate::feed::FailureKind::Fetch),
1769            "{outcome:?}"
1770        );
1771    }
1772
1773    /// Review of #225: answers that ARRIVED were filed as `Fetch` ("the request
1774    /// never produced a response"), putting deleted and deactivated accounts in
1775    /// the network bucket.
1776    #[tokio::test]
1777    async fn a_publication_failure_is_filed_under_what_happened() {
1778        let zero = std::time::Duration::ZERO;
1779        const GONE: &str = "did:plc:tombstonedaaaaaaaaaaaaaa";
1780        let plc = serve_answering(GONE, 404, (200, ""), zero, false).await;
1781        let outcome = poll_publication_at(GONE, plc, None).await;
1782        assert_eq!(
1783            kind_of(&outcome),
1784            Some(crate::feed::FailureKind::Status),
1785            "PLC 404: {outcome:?}"
1786        );
1787
1788        const NOREPO: &str = "did:plc:norepoaaaaaaaaaaaaaaaaaa";
1789        let plc = serve_answering(
1790            NOREPO,
1791            200,
1792            (400, r#"{"error":"RepoNotFound"}"#),
1793            zero,
1794            false,
1795        )
1796        .await;
1797        let outcome = poll_publication_at(NOREPO, plc, None).await;
1798        assert_eq!(
1799            kind_of(&outcome),
1800            Some(crate::feed::FailureKind::Status),
1801            "RepoNotFound: {outcome:?}"
1802        );
1803
1804        const GARBLED: &str = "did:plc:garbledaaaaaaaaaaaaaaaaa";
1805        let plc = serve_answering(GARBLED, 200, (200, r#"{"records":"x"}"#), zero, false).await;
1806        let outcome = poll_publication_at(GARBLED, plc, None).await;
1807        assert_eq!(
1808            kind_of(&outcome),
1809            Some(crate::feed::FailureKind::Parse),
1810            "garbled body: {outcome:?}"
1811        );
1812    }
1813
1814    /// Two of 17 measured publications had no documents. That is a healthy,
1815    /// empty feed — not a failure, and not a reason to back off.
1816    #[tokio::test]
1817    async fn an_empty_publication_is_a_healthy_poll() {
1818        const EMPTY: &str = "did:plc:emptypubaaaaaaaaaaaaaaaa";
1819        let site = format!("at://{EMPTY}/{}/quiet", nsid::STANDARD_PUBLICATION);
1820        let (plc, _) = serve_repo(
1821            EMPTY,
1822            vec![(
1823                nsid::STANDARD_PUBLICATION,
1824                "quiet",
1825                json!({ "name": "Quiet", "url": "https://quiet.example" }),
1826            )],
1827        )
1828        .await;
1829        let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
1830        crate::store::upsert_feed(
1831            &pool,
1832            &crate::store::NewFeed {
1833                url: site.clone(),
1834                ..Default::default()
1835            },
1836        )
1837        .await
1838        .unwrap();
1839        let feed = crate::store::get_feed_by_url(&pool, &site)
1840            .await
1841            .unwrap()
1842            .unwrap();
1843        let mut config = crate::config::Config::default();
1844        config.oauth.plc_directory = plc;
1845        let client = crate::feed::build_client().unwrap();
1846        let outcome = crate::feed::poll_feed_by_kind(&pool, &client, &config, &feed)
1847            .await
1848            .unwrap();
1849        assert!(
1850            matches!(
1851                outcome,
1852                crate::feed::PollOutcome::Updated { new_entries: 0 }
1853            ),
1854            "an empty publication was not a healthy poll: {outcome:?}"
1855        );
1856    }
1857
1858    /// **#177 through the real fetch path.** A stranger's repo with one
1859    /// malformed publication record and one malformed document, each beside
1860    /// good ones: both are skipped, and the publication still reads.
1861    #[tokio::test]
1862    async fn fetch_skips_malformed_records_beside_good_ones() {
1863        const OWN: &str = "did:plc:malformedrepo";
1864        let site = format!("at://{OWN}/{}/mine", nsid::STANDARD_PUBLICATION);
1865        let doc = |title: &str| {
1866            json!({ "title": title, "publishedAt": "2026-07-11T00:00:00Z",
1867                    "path": "/p", "site": site })
1868        };
1869        let (plc, _) = serve_repo(
1870            OWN,
1871            vec![
1872                (
1873                    nsid::STANDARD_PUBLICATION,
1874                    "",
1875                    json!({ "name": "Broken", "url": "https://x.example" }),
1876                ),
1877                (
1878                    nsid::STANDARD_PUBLICATION,
1879                    "mine",
1880                    json!({ "name": "Mine", "url": "https://mine.example" }),
1881                ),
1882                (nsid::STANDARD_DOCUMENT, "", doc("unreadable")),
1883                (nsid::STANDARD_DOCUMENT, "3l2mfaaaaaa2a", doc("kept")),
1884            ],
1885        )
1886        .await;
1887        let client = crate::feed::build_client().unwrap();
1888        let read = fetch(&client, &plc, &AtUri::parse(&site).unwrap())
1889            .await
1890            .expect("a malformed record stalled a stranger's publication");
1891        assert!(read.complete);
1892        let titles: Vec<&str> = read.entries.iter().map(|e| e.title.as_str()).collect();
1893        assert_eq!(titles, vec!["kept"]);
1894    }
1895
1896    // ---- 0.4.0 step 2b: one walk per repo ------------------------------------
1897
1898    const SHARED: &str = "did:plc:sharedrepoaaaaaaaaaaaaaa";
1899
1900    fn shared_doc(site_rkey: &str, title: &str, at: &str) -> serde_json::Value {
1901        json!({ "title": title, "publishedAt": at, "path": format!("/{title}"),
1902                "site": format!("at://{SHARED}/{}/{site_rkey}", nsid::STANDARD_PUBLICATION) })
1903    }
1904
1905    /// Cost scales with the repo, not the publication: nine publications in one
1906    /// repo were nine full walks of its documents. Two due publications in one
1907    /// repo now cost one walk, and each still gets only its own documents.
1908    #[tokio::test]
1909    async fn two_publications_in_one_repo_cost_one_walk() {
1910        let (plc, hits) = serve_repo(
1911            SHARED,
1912            vec![
1913                (
1914                    nsid::STANDARD_PUBLICATION,
1915                    "alpha",
1916                    json!({ "name": "Alpha", "url": "https://alpha.example" }),
1917                ),
1918                (
1919                    nsid::STANDARD_PUBLICATION,
1920                    "beta",
1921                    json!({ "name": "Beta", "url": "https://beta.example" }),
1922                ),
1923                (
1924                    nsid::STANDARD_DOCUMENT,
1925                    "3l2shaaaaaa2a",
1926                    shared_doc("alpha", "a1", "2026-07-11T00:00:00Z"),
1927                ),
1928                (
1929                    nsid::STANDARD_DOCUMENT,
1930                    "3l2shaaaaaa2b",
1931                    shared_doc("beta", "b1", "2026-07-10T00:00:00Z"),
1932                ),
1933                (
1934                    nsid::STANDARD_DOCUMENT,
1935                    "3l2shaaaaaa2c",
1936                    shared_doc("alpha", "a2", "2026-07-09T00:00:00Z"),
1937                ),
1938            ],
1939        )
1940        .await;
1941        let client = crate::feed::build_client().unwrap();
1942        let reads = fetch_repo(
1943            &client,
1944            &plc,
1945            SHARED,
1946            &["alpha".to_string(), "beta".to_string()],
1947        )
1948        .await
1949        .unwrap();
1950        assert_eq!(
1951            hits.load(std::sync::atomic::Ordering::SeqCst),
1952            4,
1953            "two collections walked once each (a page then an empty page), not once per publication"
1954        );
1955        let titles = |r: &anyhow::Result<PublicationRead>| {
1956            let mut t: Vec<String> = r
1957                .as_ref()
1958                .unwrap()
1959                .entries
1960                .iter()
1961                .map(|e| e.title.clone())
1962                .collect();
1963            t.sort();
1964            t
1965        };
1966        assert_eq!(titles(&reads[0]), vec!["a1", "a2"]);
1967        assert_eq!(titles(&reads[1]), vec!["b1"]);
1968    }
1969
1970    /// One walk for several publications must not let a busy one fill the
1971    /// window and starve a quiet sibling: each publication has its own cap.
1972    #[tokio::test]
1973    async fn a_busy_publication_does_not_starve_its_quiet_sibling() {
1974        let mut records = vec![
1975            (
1976                nsid::STANDARD_PUBLICATION,
1977                "busy",
1978                json!({ "name": "Busy", "url": "https://busy.example" }),
1979            ),
1980            (
1981                nsid::STANDARD_PUBLICATION,
1982                "quiet",
1983                json!({ "name": "Quiet", "url": "https://quiet.example" }),
1984            ),
1985        ];
1986        for (rkey, title) in [
1987            ("3l2bsaaaaaa2a", "b1"),
1988            ("3l2bsaaaaaa2b", "b2"),
1989            ("3l2bsaaaaaa2c", "b3"),
1990            ("3l2bsaaaaaa2d", "b4"),
1991            ("3l2bsaaaaaa2e", "b5"),
1992        ] {
1993            records.push((
1994                nsid::STANDARD_DOCUMENT,
1995                rkey,
1996                shared_doc("busy", title, "2026-07-11T00:00:00Z"),
1997            ));
1998        }
1999        records.push((
2000            nsid::STANDARD_DOCUMENT,
2001            "3l2bsaaaaaa2f",
2002            shared_doc("quiet", "q1", "2026-01-01T00:00:00Z"),
2003        ));
2004        let (plc, _) = serve_repo(SHARED, records).await;
2005        let client = crate::feed::build_client().unwrap();
2006        let reads = fetch_repo_capped(
2007            &client,
2008            &plc,
2009            SHARED,
2010            &["busy".to_string(), "quiet".to_string()],
2011            2,
2012            crate::atproto::MAX_LIST_BYTES,
2013        )
2014        .await
2015        .unwrap();
2016        let busy = reads[0].as_ref().unwrap();
2017        let quiet = reads[1].as_ref().unwrap();
2018        assert_eq!(
2019            busy.entries.len(),
2020            2,
2021            "the busy publication was not capped at its own cap"
2022        );
2023        assert!(!busy.complete, "a capped publication was reported complete");
2024        assert_eq!(quiet.entries.len(), 1, "the quiet sibling was starved");
2025        assert!(quiet.complete);
2026    }
2027
2028    /// **A known limitation, kept as a test.** A one-repo group shares one byte
2029    /// budget, so big siblings can spend it and the walk stops before a quiet
2030    /// publication's documents. A static per-publication share fixed this and
2031    /// broke two worse things (a busy publication beside idle siblings was cut
2032    /// to a fraction, and one large document failed its publication every poll),
2033    /// so it was removed. Unreachable at measured scale.
2034    #[tokio::test]
2035    #[ignore = "known limitation: a one-repo group shares one byte budget (#229)"]
2036    async fn big_siblings_do_not_spend_a_quiet_publications_share() {
2037        let body = "w".repeat(20 * 1024);
2038        let mut records = vec![
2039            (
2040                nsid::STANDARD_PUBLICATION,
2041                "big1",
2042                json!({ "name": "Big 1", "url": "https://b1.example" }),
2043            ),
2044            (
2045                nsid::STANDARD_PUBLICATION,
2046                "big2",
2047                json!({ "name": "Big 2", "url": "https://b2.example" }),
2048            ),
2049            (
2050                nsid::STANDARD_PUBLICATION,
2051                "quiet",
2052                json!({ "name": "Quiet", "url": "https://q.example" }),
2053            ),
2054        ];
2055        for i in 0..200 {
2056            let rkey: &'static str = Box::leak(format!("3l2big{i:06}").into_boxed_str());
2057            let site = if i % 2 == 0 { "big1" } else { "big2" };
2058            let mut doc = shared_doc(site, &format!("d{i}"), "2026-07-11T00:00:00Z");
2059            doc["textContent"] = json!(body);
2060            records.push((nsid::STANDARD_DOCUMENT, rkey, doc));
2061        }
2062        records.push((
2063            nsid::STANDARD_DOCUMENT,
2064            "3l2zzzzzzzzzz",
2065            shared_doc("quiet", "q1", "2026-01-01T00:00:00Z"),
2066        ));
2067        let (plc, _) = serve_repo(SHARED, records).await;
2068        let client = crate::feed::build_client().unwrap();
2069        let rkeys: Vec<String> = ["big1", "big2", "quiet"]
2070            .iter()
2071            .map(|s| s.to_string())
2072            .collect();
2073        let reads = fetch_repo_capped(&client, &plc, SHARED, &rkeys, 2_000, 4 * 1024 * 1024)
2074            .await
2075            .unwrap();
2076        let quiet = reads[2].as_ref().unwrap();
2077        assert_eq!(
2078            quiet.entries.len(),
2079            1,
2080            "the quiet publication was starved of bytes by its siblings"
2081        );
2082        assert!(
2083            !reads[0].as_ref().unwrap().complete,
2084            "a publication over its share was reported complete"
2085        );
2086    }
2087
2088    /// Second review of #228: a static per-publication byte share cut a busy
2089    /// publication beside idle siblings to a fraction of what it reads alone
2090    /// (376 of 1,200 at production scale). Grouping must never read a
2091    /// publication worse than reading it alone.
2092    #[tokio::test]
2093    async fn a_busy_publication_beside_idle_siblings_reads_as_it_would_alone() {
2094        let body = "w".repeat(20 * 1024);
2095        let mut records: Vec<(&'static str, &'static str, serde_json::Value)> = Vec::new();
2096        let rkeys: Vec<String> = (0..16).map(|i| format!("p{i:02}")).collect();
2097        for r in &rkeys {
2098            let r: &'static str = Box::leak(r.clone().into_boxed_str());
2099            records.push((
2100                nsid::STANDARD_PUBLICATION,
2101                r,
2102                json!({ "name": r, "url": "https://p.example" }),
2103            ));
2104        }
2105        for i in 0..100 {
2106            let rkey: &'static str = Box::leak(format!("3l2bus{i:06}").into_boxed_str());
2107            let mut doc = shared_doc("p00", &format!("d{i}"), "2026-07-11T00:00:00Z");
2108            doc["textContent"] = json!(body);
2109            records.push((nsid::STANDARD_DOCUMENT, rkey, doc));
2110        }
2111        let (plc, _) = serve_repo(SHARED, records).await;
2112        let client = crate::feed::build_client().unwrap();
2113        let budget = 4 * 1024 * 1024;
2114        let alone = fetch_repo_capped(&client, &plc, SHARED, &rkeys[..1], 2_000, budget)
2115            .await
2116            .unwrap();
2117        let grouped = fetch_repo_capped(&client, &plc, SHARED, &rkeys, 2_000, budget)
2118            .await
2119            .unwrap();
2120        let (a, g) = (alone[0].as_ref().unwrap(), grouped[0].as_ref().unwrap());
2121        assert_eq!((a.entries.len(), a.complete), (100, true));
2122        assert_eq!(
2123            (g.entries.len(), g.complete),
2124            (100, true),
2125            "grouping read it worse than alone"
2126        );
2127    }
2128
2129    /// Second review of #228: one document larger than a static share made its
2130    /// publication fail every poll ("stopped before its first document").
2131    #[tokio::test]
2132    async fn one_large_document_does_not_fail_its_publication_in_a_group() {
2133        let mut big = shared_doc("big", "huge", "2026-07-11T00:00:00Z");
2134        big["textContent"] = json!("w".repeat(500 * 1024));
2135        let records = vec![
2136            (
2137                nsid::STANDARD_PUBLICATION,
2138                "big",
2139                json!({ "name": "Big", "url": "https://b.example" }),
2140            ),
2141            (
2142                nsid::STANDARD_PUBLICATION,
2143                "other",
2144                json!({ "name": "Other", "url": "https://o.example" }),
2145            ),
2146            (nsid::STANDARD_DOCUMENT, "3l2hugeaaaa2a", big),
2147        ];
2148        let (plc, _) = serve_repo(SHARED, records).await;
2149        let client = crate::feed::build_client().unwrap();
2150        let rkeys = vec!["big".to_string(), "other".to_string()];
2151        let reads = fetch_repo_capped(&client, &plc, SHARED, &rkeys, 2_000, 1024 * 1024)
2152            .await
2153            .unwrap();
2154        assert_eq!(
2155            reads[0].as_ref().unwrap().entries.len(),
2156            1,
2157            "a large document was dropped"
2158        );
2159    }
2160
2161    /// Second review of #228: one result per REQUESTED rkey, even when none is
2162    /// in the repo and they repeat.
2163    #[tokio::test]
2164    async fn every_requested_rkey_gets_a_result() {
2165        let (plc, _) = serve_repo(
2166            SHARED,
2167            vec![(
2168                nsid::STANDARD_PUBLICATION,
2169                "alpha",
2170                json!({ "name": "Alpha", "url": "https://alpha.example" }),
2171            )],
2172        )
2173        .await;
2174        let client = crate::feed::build_client().unwrap();
2175        let reads = fetch_repo(
2176            &client,
2177            &plc,
2178            SHARED,
2179            &["nope".to_string(), "nope".to_string()],
2180        )
2181        .await
2182        .unwrap();
2183        assert_eq!(reads.len(), 2);
2184        assert!(reads.iter().all(|r| r.is_err()));
2185    }
2186
2187    /// Review of #228: a repeated rkey read as empty-and-complete, and an empty
2188    /// request still went to the network.
2189    #[tokio::test]
2190    async fn repeated_and_empty_requests_are_handled() {
2191        let (plc, hits) = serve_repo(
2192            SHARED,
2193            vec![
2194                (
2195                    nsid::STANDARD_PUBLICATION,
2196                    "alpha",
2197                    json!({ "name": "Alpha", "url": "https://alpha.example" }),
2198                ),
2199                (
2200                    nsid::STANDARD_DOCUMENT,
2201                    "3l2rpaaaaaa2a",
2202                    shared_doc("alpha", "a1", "2026-07-11T00:00:00Z"),
2203                ),
2204            ],
2205        )
2206        .await;
2207        let client = crate::feed::build_client().unwrap();
2208        let none = fetch_repo(&client, &plc, SHARED, &[]).await.unwrap();
2209        assert!(none.is_empty());
2210        assert_eq!(
2211            hits.load(std::sync::atomic::Ordering::SeqCst),
2212            0,
2213            "an empty request reached the network"
2214        );
2215        let twice = fetch_repo(
2216            &client,
2217            &plc,
2218            SHARED,
2219            &["alpha".to_string(), "alpha".to_string()],
2220        )
2221        .await
2222        .unwrap();
2223        for read in &twice {
2224            assert_eq!(
2225                read.as_ref().unwrap().entries.len(),
2226                1,
2227                "a repeated rkey read as empty"
2228            );
2229        }
2230    }
2231
2232    /// Review of #228: nothing checked that each publication's documents are
2233    /// stored under ITS feed — swapping them passed the whole suite.
2234    #[tokio::test]
2235    async fn a_group_stores_each_publications_documents_under_its_own_feed() {
2236        let (plc, _) = serve_repo(
2237            SHARED,
2238            vec![
2239                (
2240                    nsid::STANDARD_PUBLICATION,
2241                    "alpha",
2242                    json!({ "name": "Alpha", "url": "https://alpha.example" }),
2243                ),
2244                (
2245                    nsid::STANDARD_PUBLICATION,
2246                    "beta",
2247                    json!({ "name": "Beta", "url": "https://beta.example" }),
2248                ),
2249                (
2250                    nsid::STANDARD_DOCUMENT,
2251                    "3l2grpaaaaa2a",
2252                    shared_doc("alpha", "only-alpha", "2026-07-11T00:00:00Z"),
2253                ),
2254                (
2255                    nsid::STANDARD_DOCUMENT,
2256                    "3l2grpaaaaa2b",
2257                    shared_doc("beta", "only-beta", "2026-07-10T00:00:00Z"),
2258                ),
2259            ],
2260        )
2261        .await;
2262        let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
2263        let mut feeds = Vec::new();
2264        for rkey in ["alpha", "beta"] {
2265            let url = format!("at://{SHARED}/{}/{rkey}", nsid::STANDARD_PUBLICATION);
2266            crate::store::upsert_feed(
2267                &pool,
2268                &crate::store::NewFeed {
2269                    url: url.clone(),
2270                    ..Default::default()
2271                },
2272            )
2273            .await
2274            .unwrap();
2275            feeds.push(
2276                crate::store::get_feed_by_url(&pool, &url)
2277                    .await
2278                    .unwrap()
2279                    .unwrap(),
2280            );
2281        }
2282        let mut config = crate::config::Config::default();
2283        config.oauth.plc_directory = plc;
2284        let client = crate::feed::build_client().unwrap();
2285        let outcomes = crate::feed::poll_publication_group(&pool, &client, &config, &feeds).await;
2286        assert_eq!(outcomes.len(), 2);
2287        for (feed, want) in feeds.iter().zip(["only-alpha", "only-beta"]) {
2288            let titles: Vec<String> =
2289                sqlx::query_scalar("SELECT title FROM entries WHERE feed_id = ?")
2290                    .bind(feed.id)
2291                    .fetch_all(&pool)
2292                    .await
2293                    .unwrap();
2294            assert_eq!(
2295                titles,
2296                vec![want.to_string()],
2297                "{} got another publication's documents",
2298                feed.url
2299            );
2300        }
2301    }
2302
2303    /// **A0 — the 0.4.0 acceptance test.** A subscribed publication is due, is
2304    /// polled by the standard.site reader, and its documents land as entries.
2305    ///
2306    /// It starts from the stored row the subscribe form will create; step 3 of
2307    /// `design/STANDARD-SITE-0.4.0.md` extends it through the form itself.
2308    #[tokio::test]
2309    async fn a_publication_subscription_delivers_entries_end_to_end() {
2310        const A0: &str = "did:plc:acceptanceaaaaaaaaaaaaaa";
2311        let site = format!("at://{A0}/{}/a0pub", nsid::STANDARD_PUBLICATION);
2312        let document = |title: &str, path: &str| {
2313            json!({ "title": title, "publishedAt": "2026-07-11T00:00:00Z",
2314                    "path": path, "site": site, "textContent": "body" })
2315        };
2316        let (plc, _hits) = serve_repo(
2317            A0,
2318            vec![
2319                (
2320                    nsid::STANDARD_PUBLICATION,
2321                    "a0pub",
2322                    json!({ "name": "A0 Journal", "url": "https://a0.example" }),
2323                ),
2324                (
2325                    nsid::STANDARD_DOCUMENT,
2326                    "3l2a0aaaaaa2a",
2327                    document("First post", "/first"),
2328                ),
2329                (
2330                    nsid::STANDARD_DOCUMENT,
2331                    "3l2a0aaaaaa2b",
2332                    document("Second post", "/second"),
2333                ),
2334            ],
2335        )
2336        .await;
2337
2338        let pool = crate::store::init_url("sqlite::memory:").await.unwrap();
2339        crate::store::upsert_feed(
2340            &pool,
2341            &crate::store::NewFeed {
2342                url: site.clone(),
2343                ..Default::default()
2344            },
2345        )
2346        .await
2347        .unwrap();
2348        // The flag gates storing, not reading, so A0 runs with it off.
2349        let mut config = crate::config::Config::default();
2350        config.oauth.plc_directory = plc;
2351
2352        let now = chrono::Utc::now().to_rfc3339_opts(chrono::SecondsFormat::Secs, true);
2353        // The publication loop's own selection, so A0 tests what runs.
2354        let feed =
2355            crate::store::due_feeds_of_kind(&pool, &now, crate::feed::FeedKind::Publication, 50)
2356                .await
2357                .unwrap()
2358                .into_iter()
2359                .find(|f| f.url == site)
2360                .expect("the publication is not handed to the publication poller");
2361
2362        let client = crate::feed::build_client().unwrap();
2363        let outcome = crate::feed::poll_feed_by_kind(&pool, &client, &config, &feed)
2364            .await
2365            .unwrap();
2366        assert!(
2367            matches!(
2368                outcome,
2369                crate::feed::PollOutcome::Updated { new_entries: 2 }
2370            ),
2371            "expected two new entries, got {outcome:?}"
2372        );
2373        let titles: Vec<String> =
2374            sqlx::query_scalar("SELECT title FROM entries WHERE feed_id = ? ORDER BY title")
2375                .bind(feed.id)
2376                .fetch_all(&pool)
2377                .await
2378                .unwrap();
2379        assert_eq!(titles, vec!["First post", "Second post"]);
2380    }
2381
2382    fn canonical(rkey: &str) -> String {
2383        format!("at://{DID}/{}/{rkey}", nsid::STANDARD_PUBLICATION)
2384    }
2385
2386    /// **The `site` filter keys on the URI the PDS minted, not the string the
2387    /// reader subscribed with.** Documents reference their publication by the
2388    /// canonical URI, and every measured document does. An earlier draft
2389    /// compared against the subscribed string; storage was then meant to admit
2390    /// handle-form URIs, and a handle-form subscription found nothing forever
2391    /// while the module declared the feed healthy. #164 made storage DID-only,
2392    /// so the two strings agree today — the canonical one is still the right
2393    /// key, and this pins it.
2394    #[test]
2395    fn the_site_filter_uses_the_uri_the_pds_minted() {
2396        let records = vec![publication("p", "https://scanash.com")];
2397        let (site, pubn) = publication_from_records("p", &records).expect("publication not found");
2398        assert_eq!(site, canonical("p"), "did not take the PDS's canonical URI");
2399
2400        let docs = vec![document("d1", &canonical("p"), "Hello", "/hello")];
2401        let entries = entries_from_records(&site, &pubn, &docs);
2402        assert_eq!(
2403            entries.len(),
2404            1,
2405            "a canonical-site document was not matched"
2406        );
2407    }
2408
2409    /// The mapping into the store's row is total and loses nothing the poller
2410    /// would need — so wiring the reader has nothing to invent.
2411    #[test]
2412    fn an_entry_maps_onto_the_stores_row() {
2413        let records = vec![publication("p", "https://example.com")];
2414        let (site, pubn) = publication_from_records("p", &records).unwrap();
2415        let docs = vec![document("rk1", &site, "Hello", "/hello")];
2416        let row: crate::store::NewEntry = entries_from_records(&site, &pubn, &docs)
2417            .pop()
2418            .unwrap()
2419            .into();
2420        assert_eq!(
2421            row.guid,
2422            format!("at://{DID}/{}/rk1", nsid::STANDARD_DOCUMENT)
2423        );
2424        assert_eq!(row.url.as_deref(), Some("https://example.com/hello"));
2425        assert_eq!(row.title.as_deref(), Some("Hello"));
2426        assert_eq!(row.published.as_deref(), Some("2026-07-11T00:00:00Z"));
2427        assert_eq!(row.content_html.as_deref(), Some("body"));
2428        assert_eq!(row.author, None);
2429        assert_eq!(row.fetched_at, None);
2430    }
2431
2432    /// A repo can hold several publications — measured, some do — and
2433    /// `listRecords` cannot filter server-side.
2434    #[test]
2435    fn documents_are_filtered_by_their_site_field() {
2436        let records = vec![publication("mine", "https://example.com")];
2437        let (site, pubn) = publication_from_records("mine", &records).unwrap();
2438        let docs = vec![
2439            document("a", &site, "Mine", "/a"),
2440            document("b", &canonical("theirs"), "Theirs", "/b"),
2441            document("c", &site, "Mine again", "/c"),
2442        ];
2443        let titles: Vec<String> = entries_from_records(&site, &pubn, &docs)
2444            .into_iter()
2445            .map(|e| e.title)
2446            .collect();
2447        assert_eq!(titles, ["Mine", "Mine again"]);
2448    }
2449
2450    /// **The publication URL is a stranger's string and is vetted as one.**
2451    ///
2452    /// It becomes the base of every `Entry.url`, which is an href. This repo has
2453    /// `safe_link.rs` as a type with a module-private field precisely because
2454    /// the procedural version of this guarantee failed; #143 and #138 are this
2455    /// same bug class on `siteUrl`.
2456    #[test]
2457    fn a_publication_with_a_hostile_url_is_refused() {
2458        for hostile in [
2459            "javascript:alert(1)",
2460            "data:text/html,<script>",
2461            "file:///etc/passwd",
2462            "",
2463        ] {
2464            let records = vec![publication("p", hostile)];
2465            assert!(
2466                publication_from_records("p", &records).is_none(),
2467                "accepted a publication whose url is {hostile:?}",
2468            );
2469        }
2470    }
2471
2472    /// **Entry URLs are joined, not concatenated.**
2473    ///
2474    /// String concatenation produced `https://x.com/https://evil.com/a` for an
2475    /// absolute `path`, and put the path inside the query string for a base
2476    /// carrying one.
2477    #[test]
2478    fn entry_urls_are_joined_against_the_publication_base() {
2479        let records = vec![publication("p", "https://example.com/blog")];
2480        let (site, pubn) = publication_from_records("p", &records).unwrap();
2481        let docs = vec![
2482            document("a", &site, "Relative", "/a"),
2483            document("b", &site, "Absolute-looking", "https://evil.example/x"),
2484        ];
2485        let urls: Vec<Option<String>> = entries_from_records(&site, &pubn, &docs)
2486            .into_iter()
2487            .map(|e| e.url)
2488            .collect();
2489        assert_eq!(urls[0].as_deref(), Some("https://example.com/a"));
2490        // Exactly the publication's base, not merely "not evil": a mutation
2491        // that returned the raw path, or an empty string, passed the weaker
2492        // negative assertion this used to be.
2493        // Off-origin is dropped, not rewritten to the base. This assertion has
2494        // moved twice: it began as "not evil" (a mutant returning the raw path
2495        // passed it), was tightened to the base fallback, and is now `None` —
2496        // the base gave every affected entry the same homepage href.
2497        assert_eq!(
2498            urls[1], None,
2499            "a document path that escapes its publication's origin must yield no URL"
2500        );
2501    }
2502
2503    // ---- #205: a publisher's strings are bounded before they are stored ----
2504
2505    #[test]
2506    fn a_documents_text_fields_are_bounded_before_they_are_stored() {
2507        let site = canonical("pub");
2508        let big = "x".repeat(8 * 1024 * 1024);
2509        let records = vec![
2510            publication("pub", "https://scanash.com"),
2511            rec(
2512                nsid::STANDARD_DOCUMENT,
2513                "3l2bigaaaaa2a",
2514                json!({ "title": big, "publishedAt": "2026-07-11T00:00:00Z",
2515                        "path": format!("/{}", "p".repeat(20_000)), "site": site,
2516                        "textContent": "<".repeat(3 * 1024 * 1024) }),
2517            ),
2518        ];
2519        let (_, publication) = publication_from_records("pub", &records).unwrap();
2520        let entries = entries_from_records(&site, &publication, &records);
2521        let e = &entries[0];
2522        assert!(
2523            e.title.len() <= crate::feed::MAX_TITLE_BYTES,
2524            "title: {}",
2525            e.title.len()
2526        );
2527        let url = e
2528            .url
2529            .as_ref()
2530            .expect("an overlong path is truncated, not dropped");
2531        assert!(
2532            url.len() <= crate::feed::MAX_URL_BYTES,
2533            "url: {}",
2534            url.len()
2535        );
2536        let summary = e.summary.as_ref().unwrap();
2537        assert!(
2538            summary.len() <= crate::feed::MAX_CONTENT_HTML_BYTES,
2539            "the ESCAPED summary is what is stored: {}",
2540            summary.len()
2541        );
2542    }
2543
2544    /// Review of #224: the document's URI is the entry id, chosen by the
2545    /// publisher's PDS, and it went into the same UNIQUE-indexed column the
2546    /// RSS path bounds.
2547    #[test]
2548    fn a_documents_uri_is_bounded_as_an_entry_id() {
2549        let site = canonical("pub");
2550        let records = vec![
2551            publication("pub", "https://scanash.com"),
2552            rec(
2553                nsid::STANDARD_DOCUMENT,
2554                &"k".repeat(100_000),
2555                json!({ "title": "t", "publishedAt": "2026-07-11T00:00:00Z",
2556                        "path": "/p", "site": site }),
2557            ),
2558        ];
2559        let (_, publication) = publication_from_records("pub", &records).unwrap();
2560        let entries = entries_from_records(&site, &publication, &records);
2561        let stored: crate::store::NewEntry = entries[0].clone().into();
2562        assert!(
2563            stored.guid.len() <= crate::feed::MAX_GUID_BYTES,
2564            "guid: {}",
2565            stored.guid.len()
2566        );
2567    }
2568
2569    #[test]
2570    fn a_publications_own_name_is_bounded() {
2571        let records = vec![rec(
2572            nsid::STANDARD_PUBLICATION,
2573            "pub",
2574            json!({ "name": "n".repeat(100_000), "url": "https://scanash.com" }),
2575        )];
2576        let (_, publication) = publication_from_records("pub", &records).unwrap();
2577        assert!(publication.name.unwrap().len() <= crate::feed::MAX_TITLE_BYTES);
2578    }
2579
2580    /// 8% of measured documents (37 of 449) carry neither summary field.
2581    #[test]
2582    fn a_document_with_neither_summary_field_still_yields_an_entry() {
2583        let records = vec![publication("p", "https://example.com/")];
2584        let (site, pubn) = publication_from_records("p", &records).unwrap();
2585        let bare = rec(
2586            nsid::STANDARD_DOCUMENT,
2587            "bare",
2588            json!({
2589                "title": "Bare",
2590                "publishedAt": "2026-07-11T00:00:00Z",
2591                "path": "/bare",
2592                "site": site,
2593            }),
2594        );
2595        let entries = entries_from_records(&site, &pubn, &[bare]);
2596        assert_eq!(entries.len(), 1);
2597        assert_eq!(entries[0].summary, None);
2598        assert_eq!(entries[0].url.as_deref(), Some("https://example.com/bare"));
2599    }
2600
2601    /// `description` is the authored summary; `textContent` is the whole body.
2602    /// An EMPTY description must not shadow a real one.
2603    #[test]
2604    fn an_empty_description_does_not_shadow_the_body() {
2605        let records = vec![publication("p", "https://example.com")];
2606        let (site, pubn) = publication_from_records("p", &records).unwrap();
2607        let doc = rec(
2608            nsid::STANDARD_DOCUMENT,
2609            "d",
2610            json!({
2611                "title": "T",
2612                "publishedAt": "2026-07-11T00:00:00Z",
2613                "path": "/d",
2614                "site": site,
2615                "description": "   ",
2616                "textContent": "the real body",
2617            }),
2618        );
2619        let entries = entries_from_records(&site, &pubn, &[doc]);
2620        assert_eq!(entries[0].summary.as_deref(), Some("the real body"));
2621    }
2622
2623    /// **`publishedAt` is parsed, not passed through.** The store's `published`
2624    /// column is the RFC3339 shape `feed::fmt_time` writes, and the reading
2625    /// order sorts on it as a string. A publisher's string went in verbatim —
2626    /// a garbage value would have sorted arbitrarily among real ones, and a
2627    /// valid-but-differently-spelled one (`+00:00`, fractional seconds) would
2628    /// not have matched the RSS path's spelling for the same instant.
2629    #[test]
2630    fn published_at_is_normalised_or_dropped() {
2631        let records = vec![publication("p", "https://example.com")];
2632        let (site, pubn) = publication_from_records("p", &records).unwrap();
2633        let with = |rkey: &str, published_at: serde_json::Value| {
2634            rec(
2635                nsid::STANDARD_DOCUMENT,
2636                rkey,
2637                json!({ "title": "T", "publishedAt": published_at, "path": "/x", "site": site }),
2638            )
2639        };
2640        // Single-character rkeys on purpose: they are not TIDs, so the
2641        // rkey-derived fallback below does not apply and "unparseable" really
2642        // does mean undated here.
2643        let docs = vec![
2644            with("a", json!("2026-07-11T09:30:00.123+02:00")),
2645            with("b", json!("yesterday-ish")),
2646            with("c", json!("2026-07-11T00:00:00Z")),
2647        ];
2648        let published: Vec<Option<String>> = entries_from_records(&site, &pubn, &docs)
2649            .into_iter()
2650            .map(|e| e.published)
2651            .collect();
2652        assert_eq!(
2653            published,
2654            vec![
2655                Some("2026-07-11T07:30:00Z".to_string()),
2656                None,
2657                Some("2026-07-11T00:00:00Z".to_string()),
2658            ],
2659            "publishedAt was not normalised to the store's spelling"
2660        );
2661    }
2662
2663    /// A document with no usable `publishedAt` is dated from its rkey.
2664    ///
2665    /// **An undated entry is not merely untidy, it is immortal-and-mortal at
2666    /// once.** The store's retention sweep and its per-feed cap both order on
2667    /// `COALESCE(published, fetched_at)`, and an entry inserted with no
2668    /// `published` gets `fetched_at` stamped at insertion. So the sweep deletes
2669    /// it once it is `retention_days` old, the next poll re-inserts it with a
2670    /// fresh `fetched_at` and a new `entries.id`, its read state is gone with
2671    /// the cascade, and it arrives unread — again, on the same cycle, forever.
2672    /// The same reset also sorts it newest in the per-feed cap, where it evicts
2673    /// entries that really are newer.
2674    #[test]
2675    fn an_undated_document_is_dated_from_its_tid_rkey() {
2676        let records = vec![publication("p", "https://example.com")];
2677        let (site, pubn) = publication_from_records("p", &records).unwrap();
2678        let docs = vec![rec(
2679            nsid::STANDARD_DOCUMENT,
2680            PAST_TID,
2681            json!({ "title": "T", "path": "/x", "site": site }),
2682        )];
2683        assert_eq!(
2684            entries_from_records(&site, &pubn, &docs)
2685                .into_iter()
2686                .next()
2687                .expect("the document is an entry")
2688                .published,
2689            Some(PAST_TID_WRITTEN_AT.to_string()),
2690            "the date must come from the record key, and be spelled the way the store spells dates"
2691        );
2692    }
2693
2694    #[test]
2695    fn an_unparseable_published_at_falls_back_to_the_tid_rkey() {
2696        let records = vec![publication("p", "https://example.com")];
2697        let (site, pubn) = publication_from_records("p", &records).unwrap();
2698        let docs = vec![rec(
2699            nsid::STANDARD_DOCUMENT,
2700            PAST_TID,
2701            json!({ "title": "T", "publishedAt": "yesterday-ish", "path": "/x", "site": site }),
2702        )];
2703        assert_eq!(
2704            entries_from_records(&site, &pubn, &docs)
2705                .into_iter()
2706                .next()
2707                .expect("the document is an entry")
2708                .published,
2709            Some(PAST_TID_WRITTEN_AT.to_string()),
2710            "a date the parser cannot read is no date at all, so the rkey must stand in"
2711        );
2712    }
2713
2714    #[test]
2715    fn a_stated_date_outranks_the_rkey() {
2716        let records = vec![publication("p", "https://example.com")];
2717        let (site, pubn) = publication_from_records("p", &records).unwrap();
2718        // Both candidate dates are historical, so nothing here depends on
2719        // what the machine's clock reads.
2720        let docs = vec![rec(
2721            nsid::STANDARD_DOCUMENT,
2722            PAST_TID,
2723            json!({ "title": "T", "publishedAt": "2020-01-02T00:00:00Z", "path": "/x", "site": site }),
2724        )];
2725        assert_eq!(
2726            entries_from_records(&site, &pubn, &docs)
2727                .into_iter()
2728                .next()
2729                .expect("the document is an entry")
2730                .published,
2731            Some("2020-01-02T00:00:00Z".to_string()),
2732            "the rkey records when the file was written, which is not when the post was published"
2733        );
2734    }
2735
2736    /// **A stated date in the future is discarded, not clamped.**
2737    ///
2738    /// Clamping it to "now" looks safe and is not. The store refreshes
2739    /// `published` on every poll, so the row would be re-dated to the current
2740    /// hour forever: never older than the retention cutoff, never outranked in
2741    /// the per-feed cap, permanently first in the reading list. The date must
2742    /// not depend on when the mapping ran, which is why both of these name the
2743    /// exact value they expect rather than comparing against the clock.
2744    #[test]
2745    fn a_future_dated_document_falls_back_to_its_rkey() {
2746        let records = vec![publication("p", "https://example.com")];
2747        let (site, pubn) = publication_from_records("p", &records).unwrap();
2748        let docs = vec![rec(
2749            nsid::STANDARD_DOCUMENT,
2750            PAST_TID,
2751            json!({ "title": "T", "publishedAt": "2999-01-01T00:00:00Z", "path": "/x", "site": site }),
2752        )];
2753        assert_eq!(
2754            entries_from_records(&site, &pubn, &docs)
2755                .into_iter()
2756                .next()
2757                .expect("the document is an entry")
2758                .published,
2759            Some(PAST_TID_WRITTEN_AT.to_string()),
2760            "the date must be the record's write time, not the hour the poll happened to run"
2761        );
2762    }
2763
2764    #[test]
2765    fn a_future_dated_document_without_a_tid_rkey_is_undated() {
2766        let records = vec![publication("p", "https://example.com")];
2767        let (site, pubn) = publication_from_records("p", &records).unwrap();
2768        let docs = vec![rec(
2769            nsid::STANDARD_DOCUMENT,
2770            "self",
2771            json!({ "title": "T", "publishedAt": "2999-01-01T00:00:00Z", "path": "/x", "site": site }),
2772        )];
2773        assert_eq!(
2774            entries_from_records(&site, &pubn, &docs)
2775                .into_iter()
2776                .next()
2777                .expect("the document is an entry")
2778                .published,
2779            None,
2780            "with nothing credible to date it by, the row falls to fetched_at, which holds still"
2781        );
2782    }
2783
2784    /// **A little ahead of our clock is skew, not a lie.**
2785    ///
2786    /// The rkey here is deliberately not a TID, so nothing masks a wrongly
2787    /// discarded date: if the stated one is thrown away the entry is undated,
2788    /// and an undated entry sorts to the bottom of a list ordered on a bare
2789    /// `published DESC`. A publisher a few seconds fast would have had their
2790    /// newest post buried.
2791    #[test]
2792    fn a_stated_date_a_little_ahead_of_our_clock_is_still_believed() {
2793        let records = vec![publication("p", "https://example.com")];
2794        let (site, pubn) = publication_from_records("p", &records).unwrap();
2795        let slightly_ahead =
2796            crate::feed::fmt_time(chrono::Utc::now() + chrono::Duration::seconds(10));
2797        let docs = vec![rec(
2798            nsid::STANDARD_DOCUMENT,
2799            "self",
2800            json!({ "title": "T", "publishedAt": slightly_ahead, "path": "/x", "site": site }),
2801        )];
2802        assert_eq!(
2803            entries_from_records(&site, &pubn, &docs)
2804                .into_iter()
2805                .next()
2806                .expect("the document is an entry")
2807                .published,
2808            Some(slightly_ahead),
2809            "a few seconds of clock skew must not cost the entry its date"
2810        );
2811    }
2812
2813    #[test]
2814    fn a_document_with_neither_a_date_nor_a_tid_rkey_stays_undated() {
2815        let records = vec![publication("p", "https://example.com")];
2816        let (site, pubn) = publication_from_records("p", &records).unwrap();
2817        // Two shapes, rejected by two different checks. "my-first-post" is 13
2818        // characters but carries a `-`, so it never reaches the alphabet's
2819        // arithmetic at all. "abcdefghijklm" is 13 valid s32 characters and
2820        // decodes perfectly well — to the year 2192 — which is the case the
2821        // bound in `tid_timestamp` exists for, and the one a publisher naming
2822        // files by slug actually produces.
2823        let docs = vec![
2824            rec(
2825                nsid::STANDARD_DOCUMENT,
2826                "my-first-post",
2827                json!({ "title": "T", "path": "/x", "site": site }),
2828            ),
2829            rec(
2830                nsid::STANDARD_DOCUMENT,
2831                "abcdefghijklm",
2832                json!({ "title": "T", "path": "/y", "site": site }),
2833            ),
2834        ];
2835        assert_eq!(
2836            entries_from_records(&site, &pubn, &docs)
2837                .into_iter()
2838                .map(|e| e.published)
2839                .collect::<Vec<_>>(),
2840            vec![None, None],
2841            "an invented date is worse than no date; the store decides what to do with undated rows"
2842        );
2843    }
2844
2845    /// **Summaries are plain text and are escaped, not sanitised.**
2846    ///
2847    /// The lexicon defines `textContent` and `description` as plain text, and
2848    /// the store's `content_html` is rendered as HTML, so the text must be
2849    /// escaped on the way in. The first version of this ran `ammonia::clean`
2850    /// over them — the RSS body function — which parses its input as markup
2851    /// and deletes everything after a bare `<`. 321 of 449 measured documents
2852    /// use `textContent` as their summary; any post mentioning `Vec<T>` lost
2853    /// the rest of its summary, silently.
2854    #[test]
2855    fn summaries_are_escaped_as_plain_text_not_sanitised_as_markup() {
2856        let records = vec![publication("p", "https://example.com")];
2857        let (site, pubn) = publication_from_records("p", &records).unwrap();
2858        let doc = |rkey: &str, body: &str| {
2859            rec(
2860                nsid::STANDARD_DOCUMENT,
2861                rkey,
2862                json!({
2863                    "title": "T",
2864                    "publishedAt": "2026-07-11T00:00:00Z",
2865                    "path": "/d",
2866                    "site": site,
2867                    "textContent": body,
2868                }),
2869            )
2870        };
2871        let summaries: Vec<String> = entries_from_records(
2872            &site,
2873            &pubn,
2874            &[
2875                doc("a", "Vec<String> is a type"),
2876                doc("b", "<script>alert(1)</script>"),
2877            ],
2878        )
2879        .into_iter()
2880        .filter_map(|e| e.summary)
2881        .collect();
2882        assert_eq!(
2883            summaries[0], "Vec&lt;String&gt; is a type",
2884            "prose was eaten by an HTML parser"
2885        );
2886        assert!(
2887            !summaries[1].contains("<script"),
2888            "escaping failed: {}",
2889            summaries[1]
2890        );
2891    }
2892
2893    /// **`Entry.url` is a vetted href or nothing.** The scheme guarantee lived
2894    /// only inside `publication_from_records`; `entries_from_records` and
2895    /// `Publication` are both `pub`, so any other constructor — step 3 building
2896    /// one from the stored `feeds` row, say — gave an unparseable base, and the
2897    /// no-base branch then emitted the document's `path` verbatim. A
2898    /// `javascript:` path became the entry link. This is the class `safe_link`
2899    /// exists for.
2900    #[test]
2901    fn an_entry_url_is_never_an_unvetted_path() {
2902        let pubn = Publication {
2903            name: None,
2904            // What a caller that did not go through `publication_from_records`
2905            // can hand this function.
2906            url: "not a url".to_string(),
2907        };
2908        let site = canonical("p");
2909        let docs = vec![
2910            document("a", &site, "Hostile", "javascript:alert(1)"),
2911            document("b", &site, "Fine", "https://example.com/ok"),
2912        ];
2913        let entries = entries_from_records(&site, &pubn, &docs);
2914        assert_eq!(
2915            entries[0].url, None,
2916            "an unvetted path became an entry link"
2917        );
2918        // Contract change: with no parseable base there is no origin to check,
2919        // so a well-formed absolute URL is refused too. `safe_link` alone vets
2920        // the SCHEME; it would have published a publisher-controlled host under
2921        // this publication's name. See `no_parseable_base_means_no_url_not_any_url`.
2922        assert_eq!(
2923            entries[1].url, None,
2924            "an off-origin absolute URL was published under the publication's name"
2925        );
2926    }
2927
2928    /// A document with no `publishedAt` is an entry with no date — the same
2929    /// answer a garbage one gets. The field the module is willing to discard
2930    /// must not be the one whose absence is fatal.
2931    #[test]
2932    fn a_document_without_published_at_is_still_an_entry() {
2933        let records = vec![publication("p", "https://example.com")];
2934        let (site, pubn) = publication_from_records("p", &records).unwrap();
2935        let doc = rec(
2936            nsid::STANDARD_DOCUMENT,
2937            "d",
2938            json!({ "title": "T", "path": "/d", "site": site }),
2939        );
2940        let entries = entries_from_records(&site, &pubn, &[doc]);
2941        assert_eq!(
2942            entries.len(),
2943            1,
2944            "a missing publishedAt dropped the document"
2945        );
2946        assert_eq!(entries[0].published, None);
2947    }
2948
2949    /// **`fetch` refuses a URI naming another collection, before the network.**
2950    /// It lists publications and matches on rkey alone, so without this an
2951    /// `app.bsky.feed.post` URI would be "read as a publication" whenever a
2952    /// publication in that repo shares the rkey. Storage enforces the
2953    /// collection today, but this function is `pub`.
2954    #[tokio::test]
2955    async fn fetch_refuses_a_uri_for_another_collection() {
2956        let uri = AtUri::parse(&format!("at://{DID}/app.bsky.feed.post/3lab")).unwrap();
2957        let err = fetch(&reqwest::Client::new(), "https://plc.example", &uri)
2958            .await
2959            .expect_err("read a feed post as a publication");
2960        assert!(
2961            format!("{err:#}").contains(nsid::STANDARD_PUBLICATION),
2962            "failed for the wrong reason: {err:#}"
2963        );
2964    }
2965
2966    /// **A publication on a subpath keeps it.** `Url::join` is RFC-3986, so a
2967    /// relative `posts/a` against `https://example.com/blog` resolves to
2968    /// `/posts/a` — dropping the subpath every permalink needs, while still
2969    /// passing the origin check. The base is normalised to a directory.
2970    #[test]
2971    fn a_subpath_publication_keeps_its_base_path() {
2972        let records = vec![publication("p", "https://example.com/blog")];
2973        let (site, pubn) = publication_from_records("p", &records).unwrap();
2974        let docs = vec![document("a", &site, "Relative", "posts/a")];
2975        let urls: Vec<Option<String>> = entries_from_records(&site, &pubn, &docs)
2976            .into_iter()
2977            .map(|e| e.url)
2978            .collect();
2979        assert_eq!(urls[0].as_deref(), Some("https://example.com/blog/posts/a"));
2980    }
2981
2982    /// With no parseable base there is no origin to check, so there is no URL
2983    /// — `safe_link` alone vets the scheme and would pass any absolute URL a
2984    /// publisher chose, under this publication's name.
2985    #[test]
2986    fn no_parseable_base_means_no_url_not_any_url() {
2987        let pubn = Publication {
2988            name: None,
2989            url: "not a url".to_string(),
2990        };
2991        let site = canonical("p");
2992        let docs = vec![document("a", &site, "Absolute", "https://evil.example/x")];
2993        let entries = entries_from_records(&site, &pubn, &docs);
2994        assert_eq!(
2995            entries[0].url, None,
2996            "an off-origin absolute URL was published"
2997        );
2998    }
2999
3000    /// **An off-origin path is dropped, not rewritten to the homepage.**
3001    /// Returning the base gave every affected entry the SAME href pointing at
3002    /// the site root — realistic whenever a publication's `url` is the apex
3003    /// and its documents sit on `www.` or a CDN domain. `None` is the honest
3004    /// answer, and the template already has a no-URL branch.
3005    #[test]
3006    fn an_off_origin_path_yields_no_url_rather_than_the_homepage() {
3007        let records = vec![publication("p", "https://example.com/blog")];
3008        let (site, pubn) = publication_from_records("p", &records).unwrap();
3009        let docs = vec![
3010            document("a", &site, "Elsewhere", "https://www.example.com/post"),
3011            document("b", &site, "Home", "/ok"),
3012        ];
3013        let urls: Vec<Option<String>> = entries_from_records(&site, &pubn, &docs)
3014            .into_iter()
3015            .map(|e| e.url)
3016            .collect();
3017        assert_eq!(
3018            urls[0], None,
3019            "an off-origin path was rewritten to the base"
3020        );
3021        assert_eq!(urls[1].as_deref(), Some("https://example.com/ok"));
3022    }
3023
3024    /// **A blank `path` is no URL, not the homepage** — the same answer the
3025    /// off-origin branch now gives, and for the same reason: several
3026    /// documents with an empty `path` otherwise became several entries all
3027    /// linking to the site root.
3028    #[test]
3029    fn a_blank_path_yields_no_url() {
3030        let records = vec![publication("p", "https://example.com/blog")];
3031        let (site, pubn) = publication_from_records("p", &records).unwrap();
3032        let docs = vec![
3033            document("a", &site, "Blank", ""),
3034            document("b", &site, "Spaces", "   "),
3035            document("c", &site, "Real", "/real"),
3036        ];
3037        let urls: Vec<Option<String>> = entries_from_records(&site, &pubn, &docs)
3038            .into_iter()
3039            .map(|e| e.url)
3040            .collect();
3041        assert_eq!(urls[0], None, "a blank path became the homepage");
3042        assert_eq!(urls[1], None, "a whitespace path became the homepage");
3043        assert_eq!(urls[2].as_deref(), Some("https://example.com/real"));
3044    }
3045
3046    /// A document without `path` keeps its title, date and summary — the
3047    /// policy `publishedAt` and `Entry.url` already follow.
3048    #[test]
3049    fn a_document_without_a_path_is_still_an_entry() {
3050        let records = vec![publication("p", "https://example.com")];
3051        let (site, pubn) = publication_from_records("p", &records).unwrap();
3052        let doc = rec(
3053            nsid::STANDARD_DOCUMENT,
3054            "d",
3055            json!({ "title": "T", "publishedAt": "2026-07-11T00:00:00Z", "site": site }),
3056        );
3057        let entries = entries_from_records(&site, &pubn, &[doc]);
3058        assert_eq!(
3059            entries.len(),
3060            1,
3061            "a missing path dropped the whole document"
3062        );
3063        assert_eq!(entries[0].title, "T");
3064        assert_eq!(entries[0].url, None);
3065    }
3066
3067    /// **A sibling is not an orphan.** A repo with an empty publication A and
3068    /// a busy publication B is the exact shape the `site` filter exists for;
3069    /// treating B's documents as a signal would warn on every poll of A
3070    /// forever. The signal is a document referencing a publication this repo
3071    /// does not have — a spelling nothing can ever match.
3072    #[test]
3073    fn documents_are_classified_keep_sibling_or_orphan() {
3074        let pubs = [
3075            publication("a", "https://example.com"),
3076            publication("b", "https://b.example"),
3077        ];
3078        let known: std::collections::HashSet<&str> = pubs.iter().map(|p| p.uri.as_str()).collect();
3079        let mine = canonical("a");
3080        let wanted: std::collections::HashMap<String, usize> = [(mine.clone(), 0)].into();
3081        let fate = |d: &crate::atproto::RecordEntry| classify_document(d, &wanted, &known);
3082
3083        assert_eq!(
3084            fate(&document("d1", &mine, "Mine", "/1")),
3085            DocumentFate::Keep(0)
3086        );
3087        assert_eq!(
3088            fate(&document("d2", &canonical("b"), "B's", "/2")),
3089            DocumentFate::Sibling,
3090            "a sibling publication's document is not an orphan"
3091        );
3092        assert_eq!(
3093            fate(&document(
3094                "d3",
3095                "at://did:plc:other/site.standard.publication/x",
3096                "?",
3097                "/3"
3098            )),
3099            DocumentFate::Orphan
3100        );
3101        assert_eq!(
3102            fate(&rec(
3103                nsid::STANDARD_DOCUMENT,
3104                "d4",
3105                json!({"title": "no rest"})
3106            )),
3107            DocumentFate::Malformed
3108        );
3109    }
3110
3111    /// `AtUri` uses the crate's one spelling of the prefix.
3112    #[test]
3113    fn at_uri_parsing_uses_the_shared_prefix() {
3114        let uri = format!(
3115            "{}{DID}/{}/abc",
3116            crate::atproto::AT_URI_PREFIX,
3117            nsid::STANDARD_PUBLICATION
3118        );
3119        assert!(AtUri::parse(&uri).is_some());
3120    }
3121
3122    /// The guid is the record's own URI. `path` is mutable; dedup is
3123    /// `UNIQUE (feed_id, guid)`.
3124    #[test]
3125    fn the_guid_is_the_record_uri_not_the_path() {
3126        let records = vec![publication("p", "https://example.com")];
3127        let (site, pubn) = publication_from_records("p", &records).unwrap();
3128        let docs = vec![document("rk1", &site, "T", "/moved")];
3129        let entries = entries_from_records(&site, &pubn, &docs);
3130        assert_eq!(
3131            entries[0].guid,
3132            format!("at://{DID}/{}/rk1", nsid::STANDARD_DOCUMENT)
3133        );
3134    }
3135
3136    /// One unreadable record must not cost a publisher the whole feed.
3137    #[test]
3138    fn a_malformed_document_is_skipped_rather_than_fatal() {
3139        let records = vec![publication("p", "https://example.com")];
3140        let (site, pubn) = publication_from_records("p", &records).unwrap();
3141        let docs = vec![
3142            rec(
3143                nsid::STANDARD_DOCUMENT,
3144                "bad",
3145                json!({ "title": "no rest" }),
3146            ),
3147            document("ok", &site, "Good", "/good"),
3148        ];
3149        let entries = entries_from_records(&site, &pubn, &docs);
3150        assert_eq!(entries.len(), 1);
3151        assert_eq!(entries[0].title, "Good");
3152    }
3153
3154    /// An rkey that is not in the repo is absent, not an error.
3155    #[test]
3156    fn a_missing_publication_is_none() {
3157        let records = vec![publication("other", "https://example.com")];
3158        assert!(publication_from_records("p", &records).is_none());
3159    }
3160
3161    #[test]
3162    fn at_uris_parse_in_both_forms_and_reject_malformed_ones() {
3163        let did = AtUri::parse(&format!("at://{DID}/site.standard.publication/abc")).unwrap();
3164        assert_eq!(did.authority, DID);
3165        assert_eq!(did.rkey, "abc");
3166        assert_eq!(
3167            did.to_string(),
3168            format!("at://{DID}/site.standard.publication/abc")
3169        );
3170        assert!(AtUri::parse("at://alice.example.com/site.standard.publication/abc").is_some());
3171        for bad in [
3172            "at://",
3173            "at://only-authority",
3174            "at://authority/collection",
3175            "at://authority/collection/",
3176            "at:///collection/rkey",
3177            "at://authority/collection/rkey/extra",
3178            "https://example.com/feed.xml",
3179            "at:authority/collection/rkey",
3180        ] {
3181            assert!(AtUri::parse(bad).is_none(), "parsed {bad:?}");
3182        }
3183    }
3184}