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