use serde::Deserialize;
use crate::lexicon::nsid;
const DOCUMENT_PAGE_SIZE: u32 = 25;
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct AtUri {
pub authority: String,
pub collection: String,
pub rkey: String,
}
impl AtUri {
pub fn parse(uri: &str) -> Option<Self> {
let rest = uri.strip_prefix(crate::atproto::AT_URI_PREFIX)?;
let mut parts = rest.split('/');
let (authority, collection, rkey) = (parts.next()?, parts.next()?, parts.next()?);
if parts.next().is_some()
|| authority.is_empty()
|| collection.is_empty()
|| rkey.is_empty()
{
return None;
}
Some(Self {
authority: authority.to_string(),
collection: collection.to_string(),
rkey: rkey.to_string(),
})
}
}
impl std::fmt::Display for AtUri {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(
f,
"at://{}/{}/{}",
self.authority, self.collection, self.rkey
)
}
}
#[derive(Debug)]
pub struct PublicationRead {
pub publication: Publication,
pub entries: Vec<Entry>,
pub complete: bool,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct Publication {
pub name: Option<String>,
pub url: String,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct Entry {
pub guid: String,
pub title: String,
pub published: Option<String>,
pub url: Option<String>,
pub summary: Option<String>,
}
impl From<Entry> for crate::store::NewEntry {
fn from(e: Entry) -> Self {
crate::store::NewEntry {
guid: e.guid,
url: e.url,
title: Some(e.title),
author: None,
published: e.published,
content_html: e.summary,
fetched_at: None,
}
}
}
#[derive(Debug, Deserialize)]
struct PublicationValue {
name: Option<String>,
url: String,
}
#[derive(Debug, Deserialize)]
struct DocumentValue {
title: String,
#[serde(rename = "publishedAt")]
published_at: Option<String>,
path: Option<String>,
site: String,
#[serde(rename = "textContent")]
text_content: Option<String>,
description: Option<String>,
}
pub fn publication_from_records(
rkey: &str,
records: &[crate::atproto::RecordEntry],
) -> Option<(String, Publication)> {
let entry = records
.iter()
.find(|r| AtUri::parse(&r.uri).is_some_and(|u| u.rkey == rkey))?;
let value: PublicationValue = serde_json::from_value(entry.value.clone()).ok()?;
let url = crate::net::safe_link(&value.url)?;
Some((
entry.uri.clone(),
Publication {
name: value.name,
url,
},
))
}
pub fn entries_from_records(
canonical_site: &str,
publication: &Publication,
records: &[crate::atproto::RecordEntry],
) -> Vec<Entry> {
let base = url::Url::parse(&publication.url).ok().map(|mut u| {
if !u.path().ends_with('/') {
u.set_path(&format!("{}/", u.path()));
}
u
});
let now = chrono::Utc::now();
let ceiling = now + chrono::Duration::seconds(crate::atproto::CLOCK_SKEW_GRACE_SECS);
records
.iter()
.filter_map(|record| {
let doc: DocumentValue = serde_json::from_value(record.value.clone()).ok()?;
if doc.site != canonical_site {
return None;
}
Some(Entry {
guid: record.uri.clone(),
title: doc.title,
published: doc
.published_at
.as_deref()
.and_then(|raw| chrono::DateTime::parse_from_rfc3339(raw).ok())
.map(|d| d.with_timezone(&chrono::Utc))
.filter(|d| *d <= ceiling)
.or_else(|| {
AtUri::parse(&record.uri)
.and_then(|uri| crate::atproto::tid_timestamp(&uri.rkey))
})
.map(crate::feed::fmt_time),
url: non_blank(doc.path)
.as_deref()
.and_then(|path| join_path(base.as_ref(), path)),
summary: non_blank(doc.description)
.or_else(|| non_blank(doc.text_content))
.map(|raw| crate::feed::plain_text_to_html(&raw)),
})
})
.collect()
}
fn ingest_floor(
retention_days: u32,
retention_hard_days: u32,
now: chrono::DateTime<chrono::Utc>,
) -> Option<String> {
let days = if retention_days > 0 {
retention_days
} else if retention_hard_days > 0 {
retention_hard_days
} else {
return None;
};
let window = chrono::Duration::try_days(days.into())?;
now.checked_sub_signed(window).map(crate::feed::fmt_time)
}
pub async fn store_publication(
pool: &sqlx::SqlitePool,
url: &str,
read: PublicationRead,
max_entries_per_feed: i64,
retention_days: u32,
retention_hard_days: u32,
) -> anyhow::Result<crate::feed::PollOutcome> {
let offered = read.entries.len();
if !read.complete && offered == 0 {
return Ok(crate::feed::PollOutcome::Failed {
backoff: crate::feed::backoff_for(1),
kind: crate::feed::FailureKind::Body,
detail: crate::feed::failure_detail(
"the publication read stopped before its first document",
),
});
}
let floor = ingest_floor(retention_days, retention_hard_days, chrono::Utc::now());
let rows: Vec<crate::store::NewEntry> = read
.entries
.into_iter()
.filter(|e| match (&floor, &e.published) {
(Some(floor), Some(published)) => published.as_str() >= floor.as_str(),
_ => true,
})
.map(Into::into)
.collect();
if offered > rows.len() {
tracing::info!(
feed = %url,
offered,
stored = rows.len(),
"the retention floor dropped documents older than the window"
);
}
let feed_id = crate::store::upsert_feed(
pool,
&crate::store::NewFeed {
url: url.to_string(),
title: read.publication.name.clone(),
site_url: Some(read.publication.url.clone()),
last_polled: Some(crate::feed::fmt_time(chrono::Utc::now())),
..Default::default()
},
)
.await?;
let new_entries =
crate::store::insert_entries(pool, feed_id, &rows, max_entries_per_feed).await?;
Ok(crate::feed::PollOutcome::Updated { new_entries })
}
fn non_blank(s: Option<String>) -> Option<String> {
s.filter(|v| !v.trim().is_empty())
}
fn join_path(base: Option<&url::Url>, path: &str) -> Option<String> {
let base = base?;
match base.join(path) {
Ok(joined) if joined.origin() == base.origin() => crate::net::safe_link(joined.as_str()),
_ => None,
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum DocumentFate {
Keep,
Sibling,
Orphan,
Malformed,
}
fn classify_document(
record: &crate::atproto::RecordEntry,
canonical_site: &str,
known: &std::collections::HashSet<&str>,
) -> DocumentFate {
match serde_json::from_value::<DocumentValue>(record.value.clone()) {
Ok(doc) if doc.site == canonical_site => DocumentFate::Keep,
Ok(doc) if known.contains(doc.site.as_str()) => DocumentFate::Sibling,
Ok(_) => DocumentFate::Orphan,
Err(_) => DocumentFate::Malformed,
}
}
pub async fn fetch(
http: &reqwest::Client,
plc_directory: &str,
uri: &AtUri,
) -> anyhow::Result<PublicationRead> {
use anyhow::Context;
anyhow::ensure!(
uri.collection == nsid::STANDARD_PUBLICATION,
"{uri} is not a {} URI",
nsid::STANDARD_PUBLICATION
);
let pds = crate::atproto::resolve_did_to_pds(http, plc_directory, &uri.authority)
.await
.with_context(|| format!("resolving the PDS for {}", uri.authority))?;
let client = crate::atproto::PdsClient::anonymous(http.clone(), pds, uri.authority.clone());
let mut budget = crate::atproto::ByteBudget::new(crate::atproto::MAX_LIST_BYTES);
let publications = client
.list_all_records_within(nsid::STANDARD_PUBLICATION, &mut budget)
.await
.with_context(|| format!("listing publications for {}", uri.authority))?;
let (canonical_site, publication) = publication_from_records(&uri.rkey, &publications)
.with_context(|| format!("{uri} is not a readable site.standard.publication"))?;
let known: std::collections::HashSet<&str> =
publications.iter().map(|p| p.uri.as_str()).collect();
let mut orphaned = 0usize;
let documents = client
.list_recent_matching_within(
nsid::STANDARD_DOCUMENT,
crate::atproto::MAX_LARGE_RECORDS,
&mut budget,
DOCUMENT_PAGE_SIZE,
|record| match classify_document(record, &canonical_site, &known) {
DocumentFate::Keep => true,
DocumentFate::Orphan => {
orphaned += 1;
false
}
DocumentFate::Sibling | DocumentFate::Malformed => false,
},
)
.await
.with_context(|| format!("listing documents for {canonical_site}"))?;
Ok(read_from(publication, &canonical_site, documents, orphaned))
}
fn read_from(
publication: Publication,
canonical_site: &str,
documents: crate::atproto::RecordWalk,
orphaned: usize,
) -> PublicationRead {
let entries = entries_from_records(canonical_site, &publication, &documents.records);
if !documents.complete {
tracing::warn!(
site = %canonical_site,
kept = entries.len(),
"stopped reading this publication before its documents ran out"
);
}
if orphaned > 0 {
tracing::warn!(
site = %canonical_site,
orphaned,
"documents in this repo reference no publication in it — a `site` spelling nothing matches"
);
}
PublicationRead {
publication,
entries,
complete: documents.complete,
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::atproto::RecordEntry;
use serde_json::json;
const DID: &str = "did:plc:ohutz6x5acjmpuulp3x7wxxc";
#[test]
fn a_read_carries_whether_the_walk_finished() {
let publication = Publication {
name: Some("Scan's Lab".to_string()),
url: "https://example.com/blog/".to_string(),
};
for complete in [true, false] {
let walk = crate::atproto::RecordWalk {
records: Vec::new(),
complete,
};
let read = read_from(publication.clone(), "at://d/c/r", walk, 0);
assert_eq!(
read.complete, complete,
"the walk said complete={complete} and the read said {}",
read.complete
);
}
}
fn read_of(entries: Vec<Entry>, complete: bool) -> PublicationRead {
PublicationRead {
publication: Publication {
name: Some("Scan's Lab".to_string()),
url: "https://example.com/blog/".to_string(),
},
entries,
complete,
}
}
fn entry_dated(guid: &str, published: Option<&str>) -> Entry {
Entry {
guid: guid.to_string(),
title: "T".to_string(),
published: published.map(str::to_string),
url: Some("https://example.com/blog/a".to_string()),
summary: None,
}
}
fn days_ago(n: i64) -> String {
crate::feed::fmt_time(chrono::Utc::now() - chrono::Duration::days(n))
}
const PUB_URL: &str = "at://did:plc:ohutz6x5acjmpuulp3x7wxxc/site.standard.publication/3lab";
async fn pool() -> sqlx::SqlitePool {
crate::store::init_url("sqlite::memory:").await.unwrap()
}
#[tokio::test]
async fn a_complete_read_stores_its_entries_and_reports_them() {
let pool = pool().await;
let read = read_of(
vec![
entry_dated("at://d/c/1", Some(&days_ago(1))),
entry_dated("at://d/c/2", Some(&days_ago(2))),
],
true,
);
let outcome = store_publication(&pool, PUB_URL, read, 0, 14, 180)
.await
.unwrap();
assert!(
matches!(
outcome,
crate::feed::PollOutcome::Updated { new_entries: 2 }
),
"expected two new entries, got {outcome:?}"
);
let n: i64 = sqlx::query_scalar("SELECT COUNT(*) FROM entries")
.fetch_one(&pool)
.await
.unwrap();
assert_eq!(n, 2, "the entries were not stored");
}
#[tokio::test]
async fn an_incomplete_read_that_offered_nothing_is_a_failure() {
let pool = pool().await;
let outcome = store_publication(&pool, PUB_URL, read_of(vec![], false), 0, 14, 180)
.await
.unwrap();
assert!(
matches!(outcome, crate::feed::PollOutcome::Failed { .. }),
"a truncated read that produced nothing is not a healthy poll: {outcome:?}"
);
}
#[tokio::test]
async fn an_incomplete_read_whose_entries_the_floor_dropped_is_not_a_failure() {
let pool = pool().await;
let read = read_of(
vec![entry_dated("at://d/c/old", Some(&days_ago(900)))],
false,
);
let outcome = store_publication(&pool, PUB_URL, read, 0, 14, 180)
.await
.unwrap();
assert!(
!matches!(outcome, crate::feed::PollOutcome::Failed { .. }),
"the read offered an entry; the floor dropping it is not a failed poll: {outcome:?}"
);
}
#[tokio::test]
async fn a_failed_read_does_not_stamp_last_polled() {
let pool = pool().await;
let outcome = store_publication(&pool, PUB_URL, read_of(vec![], false), 0, 14, 180)
.await
.unwrap();
assert!(matches!(outcome, crate::feed::PollOutcome::Failed { .. }));
let stamped: Option<String> =
sqlx::query_scalar("SELECT last_polled FROM feeds WHERE url = ?1")
.bind(PUB_URL)
.fetch_optional(&pool)
.await
.unwrap()
.flatten();
assert_eq!(
stamped, None,
"a failed poll stamped last_polled, so the feed reads as freshly polled"
);
}
#[tokio::test]
async fn an_entry_already_older_than_the_window_is_not_stored() {
let pool = pool().await;
let read = read_of(
vec![
entry_dated("at://d/c/fresh", Some(&days_ago(1))),
entry_dated("at://d/c/ancient", Some(&days_ago(900))),
],
true,
);
store_publication(&pool, PUB_URL, read, 0, 14, 180)
.await
.unwrap();
let guids: Vec<String> = sqlx::query_scalar("SELECT guid FROM entries ORDER BY guid")
.fetch_all(&pool)
.await
.unwrap();
assert_eq!(
guids,
vec!["at://d/c/fresh".to_string()],
"an entry the next sweep would delete was stored anyway"
);
}
#[tokio::test]
async fn the_floor_follows_the_hard_ceiling_when_the_window_is_disabled() {
let pool = pool().await;
let read = read_of(
vec![entry_dated("at://d/c/ancient", Some(&days_ago(900)))],
true,
);
store_publication(&pool, PUB_URL, read, 0, 0, 180)
.await
.unwrap();
let n: i64 = sqlx::query_scalar("SELECT COUNT(*) FROM entries")
.fetch_one(&pool)
.await
.unwrap();
assert_eq!(
n, 0,
"an entry the hard ceiling will delete was stored, so it will resurrect"
);
}
#[tokio::test]
async fn the_floor_follows_the_window_when_both_are_set() {
let pool = pool().await;
let read = read_of(
vec![entry_dated("at://d/c/hundred", Some(&days_ago(100)))],
true,
);
store_publication(&pool, PUB_URL, read, 0, 14, 180)
.await
.unwrap();
let n: i64 = sqlx::query_scalar("SELECT COUNT(*) FROM entries")
.fetch_one(&pool)
.await
.unwrap();
assert_eq!(
n, 0,
"an entry inside the ceiling but outside the window was stored, so it will cycle"
);
}
#[test]
fn a_publications_floor_is_its_archive_ceiling_not_the_rss_window() {
let config = crate::config::Config::default();
let now = chrono::DateTime::parse_from_rfc3339("2026-09-29T12:00:00Z")
.unwrap()
.with_timezone(&chrono::Utc);
let (days, hard) = config.retention_for(crate::feed::FeedKind::Publication);
let floor = ingest_floor(days, hard, now).expect("a publication has an ingest floor");
let expected = crate::feed::fmt_time(
now - chrono::Duration::days(config.publication_retention_days.into()),
);
assert_eq!(
floor, expected,
"a publication's floor must be its archive ceiling ({} days), because \
that is the only sweep pass that can delete its rows",
config.publication_retention_days,
);
let rss_floor = ingest_floor(config.retention_days, config.retention_hard_days, now)
.expect("an RSS feed has an ingest floor");
assert!(
floor < rss_floor,
"the publication floor ({floor}) is no older than the RSS one \
({rss_floor}), so a months-old document would still be dropped",
);
let a_real_publications_newest_document =
crate::feed::fmt_time(now - chrono::Duration::days(109));
assert!(
a_real_publications_newest_document.as_str() >= floor.as_str(),
"the newest document a real publication offered would be refused at \
ingest: {a_real_publications_newest_document} against a floor of {floor}",
);
assert!(
a_real_publications_newest_document.as_str() < rss_floor.as_str(),
"this assertion is only meaningful while the RSS window WOULD have \
dropped it, and it no longer does",
);
}
#[test]
fn the_ingest_floor_is_the_window_the_sweep_would_use() {
let now = chrono::DateTime::parse_from_rfc3339("2026-09-26T12:00:00Z")
.unwrap()
.with_timezone(&chrono::Utc);
assert_eq!(
ingest_floor(14, 180, now).as_deref(),
Some("2026-09-12T12:00:00Z"),
"with both set, the floor is the WINDOW — the thing that deletes first",
);
assert_eq!(
ingest_floor(180, 30, now).as_deref(),
Some("2026-03-30T12:00:00Z"),
"a ceiling INSIDE the window is one `prune_old_entries` ignores, so it \
must not lower the floor — `min` here discarded five months of archive \
that nothing would have deleted",
);
assert_eq!(
ingest_floor(0, 30, now).as_deref(),
Some("2026-08-27T12:00:00Z"),
"with no window the ceiling stands alone, and it still deletes",
);
assert_eq!(
ingest_floor(0, 0, now),
None,
"with no retention at all there is nothing to floor against",
);
assert_eq!(
ingest_floor(u32::MAX, u32::MAX, now),
None,
"an unrepresentable window produced a floor, so the comparison is \
resting on how `fmt_time` renders an out-of-range year",
);
}
#[tokio::test]
async fn a_ceiling_inside_the_window_does_not_lower_the_floor() {
let pool = pool().await;
let read = read_of(
vec![entry_dated("at://d/c/hundred", Some(&days_ago(100)))],
true,
);
store_publication(&pool, PUB_URL, read, 0, 180, 30)
.await
.unwrap();
let n: i64 = sqlx::query_scalar("SELECT COUNT(*) FROM entries")
.fetch_one(&pool)
.await
.unwrap();
assert_eq!(
n, 1,
"an entry inside the 180-day window was dropped because of a 30-day \
ceiling the sweep ignores — five months of archive discarded at ingest \
that nothing would have deleted",
);
}
#[tokio::test]
async fn a_complete_read_of_nothing_is_not_a_failure() {
let pool = pool().await;
let outcome = store_publication(&pool, PUB_URL, read_of(vec![], true), 0, 14, 180)
.await
.unwrap();
assert!(
matches!(
outcome,
crate::feed::PollOutcome::Updated { new_entries: 0 }
),
"a complete read of an empty publication was not a healthy poll: {outcome:?}"
);
}
#[tokio::test]
async fn a_successful_read_stamps_last_polled() {
let pool = pool().await;
let read = read_of(vec![entry_dated("at://d/c/1", Some(&days_ago(1)))], true);
store_publication(&pool, PUB_URL, read, 0, 14, 180)
.await
.unwrap();
let stamped: Option<String> =
sqlx::query_scalar("SELECT last_polled FROM feeds WHERE url = ?1")
.bind(PUB_URL)
.fetch_one(&pool)
.await
.unwrap();
assert!(
stamped.is_some(),
"a successful read left `last_polled` NULL, so the publication stays \
due forever and `/stats` never shows it as polled",
);
}
#[tokio::test]
async fn an_absurd_retention_window_keeps_everything_rather_than_nothing() {
let pool = pool().await;
let read = read_of(
vec![entry_dated("at://d/c/ancient", Some(&days_ago(10_000)))],
true,
);
store_publication(&pool, PUB_URL, read, 0, u32::MAX, u32::MAX)
.await
.expect("a huge window is a wide floor, not a crash");
let n: i64 = sqlx::query_scalar("SELECT COUNT(*) FROM entries")
.fetch_one(&pool)
.await
.unwrap();
assert_eq!(
n, 1,
"a window of u32::MAX days dropped a 27-year-old entry, so the \
saturation went the wrong way",
);
}
#[tokio::test]
async fn the_per_feed_cap_is_the_one_the_caller_passed() {
let pool = pool().await;
let read = read_of(
(0..5)
.map(|i| entry_dated(&format!("at://d/c/{i}"), Some(&days_ago(i + 1))))
.collect(),
true,
);
store_publication(&pool, PUB_URL, read, 2, 0, 0)
.await
.unwrap();
let n: i64 = sqlx::query_scalar("SELECT COUNT(*) FROM entries")
.fetch_one(&pool)
.await
.unwrap();
assert_eq!(
n, 2,
"five entries under a cap of two left {n} rows, so the caller's cap is \
not the one being applied",
);
}
#[tokio::test]
async fn no_retention_at_all_means_no_ingest_floor() {
let pool = pool().await;
let read = read_of(
vec![entry_dated("at://d/c/ancient", Some(&days_ago(900)))],
true,
);
store_publication(&pool, PUB_URL, read, 0, 0, 0)
.await
.unwrap();
let n: i64 = sqlx::query_scalar("SELECT COUNT(*) FROM entries")
.fetch_one(&pool)
.await
.unwrap();
assert_eq!(
n, 1,
"with nothing deleting it, an old entry is worth keeping"
);
}
#[tokio::test]
async fn an_absurd_retention_window_does_not_panic() {
let pool = pool().await;
let read = read_of(vec![entry_dated("at://d/c/x", Some(&days_ago(1)))], true);
store_publication(&pool, PUB_URL, read, 0, u32::MAX, u32::MAX)
.await
.expect("a huge window is a wide floor, not a crash");
}
#[tokio::test]
async fn the_feed_row_learns_the_publications_name_and_site() {
let pool = pool().await;
let read = read_of(vec![entry_dated("at://d/c/1", Some(&days_ago(1)))], true);
store_publication(&pool, PUB_URL, read, 0, 14, 180)
.await
.unwrap();
let (title, site): (Option<String>, Option<String>) =
sqlx::query_as("SELECT title, site_url FROM feeds WHERE url = ?1")
.bind(PUB_URL)
.fetch_one(&pool)
.await
.unwrap();
assert_eq!(
title.as_deref(),
Some("Scan's Lab"),
"the name never reached the row"
);
assert_eq!(
site.as_deref(),
Some("https://example.com/blog/"),
"the homepage never reached the row"
);
}
#[tokio::test]
async fn an_undated_entry_is_stored_rather_than_dropped() {
let pool = pool().await;
let read = read_of(vec![entry_dated("at://d/c/undated", None)], true);
store_publication(&pool, PUB_URL, read, 0, 14, 180)
.await
.unwrap();
let n: i64 = sqlx::query_scalar("SELECT COUNT(*) FROM entries")
.fetch_one(&pool)
.await
.unwrap();
assert_eq!(n, 1, "an entry with no date was discarded");
}
const PAST_TID: &str = "3jzfcijpj2z2a";
const PAST_TID_WRITTEN_AT: &str = "2023-06-30T15:03:01Z";
fn rec(collection: &str, rkey: &str, value: serde_json::Value) -> RecordEntry {
RecordEntry {
uri: format!("at://{DID}/{collection}/{rkey}"),
cid: None,
value,
}
}
fn publication(rkey: &str, url: &str) -> RecordEntry {
rec(
nsid::STANDARD_PUBLICATION,
rkey,
json!({ "name": "Scan's Lab", "url": url }),
)
}
fn document(rkey: &str, site: &str, title: &str, path: &str) -> RecordEntry {
rec(
nsid::STANDARD_DOCUMENT,
rkey,
json!({
"title": title,
"publishedAt": "2026-07-11T00:00:00Z",
"path": path,
"site": site,
"textContent": "body",
}),
)
}
fn canonical(rkey: &str) -> String {
format!("at://{DID}/{}/{rkey}", nsid::STANDARD_PUBLICATION)
}
#[test]
fn the_site_filter_uses_the_uri_the_pds_minted() {
let records = vec![publication("p", "https://scanash.com")];
let (site, pubn) = publication_from_records("p", &records).expect("publication not found");
assert_eq!(site, canonical("p"), "did not take the PDS's canonical URI");
let docs = vec![document("d1", &canonical("p"), "Hello", "/hello")];
let entries = entries_from_records(&site, &pubn, &docs);
assert_eq!(
entries.len(),
1,
"a canonical-site document was not matched"
);
}
#[test]
fn an_entry_maps_onto_the_stores_row() {
let records = vec![publication("p", "https://example.com")];
let (site, pubn) = publication_from_records("p", &records).unwrap();
let docs = vec![document("rk1", &site, "Hello", "/hello")];
let row: crate::store::NewEntry = entries_from_records(&site, &pubn, &docs)
.pop()
.unwrap()
.into();
assert_eq!(
row.guid,
format!("at://{DID}/{}/rk1", nsid::STANDARD_DOCUMENT)
);
assert_eq!(row.url.as_deref(), Some("https://example.com/hello"));
assert_eq!(row.title.as_deref(), Some("Hello"));
assert_eq!(row.published.as_deref(), Some("2026-07-11T00:00:00Z"));
assert_eq!(row.content_html.as_deref(), Some("body"));
assert_eq!(row.author, None);
assert_eq!(row.fetched_at, None);
}
#[test]
fn documents_are_filtered_by_their_site_field() {
let records = vec![publication("mine", "https://example.com")];
let (site, pubn) = publication_from_records("mine", &records).unwrap();
let docs = vec![
document("a", &site, "Mine", "/a"),
document("b", &canonical("theirs"), "Theirs", "/b"),
document("c", &site, "Mine again", "/c"),
];
let titles: Vec<String> = entries_from_records(&site, &pubn, &docs)
.into_iter()
.map(|e| e.title)
.collect();
assert_eq!(titles, ["Mine", "Mine again"]);
}
#[test]
fn a_publication_with_a_hostile_url_is_refused() {
for hostile in [
"javascript:alert(1)",
"data:text/html,<script>",
"file:///etc/passwd",
"",
] {
let records = vec![publication("p", hostile)];
assert!(
publication_from_records("p", &records).is_none(),
"accepted a publication whose url is {hostile:?}",
);
}
}
#[test]
fn entry_urls_are_joined_against_the_publication_base() {
let records = vec![publication("p", "https://example.com/blog")];
let (site, pubn) = publication_from_records("p", &records).unwrap();
let docs = vec![
document("a", &site, "Relative", "/a"),
document("b", &site, "Absolute-looking", "https://evil.example/x"),
];
let urls: Vec<Option<String>> = entries_from_records(&site, &pubn, &docs)
.into_iter()
.map(|e| e.url)
.collect();
assert_eq!(urls[0].as_deref(), Some("https://example.com/a"));
assert_eq!(
urls[1], None,
"a document path that escapes its publication's origin must yield no URL"
);
}
#[test]
fn a_document_with_neither_summary_field_still_yields_an_entry() {
let records = vec![publication("p", "https://example.com/")];
let (site, pubn) = publication_from_records("p", &records).unwrap();
let bare = rec(
nsid::STANDARD_DOCUMENT,
"bare",
json!({
"title": "Bare",
"publishedAt": "2026-07-11T00:00:00Z",
"path": "/bare",
"site": site,
}),
);
let entries = entries_from_records(&site, &pubn, &[bare]);
assert_eq!(entries.len(), 1);
assert_eq!(entries[0].summary, None);
assert_eq!(entries[0].url.as_deref(), Some("https://example.com/bare"));
}
#[test]
fn an_empty_description_does_not_shadow_the_body() {
let records = vec![publication("p", "https://example.com")];
let (site, pubn) = publication_from_records("p", &records).unwrap();
let doc = rec(
nsid::STANDARD_DOCUMENT,
"d",
json!({
"title": "T",
"publishedAt": "2026-07-11T00:00:00Z",
"path": "/d",
"site": site,
"description": " ",
"textContent": "the real body",
}),
);
let entries = entries_from_records(&site, &pubn, &[doc]);
assert_eq!(entries[0].summary.as_deref(), Some("the real body"));
}
#[test]
fn published_at_is_normalised_or_dropped() {
let records = vec![publication("p", "https://example.com")];
let (site, pubn) = publication_from_records("p", &records).unwrap();
let with = |rkey: &str, published_at: serde_json::Value| {
rec(
nsid::STANDARD_DOCUMENT,
rkey,
json!({ "title": "T", "publishedAt": published_at, "path": "/x", "site": site }),
)
};
let docs = vec![
with("a", json!("2026-07-11T09:30:00.123+02:00")),
with("b", json!("yesterday-ish")),
with("c", json!("2026-07-11T00:00:00Z")),
];
let published: Vec<Option<String>> = entries_from_records(&site, &pubn, &docs)
.into_iter()
.map(|e| e.published)
.collect();
assert_eq!(
published,
vec![
Some("2026-07-11T07:30:00Z".to_string()),
None,
Some("2026-07-11T00:00:00Z".to_string()),
],
"publishedAt was not normalised to the store's spelling"
);
}
#[test]
fn an_undated_document_is_dated_from_its_tid_rkey() {
let records = vec![publication("p", "https://example.com")];
let (site, pubn) = publication_from_records("p", &records).unwrap();
let docs = vec![rec(
nsid::STANDARD_DOCUMENT,
PAST_TID,
json!({ "title": "T", "path": "/x", "site": site }),
)];
assert_eq!(
entries_from_records(&site, &pubn, &docs)
.into_iter()
.next()
.expect("the document is an entry")
.published,
Some(PAST_TID_WRITTEN_AT.to_string()),
"the date must come from the record key, and be spelled the way the store spells dates"
);
}
#[test]
fn an_unparseable_published_at_falls_back_to_the_tid_rkey() {
let records = vec![publication("p", "https://example.com")];
let (site, pubn) = publication_from_records("p", &records).unwrap();
let docs = vec![rec(
nsid::STANDARD_DOCUMENT,
PAST_TID,
json!({ "title": "T", "publishedAt": "yesterday-ish", "path": "/x", "site": site }),
)];
assert_eq!(
entries_from_records(&site, &pubn, &docs)
.into_iter()
.next()
.expect("the document is an entry")
.published,
Some(PAST_TID_WRITTEN_AT.to_string()),
"a date the parser cannot read is no date at all, so the rkey must stand in"
);
}
#[test]
fn a_stated_date_outranks_the_rkey() {
let records = vec![publication("p", "https://example.com")];
let (site, pubn) = publication_from_records("p", &records).unwrap();
let docs = vec![rec(
nsid::STANDARD_DOCUMENT,
PAST_TID,
json!({ "title": "T", "publishedAt": "2020-01-02T00:00:00Z", "path": "/x", "site": site }),
)];
assert_eq!(
entries_from_records(&site, &pubn, &docs)
.into_iter()
.next()
.expect("the document is an entry")
.published,
Some("2020-01-02T00:00:00Z".to_string()),
"the rkey records when the file was written, which is not when the post was published"
);
}
#[test]
fn a_future_dated_document_falls_back_to_its_rkey() {
let records = vec![publication("p", "https://example.com")];
let (site, pubn) = publication_from_records("p", &records).unwrap();
let docs = vec![rec(
nsid::STANDARD_DOCUMENT,
PAST_TID,
json!({ "title": "T", "publishedAt": "2999-01-01T00:00:00Z", "path": "/x", "site": site }),
)];
assert_eq!(
entries_from_records(&site, &pubn, &docs)
.into_iter()
.next()
.expect("the document is an entry")
.published,
Some(PAST_TID_WRITTEN_AT.to_string()),
"the date must be the record's write time, not the hour the poll happened to run"
);
}
#[test]
fn a_future_dated_document_without_a_tid_rkey_is_undated() {
let records = vec![publication("p", "https://example.com")];
let (site, pubn) = publication_from_records("p", &records).unwrap();
let docs = vec![rec(
nsid::STANDARD_DOCUMENT,
"self",
json!({ "title": "T", "publishedAt": "2999-01-01T00:00:00Z", "path": "/x", "site": site }),
)];
assert_eq!(
entries_from_records(&site, &pubn, &docs)
.into_iter()
.next()
.expect("the document is an entry")
.published,
None,
"with nothing credible to date it by, the row falls to fetched_at, which holds still"
);
}
#[test]
fn a_stated_date_a_little_ahead_of_our_clock_is_still_believed() {
let records = vec![publication("p", "https://example.com")];
let (site, pubn) = publication_from_records("p", &records).unwrap();
let slightly_ahead =
crate::feed::fmt_time(chrono::Utc::now() + chrono::Duration::seconds(10));
let docs = vec![rec(
nsid::STANDARD_DOCUMENT,
"self",
json!({ "title": "T", "publishedAt": slightly_ahead, "path": "/x", "site": site }),
)];
assert_eq!(
entries_from_records(&site, &pubn, &docs)
.into_iter()
.next()
.expect("the document is an entry")
.published,
Some(slightly_ahead),
"a few seconds of clock skew must not cost the entry its date"
);
}
#[test]
fn a_document_with_neither_a_date_nor_a_tid_rkey_stays_undated() {
let records = vec![publication("p", "https://example.com")];
let (site, pubn) = publication_from_records("p", &records).unwrap();
let docs = vec![
rec(
nsid::STANDARD_DOCUMENT,
"my-first-post",
json!({ "title": "T", "path": "/x", "site": site }),
),
rec(
nsid::STANDARD_DOCUMENT,
"abcdefghijklm",
json!({ "title": "T", "path": "/y", "site": site }),
),
];
assert_eq!(
entries_from_records(&site, &pubn, &docs)
.into_iter()
.map(|e| e.published)
.collect::<Vec<_>>(),
vec![None, None],
"an invented date is worse than no date; the store decides what to do with undated rows"
);
}
#[test]
fn summaries_are_escaped_as_plain_text_not_sanitised_as_markup() {
let records = vec![publication("p", "https://example.com")];
let (site, pubn) = publication_from_records("p", &records).unwrap();
let doc = |rkey: &str, body: &str| {
rec(
nsid::STANDARD_DOCUMENT,
rkey,
json!({
"title": "T",
"publishedAt": "2026-07-11T00:00:00Z",
"path": "/d",
"site": site,
"textContent": body,
}),
)
};
let summaries: Vec<String> = entries_from_records(
&site,
&pubn,
&[
doc("a", "Vec<String> is a type"),
doc("b", "<script>alert(1)</script>"),
],
)
.into_iter()
.filter_map(|e| e.summary)
.collect();
assert_eq!(
summaries[0], "Vec<String> is a type",
"prose was eaten by an HTML parser"
);
assert!(
!summaries[1].contains("<script"),
"escaping failed: {}",
summaries[1]
);
}
#[test]
fn an_entry_url_is_never_an_unvetted_path() {
let pubn = Publication {
name: None,
url: "not a url".to_string(),
};
let site = canonical("p");
let docs = vec![
document("a", &site, "Hostile", "javascript:alert(1)"),
document("b", &site, "Fine", "https://example.com/ok"),
];
let entries = entries_from_records(&site, &pubn, &docs);
assert_eq!(
entries[0].url, None,
"an unvetted path became an entry link"
);
assert_eq!(
entries[1].url, None,
"an off-origin absolute URL was published under the publication's name"
);
}
#[test]
fn a_document_without_published_at_is_still_an_entry() {
let records = vec![publication("p", "https://example.com")];
let (site, pubn) = publication_from_records("p", &records).unwrap();
let doc = rec(
nsid::STANDARD_DOCUMENT,
"d",
json!({ "title": "T", "path": "/d", "site": site }),
);
let entries = entries_from_records(&site, &pubn, &[doc]);
assert_eq!(
entries.len(),
1,
"a missing publishedAt dropped the document"
);
assert_eq!(entries[0].published, None);
}
#[tokio::test]
async fn fetch_refuses_a_uri_for_another_collection() {
let uri = AtUri::parse(&format!("at://{DID}/app.bsky.feed.post/3lab")).unwrap();
let err = fetch(&reqwest::Client::new(), "https://plc.example", &uri)
.await
.expect_err("read a feed post as a publication");
assert!(
format!("{err:#}").contains(nsid::STANDARD_PUBLICATION),
"failed for the wrong reason: {err:#}"
);
}
#[test]
fn a_subpath_publication_keeps_its_base_path() {
let records = vec![publication("p", "https://example.com/blog")];
let (site, pubn) = publication_from_records("p", &records).unwrap();
let docs = vec![document("a", &site, "Relative", "posts/a")];
let urls: Vec<Option<String>> = entries_from_records(&site, &pubn, &docs)
.into_iter()
.map(|e| e.url)
.collect();
assert_eq!(urls[0].as_deref(), Some("https://example.com/blog/posts/a"));
}
#[test]
fn no_parseable_base_means_no_url_not_any_url() {
let pubn = Publication {
name: None,
url: "not a url".to_string(),
};
let site = canonical("p");
let docs = vec![document("a", &site, "Absolute", "https://evil.example/x")];
let entries = entries_from_records(&site, &pubn, &docs);
assert_eq!(
entries[0].url, None,
"an off-origin absolute URL was published"
);
}
#[test]
fn an_off_origin_path_yields_no_url_rather_than_the_homepage() {
let records = vec![publication("p", "https://example.com/blog")];
let (site, pubn) = publication_from_records("p", &records).unwrap();
let docs = vec![
document("a", &site, "Elsewhere", "https://www.example.com/post"),
document("b", &site, "Home", "/ok"),
];
let urls: Vec<Option<String>> = entries_from_records(&site, &pubn, &docs)
.into_iter()
.map(|e| e.url)
.collect();
assert_eq!(
urls[0], None,
"an off-origin path was rewritten to the base"
);
assert_eq!(urls[1].as_deref(), Some("https://example.com/ok"));
}
#[test]
fn a_blank_path_yields_no_url() {
let records = vec![publication("p", "https://example.com/blog")];
let (site, pubn) = publication_from_records("p", &records).unwrap();
let docs = vec![
document("a", &site, "Blank", ""),
document("b", &site, "Spaces", " "),
document("c", &site, "Real", "/real"),
];
let urls: Vec<Option<String>> = entries_from_records(&site, &pubn, &docs)
.into_iter()
.map(|e| e.url)
.collect();
assert_eq!(urls[0], None, "a blank path became the homepage");
assert_eq!(urls[1], None, "a whitespace path became the homepage");
assert_eq!(urls[2].as_deref(), Some("https://example.com/real"));
}
#[test]
fn a_document_without_a_path_is_still_an_entry() {
let records = vec![publication("p", "https://example.com")];
let (site, pubn) = publication_from_records("p", &records).unwrap();
let doc = rec(
nsid::STANDARD_DOCUMENT,
"d",
json!({ "title": "T", "publishedAt": "2026-07-11T00:00:00Z", "site": site }),
);
let entries = entries_from_records(&site, &pubn, &[doc]);
assert_eq!(
entries.len(),
1,
"a missing path dropped the whole document"
);
assert_eq!(entries[0].title, "T");
assert_eq!(entries[0].url, None);
}
#[test]
fn documents_are_classified_keep_sibling_or_orphan() {
let pubs = [
publication("a", "https://example.com"),
publication("b", "https://b.example"),
];
let known: std::collections::HashSet<&str> = pubs.iter().map(|p| p.uri.as_str()).collect();
let mine = canonical("a");
let fate = |d: &crate::atproto::RecordEntry| classify_document(d, &mine, &known);
assert_eq!(
fate(&document("d1", &mine, "Mine", "/1")),
DocumentFate::Keep
);
assert_eq!(
fate(&document("d2", &canonical("b"), "B's", "/2")),
DocumentFate::Sibling,
"a sibling publication's document is not an orphan"
);
assert_eq!(
fate(&document(
"d3",
"at://did:plc:other/site.standard.publication/x",
"?",
"/3"
)),
DocumentFate::Orphan
);
assert_eq!(
fate(&rec(
nsid::STANDARD_DOCUMENT,
"d4",
json!({"title": "no rest"})
)),
DocumentFate::Malformed
);
}
#[test]
fn at_uri_parsing_uses_the_shared_prefix() {
let uri = format!(
"{}{DID}/{}/abc",
crate::atproto::AT_URI_PREFIX,
nsid::STANDARD_PUBLICATION
);
assert!(AtUri::parse(&uri).is_some());
}
#[test]
fn the_guid_is_the_record_uri_not_the_path() {
let records = vec![publication("p", "https://example.com")];
let (site, pubn) = publication_from_records("p", &records).unwrap();
let docs = vec![document("rk1", &site, "T", "/moved")];
let entries = entries_from_records(&site, &pubn, &docs);
assert_eq!(
entries[0].guid,
format!("at://{DID}/{}/rk1", nsid::STANDARD_DOCUMENT)
);
}
#[test]
fn a_malformed_document_is_skipped_rather_than_fatal() {
let records = vec![publication("p", "https://example.com")];
let (site, pubn) = publication_from_records("p", &records).unwrap();
let docs = vec![
rec(
nsid::STANDARD_DOCUMENT,
"bad",
json!({ "title": "no rest" }),
),
document("ok", &site, "Good", "/good"),
];
let entries = entries_from_records(&site, &pubn, &docs);
assert_eq!(entries.len(), 1);
assert_eq!(entries[0].title, "Good");
}
#[test]
fn a_missing_publication_is_none() {
let records = vec![publication("other", "https://example.com")];
assert!(publication_from_records("p", &records).is_none());
}
#[test]
fn at_uris_parse_in_both_forms_and_reject_malformed_ones() {
let did = AtUri::parse(&format!("at://{DID}/site.standard.publication/abc")).unwrap();
assert_eq!(did.authority, DID);
assert_eq!(did.rkey, "abc");
assert_eq!(
did.to_string(),
format!("at://{DID}/site.standard.publication/abc")
);
assert!(AtUri::parse("at://alice.example.com/site.standard.publication/abc").is_some());
for bad in [
"at://",
"at://only-authority",
"at://authority/collection",
"at://authority/collection/",
"at:///collection/rkey",
"at://authority/collection/rkey/extra",
"https://example.com/feed.xml",
"at:authority/collection/rkey",
] {
assert!(AtUri::parse(bad).is_none(), "parsed {bad:?}");
}
}
}