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